1
0

spawn-runner.ts 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363
  1. /** Native managed-range runner for ordinary local subprocesses. */
  2. import { spawn } from 'node:child_process'
  3. import {
  4. closeHandleChecked,
  5. isJobEmpty,
  6. loadWin32ProcessBindings,
  7. openNamedPipeForStdio,
  8. pollProcessExit,
  9. spawnCurrentTokenJobProcess,
  10. terminateJob,
  11. waitForProcessExit,
  12. Win32Error,
  13. } from '@deepseek-ai/dsh-win32-process'
  14. import type { ChildStdioHandles, NativePtr } from '@deepseek-ai/dsh-win32-process'
  15. import {
  16. appendRunnerEvent,
  17. consumeRunnerRequest,
  18. serializeSpawnError,
  19. } from './runner-protocol.ts'
  20. import type { RunnerRequest, SerializedSpawnError } from './runner-protocol.ts'
  21. type RunnerArgs =
  22. | { mode: 'probe-node' }
  23. | { mode: 'probe-win32' }
  24. | { mode: 'node'; requestPath: string; eventsPath: string }
  25. | {
  26. mode: 'win32'
  27. requestPath: string
  28. eventsPath: string
  29. stdinPipe?: string
  30. stdoutPipe?: string
  31. stderrPipe?: string
  32. }
  33. type RunnerHost = Pick<
  34. NodeJS.Process,
  35. 'env' | 'exitCode' | 'connected' | 'cwd' | 'chdir' | 'on' | 'off' | 'disconnect'
  36. >
  37. interface RunnerInternals {
  38. spawn: typeof spawn
  39. loadWin32ProcessBindings: typeof loadWin32ProcessBindings
  40. openNamedPipeForStdio: typeof openNamedPipeForStdio
  41. spawnCurrentTokenJobProcess: typeof spawnCurrentTokenJobProcess
  42. pollProcessExit: typeof pollProcessExit
  43. isJobEmpty: typeof isJobEmpty
  44. terminateJob: typeof terminateJob
  45. waitForProcessExit: typeof waitForProcessExit
  46. closeHandleChecked: typeof closeHandleChecked
  47. }
  48. const defaultRunnerInternals: RunnerInternals = {
  49. spawn,
  50. loadWin32ProcessBindings,
  51. openNamedPipeForStdio,
  52. spawnCurrentTokenJobProcess,
  53. pollProcessExit,
  54. isJobEmpty,
  55. terminateJob,
  56. waitForProcessExit,
  57. closeHandleChecked,
  58. }
  59. function parseArgs(argv: string[]): RunnerArgs {
  60. let mode: string | undefined
  61. let requestPath: string | undefined
  62. let eventsPath: string | undefined
  63. let stdinPipe: string | undefined
  64. let stdoutPipe: string | undefined
  65. let stderrPipe: string | undefined
  66. for (let index = 0; index < argv.length; index += 2) {
  67. const key = argv[index]
  68. const value = argv[index + 1]
  69. if (value === undefined) throw new Error(`subprocess runner missing value after ${String(key)}`)
  70. if (key === '--mode') mode = value
  71. else if (key === '--request') requestPath = value
  72. else if (key === '--events') eventsPath = value
  73. else if (key === '--stdin-pipe') stdinPipe = value
  74. else if (key === '--stdout-pipe') stdoutPipe = value
  75. else if (key === '--stderr-pipe') stderrPipe = value
  76. else throw new Error(`subprocess runner unknown argument: ${String(key)}`)
  77. }
  78. if (mode === 'probe-node' || mode === 'probe-win32') return { mode }
  79. if (mode !== 'node' && mode !== 'win32') throw new Error(`subprocess runner unknown mode: ${String(mode)}`)
  80. if (requestPath === undefined || eventsPath === undefined) throw new Error('subprocess runner requires request and event paths')
  81. if (mode === 'node') return { mode, requestPath, eventsPath }
  82. return {
  83. mode,
  84. requestPath,
  85. eventsPath,
  86. ...stdinPipe === undefined ? {} : { stdinPipe },
  87. ...stdoutPipe === undefined ? {} : { stdoutPipe },
  88. ...stderrPipe === undefined ? {} : { stderrPipe },
  89. }
  90. }
  91. function win32SpawnError(error: unknown, request: RunnerRequest): SerializedSpawnError {
  92. const serialized = serializeSpawnError(error)
  93. const code = error instanceof Win32Error
  94. ? error.win32Code === 2 || error.win32Code === 3 || error.win32Code === 267
  95. ? 'ENOENT'
  96. : error.win32Code === 5
  97. ? 'EACCES'
  98. : error.win32Code === 193
  99. ? 'EFTYPE'
  100. : 'UNKNOWN'
  101. : serialized.code
  102. if (code === undefined) return serialized
  103. const program = request.argv[0] as string
  104. return {
  105. ...serialized,
  106. message: `spawn ${program} ${code}: ${serialized.message}`,
  107. code,
  108. syscall: `spawn ${program}`,
  109. path: program,
  110. spawnargs: request.argv.slice(1),
  111. }
  112. }
  113. async function runNode(
  114. request: RunnerRequest,
  115. eventsPath: string,
  116. host: RunnerHost,
  117. internals: RunnerInternals,
  118. ): Promise<void> {
  119. const ignoreScopeSignal = (): void => { /* The target receives the scope signal; the runner reports its outcome. */ }
  120. for (const signal of ['SIGTERM', 'SIGINT', 'SIGHUP'] as const) {
  121. host.on(signal, ignoreScopeSignal)
  122. }
  123. const [program, ...args] = request.argv
  124. const child = internals.spawn(program as string, args, {
  125. cwd: request.cwd,
  126. env: request.env,
  127. stdio: 'inherit',
  128. })
  129. await new Promise<void>((resolve) => {
  130. let started = false
  131. let failed = false
  132. let settled = false
  133. const finish = (): void => {
  134. if (settled) return
  135. settled = true
  136. for (const signal of ['SIGTERM', 'SIGINT', 'SIGHUP'] as const) host.off(signal, ignoreScopeSignal)
  137. resolve()
  138. }
  139. child.once('spawn', () => {
  140. started = true
  141. appendRunnerEvent(eventsPath, { type: 'started', pid: child.pid as number })
  142. })
  143. child.once('error', (error) => {
  144. failed = true
  145. if (!started) appendRunnerEvent(eventsPath, { type: 'spawn-error', error: serializeSpawnError(error) })
  146. else appendRunnerEvent(eventsPath, { type: 'runner-error', error: serializeSpawnError(error) })
  147. host.exitCode = 127
  148. finish()
  149. })
  150. child.once('exit', (exitCode, signal) => {
  151. if (!failed) {
  152. appendRunnerEvent(eventsPath, { type: 'exit', exitCode, signal })
  153. host.exitCode = exitCode ?? 1
  154. }
  155. finish()
  156. })
  157. })
  158. }
  159. function replaceEnvironment(target: NodeJS.ProcessEnv, env: Record<string, string>): void {
  160. for (const key of Object.keys(target)) Reflect.deleteProperty(target, key)
  161. Object.assign(target, env)
  162. }
  163. function closeStdioHandles(
  164. api: ReturnType<typeof loadWin32ProcessBindings>,
  165. handles: Array<{ handle: NativePtr; label: string }>,
  166. reportFailure: boolean,
  167. internals: RunnerInternals,
  168. ): void {
  169. let failure: Error | undefined
  170. for (const owned of handles.splice(0)) {
  171. try {
  172. internals.closeHandleChecked(api, owned.handle, owned.label)
  173. } catch (error) {
  174. handles.push(owned)
  175. failure ??= error instanceof Error ? error : new Error(serializeSpawnError(error).message)
  176. }
  177. }
  178. if (reportFailure && failure !== undefined) throw failure
  179. }
  180. async function runWin32(
  181. request: RunnerRequest,
  182. eventsPath: string,
  183. pipes: Pick<Extract<RunnerArgs, { mode: 'win32' }>, 'stdinPipe' | 'stdoutPipe' | 'stderrPipe'>,
  184. host: RunnerHost,
  185. internals: RunnerInternals,
  186. ): Promise<void> {
  187. replaceEnvironment(host.env, request.env)
  188. const api = internals.loadWin32ProcessBindings()
  189. let processHandle: NativePtr | undefined
  190. let jobHandle: NativePtr | undefined
  191. const openedStdio: Array<{ handle: NativePtr; label: string }> = []
  192. try {
  193. const stdio: ChildStdioHandles = {}
  194. for (const [key, path, access] of [
  195. ['stdin', pipes.stdinPipe, 'read'],
  196. ['stdout', pipes.stdoutPipe, 'write'],
  197. ['stderr', pipes.stderrPipe, 'write'],
  198. ] as const) {
  199. if (path === undefined) continue
  200. const handle = internals.openNamedPipeForStdio(api, path, access)
  201. stdio[key] = handle
  202. openedStdio.push({ handle, label: `ordinary target ${key} pipe` })
  203. }
  204. // Match Node's cwd-relative executable lookup and spawn-error attribution.
  205. const runnerCwd = host.cwd()
  206. host.chdir(request.cwd)
  207. try {
  208. const [command, ...args] = request.argv
  209. const spawned = internals.spawnCurrentTokenJobProcess(
  210. api,
  211. { command: command as string, args, cwd: host.cwd() },
  212. stdio,
  213. )
  214. processHandle = spawned.process
  215. jobHandle = spawned.job
  216. appendRunnerEvent(eventsPath, { type: 'started', pid: spawned.pid })
  217. } finally {
  218. host.chdir(runnerCwd)
  219. }
  220. closeStdioHandles(api, openedStdio, true, internals)
  221. await new Promise<void>((resolve, reject) => {
  222. let settled = false
  223. let terminationRequested = false
  224. const settle = (error?: unknown): void => {
  225. if (settled) return
  226. settled = true
  227. clearInterval(timer)
  228. host.off('message', onMessage)
  229. host.off('disconnect', onDisconnect)
  230. if (error === undefined) resolve()
  231. else reject(error instanceof Error ? error : new Error(serializeSpawnError(error).message))
  232. }
  233. const terminate = (): void => {
  234. if (terminationRequested || jobHandle === undefined) return
  235. terminationRequested = true
  236. try {
  237. internals.terminateJob(api, jobHandle, 1)
  238. } catch (error) {
  239. settle(error)
  240. }
  241. }
  242. const onMessage = (message: unknown): void => {
  243. if (message !== null && typeof message === 'object' && (message as { type?: unknown }).type === 'terminate') {
  244. terminate()
  245. }
  246. }
  247. const onDisconnect = (): void => { terminate() }
  248. host.on('message', onMessage)
  249. host.on('disconnect', onDisconnect)
  250. const timer = setInterval(() => {
  251. try {
  252. if (processHandle !== undefined) {
  253. const exitCode = internals.pollProcessExit(api, processHandle)
  254. if (exitCode !== undefined) {
  255. appendRunnerEvent(eventsPath, { type: 'exit', exitCode, signal: null })
  256. internals.closeHandleChecked(api, processHandle, 'ordinary direct process')
  257. processHandle = undefined
  258. }
  259. }
  260. if (processHandle === undefined && jobHandle !== undefined && internals.isJobEmpty(api, jobHandle)) {
  261. internals.closeHandleChecked(api, jobHandle, 'ordinary process Job')
  262. jobHandle = undefined
  263. settle()
  264. }
  265. } catch (error) {
  266. settle(error)
  267. }
  268. }, 10)
  269. })
  270. } catch (error) {
  271. const targetSpawnFailed = (error instanceof Win32Error && error.api === 'CreateProcessW')
  272. || (processHandle === undefined
  273. && error instanceof Error
  274. && (error as NodeJS.ErrnoException).syscall === 'chdir')
  275. appendRunnerEvent(eventsPath, {
  276. type: targetSpawnFailed ? 'spawn-error' : 'runner-error',
  277. error: targetSpawnFailed ? win32SpawnError(error, request) : serializeSpawnError(error),
  278. })
  279. if (!targetSpawnFailed) host.exitCode = 127
  280. } finally {
  281. closeStdioHandles(api, openedStdio, false, internals)
  282. if (processHandle !== undefined) {
  283. try { internals.closeHandleChecked(api, processHandle, 'ordinary direct process cleanup') } catch { /* best effort after reported failure */ }
  284. }
  285. if (jobHandle !== undefined) {
  286. try { internals.closeHandleChecked(api, jobHandle, 'ordinary process Job cleanup') } catch { /* best effort after reported failure */ }
  287. }
  288. }
  289. }
  290. function probeWin32Job(host: RunnerHost, internals: RunnerInternals): void {
  291. const command = host.env.ComSpec ?? host.env.COMSPEC
  292. if (command === undefined) throw new Error('subprocess runner cannot probe a Windows Job without ComSpec')
  293. const api = internals.loadWin32ProcessBindings()
  294. const spawned = internals.spawnCurrentTokenJobProcess(api, {
  295. command,
  296. args: ['/d', '/s', '/c', 'exit 0'],
  297. cwd: host.cwd(),
  298. })
  299. try {
  300. const exitCode = internals.waitForProcessExit(api, spawned.process)
  301. if (exitCode !== 0) throw new Error(`subprocess Windows Job probe exited with code ${String(exitCode)}`)
  302. } finally {
  303. internals.closeHandleChecked(api, spawned.job, 'subprocess Windows Job probe')
  304. }
  305. }
  306. /**
  307. * Execute one parsed private-runner request.
  308. * @param argv - runner arguments after the executable and entry path.
  309. * @param host - process operations; tests provide an isolated host facade.
  310. * @param internals - platform operations; tests replace native Win32 calls.
  311. * @returns after the requested probe or target lifecycle completes.
  312. */
  313. export async function runSpawnRunner(
  314. argv: string[],
  315. host: RunnerHost = process,
  316. internals: RunnerInternals = defaultRunnerInternals,
  317. ): Promise<void> {
  318. const args = parseArgs(argv)
  319. if (args.mode === 'probe-node') return
  320. if (args.mode === 'probe-win32') {
  321. probeWin32Job(host, internals)
  322. return
  323. }
  324. const request = consumeRunnerRequest(args.requestPath)
  325. if (args.mode === 'node') await runNode(request, args.eventsPath, host, internals)
  326. else {
  327. try {
  328. await runWin32(request, args.eventsPath, args, host, internals)
  329. } finally {
  330. if (host.connected) host.disconnect()
  331. }
  332. }
  333. }
  334. /**
  335. * Publish an infrastructure failure when runner arguments still identify an event file.
  336. * @param argv - original runner arguments.
  337. * @param error - uncaught runner failure.
  338. */
  339. export function reportSpawnRunnerFailure(argv: string[], error: unknown): void {
  340. try {
  341. const args = parseArgs(argv)
  342. if (args.mode !== 'probe-node' && args.mode !== 'probe-win32') {
  343. appendRunnerEvent(args.eventsPath, { type: 'runner-error', error: serializeSpawnError(error) })
  344. }
  345. } catch {
  346. // No trustworthy transport remains; the parent reports the missing result.
  347. }
  348. }