spawn.ts 26 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687
  1. /**
  2. * Process plumbing for the local subprocess service: ordinary process launch
  3. * with per-stream stdio dispositions, tail-keep collection with spill
  4. * files, provider-owned range signalling, and common termination scheduling.
  5. * POSIX owners stage TERM before KILL; Windows owners terminate immediately.
  6. * This layer reacts to an abort signal; callers own deadlines, teardown
  7. * ladders, and cause classification.
  8. * @module dsh-subprocess-local/spawn
  9. */
  10. import { type ChildProcess, type SpawnOptions, spawn, spawnSync } from 'node:child_process'
  11. import type { Readable } from 'node:stream'
  12. import { randomBytes } from 'node:crypto'
  13. import { closeSync, mkdtempSync, openSync, rmdirSync, unlinkSync, writeSync } from 'node:fs'
  14. import { tmpdir } from 'node:os'
  15. import { join } from 'node:path'
  16. import { setTimeout as sleepMs } from 'node:timers/promises'
  17. import { scrubbedParentEnv } from '@deepseek-ai/dsh-subprocess'
  18. import { MAX_TIMER_DELAY_MS } from '@deepseek-ai/dsh-timeout'
  19. import type {
  20. CollectedOutput,
  21. SubprocessCollect,
  22. SubprocessHandle,
  23. SubprocessOutcome,
  24. SubprocessOutputMode,
  25. SubprocessSpawnSpec,
  26. } from '@deepseek-ai/dsh-subprocess'
  27. import type { BoundProcessOwner, ManagedProcessLaunch } from './managed-owner.ts'
  28. import { waitWithAbort } from './managed-owner.ts'
  29. import { controlEnvironment, controlPipe } from './control-spawn.ts'
  30. import { SUBPROCESS_CONTROL_FD } from '@deepseek-ai/dsh-subprocess/control'
  31. import { linuxProcessGroupHasLiveMembers } from './process-inspector.ts'
  32. type SpawnProcess = (
  33. program: string,
  34. args: readonly string[],
  35. options: SpawnOptions,
  36. ) => ChildProcess
  37. /**
  38. * Build a child environment: explicit caller entries override the scrubbed
  39. * parent base using the target platform's environment-key semantics. A string
  40. * deliberately restores or overrides an entry; an explicit `undefined`
  41. * tombstone removes an ordinary ambient entry.
  42. * @param extra - explicit caller entries and tombstones, merged after the scrub.
  43. * @returns the environment to hand to `spawn` for the child process.
  44. */
  45. export function childEnv(extra?: Readonly<NodeJS.ProcessEnv>): NodeJS.ProcessEnv {
  46. const env = scrubbedParentEnv()
  47. if (process.platform !== 'win32') return { ...env, ...extra }
  48. let entries: [string, string | undefined][] = Object.entries(env)
  49. for (const [key, value] of Object.entries(extra ?? {})) {
  50. const normalized = key.toUpperCase()
  51. entries = entries.filter(([inherited]) => inherited.toUpperCase() !== normalized)
  52. entries.push([key, value])
  53. }
  54. return Object.fromEntries(entries)
  55. }
  56. /** Injectable process, spill, and platform operations. */
  57. export interface SpawnInternals {
  58. /** Process spawner (defaults to `node:child_process` `spawn`). */
  59. spawn?: SpawnProcess
  60. /** Directory for spill files (defaults to the OS temp dir). */
  61. spillDir?: string
  62. /** Windows tree-termination runner (defaults to `taskkill /PID <pid> /T /F`). */
  63. taskkill?: (pid: number) => void
  64. /** Host platform override for signalling decisions. */
  65. platform?: NodeJS.Platform
  66. /** Linux process-group member probe (defaults to `/proc` inspection). */
  67. linuxProcessGroupHasLiveMembers?: (processGroupId: number) => boolean | undefined
  68. }
  69. /**
  70. * Local-only synchronous final termination used by the owning service during
  71. * host exit and as the last fallback after failed normal disposal. It is
  72. * intentionally absent from the public subprocess seam.
  73. */
  74. export interface LocalSubprocessHandle extends SubprocessHandle {
  75. /** Force-terminate the current tree synchronously without starting timers or waits. */
  76. terminateForHostExit(): void
  77. }
  78. /**
  79. * Liveness-poll cadence for tree-exit waits. The timer stays ref'd: an
  80. * awaited teardown must keep the event loop alive until the tree really
  81. * exits, or the parent can exit while claiming quiescence and orphan the
  82. * survivors it promised to reap.
  83. */
  84. function sleepTick(): Promise<void> {
  85. return sleepMs(15)
  86. }
  87. let spillCounter = 0
  88. let defaultSpillDir: string | undefined
  89. /**
  90. * The default spill location: a private (0700) per-process directory under
  91. * the OS tmpdir, created lazily. Predictable world-readable paths would let
  92. * other local users read command output or pre-create symlinks. At a
  93. * JavaScript-observable process exit the directory is removed only when it
  94. * holds no completed spill file (spill files are retained as full-output
  95. * recovery artifacts until an external cleanup).
  96. */
  97. function privateSpillDir(): string {
  98. defaultSpillDir ??= mkdtempSync(join(tmpdir(), 'dsh-subprocess-'))
  99. return defaultSpillDir
  100. }
  101. // The per-process spill directory is removed at process exit when it holds no
  102. // completed spill file: a directory that never spilled is empty and is safe to
  103. // remove, while a directory holding completed spill files keeps them (their
  104. // content is retained until an external cleanup). A SIGKILLed process cannot
  105. // run this at all; its residue is left to OS temp hygiene.
  106. /* v8 ignore next 4 -- exit listeners run after the coverage dump; removal is verified by the CI /tmp residue measurement. */
  107. process.once('exit', () => {
  108. if (defaultSpillDir === undefined) return
  109. try { rmdirSync(defaultSpillDir) } catch { /* best-effort: ENOENT/ENOTEMPTY/EBUSY/EPERM must not change the exit code. */ }
  110. })
  111. /**
  112. * Prepare fallible output storage before starting a managed native process.
  113. * @param internals - optional caller-owned spill directory.
  114. * @returns binding inputs whose spill directory is ready for use.
  115. */
  116. export function prepareManagedProcessBinding(
  117. internals: Pick<SpawnInternals, 'spillDir'> = {},
  118. ): { spillDir: string } {
  119. return { spillDir: internals.spillDir ?? privateSpillDir() }
  120. }
  121. /**
  122. * Collects one stream with a bounded in-memory tail. With a spill cap, on
  123. * first overflow a spill file is created and every chunk (including those
  124. * already collected) is appended there while the full stream remains within
  125. * the cap; without one, only the in-memory tail is ever retained (the
  126. * diagnostic-tail shape — a language server's stderr).
  127. *
  128. * Tail-keep rationale (pi/OpenCode): errors and final results cluster at the
  129. * end of command output; the spill file covers the head.
  130. */
  131. export class OutputCollector {
  132. private chunks: Buffer[] = []
  133. private bytes = 0
  134. private dropped = false
  135. private spillFd: number | undefined
  136. private spillFile: string | undefined
  137. private spillDisabled: boolean
  138. /** Total bytes ever pushed (not just retained). */
  139. private total = 0
  140. constructor(
  141. private readonly maxBytes: number,
  142. private readonly maxSpillBytes: number | undefined,
  143. private readonly label: string,
  144. private readonly spillDir: string,
  145. ) {
  146. this.spillDisabled = maxSpillBytes === undefined
  147. }
  148. /**
  149. * Ingest one stream chunk, counting it toward the whole-stream total. On
  150. * first overflow of the in-memory cap a spill file is opened (when spilling
  151. * is enabled) and every chunk (already-collected ones included) is appended
  152. * there from then on; the in-memory tail then drops whole chunks from its
  153. * head (or the head of a single over-cap chunk) until it fits the cap again.
  154. * @param chunk - the raw bytes from one stream 'data' event.
  155. */
  156. push(chunk: Buffer): void {
  157. this.total += chunk.length
  158. const overflows = this.bytes + chunk.length > this.maxBytes
  159. if (!this.spillDisabled && (overflows || this.spillFd !== undefined)) this.spillAll(chunk)
  160. this.chunks.push(chunk)
  161. this.bytes += chunk.length
  162. while (this.bytes > this.maxBytes) {
  163. const head = this.chunks[0] as Buffer
  164. const excess = this.bytes - this.maxBytes
  165. if (head.length <= excess) {
  166. // Drop the whole head chunk (length ≥ 1 is guaranteed while over cap).
  167. this.chunks.shift()
  168. this.bytes -= head.length
  169. } else {
  170. // Trim the head so the retained window is byte-exact at the cap — a
  171. // diagnostic tail (an LSP server's stderr) must hold the LAST
  172. // maxBytes regardless of how the stream was chunked.
  173. this.chunks[0] = head.subarray(excess)
  174. this.bytes -= excess
  175. }
  176. this.dropped = true
  177. }
  178. }
  179. /** Open the spill file lazily and append `chunk` (and any prior chunks once). */
  180. private spillAll(chunk: Buffer): void {
  181. if (this.maxSpillBytes !== undefined && this.total > this.maxSpillBytes) {
  182. this.discardSpill()
  183. return
  184. }
  185. if (this.spillFd === undefined) {
  186. // Random suffix + O_EXCL + no-follow-equivalent ('wx' fails on any
  187. // existing path, symlink or not) + owner-only mode: defeats spill-path
  188. // prediction and symlink planting in shared tmp dirs.
  189. this.spillFile = join(
  190. this.spillDir,
  191. `dsh-subprocess-${process.pid}-${++spillCounter}-${randomBytes(6).toString('hex')}-${this.label}.log`,
  192. )
  193. this.spillFd = openSync(this.spillFile, 'wx', 0o600)
  194. for (const prior of this.chunks) writeSync(this.spillFd, prior)
  195. }
  196. writeSync(this.spillFd, chunk)
  197. }
  198. /** Stop spilling and remove the file once it can no longer hold the complete stream. */
  199. private discardSpill(): void {
  200. const fd = this.spillFd
  201. const file = this.spillFile
  202. this.spillFd = undefined
  203. this.spillFile = undefined
  204. this.spillDisabled = true
  205. if (fd !== undefined) {
  206. try {
  207. closeSync(fd)
  208. } catch {
  209. // Retain the descriptor so finalize can retry the failed close.
  210. this.spillFd = fd
  211. }
  212. }
  213. if (file !== undefined) {
  214. try {
  215. unlinkSync(file)
  216. } catch {
  217. // A failed unlink leaves at most maxSpillBytes behind, never an unbounded file.
  218. }
  219. }
  220. }
  221. /**
  222. * Incremental read in whole-stream byte coordinates: returns everything
  223. * pushed since `fromByte`. When `fromByte` has already slid out of the
  224. * in-memory tail window, the read is `lossy` — it returns the whole
  225. * retained tail and the gap is only recoverable from the spill file.
  226. * @param fromByte - whole-stream offset to resume from (a prior read's `nextOffset`; 0 for the first read).
  227. * @returns the delta text, the offset for the next read, the `lossy` flag, and the spill path when one was created.
  228. */
  229. readFrom(fromByte: number): { text: string; nextOffset: number; lossy: boolean; spillPath?: string } {
  230. const windowStart = this.total - this.bytes
  231. const buffer = Buffer.concat(this.chunks)
  232. const lossy = fromByte < windowStart
  233. const slice = lossy ? buffer : buffer.subarray(fromByte - windowStart)
  234. return {
  235. text: slice.toString('utf8'),
  236. nextOffset: this.total,
  237. lossy,
  238. ...this.spillFile !== undefined ? { spillPath: this.spillFile } : {},
  239. }
  240. }
  241. /**
  242. * Close the spill file once the stream has ended. A failed close (delayed
  243. * writeback fault) stops advertising the spill path — the file may be
  244. * missing its tail — while every in-memory read keeps working. Idempotent;
  245. * the spawn path seals both collectors at settlement so reads after exit
  246. * never point at a still-open file.
  247. */
  248. seal(): void {
  249. if (this.spillFd === undefined) return
  250. try {
  251. closeSync(this.spillFd)
  252. } catch {
  253. // A delayed writeback failure makes the spill unreliable; keep the
  254. // in-memory result but stop advertising that file.
  255. this.spillFile = undefined
  256. }
  257. this.spillFd = undefined
  258. }
  259. /**
  260. * Seal the spill file and return the final output.
  261. * @returns the final collected output: tail text, truncation flag, and the spill path when intact.
  262. */
  263. finalize(): CollectedOutput {
  264. this.seal()
  265. return {
  266. text: Buffer.concat(this.chunks).toString('utf8'),
  267. truncated: this.dropped,
  268. ...this.spillFile !== undefined ? { spillPath: this.spillFile } : {},
  269. }
  270. }
  271. }
  272. /**
  273. * Send `sig` to a detached POSIX process group. Never throws: delivery races
  274. * process exit and may run in a timer callback, so failures are contained and
  275. * a missing pid is a no-op.
  276. * @param pid - the group leader's pid, when the spawn published one.
  277. * @param sig - the signal to deliver to the whole group.
  278. */
  279. export function killGroup(pid: number | undefined, sig: NodeJS.Signals): void {
  280. if (pid === undefined) return
  281. try {
  282. process.kill(-pid, sig)
  283. } catch {
  284. // Swallow: see contract above.
  285. }
  286. }
  287. /**
  288. * Terminate one Windows process tree with `taskkill /T /F`. Contained like
  289. * POSIX group signalling — delivery races tree exit, so an absent tree, a
  290. * nonzero status, or a missing taskkill binary must not break idempotent
  291. * teardown.
  292. * @param pid - root process id, when the spawn published one.
  293. */
  294. export function taskkillProcessTree(pid: number | undefined): void {
  295. if (pid === undefined || pid <= 0) return
  296. // Outcome deliberately unchecked: an already-absent tree (status 128), exit
  297. // races, and a missing taskkill binary (spawnSync reports, never throws) are
  298. // as tolerable here as ESRCH is for a POSIX group signal.
  299. spawnSync('taskkill', ['/PID', String(pid), '/T', '/F'], {
  300. stdio: 'ignore',
  301. windowsHide: true,
  302. })
  303. }
  304. /**
  305. * Signal a detached process tree with platform-correct semantics: POSIX
  306. * signals the negative process-group id and falls back to the direct child
  307. * when the group is gone; Windows terminates the tree via taskkill (any
  308. * signal value force-terminates — Node maps signals to TerminateProcess).
  309. */
  310. function signalTree(
  311. platform: NodeJS.Platform,
  312. pid: number | undefined,
  313. sig: NodeJS.Signals,
  314. child: ChildProcess,
  315. taskkill: (pid: number) => void,
  316. ): void {
  317. /* v8 ignore next -- kill/terminate gate on treeAlive(), which is false without a pid; this guard protects direct callers only. */
  318. if (pid === undefined) return
  319. if (platform === 'win32') {
  320. taskkill(pid)
  321. return
  322. }
  323. try {
  324. process.kill(-pid, sig)
  325. } catch {
  326. /* v8 ignore start -- the fallback needs a live child whose group signal fails
  327. (EPERM-style), which POSIX CI cannot stage; the swallow keeps teardown idempotent. */
  328. try {
  329. child.kill(sig)
  330. } catch {
  331. // The direct child already exited; teardown remains idempotent.
  332. }
  333. /* v8 ignore stop */
  334. }
  335. }
  336. /**
  337. * Validate the synchronous portion of one ordinary spawn request.
  338. * @param spec - exact target request.
  339. * @throws when grace, cancellation, or argv is invalid before launch.
  340. */
  341. export function validateSubprocessSpec(spec: SubprocessSpawnSpec): void {
  342. if (!Number.isFinite(spec.graceMs) || spec.graceMs <= 0 || spec.graceMs > MAX_TIMER_DELAY_MS) {
  343. throw new Error(`subprocess graceMs must be a positive finite number no greater than ${MAX_TIMER_DELAY_MS}`)
  344. }
  345. if (spec.signal?.aborted) {
  346. let reason = 'aborted'
  347. try {
  348. reason = String(spec.signal.reason ?? reason)
  349. } catch {
  350. // Arbitrary caller-owned reasons cannot escape the stable Error boundary.
  351. }
  352. throw new Error(`aborted before spawn: ${reason}`)
  353. }
  354. const [program] = spec.argv
  355. if (program === undefined || program.length === 0) {
  356. throw new Error('invalid argv: expected a non-empty program name at argv[0]')
  357. }
  358. }
  359. function directChildResult(child: ChildProcess): Promise<SubprocessOutcome> {
  360. return new Promise((resolve, reject) => {
  361. let completed = false
  362. child.once('error', (error) => {
  363. /* v8 ignore next -- ChildProcess may report a later operational error after its
  364. terminal exit event; the first terminal event owns the result. */
  365. if (completed) return
  366. completed = true
  367. reject(error)
  368. })
  369. child.once('exit', (exitCode, signal) => {
  370. /* v8 ignore next -- a spawn/kill error may be followed by exit; a Promise can publish only the first terminal event. */
  371. if (completed) return
  372. completed = true
  373. resolve({ exitCode, signal })
  374. })
  375. })
  376. }
  377. function fallbackOwner(
  378. platform: NodeJS.Platform,
  379. pid: number | undefined,
  380. child: ChildProcess,
  381. taskkill: (pid: number) => void,
  382. linuxGroupHasLiveMembers: (processGroupId: number) => boolean | undefined,
  383. direct: Promise<SubprocessOutcome>,
  384. ): BoundProcessOwner {
  385. let stopped = false
  386. let directSettled = false
  387. let observation: Promise<void> | undefined
  388. void direct.then(
  389. () => { directSettled = true },
  390. () => { directSettled = true },
  391. )
  392. const alive = (): boolean => {
  393. if (stopped || pid === undefined) return false
  394. if (platform === 'win32') return child.exitCode === null && child.signalCode === null
  395. try {
  396. process.kill(-pid, 0)
  397. if (directSettled && platform === 'linux' && linuxGroupHasLiveMembers(pid) === false) return false
  398. return true
  399. } catch (error) {
  400. const code = (error as NodeJS.ErrnoException).code
  401. if (code === 'ESRCH') return false
  402. /* v8 ignore start -- EPERM and non-POSIX negative-pid failures are platform defenses. */
  403. if (code === 'EPERM') return true
  404. return child.exitCode === null && child.signalCode === null
  405. /* v8 ignore stop */
  406. }
  407. }
  408. return {
  409. signal: (signal) => {
  410. if (!alive()) {
  411. stopped = true
  412. return
  413. }
  414. signalTree(platform, pid, signal, child, taskkill)
  415. },
  416. waitForExit: async () => {
  417. /* v8 ignore next -- bindManagedProcess memoizes this owner wait; the guard only
  418. protects direct internal re-entry after signal() observed absence. */
  419. if (stopped) return
  420. observation ??= (async () => {
  421. while (alive()) await sleepTick()
  422. stopped = true
  423. })()
  424. await observation
  425. },
  426. terminateForHostExit: () => {
  427. if (stopped) return
  428. signalTree(platform, pid, 'SIGKILL', child, taskkill)
  429. },
  430. }
  431. }
  432. /**
  433. * Bind platform launch facts to the existing stdio, outcome, abort, and termination lifecycle.
  434. * @param spec - fully resolved argv, cwd, stdio, grace, cancellation, environment.
  435. * @param launch - platform streams, direct outcome, and managed-range owner.
  436. * @param internals - test-only spill-directory override.
  437. * @returns live subprocess handle.
  438. */
  439. export function bindManagedProcess(
  440. spec: SubprocessSpawnSpec,
  441. launch: ManagedProcessLaunch,
  442. internals: Pick<SpawnInternals, 'spillDir'> = {},
  443. ): LocalSubprocessHandle {
  444. const { spillDir } = prepareManagedProcessBinding(internals)
  445. const { stdin, stdout, stderr } = launch
  446. const isCollect = (mode: SubprocessOutputMode): mode is SubprocessCollect =>
  447. mode !== 'pipe' && mode !== 'inherit'
  448. const outMode = spec.stdio.stdout
  449. const errMode = spec.stdio.stderr
  450. const stdinMode = spec.stdio.stdin
  451. const collectStream = (mode: SubprocessOutputMode, stream: Readable | null, label: string): OutputCollector | undefined => {
  452. if (!isCollect(mode) || stream === null) return undefined
  453. const collector = new OutputCollector(mode.maxBytes, mode.spill?.maxBytes, label, spillDir)
  454. stream.on('data', (chunk: Buffer) => { collector.push(chunk) })
  455. return collector
  456. }
  457. const stdoutCollector = collectStream(outMode, stdout, 'stdout')
  458. const stderrCollector = collectStream(errMode, stderr, 'stderr')
  459. const observeOutputStream = (mode: SubprocessOutputMode, stream: Readable | null): Promise<void> | undefined => {
  460. if (mode === 'inherit' || stream === null || stream.readableEnded || stream.destroyed) return undefined
  461. return new Promise((resolve) => {
  462. const settle = (): void => {
  463. stream.off('end', settle)
  464. stream.off('close', settle)
  465. stream.off('error', settle)
  466. resolve()
  467. }
  468. stream.once('end', settle)
  469. stream.once('close', settle)
  470. stream.once('error', settle)
  471. })
  472. }
  473. const stdoutClosed = observeOutputStream(outMode, stdout)
  474. const stderrClosed = observeOutputStream(errMode, stderr)
  475. const outputStreamsClosed = Promise.all([stdoutClosed, stderrClosed])
  476. const stopCollectors = (): void => {
  477. if (stdoutCollector !== undefined) stdout?.destroy()
  478. if (stderrCollector !== undefined) stderr?.destroy()
  479. stdoutCollector?.seal()
  480. stderrCollector?.seal()
  481. }
  482. let graceTimer: ReturnType<typeof setTimeout> | undefined
  483. let terminationStarted = false
  484. let rangeExitObserved = false
  485. let rangeExitObservation: Promise<void> | undefined
  486. let settled = false
  487. const scheduleOwnerCleanup = (): boolean => {
  488. if (launch.owner.cleanup === undefined) return false
  489. queueMicrotask(() => { void done.finally(() => { launch.owner.cleanup?.() }).catch(() => {}) })
  490. return true
  491. }
  492. /**
  493. * Start or reuse the handle's managed-range exit observer. A failed read
  494. * before direct settlement can be retried. Once direct settlement permits
  495. * cleanup, retain a failed observation because removing its private evidence
  496. * must not turn a later wait into a false success. The first confirmed
  497. * absence is the permanent no-more-signals boundary and cancels pending
  498. * escalation before stale identity can be used.
  499. */
  500. const observeRangeExit = (): Promise<void> => {
  501. rangeExitObservation ??= (async () => {
  502. await launch.owner.waitForExit()
  503. rangeExitObserved = true
  504. if (graceTimer !== undefined) clearTimeout(graceTimer)
  505. graceTimer = undefined
  506. spec.signal?.removeEventListener('abort', onAbort)
  507. scheduleOwnerCleanup()
  508. })().catch((error: unknown) => {
  509. if (!settled || !scheduleOwnerCleanup()) rangeExitObservation = undefined
  510. throw error
  511. })
  512. return rangeExitObservation
  513. }
  514. const kill = (sig: 'SIGTERM' | 'SIGKILL', cancellationReason?: unknown): void => {
  515. if (rangeExitObserved) return
  516. launch.owner.signal(sig, cancellationReason)
  517. }
  518. const terminateWithReason = (cancellationReason: unknown): void => {
  519. if (rangeExitObserved || terminationStarted) return
  520. terminationStarted = true
  521. // Keep the shared observation rejection available to waitForExit() without
  522. // leaking an unhandled rejection when a caller only invokes terminate().
  523. void observeRangeExit().catch(() => {})
  524. kill('SIGTERM', cancellationReason)
  525. graceTimer = setTimeout(() => {
  526. graceTimer = undefined
  527. kill('SIGKILL')
  528. }, spec.graceMs)
  529. }
  530. const terminate = (): void => {
  531. terminateWithReason(new Error('subprocess terminated before target start'))
  532. }
  533. const terminateForHostExit = (): void => {
  534. launch.owner.terminateForHostExit()
  535. }
  536. // The caller owns timeout classification; this layer only reacts to abort.
  537. const onAbort = (): void => { terminateWithReason(spec.signal?.reason) }
  538. spec.signal?.addEventListener('abort', onAbort, { once: true })
  539. // Batch stdin is written and closed up front; process exit and captured
  540. // output remain authoritative, so write errors (EPIPE) are best-effort.
  541. if (typeof stdinMode === 'object' && stdin !== null) {
  542. stdin.on('error', () => { /* stdin write is best-effort; outcome rides on exit/output. */ })
  543. stdin.end(stdinMode.data)
  544. }
  545. const done = new Promise<SubprocessOutcome>((resolve, reject) => {
  546. let pipeDrainTimer: ReturnType<typeof setTimeout> | undefined
  547. const settle = (outcome: SubprocessOutcome): void => {
  548. if (settled) return
  549. settled = true
  550. // Only harness-collected pipes are force-closed at the drain boundary;
  551. // a 'pipe'-mode stream belongs to the caller and closes with the child.
  552. stopCollectors()
  553. cleanup()
  554. resolve(outcome)
  555. }
  556. const fail = (error: unknown): void => {
  557. settled = true
  558. terminate()
  559. stopCollectors()
  560. cleanup()
  561. // Preserve the exact parent-local AbortSignal reason, including null or undefined.
  562. // oxlint-disable-next-line typescript/prefer-promise-reject-errors -- Exact cancellation reason is the contract.
  563. reject(error)
  564. }
  565. launch.direct.then((outcome) => {
  566. if (stdoutClosed === undefined && stderrClosed === undefined) {
  567. settle(outcome)
  568. return
  569. }
  570. pipeDrainTimer = setTimeout(() => { settle(outcome) }, spec.graceMs)
  571. void outputStreamsClosed.then(() => { settle(outcome) })
  572. }, fail)
  573. function cleanup(): void {
  574. // graceTimer deliberately NOT cleared: forced termination must still
  575. // reach range survivors after the spawned command settles.
  576. if (pipeDrainTimer !== undefined) clearTimeout(pipeDrainTimer)
  577. }
  578. })
  579. const waitForExit = async (signal?: AbortSignal): Promise<boolean> => {
  580. if (rangeExitObserved) return true
  581. return waitWithAbort(observeRangeExit(), signal)
  582. }
  583. return {
  584. /* v8 ignore start -- pipe-mode streams exist on every conforming launch;
  585. the null-coalesces guard an internal adapter defect only. */
  586. stdin: stdinMode === 'pipe' ? stdin ?? undefined : undefined,
  587. stdout: outMode === 'pipe' ? stdout ?? undefined : undefined,
  588. stderr: errMode === 'pipe' ? stderr ?? undefined : undefined,
  589. control: launch.control,
  590. /* v8 ignore stop */
  591. collected: {
  592. ...stdoutCollector !== undefined ? { stdout: stdoutCollector } : {},
  593. ...stderrCollector !== undefined ? { stderr: stderrCollector } : {},
  594. },
  595. done,
  596. terminate,
  597. terminateForHostExit,
  598. waitForExit,
  599. }
  600. }
  601. /**
  602. * Spawn one detached PGID/taskkill fallback and bind the common lifecycle.
  603. * @param spec - fully resolved argv, cwd, stdio, grace, cancellation, environment.
  604. * @param internals - test-only spill-directory, platform, and taskkill overrides.
  605. * @returns live subprocess handle.
  606. */
  607. export function spawnSubprocess(spec: SubprocessSpawnSpec, internals: SpawnInternals = {}): LocalSubprocessHandle {
  608. const binding = prepareManagedProcessBinding(internals)
  609. const platform = internals.platform ?? process.platform
  610. const [program, ...args] = spec.argv
  611. const stdio: import('node:child_process').StdioOptions = [
  612. spec.stdio.stdin === 'ignore' ? 'ignore' : 'pipe',
  613. spec.stdio.stdout === 'inherit' ? 'inherit' : 'pipe',
  614. spec.stdio.stderr === 'inherit' ? 'inherit' : 'pipe',
  615. ]
  616. if (spec.stdio.control === 'pipe') {
  617. while (stdio.length < SUBPROCESS_CONTROL_FD) stdio.push('ignore')
  618. stdio.push('overlapped')
  619. }
  620. const child = (internals.spawn ?? spawn)(program as string, args, {
  621. cwd: spec.cwd,
  622. env: controlEnvironment(childEnv(spec.env), spec.stdio.control),
  623. stdio,
  624. detached: platform !== 'win32',
  625. windowsHide: platform === 'win32',
  626. })
  627. const direct = directChildResult(child)
  628. const pid = child.pid
  629. const owner = fallbackOwner(
  630. platform,
  631. pid,
  632. child,
  633. internals.taskkill ?? taskkillProcessTree,
  634. internals.linuxProcessGroupHasLiveMembers ?? linuxProcessGroupHasLiveMembers,
  635. direct,
  636. )
  637. return bindManagedProcess(spec, {
  638. stdin: child.stdin,
  639. stdout: child.stdout,
  640. stderr: child.stderr,
  641. control: controlPipe(child, spec.stdio.control),
  642. direct,
  643. owner,
  644. }, binding)
  645. }