spawn-runner.ts 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429
  1. /** One-shot Linux exec bootstrap and Windows Job-owning subprocess runner. */
  2. import { delimiter, resolve } from 'node:path'
  3. import {
  4. closeCurrentProcessStandardHandles,
  5. closeHandleChecked,
  6. isJobEmpty,
  7. loadWin32ProcessBindings,
  8. pollProcessExit,
  9. spawnCurrentTokenJobProcess,
  10. terminateJob,
  11. Win32Error,
  12. } from '@deepseek-ai/dsh-win32-process'
  13. import type { NativePtr, Win32ProcessBindings } from '@deepseek-ai/dsh-win32-process'
  14. import {
  15. consumeLinuxLaunchRequest,
  16. isWindowsTerminateRequest,
  17. linuxLaunchFilesFromLocator,
  18. parseWindowsStartRequest,
  19. serializeRunnerError,
  20. writeLinuxStartupError,
  21. } from './runner-protocol.ts'
  22. import type {
  23. LinuxLaunchFiles,
  24. SerializedRunnerError,
  25. WindowsRunnerResult,
  26. WindowsStartRequest,
  27. } from './runner-protocol.ts'
  28. import {
  29. parseRunnerTargetArgv,
  30. SUBPROCESS_RUNNER_ENV,
  31. WINDOWS_RUNNER_SELECTION,
  32. } from './runner-launch.ts'
  33. type RunnerHost = Pick<NodeJS.Process, 'env' | 'exitCode' | 'connected' | 'cwd' | 'chdir' | 'on' | 'off' | 'once' | 'disconnect'> & {
  34. send?: NodeJS.Process['send']
  35. }
  36. /** Injectable operations used by the protocol-owner tests. */
  37. export interface SpawnRunnerInternals {
  38. execve(file: string, argv: string[], env: Record<string, string>): never
  39. loadWin32ProcessBindings(): Win32ProcessBindings
  40. spawnCurrentTokenJobProcess: typeof spawnCurrentTokenJobProcess
  41. closeCurrentProcessStandardHandles: typeof closeCurrentProcessStandardHandles
  42. pollProcessExit: typeof pollProcessExit
  43. isJobEmpty: typeof isJobEmpty
  44. terminateJob: typeof terminateJob
  45. closeHandleChecked: typeof closeHandleChecked
  46. }
  47. const defaultInternals: SpawnRunnerInternals = {
  48. /* v8 ignore next -- source/built/packaged subprocess smoke executes this only in a replaceable child process. */
  49. execve: (file, argv, env) => (process.execve as NonNullable<typeof process.execve>)(file, argv, env),
  50. loadWin32ProcessBindings,
  51. spawnCurrentTokenJobProcess,
  52. closeCurrentProcessStandardHandles,
  53. pollProcessExit,
  54. isJobEmpty,
  55. terminateJob,
  56. closeHandleChecked,
  57. }
  58. function replaceEnvironment(target: NodeJS.ProcessEnv, env: Record<string, string>): void {
  59. for (const key of Object.keys(target)) Reflect.deleteProperty(target, key)
  60. Object.assign(target, env)
  61. }
  62. function asSpawnError(error: unknown, program: string, args: readonly string[]): SerializedRunnerError {
  63. const serialized = serializeRunnerError(error)
  64. const code = error instanceof Win32Error
  65. ? error.win32Code === 2 || error.win32Code === 3 || error.win32Code === 267
  66. ? 'ENOENT'
  67. : error.win32Code === 5
  68. ? 'EACCES'
  69. : error.win32Code === 193
  70. ? 'EFTYPE'
  71. : 'UNKNOWN'
  72. : serialized.code
  73. if (code === undefined) return serialized
  74. return {
  75. ...serialized,
  76. message: `spawn ${program} ${code}: ${serialized.message}`,
  77. code,
  78. syscall: `spawn ${program}`,
  79. path: program,
  80. spawnargs: [...args],
  81. }
  82. }
  83. function linuxPathNotFoundError(program: string): NodeJS.ErrnoException {
  84. return Object.assign(new Error(`spawn ${program} ENOENT`), {
  85. code: 'ENOENT',
  86. errno: -2,
  87. syscall: `spawn ${program}`,
  88. path: program,
  89. spawnargs: [] as string[],
  90. })
  91. }
  92. function execLinuxTarget(
  93. request: { cwd: string; env: Record<string, string> },
  94. argv: string[],
  95. internals: SpawnRunnerInternals,
  96. ): never {
  97. const program = argv[0] as string
  98. if (program.includes('/')) return internals.execve(program, argv, request.env)
  99. const path = request.env.PATH ?? '/usr/bin:/bin'
  100. let permissionFailure: Error | undefined
  101. for (const directory of path.split(delimiter)) {
  102. const candidate = resolve(request.cwd, directory, program)
  103. try {
  104. return internals.execve(candidate, argv, request.env)
  105. } catch (error) {
  106. const code = (error as NodeJS.ErrnoException).code
  107. if (code === 'EACCES') {
  108. permissionFailure ??= error as Error
  109. continue
  110. }
  111. if (code === 'ENOENT' || code === 'ENOTDIR') continue
  112. throw error
  113. }
  114. }
  115. throw permissionFailure ?? linuxPathNotFoundError(program)
  116. }
  117. function runLinux(
  118. locator: string,
  119. argv: string[],
  120. host: RunnerHost,
  121. internals: SpawnRunnerInternals,
  122. ): void {
  123. const files = linuxLaunchFilesFromLocator(locator)
  124. let request: ReturnType<typeof consumeLinuxLaunchRequest>
  125. try {
  126. request = consumeLinuxLaunchRequest(files.requestPath)
  127. } catch (error) {
  128. writeLinuxStartupError(files, { type: 'runner-error', error: serializeRunnerError(error) })
  129. host.exitCode = 127
  130. return
  131. }
  132. try {
  133. host.chdir(request.cwd)
  134. execLinuxTarget(request, argv, internals)
  135. } catch (error) {
  136. writeLinuxStartupError(files, {
  137. type: 'spawn-error',
  138. error: asSpawnError(error, argv[0] as string, argv.slice(1)),
  139. })
  140. host.exitCode = 127
  141. }
  142. }
  143. function sendMessage(host: RunnerHost, result: WindowsRunnerResult): Promise<void> {
  144. return new Promise((resolve, reject) => {
  145. if (!host.connected || host.send === undefined) {
  146. reject(new Error('subprocess runner IPC is not connected'))
  147. return
  148. }
  149. try {
  150. host.send(result, (error) => {
  151. if (error === null) resolve()
  152. else reject(error)
  153. })
  154. } catch (error) {
  155. /* v8 ignore next -- process.send throws Error instances. */
  156. const failure = error instanceof Error ? error : new Error(String(error))
  157. reject(failure)
  158. }
  159. })
  160. }
  161. class WindowsJobRunner {
  162. private api: Win32ProcessBindings | undefined
  163. private processHandle: NativePtr | undefined
  164. private jobHandle: NativePtr | undefined
  165. private pollTimer: ReturnType<typeof setInterval> | undefined
  166. private startSeen = false
  167. private committed = false
  168. private terminateRequested = false
  169. private resultStarted = false
  170. private resultDelivered = false
  171. private jobEmpty = false
  172. private finished = false
  173. private readonly completion = Promise.withResolvers<void>()
  174. constructor(
  175. private readonly argv: string[],
  176. private readonly host: RunnerHost,
  177. private readonly internals: SpawnRunnerInternals,
  178. ) {}
  179. run(): Promise<void> {
  180. if (!this.host.connected || this.host.send === undefined) {
  181. this.finish(127)
  182. return this.completion.promise
  183. }
  184. this.host.on('message', this.onMessage)
  185. this.host.once('disconnect', this.onDisconnect)
  186. return this.completion.promise
  187. }
  188. private readonly onMessage = (value: unknown): void => {
  189. if (this.finished) return
  190. if (isWindowsTerminateRequest(value)) {
  191. this.requestTermination()
  192. return
  193. }
  194. if (this.startSeen) {
  195. void this.runnerFailure(new Error('subprocess runner received more than one Windows start request'))
  196. return
  197. }
  198. let request: WindowsStartRequest
  199. try {
  200. request = parseWindowsStartRequest(value)
  201. } catch (error) {
  202. void this.runnerFailure(error)
  203. return
  204. }
  205. this.startSeen = true
  206. void this.start(request)
  207. }
  208. private readonly onDisconnect = (): void => {
  209. if (this.finished) return
  210. this.releaseOwnedJob()
  211. this.finish(127, false)
  212. }
  213. private async start(request: WindowsStartRequest): Promise<void> {
  214. if (this.terminateRequested) {
  215. await this.publishTerminalResult({ type: 'start-cancelled' }, 0)
  216. return
  217. }
  218. await new Promise<void>((resolveImmediate) => { setImmediate(resolveImmediate) })
  219. if (this.finished) return
  220. if (this.startCancellationPending()) {
  221. await this.publishTerminalResult({ type: 'start-cancelled' }, 0)
  222. return
  223. }
  224. try {
  225. replaceEnvironment(this.host.env, request.env)
  226. this.api = this.internals.loadWin32ProcessBindings()
  227. const [command, ...args] = this.argv
  228. const spawned = this.internals.spawnCurrentTokenJobProcess(this.api, {
  229. command: command as string,
  230. args,
  231. cwd: request.cwd,
  232. })
  233. this.processHandle = spawned.process
  234. this.jobHandle = spawned.job
  235. this.committed = true
  236. this.internals.closeCurrentProcessStandardHandles(this.api)
  237. if (this.startCancellationPending()) this.terminateOwnedJob()
  238. this.pollTimer = setInterval(() => { this.poll() }, 10)
  239. this.poll()
  240. } catch (error) {
  241. if (!this.committed && error instanceof Win32Error && error.api === 'CreateProcessW') {
  242. await this.publishTerminalResult({
  243. type: 'spawn-error',
  244. error: asSpawnError(error, this.argv[0] as string, this.argv.slice(1)),
  245. }, 0)
  246. return
  247. }
  248. await this.runnerFailure(error)
  249. }
  250. }
  251. private requestTermination(): void {
  252. if (this.terminateRequested) return
  253. this.terminateRequested = true
  254. if (this.committed) {
  255. try {
  256. this.terminateOwnedJob()
  257. } catch (error) {
  258. void this.runnerFailure(error)
  259. }
  260. }
  261. }
  262. private startCancellationPending(): boolean {
  263. return this.terminateRequested
  264. }
  265. private terminateOwnedJob(): void {
  266. const job = this.jobHandle
  267. if (job === undefined) return
  268. /* v8 ignore next -- a Job handle is assigned only after the bindings are loaded;
  269. * the guard above is the only reachable empty-owner state. */
  270. if (this.api === undefined) return
  271. this.internals.terminateJob(this.api, job, 1)
  272. }
  273. private poll(): void {
  274. if (this.finished) return
  275. /* v8 ignore next -- poll is installed only after start() stores the bindings; retained as a defensive invariant guard. */
  276. if (this.api === undefined) return
  277. try {
  278. if (this.processHandle !== undefined) {
  279. const exitCode = this.internals.pollProcessExit(this.api, this.processHandle)
  280. if (exitCode !== undefined) {
  281. this.internals.closeHandleChecked(this.api, this.processHandle, 'ordinary direct process')
  282. this.processHandle = undefined
  283. void this.publishTerminalResult({ type: 'target-exit', exitCode, signal: null })
  284. }
  285. }
  286. if (this.jobHandle !== undefined && this.internals.isJobEmpty(this.api, this.jobHandle)) {
  287. this.internals.closeHandleChecked(this.api, this.jobHandle, 'ordinary process Job')
  288. this.jobHandle = undefined
  289. this.jobEmpty = true
  290. if (this.resultDelivered) this.finish(0)
  291. }
  292. } catch (error) {
  293. void this.runnerFailure(error)
  294. }
  295. }
  296. private async publishTerminalResult(result: WindowsRunnerResult, exitCode?: number): Promise<void> {
  297. /* v8 ignore next -- each state transition has a single result call site; the guard contains only re-entrant internal defects. */
  298. if (this.finished || this.resultStarted) return
  299. this.resultStarted = true
  300. try {
  301. await sendMessage(this.host, result)
  302. this.resultDelivered = true
  303. } catch {
  304. this.releaseOwnedJob()
  305. this.finish(127, false)
  306. return
  307. }
  308. if (exitCode !== undefined) {
  309. this.finish(exitCode)
  310. return
  311. }
  312. if (this.jobEmpty) this.finish(0)
  313. }
  314. private async runnerFailure(error: unknown): Promise<void> {
  315. /* v8 ignore next -- callers stop/detach on finish; this guard contains only an already-queued internal callback. */
  316. if (this.finished) return
  317. if (!this.resultStarted) {
  318. this.resultStarted = true
  319. try {
  320. await sendMessage(this.host, { type: 'runner-error', error: serializeRunnerError(error) })
  321. this.resultDelivered = true
  322. } catch {
  323. // The disconnected parent observes runner infrastructure failure.
  324. }
  325. }
  326. this.releaseOwnedJob()
  327. this.finish(127)
  328. }
  329. private releaseOwnedJob(): void {
  330. if (this.pollTimer !== undefined) clearInterval(this.pollTimer)
  331. this.pollTimer = undefined
  332. const api = this.api
  333. if (api === undefined) return
  334. if (this.jobHandle !== undefined) {
  335. try { this.internals.terminateJob(api, this.jobHandle, 1) } catch { /* Continue to kill-on-close. */ }
  336. try { this.internals.closeHandleChecked(api, this.jobHandle, 'ordinary process Job cleanup') } catch { /* Best effort after failure. */ }
  337. this.jobHandle = undefined
  338. }
  339. if (this.processHandle !== undefined) {
  340. try { this.internals.closeHandleChecked(api, this.processHandle, 'ordinary direct process cleanup') } catch { /* Best effort after failure. */ }
  341. this.processHandle = undefined
  342. }
  343. }
  344. private finish(exitCode: number, disconnect = true): void {
  345. if (this.finished) return
  346. this.finished = true
  347. if (this.pollTimer !== undefined) clearInterval(this.pollTimer)
  348. this.pollTimer = undefined
  349. this.host.off('message', this.onMessage)
  350. this.host.off('disconnect', this.onDisconnect)
  351. this.host.exitCode = exitCode
  352. if (disconnect && this.host.connected) this.host.disconnect()
  353. this.completion.resolve()
  354. }
  355. }
  356. /**
  357. * Execute the selected Linux bootstrap or Windows Job runner.
  358. * @param selection - Windows sentinel or Linux launch-request locator.
  359. * @param argv - private runner arguments beginning with the target delimiter.
  360. * @param host - process transport and lifecycle host.
  361. * @param internals - native and filesystem operations used by the runner.
  362. */
  363. export async function runSpawnRunner(
  364. selection: string,
  365. argv: readonly string[],
  366. host: RunnerHost = process,
  367. internals: SpawnRunnerInternals = defaultInternals,
  368. ): Promise<void> {
  369. Reflect.deleteProperty(host.env, SUBPROCESS_RUNNER_ENV)
  370. const targetArgv = parseRunnerTargetArgv(argv)
  371. if (selection === WINDOWS_RUNNER_SELECTION) {
  372. await new WindowsJobRunner(targetArgv, host, internals).run()
  373. return
  374. }
  375. runLinux(selection, targetArgv, host, internals)
  376. }
  377. /**
  378. * Best-effort reporting for failures before the selected runner established its owner.
  379. * @param selection - Windows sentinel, Linux launch-request locator, or no selection.
  380. * @param error - failure raised before normal runner settlement.
  381. * @param host - process transport and lifecycle host.
  382. */
  383. export async function reportSpawnRunnerFailure(
  384. selection: string | undefined,
  385. error: unknown,
  386. host: RunnerHost = process,
  387. ): Promise<void> {
  388. if (selection === WINDOWS_RUNNER_SELECTION) {
  389. try { await sendMessage(host, { type: 'runner-error', error: serializeRunnerError(error) }) } catch { /* No transport remains. */ }
  390. host.exitCode = 127
  391. if (host.connected) host.disconnect()
  392. return
  393. }
  394. if (selection !== undefined) {
  395. try {
  396. const files: LinuxLaunchFiles = linuxLaunchFilesFromLocator(selection)
  397. writeLinuxStartupError(files, { type: 'runner-error', error: serializeRunnerError(error) })
  398. } catch {
  399. // The parent will report an unconsumed request or missing runner result.
  400. }
  401. }
  402. host.exitCode = 127
  403. }