| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522 |
- /**
- * Process plumbing for the local subprocess service: detached process-tree
- * spawn with per-stream stdio dispositions, tail-keep collection with spill
- * files, tree-scoped signalling (POSIX groups; Windows taskkill), the
- * SIGTERM→SIGKILL escalation, and the cooperative EOF-first dispose ladder.
- * This layer reacts to an abort signal; callers own deadlines and classify
- * causes.
- * @module dsh-subprocess-local/spawn
- */
- import { type ChildProcess, spawn, spawnSync } from 'node:child_process'
- import type { Readable } from 'node:stream'
- import { randomBytes } from 'node:crypto'
- import { closeSync, mkdtempSync, openSync, unlinkSync, writeSync } from 'node:fs'
- import { tmpdir } from 'node:os'
- import { join } from 'node:path'
- import { setTimeout as sleepMs } from 'node:timers/promises'
- import { deadline } from '@deepseek-ai/dsh-timeout'
- import { DSH_ENV_PREFIX, scrubbedParentEnv } from '@deepseek-ai/dsh-subprocess'
- import type {
- CollectedOutput,
- DshEnvironment,
- SubprocessCollect,
- SubprocessDisposeGraces,
- SubprocessHandle,
- SubprocessOutcome,
- SubprocessOutputMode,
- SubprocessSpawnSpec,
- } from '@deepseek-ai/dsh-subprocess'
- /**
- * Build a child environment from the scrubbed parent base, ordinary caller
- * entries, and a managed `DSH_*` snapshot. Ordinary and managed entries
- * reject the other channel's namespace before `dshEnv` merges last.
- * @param extra - caller entries; `DSH_*` names are rejected.
- * @param dshEnv - managed entries; non-`DSH_*` names are rejected.
- * @returns the environment to hand to `spawn` for the child process.
- */
- export function childEnv(
- extra?: Readonly<Record<string, string>>,
- dshEnv?: DshEnvironment,
- ): NodeJS.ProcessEnv {
- for (const key of Object.keys(extra ?? {})) {
- if (key.startsWith(DSH_ENV_PREFIX)) {
- throw new Error(`ordinary child env cannot set reserved variable "${key}"; use dshEnv`)
- }
- }
- for (const key of Object.keys(dshEnv ?? {})) {
- if (!key.startsWith(DSH_ENV_PREFIX)) {
- throw new Error(`managed child env cannot set ordinary variable "${key}"; use env`)
- }
- }
- return { ...scrubbedParentEnv(), ...extra, ...dshEnv }
- }
- /** Injectable knobs so tests can exercise spill and platform behavior deterministically. */
- export interface SpawnInternals {
- /** Directory for spill files (defaults to the OS temp dir). */
- spillDir?: string
- /** Windows tree-termination runner (defaults to `taskkill /PID <pid> /T /F`). */
- taskkill?: (pid: number) => void
- /** Host platform override for signalling decisions. */
- platform?: NodeJS.Platform
- }
- /** Timeout code marking a dispose-ladder tier bound (vs an external abort). */
- const DISPOSE_TIER_TIMEOUT = 'SUBPROCESS_DISPOSE_TIER'
- /**
- * Liveness-poll cadence for tree-exit waits. The timer stays ref'd: an
- * awaited teardown must keep the event loop alive until the tree really
- * exits, or the parent can exit while claiming quiescence and orphan the
- * survivors it promised to reap.
- */
- function sleepTick(): Promise<void> {
- return sleepMs(15)
- }
- let spillCounter = 0
- let defaultSpillDir: string | undefined
- /**
- * The default spill location: a private (0700) per-process directory under
- * the OS tmpdir, created lazily. Predictable world-readable paths would let
- * other local users read command output or pre-create symlinks.
- */
- function privateSpillDir(): string {
- defaultSpillDir ??= mkdtempSync(join(tmpdir(), 'dsh-subprocess-'))
- return defaultSpillDir
- }
- /**
- * Collects one stream with a bounded in-memory tail. With a spill cap, on
- * first overflow a spill file is created and every chunk (including those
- * already collected) is appended there while the full stream remains within
- * the cap; without one, only the in-memory tail is ever retained (the
- * diagnostic-tail shape — a language server's stderr).
- *
- * Tail-keep rationale (pi/OpenCode): errors and final results cluster at the
- * end of command output; the spill file covers the head.
- */
- export class OutputCollector {
- private chunks: Buffer[] = []
- private bytes = 0
- private dropped = false
- private spillFd: number | undefined
- private spillFile: string | undefined
- private spillDisabled: boolean
- /** Total bytes ever pushed (not just retained). */
- private total = 0
- constructor(
- private readonly maxBytes: number,
- private readonly maxSpillBytes: number | undefined,
- private readonly label: string,
- private readonly spillDir: string,
- ) {
- this.spillDisabled = maxSpillBytes === undefined
- }
- /**
- * Ingest one stream chunk, counting it toward the whole-stream total. On
- * first overflow of the in-memory cap a spill file is opened (when spilling
- * is enabled) and every chunk (already-collected ones included) is appended
- * there from then on; the in-memory tail then drops whole chunks from its
- * head (or the head of a single over-cap chunk) until it fits the cap again.
- * @param chunk - the raw bytes from one stream 'data' event.
- */
- push(chunk: Buffer): void {
- this.total += chunk.length
- const overflows = this.bytes + chunk.length > this.maxBytes
- if (!this.spillDisabled && (overflows || this.spillFd !== undefined)) this.spillAll(chunk)
- this.chunks.push(chunk)
- this.bytes += chunk.length
- while (this.bytes > this.maxBytes) {
- const head = this.chunks[0] as Buffer
- const excess = this.bytes - this.maxBytes
- if (head.length <= excess) {
- // Drop the whole head chunk (length ≥ 1 is guaranteed while over cap).
- this.chunks.shift()
- this.bytes -= head.length
- } else {
- // Trim the head so the retained window is byte-exact at the cap — a
- // diagnostic tail (an LSP server's stderr) must hold the LAST
- // maxBytes regardless of how the stream was chunked.
- this.chunks[0] = head.subarray(excess)
- this.bytes -= excess
- }
- this.dropped = true
- }
- }
- /** Open the spill file lazily and append `chunk` (and any prior chunks once). */
- private spillAll(chunk: Buffer): void {
- if (this.maxSpillBytes !== undefined && this.total > this.maxSpillBytes) {
- this.discardSpill()
- return
- }
- if (this.spillFd === undefined) {
- // Random suffix + O_EXCL + no-follow-equivalent ('wx' fails on any
- // existing path, symlink or not) + owner-only mode: defeats spill-path
- // prediction and symlink planting in shared tmp dirs.
- this.spillFile = join(
- this.spillDir,
- `dsh-subprocess-${process.pid}-${++spillCounter}-${randomBytes(6).toString('hex')}-${this.label}.log`,
- )
- this.spillFd = openSync(this.spillFile, 'wx', 0o600)
- for (const prior of this.chunks) writeSync(this.spillFd, prior)
- }
- writeSync(this.spillFd, chunk)
- }
- /** Stop spilling and remove the file once it can no longer hold the complete stream. */
- private discardSpill(): void {
- const fd = this.spillFd
- const file = this.spillFile
- this.spillFd = undefined
- this.spillFile = undefined
- this.spillDisabled = true
- if (fd !== undefined) {
- try {
- closeSync(fd)
- } catch {
- // Retain the descriptor so finalize can retry the failed close.
- this.spillFd = fd
- }
- }
- if (file !== undefined) {
- try {
- unlinkSync(file)
- } catch {
- // A failed unlink leaves at most maxSpillBytes behind, never an unbounded file.
- }
- }
- }
- /**
- * Incremental read in whole-stream byte coordinates: returns everything
- * pushed since `fromByte`. When `fromByte` has already slid out of the
- * in-memory tail window, the read is `lossy` — it returns the whole
- * retained tail and the gap is only recoverable from the spill file.
- * @param fromByte - whole-stream offset to resume from (a prior read's `nextOffset`; 0 for the first read).
- * @returns the delta text, the offset for the next read, the `lossy` flag, and the spill path when one was created.
- */
- readFrom(fromByte: number): { text: string; nextOffset: number; lossy: boolean; spillPath?: string } {
- const windowStart = this.total - this.bytes
- const buffer = Buffer.concat(this.chunks)
- const lossy = fromByte < windowStart
- const slice = lossy ? buffer : buffer.subarray(fromByte - windowStart)
- return {
- text: slice.toString('utf8'),
- nextOffset: this.total,
- lossy,
- ...this.spillFile !== undefined ? { spillPath: this.spillFile } : {},
- }
- }
- /**
- * Close the spill file once the stream has ended. A failed close (delayed
- * writeback fault) stops advertising the spill path — the file may be
- * missing its tail — while every in-memory read keeps working. Idempotent;
- * the spawn path seals both collectors at settlement so reads after exit
- * never point at a still-open file.
- */
- seal(): void {
- if (this.spillFd === undefined) return
- try {
- closeSync(this.spillFd)
- } catch {
- // A delayed writeback failure makes the spill unreliable; keep the
- // in-memory result but stop advertising that file.
- this.spillFile = undefined
- }
- this.spillFd = undefined
- }
- /**
- * Seal the spill file and return the final output.
- * @returns the final collected output: tail text, truncation flag, and the spill path when intact.
- */
- finalize(): CollectedOutput {
- this.seal()
- return {
- text: Buffer.concat(this.chunks).toString('utf8'),
- truncated: this.dropped,
- ...this.spillFile !== undefined ? { spillPath: this.spillFile } : {},
- }
- }
- }
- /**
- * Send `sig` to a detached POSIX process group. Never throws: delivery races
- * process exit and may run in a timer callback, so failures are contained and
- * a non-positive pid is a no-op.
- * @param pid - the group leader's pid; non-positive means the spawn failed and the call is a no-op.
- * @param sig - the signal to deliver to the whole group.
- */
- export function killGroup(pid: number, sig: NodeJS.Signals): void {
- if (pid <= 0) return
- try {
- process.kill(-pid, sig)
- } catch {
- // Swallow: see contract above.
- }
- }
- /**
- * Terminate one Windows process tree with `taskkill /T /F`. Contained like
- * POSIX group signalling — delivery races tree exit, so an absent tree, a
- * nonzero status, or a missing taskkill binary must not break idempotent
- * teardown.
- * @param pid - root process id; non-positive is a no-op.
- */
- export function taskkillProcessTree(pid: number): void {
- if (pid <= 0) return
- // Outcome deliberately unchecked: an already-absent tree (status 128), exit
- // races, and a missing taskkill binary (spawnSync reports, never throws) are
- // as tolerable here as ESRCH is for a POSIX group signal.
- spawnSync('taskkill', ['/PID', String(pid), '/T', '/F'], { stdio: 'ignore' })
- }
- /**
- * Signal a detached process tree with platform-correct semantics: POSIX
- * signals the negative process-group id and falls back to the direct child
- * when the group is gone; Windows terminates the tree via taskkill (any
- * signal value force-terminates — Node maps signals to TerminateProcess).
- */
- function signalTree(
- platform: NodeJS.Platform,
- pid: number,
- sig: NodeJS.Signals,
- child: ChildProcess,
- taskkill: (pid: number) => void,
- ): void {
- if (platform === 'win32') {
- taskkill(pid)
- return
- }
- /* v8 ignore next -- kill/terminate gate on treeAlive(), which is false for pid -1; this guard protects direct callers only. */
- if (pid <= 0) return
- try {
- process.kill(-pid, sig)
- } catch {
- /* v8 ignore start -- the fallback needs a live child whose group signal fails
- (EPERM-style), which POSIX CI cannot stage; the swallow keeps teardown idempotent. */
- try {
- child.kill(sig)
- } catch {
- // The direct child already exited; teardown remains idempotent.
- }
- /* v8 ignore stop */
- }
- }
- /**
- * Spawn one isolated detached process tree with the spec's per-stream stdio
- * dispositions. Runtime exits resolve `done` as {@link SubprocessOutcome};
- * only spawn failures reject.
- * @param spec - fully resolved argv, cwd, stdio, grace, cancellation, environment.
- * @param internals - test-only spill-directory, platform, and taskkill overrides.
- * @returns live subprocess handle.
- */
- export function spawnSubprocess(spec: SubprocessSpawnSpec, internals: SpawnInternals = {}): SubprocessHandle {
- const spillDir = internals.spillDir ?? privateSpillDir()
- const platform = internals.platform ?? process.platform
- const taskkill = internals.taskkill ?? taskkillProcessTree
- if (spec.signal?.aborted) {
- throw new Error(`aborted before spawn: ${String(spec.signal.reason ?? 'aborted')}`)
- }
- const [program, ...args] = spec.argv
- if (program === undefined || program.length === 0) {
- throw new Error('invalid argv: expected a non-empty program name at argv[0]')
- }
- const isCollect = (mode: SubprocessOutputMode): mode is SubprocessCollect =>
- mode !== 'pipe' && mode !== 'inherit'
- const outMode = spec.stdio.stdout
- const errMode = spec.stdio.stderr
- const stdinMode = spec.stdio.stdin
- const env = childEnv(spec.env, spec.dshEnv)
- const child = spawn(program, args, {
- cwd: spec.cwd,
- env,
- stdio: [
- stdinMode === 'ignore' ? 'ignore' : 'pipe',
- outMode === 'inherit' ? 'inherit' : 'pipe',
- errMode === 'inherit' ? 'inherit' : 'pipe',
- ],
- // `detached` gives teardown a tree root on POSIX (its own process group);
- // Windows terminates by root pid through taskkill /T instead.
- detached: platform !== 'win32',
- })
- const collectStream = (mode: SubprocessOutputMode, stream: Readable | null, label: string): OutputCollector | undefined => {
- if (!isCollect(mode) || stream === null) return undefined
- const collector = new OutputCollector(mode.maxBytes, mode.spill?.maxBytes, label, spillDir)
- stream.on('data', (chunk: Buffer) => { collector.push(chunk) })
- return collector
- }
- const stdoutCollector = collectStream(outMode, child.stdout, 'stdout')
- const stderrCollector = collectStream(errMode, child.stderr, 'stderr')
- let graceTimer: NodeJS.Timeout | undefined
- let settled = false
- // Failed spawns use pid -1 so signalling remains a no-op.
- const pid = child.pid ?? -1
- /** Whether the detached tree's root (or POSIX group) is still alive. */
- const treeAlive = (): boolean => {
- if (pid <= 0) return false
- if (platform === 'win32') {
- // Windows has no group-liveness probe; the direct child's exit is the
- // observable boundary (taskkill /T already took the tree with it).
- return child.exitCode === null && child.signalCode === null
- }
- try {
- process.kill(-pid, 0)
- return true
- } catch (error) {
- const code = (error as NodeJS.ErrnoException).code
- /* v8 ignore next 2 -- POSIX reports an absent group as ESRCH; child-reaping timing
- makes observing the other arm platform-dependent. */
- if (code === 'ESRCH') return false
- /* v8 ignore start -- EPERM and non-POSIX negative-pid failures are platform defenses; CI runs
- tree-lifecycle tests on POSIX hosts where absence reports ESRCH. */
- if (code === 'EPERM') return true
- return child.exitCode === null && child.signalCode === null
- /* v8 ignore stop */
- }
- }
- const kill = (sig: NodeJS.Signals = 'SIGTERM'): void => {
- // Guard on TREE liveness, not outcome settlement: a TERM-trapping helper
- // can outlive the settled direct child and must stay signalable, while a
- // fully-dead tree (possible pid reuse) must not be re-signalled from a
- // caller's finally block.
- if (!treeAlive()) return
- signalTree(platform, pid, sig, child, taskkill)
- }
- const terminate = (): void => {
- if (graceTimer !== undefined) return // escalation already in flight
- if (!treeAlive()) return
- signalTree(platform, pid, 'SIGTERM', child, taskkill)
- // The escalation must survive direct-child settlement — the leader dying
- // does not mean the tree died — so settle does not clear this timer, and
- // it re-probes tree liveness before force-killing. It stays ref'd: the
- // pending SIGKILL is a commitment, and a parent exiting before it fires
- // would orphan a trapped survivor. Self-bounds at graceMs.
- graceTimer = setTimeout(() => {
- if (treeAlive()) signalTree(platform, pid, 'SIGKILL', child, taskkill)
- }, spec.graceMs)
- }
- // The caller owns timeout classification; this layer only reacts to abort.
- const onAbort = (): void => { terminate() }
- spec.signal?.addEventListener('abort', onAbort, { once: true })
- // Batch stdin is written and closed up front; process exit and captured
- // output remain authoritative, so write errors (EPIPE) are best-effort.
- if (typeof stdinMode === 'object' && child.stdin !== null) {
- child.stdin.on('error', () => { /* stdin write is best-effort; outcome rides on exit/output. */ })
- child.stdin.end(stdinMode.data)
- }
- const done = new Promise<SubprocessOutcome>((resolve, reject) => {
- let pipeDrainTimer: NodeJS.Timeout | undefined
- const settle = (exitCode: number | null, signal: NodeJS.Signals | null): void => {
- if (settled) return
- settled = true
- // Only harness-collected pipes are force-closed at the drain boundary;
- // a 'pipe'-mode stream belongs to the caller and closes with the child.
- if (stdoutCollector !== undefined) child.stdout?.destroy()
- if (stderrCollector !== undefined) child.stderr?.destroy()
- stdoutCollector?.seal()
- stderrCollector?.seal()
- cleanup()
- resolve({ exitCode, signal })
- }
- child.on('error', (error) => {
- // No meaningful close outcome follows a spawn failure.
- settled = true
- cleanup()
- reject(error)
- })
- child.on('exit', (exitCode, signal) => {
- // A surviving descendant that inherited a pipe must not hold the
- // outcome open indefinitely: after exit, the same bounded grace that
- // governs kills also bounds the close wait.
- pipeDrainTimer = setTimeout(() => { settle(exitCode, signal) }, spec.graceMs)
- })
- child.on('close', settle)
- function cleanup(): void {
- // graceTimer deliberately NOT cleared: the SIGKILL escalation must be
- // able to reach tree survivors after the direct child settles.
- if (pipeDrainTimer !== undefined) clearTimeout(pipeDrainTimer)
- spec.signal?.removeEventListener('abort', onAbort)
- }
- })
- const waitForExit = async (signal?: AbortSignal): Promise<boolean> => {
- while (treeAlive()) {
- if (signal?.aborted) return false
- await sleepTick()
- }
- return true
- }
- /**
- * Wait, bounded, for whole-tree exit — the dispose ladder's quiescence test.
- * Tree liveness, not direct-child settlement: a TERM-trapping helper that
- * outlives the leader must hold the ladder on its tier until it exits.
- */
- const treeExitsWithin = async (ms: number): Promise<boolean> => {
- using bound = deadline(undefined, ms, DISPOSE_TIER_TIMEOUT)
- return await waitForExit(bound.signal)
- }
- let disposal: Promise<void> | undefined
- const dispose = (graces: SubprocessDisposeGraces): Promise<void> => (disposal ??= (async () => {
- // A spawn failure has no process to tear down; observe the rejection so
- // disposal in a finally block cannot surface it as unhandled.
- if (pid <= 0) {
- await done.catch(() => {})
- return
- }
- // 1. Close a piped stdin and allow cooperative teardown and flush.
- if (stdinMode === 'pipe') child.stdin?.end()
- if (await treeExitsWithin(graces.eofGraceMs)) return
- // 2. POSIX gets a catchable graceful signal; Windows taskkill force-terminates.
- if (platform !== 'win32') {
- kill('SIGTERM')
- if (await treeExitsWithin(graces.graceMs)) return
- }
- // 3. Force-kill the tree and await a bounded exit edge.
- kill('SIGKILL')
- if (!(await treeExitsWithin(graces.graceMs))) {
- throw new Error(`child process tree did not exit within ${graces.graceMs}ms after forced termination`)
- }
- })())
- return {
- pid,
- /* v8 ignore start -- pipe-mode fds exist on every spawn Node returns; the null-coalesces guard a nonconforming ChildProcess only. */
- stdin: stdinMode === 'pipe' ? child.stdin ?? undefined : undefined,
- stdout: outMode === 'pipe' ? child.stdout ?? undefined : undefined,
- stderr: errMode === 'pipe' ? child.stderr ?? undefined : undefined,
- /* v8 ignore stop */
- collected: {
- ...stdoutCollector !== undefined ? { stdout: stdoutCollector } : {},
- ...stderrCollector !== undefined ? { stderr: stderrCollector } : {},
- },
- done,
- kill,
- terminate,
- waitForExit,
- dispose,
- }
- }
|