spawn.ts 18 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465
  1. /**
  2. * Process plumbing for the local subprocess service: detached process-tree
  3. * spawn with per-stream stdio dispositions, tail-keep collection with spill
  4. * files, tree-scoped signalling (POSIX groups; Windows taskkill), and the
  5. * SIGTERM→SIGKILL escalation. This layer reacts to an abort signal; callers
  6. * own deadlines, teardown ladders, and cause classification.
  7. * @module dsh-subprocess-local/spawn
  8. */
  9. import { type ChildProcess, spawn, spawnSync } from 'node:child_process'
  10. import type { Readable } from 'node:stream'
  11. import { randomBytes } from 'node:crypto'
  12. import { closeSync, mkdtempSync, openSync, unlinkSync, writeSync } from 'node:fs'
  13. import { tmpdir } from 'node:os'
  14. import { join } from 'node:path'
  15. import { setTimeout as sleepMs } from 'node:timers/promises'
  16. import { scrubbedParentEnv } from '@deepseek-ai/dsh-subprocess'
  17. import type {
  18. CollectedOutput,
  19. SubprocessCollect,
  20. SubprocessHandle,
  21. SubprocessOutcome,
  22. SubprocessOutputMode,
  23. SubprocessSpawnSpec,
  24. } from '@deepseek-ai/dsh-subprocess'
  25. /**
  26. * Build a child environment: explicit caller entries merge after the scrubbed
  27. * parent base, so a deliberately supplied credential or current `DSH_*` fact
  28. * wins over the scrub that dropped its ambient namesake.
  29. * @param extra - explicit caller entries, merged verbatim after the scrub.
  30. * @returns the environment to hand to `spawn` for the child process.
  31. */
  32. export function childEnv(extra?: Readonly<Record<string, string>>): NodeJS.ProcessEnv {
  33. return { ...scrubbedParentEnv(), ...extra }
  34. }
  35. /** Injectable knobs so tests can exercise spill and platform behavior deterministically. */
  36. export interface SpawnInternals {
  37. /** Directory for spill files (defaults to the OS temp dir). */
  38. spillDir?: string
  39. /** Windows tree-termination runner (defaults to `taskkill /PID <pid> /T /F`). */
  40. taskkill?: (pid: number) => void
  41. /** Host platform override for signalling decisions. */
  42. platform?: NodeJS.Platform
  43. }
  44. /**
  45. * Liveness-poll cadence for tree-exit waits. The timer stays ref'd: an
  46. * awaited teardown must keep the event loop alive until the tree really
  47. * exits, or the parent can exit while claiming quiescence and orphan the
  48. * survivors it promised to reap.
  49. */
  50. function sleepTick(): Promise<void> {
  51. return sleepMs(15)
  52. }
  53. let spillCounter = 0
  54. let defaultSpillDir: string | undefined
  55. /**
  56. * The default spill location: a private (0700) per-process directory under
  57. * the OS tmpdir, created lazily. Predictable world-readable paths would let
  58. * other local users read command output or pre-create symlinks.
  59. */
  60. function privateSpillDir(): string {
  61. defaultSpillDir ??= mkdtempSync(join(tmpdir(), 'dsh-subprocess-'))
  62. return defaultSpillDir
  63. }
  64. /**
  65. * Collects one stream with a bounded in-memory tail. With a spill cap, on
  66. * first overflow a spill file is created and every chunk (including those
  67. * already collected) is appended there while the full stream remains within
  68. * the cap; without one, only the in-memory tail is ever retained (the
  69. * diagnostic-tail shape — a language server's stderr).
  70. *
  71. * Tail-keep rationale (pi/OpenCode): errors and final results cluster at the
  72. * end of command output; the spill file covers the head.
  73. */
  74. export class OutputCollector {
  75. private chunks: Buffer[] = []
  76. private bytes = 0
  77. private dropped = false
  78. private spillFd: number | undefined
  79. private spillFile: string | undefined
  80. private spillDisabled: boolean
  81. /** Total bytes ever pushed (not just retained). */
  82. private total = 0
  83. constructor(
  84. private readonly maxBytes: number,
  85. private readonly maxSpillBytes: number | undefined,
  86. private readonly label: string,
  87. private readonly spillDir: string,
  88. ) {
  89. this.spillDisabled = maxSpillBytes === undefined
  90. }
  91. /**
  92. * Ingest one stream chunk, counting it toward the whole-stream total. On
  93. * first overflow of the in-memory cap a spill file is opened (when spilling
  94. * is enabled) and every chunk (already-collected ones included) is appended
  95. * there from then on; the in-memory tail then drops whole chunks from its
  96. * head (or the head of a single over-cap chunk) until it fits the cap again.
  97. * @param chunk - the raw bytes from one stream 'data' event.
  98. */
  99. push(chunk: Buffer): void {
  100. this.total += chunk.length
  101. const overflows = this.bytes + chunk.length > this.maxBytes
  102. if (!this.spillDisabled && (overflows || this.spillFd !== undefined)) this.spillAll(chunk)
  103. this.chunks.push(chunk)
  104. this.bytes += chunk.length
  105. while (this.bytes > this.maxBytes) {
  106. const head = this.chunks[0] as Buffer
  107. const excess = this.bytes - this.maxBytes
  108. if (head.length <= excess) {
  109. // Drop the whole head chunk (length ≥ 1 is guaranteed while over cap).
  110. this.chunks.shift()
  111. this.bytes -= head.length
  112. } else {
  113. // Trim the head so the retained window is byte-exact at the cap — a
  114. // diagnostic tail (an LSP server's stderr) must hold the LAST
  115. // maxBytes regardless of how the stream was chunked.
  116. this.chunks[0] = head.subarray(excess)
  117. this.bytes -= excess
  118. }
  119. this.dropped = true
  120. }
  121. }
  122. /** Open the spill file lazily and append `chunk` (and any prior chunks once). */
  123. private spillAll(chunk: Buffer): void {
  124. if (this.maxSpillBytes !== undefined && this.total > this.maxSpillBytes) {
  125. this.discardSpill()
  126. return
  127. }
  128. if (this.spillFd === undefined) {
  129. // Random suffix + O_EXCL + no-follow-equivalent ('wx' fails on any
  130. // existing path, symlink or not) + owner-only mode: defeats spill-path
  131. // prediction and symlink planting in shared tmp dirs.
  132. this.spillFile = join(
  133. this.spillDir,
  134. `dsh-subprocess-${process.pid}-${++spillCounter}-${randomBytes(6).toString('hex')}-${this.label}.log`,
  135. )
  136. this.spillFd = openSync(this.spillFile, 'wx', 0o600)
  137. for (const prior of this.chunks) writeSync(this.spillFd, prior)
  138. }
  139. writeSync(this.spillFd, chunk)
  140. }
  141. /** Stop spilling and remove the file once it can no longer hold the complete stream. */
  142. private discardSpill(): void {
  143. const fd = this.spillFd
  144. const file = this.spillFile
  145. this.spillFd = undefined
  146. this.spillFile = undefined
  147. this.spillDisabled = true
  148. if (fd !== undefined) {
  149. try {
  150. closeSync(fd)
  151. } catch {
  152. // Retain the descriptor so finalize can retry the failed close.
  153. this.spillFd = fd
  154. }
  155. }
  156. if (file !== undefined) {
  157. try {
  158. unlinkSync(file)
  159. } catch {
  160. // A failed unlink leaves at most maxSpillBytes behind, never an unbounded file.
  161. }
  162. }
  163. }
  164. /**
  165. * Incremental read in whole-stream byte coordinates: returns everything
  166. * pushed since `fromByte`. When `fromByte` has already slid out of the
  167. * in-memory tail window, the read is `lossy` — it returns the whole
  168. * retained tail and the gap is only recoverable from the spill file.
  169. * @param fromByte - whole-stream offset to resume from (a prior read's `nextOffset`; 0 for the first read).
  170. * @returns the delta text, the offset for the next read, the `lossy` flag, and the spill path when one was created.
  171. */
  172. readFrom(fromByte: number): { text: string; nextOffset: number; lossy: boolean; spillPath?: string } {
  173. const windowStart = this.total - this.bytes
  174. const buffer = Buffer.concat(this.chunks)
  175. const lossy = fromByte < windowStart
  176. const slice = lossy ? buffer : buffer.subarray(fromByte - windowStart)
  177. return {
  178. text: slice.toString('utf8'),
  179. nextOffset: this.total,
  180. lossy,
  181. ...this.spillFile !== undefined ? { spillPath: this.spillFile } : {},
  182. }
  183. }
  184. /**
  185. * Close the spill file once the stream has ended. A failed close (delayed
  186. * writeback fault) stops advertising the spill path — the file may be
  187. * missing its tail — while every in-memory read keeps working. Idempotent;
  188. * the spawn path seals both collectors at settlement so reads after exit
  189. * never point at a still-open file.
  190. */
  191. seal(): void {
  192. if (this.spillFd === undefined) return
  193. try {
  194. closeSync(this.spillFd)
  195. } catch {
  196. // A delayed writeback failure makes the spill unreliable; keep the
  197. // in-memory result but stop advertising that file.
  198. this.spillFile = undefined
  199. }
  200. this.spillFd = undefined
  201. }
  202. /**
  203. * Seal the spill file and return the final output.
  204. * @returns the final collected output: tail text, truncation flag, and the spill path when intact.
  205. */
  206. finalize(): CollectedOutput {
  207. this.seal()
  208. return {
  209. text: Buffer.concat(this.chunks).toString('utf8'),
  210. truncated: this.dropped,
  211. ...this.spillFile !== undefined ? { spillPath: this.spillFile } : {},
  212. }
  213. }
  214. }
  215. /**
  216. * Send `sig` to a detached POSIX process group. Never throws: delivery races
  217. * process exit and may run in a timer callback, so failures are contained and
  218. * a non-positive pid is a no-op.
  219. * @param pid - the group leader's pid; non-positive means the spawn failed and the call is a no-op.
  220. * @param sig - the signal to deliver to the whole group.
  221. */
  222. export function killGroup(pid: number, sig: NodeJS.Signals): void {
  223. if (pid <= 0) return
  224. try {
  225. process.kill(-pid, sig)
  226. } catch {
  227. // Swallow: see contract above.
  228. }
  229. }
  230. /**
  231. * Terminate one Windows process tree with `taskkill /T /F`. Contained like
  232. * POSIX group signalling — delivery races tree exit, so an absent tree, a
  233. * nonzero status, or a missing taskkill binary must not break idempotent
  234. * teardown.
  235. * @param pid - root process id; non-positive is a no-op.
  236. */
  237. export function taskkillProcessTree(pid: number): void {
  238. if (pid <= 0) return
  239. // Outcome deliberately unchecked: an already-absent tree (status 128), exit
  240. // races, and a missing taskkill binary (spawnSync reports, never throws) are
  241. // as tolerable here as ESRCH is for a POSIX group signal.
  242. spawnSync('taskkill', ['/PID', String(pid), '/T', '/F'], { stdio: 'ignore' })
  243. }
  244. /**
  245. * Signal a detached process tree with platform-correct semantics: POSIX
  246. * signals the negative process-group id and falls back to the direct child
  247. * when the group is gone; Windows terminates the tree via taskkill (any
  248. * signal value force-terminates — Node maps signals to TerminateProcess).
  249. */
  250. function signalTree(
  251. platform: NodeJS.Platform,
  252. pid: number,
  253. sig: NodeJS.Signals,
  254. child: ChildProcess,
  255. taskkill: (pid: number) => void,
  256. ): void {
  257. if (platform === 'win32') {
  258. taskkill(pid)
  259. return
  260. }
  261. /* v8 ignore next -- kill/terminate gate on treeAlive(), which is false for pid -1; this guard protects direct callers only. */
  262. if (pid <= 0) return
  263. try {
  264. process.kill(-pid, sig)
  265. } catch {
  266. /* v8 ignore start -- the fallback needs a live child whose group signal fails
  267. (EPERM-style), which POSIX CI cannot stage; the swallow keeps teardown idempotent. */
  268. try {
  269. child.kill(sig)
  270. } catch {
  271. // The direct child already exited; teardown remains idempotent.
  272. }
  273. /* v8 ignore stop */
  274. }
  275. }
  276. /**
  277. * Spawn one isolated detached process tree with the spec's per-stream stdio
  278. * dispositions. Runtime exits resolve `done` as {@link SubprocessOutcome};
  279. * only spawn failures reject.
  280. * @param spec - fully resolved argv, cwd, stdio, grace, cancellation, environment.
  281. * @param internals - test-only spill-directory, platform, and taskkill overrides.
  282. * @returns live subprocess handle.
  283. */
  284. export function spawnSubprocess(spec: SubprocessSpawnSpec, internals: SpawnInternals = {}): SubprocessHandle {
  285. const spillDir = internals.spillDir ?? privateSpillDir()
  286. const platform = internals.platform ?? process.platform
  287. const taskkill = internals.taskkill ?? taskkillProcessTree
  288. if (spec.signal?.aborted) {
  289. throw new Error(`aborted before spawn: ${String(spec.signal.reason ?? 'aborted')}`)
  290. }
  291. const [program, ...args] = spec.argv
  292. if (program === undefined || program.length === 0) {
  293. throw new Error('invalid argv: expected a non-empty program name at argv[0]')
  294. }
  295. const isCollect = (mode: SubprocessOutputMode): mode is SubprocessCollect =>
  296. mode !== 'pipe' && mode !== 'inherit'
  297. const outMode = spec.stdio.stdout
  298. const errMode = spec.stdio.stderr
  299. const stdinMode = spec.stdio.stdin
  300. const env = childEnv(spec.env)
  301. const child = spawn(program, args, {
  302. cwd: spec.cwd,
  303. env,
  304. stdio: [
  305. stdinMode === 'ignore' ? 'ignore' : 'pipe',
  306. outMode === 'inherit' ? 'inherit' : 'pipe',
  307. errMode === 'inherit' ? 'inherit' : 'pipe',
  308. ],
  309. // `detached` gives teardown a tree root on POSIX (its own process group);
  310. // Windows terminates by root pid through taskkill /T instead.
  311. detached: platform !== 'win32',
  312. })
  313. const collectStream = (mode: SubprocessOutputMode, stream: Readable | null, label: string): OutputCollector | undefined => {
  314. if (!isCollect(mode) || stream === null) return undefined
  315. const collector = new OutputCollector(mode.maxBytes, mode.spill?.maxBytes, label, spillDir)
  316. stream.on('data', (chunk: Buffer) => { collector.push(chunk) })
  317. return collector
  318. }
  319. const stdoutCollector = collectStream(outMode, child.stdout, 'stdout')
  320. const stderrCollector = collectStream(errMode, child.stderr, 'stderr')
  321. let graceTimer: NodeJS.Timeout | undefined
  322. let settled = false
  323. // Failed spawns use pid -1 so signalling remains a no-op.
  324. const pid = child.pid ?? -1
  325. /** Whether the detached tree's root (or POSIX group) is still alive. */
  326. const treeAlive = (): boolean => {
  327. if (pid <= 0) return false
  328. if (platform === 'win32') {
  329. // Windows has no group-liveness probe; the direct child's exit is the
  330. // observable boundary (taskkill /T already took the tree with it).
  331. return child.exitCode === null && child.signalCode === null
  332. }
  333. try {
  334. process.kill(-pid, 0)
  335. return true
  336. } catch (error) {
  337. const code = (error as NodeJS.ErrnoException).code
  338. /* v8 ignore next 2 -- POSIX reports an absent group as ESRCH; child-reaping timing
  339. makes observing the other arm platform-dependent. */
  340. if (code === 'ESRCH') return false
  341. /* v8 ignore start -- EPERM and non-POSIX negative-pid failures are platform defenses; CI runs
  342. tree-lifecycle tests on POSIX hosts where absence reports ESRCH. */
  343. if (code === 'EPERM') return true
  344. return child.exitCode === null && child.signalCode === null
  345. /* v8 ignore stop */
  346. }
  347. }
  348. // The escalation's tier primitive (not on the handle — terminate() is the
  349. // only consumer-facing termination verb). Guards on TREE liveness, not
  350. // outcome settlement: a TERM-trapping helper can outlive the settled direct
  351. // child and must stay signalable, while a fully-dead tree (possible pid
  352. // reuse) must not be re-signalled by a later tier.
  353. const kill = (sig: NodeJS.Signals): void => {
  354. if (!treeAlive()) return
  355. signalTree(platform, pid, sig, child, taskkill)
  356. }
  357. const terminate = (): void => {
  358. if (graceTimer !== undefined) return // escalation already in flight
  359. if (!treeAlive()) return
  360. kill('SIGTERM')
  361. // The escalation must survive direct-child settlement — the leader dying
  362. // does not mean the tree died — so settle does not clear this timer, and
  363. // kill() re-probes tree liveness before force-killing. It stays ref'd:
  364. // the pending SIGKILL is a commitment, and a parent exiting before it
  365. // fires would orphan a trapped survivor. Self-bounds at graceMs.
  366. graceTimer = setTimeout(() => { kill('SIGKILL') }, spec.graceMs)
  367. }
  368. // The caller owns timeout classification; this layer only reacts to abort.
  369. const onAbort = (): void => { terminate() }
  370. spec.signal?.addEventListener('abort', onAbort, { once: true })
  371. // Batch stdin is written and closed up front; process exit and captured
  372. // output remain authoritative, so write errors (EPIPE) are best-effort.
  373. if (typeof stdinMode === 'object' && child.stdin !== null) {
  374. child.stdin.on('error', () => { /* stdin write is best-effort; outcome rides on exit/output. */ })
  375. child.stdin.end(stdinMode.data)
  376. }
  377. const done = new Promise<SubprocessOutcome>((resolve, reject) => {
  378. let pipeDrainTimer: NodeJS.Timeout | undefined
  379. const settle = (exitCode: number | null, signal: NodeJS.Signals | null): void => {
  380. if (settled) return
  381. settled = true
  382. // Only harness-collected pipes are force-closed at the drain boundary;
  383. // a 'pipe'-mode stream belongs to the caller and closes with the child.
  384. if (stdoutCollector !== undefined) child.stdout?.destroy()
  385. if (stderrCollector !== undefined) child.stderr?.destroy()
  386. stdoutCollector?.seal()
  387. stderrCollector?.seal()
  388. cleanup()
  389. resolve({ exitCode, signal })
  390. }
  391. child.on('error', (error) => {
  392. // No meaningful close outcome follows a spawn failure.
  393. settled = true
  394. cleanup()
  395. reject(error)
  396. })
  397. child.on('exit', (exitCode, signal) => {
  398. // A surviving descendant that inherited a pipe must not hold the
  399. // outcome open indefinitely: after exit, the same bounded grace that
  400. // governs kills also bounds the close wait.
  401. pipeDrainTimer = setTimeout(() => { settle(exitCode, signal) }, spec.graceMs)
  402. })
  403. child.on('close', settle)
  404. function cleanup(): void {
  405. // graceTimer deliberately NOT cleared: the SIGKILL escalation must be
  406. // able to reach tree survivors after the direct child settles.
  407. if (pipeDrainTimer !== undefined) clearTimeout(pipeDrainTimer)
  408. spec.signal?.removeEventListener('abort', onAbort)
  409. }
  410. })
  411. const waitForExit = async (signal?: AbortSignal): Promise<boolean> => {
  412. while (treeAlive()) {
  413. if (signal?.aborted) return false
  414. await sleepTick()
  415. }
  416. return true
  417. }
  418. return {
  419. pid,
  420. /* v8 ignore start -- pipe-mode fds exist on every spawn Node returns; the null-coalesces guard a nonconforming ChildProcess only. */
  421. stdin: stdinMode === 'pipe' ? child.stdin ?? undefined : undefined,
  422. stdout: outMode === 'pipe' ? child.stdout ?? undefined : undefined,
  423. stderr: errMode === 'pipe' ? child.stderr ?? undefined : undefined,
  424. /* v8 ignore stop */
  425. collected: {
  426. ...stdoutCollector !== undefined ? { stdout: stdoutCollector } : {},
  427. ...stderrCollector !== undefined ? { stderr: stderrCollector } : {},
  428. },
  429. done,
  430. terminate,
  431. waitForExit,
  432. }
  433. }