| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399 |
- /**
- * Process plumbing for the local bash executor: detached process-group spawn,
- * tail-keep output with spill files, and SIGTERM→SIGKILL escalation. This layer
- * reacts to an abort signal; the executor owns deadlines and classifies causes.
- * @module dsh-bash-local/run
- */
- import { type ChildProcessByStdio, spawn } from 'node:child_process'
- import type { Readable, Writable } 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 { DSH_ENV_PREFIX } from '@deepseek-ai/dsh-bash'
- import type { CollectedOutput, DshEnvironment } from '@deepseek-ai/dsh-bash'
- /**
- * Model-friendly environment overrides: disable colors, pagers, and
- * interactive terminal features that would garble tool output (the same set
- * Codex hardcodes; Claude Code achieves it via TERM=dumb).
- */
- export const ENV_OVERRIDES = {
- NO_COLOR: '1',
- TERM: 'dumb',
- PAGER: 'cat',
- GIT_PAGER: 'cat',
- } as const
- /**
- * Credential-shaped env vars are NOT forwarded to commands (the harness's
- * own DEEPSEEK_API_KEY must not leak into `env` output, tool results, or
- * spill files). Same default pattern as Codex's env policy; a future config
- * can whitelist specific vars when a workflow genuinely needs one.
- */
- export const SENSITIVE_ENV_PATTERN = /KEY|SECRET|TOKEN/i
- /**
- * Build a child environment from scrubbed ambient values, terminal overrides,
- * ordinary caller entries, and a managed `DSH_*` snapshot. Ambient managed
- * names are removed; 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 {
- const env: NodeJS.ProcessEnv = {}
- for (const [key, value] of Object.entries(process.env)) {
- if (!SENSITIVE_ENV_PATTERN.test(key) && !key.startsWith(DSH_ENV_PREFIX)) env[key] = value
- }
- for (const key of Object.keys(extra ?? {})) {
- if (key.startsWith(DSH_ENV_PREFIX)) {
- throw new Error(`ordinary bash 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 bash env cannot set ordinary variable "${key}"; use env`)
- }
- }
- return { ...env, ...ENV_OVERRIDES, ...extra, ...dshEnv }
- }
- /** What to run and under which limits (resolved — no defaults in here). */
- export interface SpawnSpec {
- command: string
- cwd: string
- /** Stdout in-memory cap; overflow spills to disk (tail kept in memory). */
- stdoutMaxBytes: number
- /** Stderr in-memory cap; overflow spills to disk (tail kept in memory). */
- stderrMaxBytes: number
- /** Per-stream spill-file cap; larger streams retain only their in-memory tail. */
- maxSpillBytes: number
- /** Grace period for kill escalation and for inherited pipes after shell exit. */
- graceMs: number
- /**
- * Abort signal — kills the process group when it fires. The executor owns
- * timing: `run()` passes a fused timeout/cancel deadline signal (see
- * `@deepseek-ai/dsh-timeout`), `start()` passes the bare upstream signal.
- * runBash only listens and kills; it does NOT classify why (the executor
- * reads the signal's reason afterward).
- */
- signal?: AbortSignal | undefined
- /**
- * Bytes to write to the child's stdin, then close it. Absent (or empty)
- * leaves stdin closed/empty. Set by in-process plugins (the hooks bridges);
- * the model-facing `dsh-tool-bash` tool does not thread model input here.
- */
- stdin?: string | undefined
- /**
- * Ordinary environment entries merged after the credential scrub and
- * terminal overrides. `DSH_*` names are rejected and belong in `dshEnv`.
- */
- env?: Record<string, string> | undefined
- /** Harness-owned entries; non-`DSH_*` names are rejected before spawn. */
- dshEnv?: DshEnvironment | undefined
- }
- /**
- * Raw outcome of one closed process (before result shaping). Deliberately
- * carries NO timeout/cancel classification: runBash kills on abort but does not
- * decide why — the executor's `run()`/`start()` reads the deadline signal it
- * owns to classify `timedOut`/`aborted` (see the package README).
- */
- export interface SpawnOutcome {
- exitCode: number | null
- signal: NodeJS.Signals | null
- stdout: CollectedOutput
- stderr: CollectedOutput
- }
- /** Injectable knobs so tests can exercise spill behavior without the OS tmpdir. */
- export interface RunInternals {
- /** Directory for spill files (defaults to the OS temp dir). */
- spillDir?: string
- }
- /** Default SIGTERM→SIGKILL grace period (the `graceMs` config; matches OpenCode's 3s). */
- export const DEFAULT_GRACE_MS = 3_000
- /** Default per-stream spill cap (the `maxSpillBytes` config). */
- export const DEFAULT_MAX_SPILL_BYTES = 64 * 1024 * 1024
- 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-bash-'))
- return defaultSpillDir
- }
- /**
- * Collects one stream with a bounded in-memory tail. On first overflow a
- * spill file is created and every chunk (including those already collected)
- * is appended there while the full stream remains within `maxSpillBytes`.
- *
- * 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 = false
- /** Total bytes ever pushed (not just retained). */
- private total = 0
- constructor(
- private readonly maxBytes: number,
- private readonly maxSpillBytes: number,
- private readonly label: string,
- private readonly spillDir: string,
- ) {}
- /**
- * Ingest one stream chunk, counting it toward the whole-stream total. On
- * first overflow of the in-memory cap a spill file is opened 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 && this.chunks.length > 1) {
- // Drop whole chunks from the head; pipe chunks are small (≤64KiB), so
- // the retained tail tracks the cap closely enough for a model-facing
- // truncation boundary. (length > 1 was just checked — shift() returns.)
- const head = this.chunks.shift() as Buffer
- this.bytes -= head.length
- this.dropped = true
- }
- if (this.bytes > this.maxBytes && this.chunks.length === 1) {
- // A single chunk larger than the cap: keep its tail.
- const only = this.chunks[0] as Buffer
- this.chunks[0] = only.subarray(only.length - this.maxBytes)
- this.bytes = this.maxBytes
- this.dropped = true
- }
- }
- /** Open the spill file lazily and append `chunk` (and any prior chunks once). */
- private spillAll(chunk: Buffer): void {
- if (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-bash-${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 (if any) and return the final output. A failed close
- * (delayed writeback fault) stops advertising the spill path — the file may
- * be missing its tail — but still returns the in-memory result.
- * @returns the final collected output: tail text, truncation flag, and the spill path when intact.
- */
- finalize(): CollectedOutput {
- if (this.spillFd !== undefined) {
- try {
- closeSync(this.spillFd)
- } catch {
- // A delayed writeback failure makes the spill unreliable; keep finalize
- // total but stop advertising that file.
- this.spillFile = undefined
- }
- this.spillFd = undefined
- }
- return {
- text: Buffer.concat(this.chunks).toString('utf8'),
- truncated: this.dropped,
- ...this.spillFile !== undefined ? { spillPath: this.spillFile } : {},
- }
- }
- }
- /**
- * Send `sig` to a detached 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.
- }
- }
- /**
- * A live bash child process: the promise resolves when the process closes;
- * `kill()` starts the SIGTERM→grace→SIGKILL escalation on its group.
- */
- export interface RunningBash {
- /** Process id (group leader); -1 when the spawn itself failed. */
- readonly pid: number
- /** stdout/stderr collectors (live — background polling reads incrementally). */
- readonly stdout: OutputCollector
- readonly stderr: OutputCollector
- /** Resolves when the process closes; rejects only for spawn-level failures. */
- readonly done: Promise<SpawnOutcome>
- /** Begin SIGTERM→grace→SIGKILL on the process group. Idempotent. */
- kill(): void
- }
- /**
- * Spawn one isolated `bash -c` process group and collect its output.
- * Runtime exits resolve as {@link SpawnOutcome}; only spawn failures reject.
- * @param spec - fully resolved command, cwd, limits, and cancellation.
- * @param internals - test-only process and spill-directory overrides.
- * @returns live process handle and outcome promise.
- */
- // XXX(stateful-shell): evaluate persistent cwd or PTY sessions when workflows require shell state.
- export function runBash(spec: SpawnSpec, internals: RunInternals = {}): RunningBash {
- const spillDir = internals.spillDir ?? privateSpillDir()
- if (spec.signal?.aborted) {
- throw new Error(`aborted before spawn: ${String(spec.signal.reason ?? 'aborted')}`)
- }
- // Keep absent stdin as /dev/null; literal tuples preserve non-null output types.
- const env = childEnv(spec.env, spec.dshEnv)
- const child: ChildProcessByStdio<Writable | null, Readable, Readable> = spec.stdin !== undefined
- ? spawn('bash', ['-c', spec.command], { cwd: spec.cwd, env, stdio: ['pipe', 'pipe', 'pipe'], detached: true })
- : spawn('bash', ['-c', spec.command], { cwd: spec.cwd, env, stdio: ['ignore', 'pipe', 'pipe'], detached: true })
- const stdout = new OutputCollector(spec.stdoutMaxBytes, spec.maxSpillBytes, 'stdout', spillDir)
- const stderr = new OutputCollector(spec.stderrMaxBytes, spec.maxSpillBytes, 'stderr', spillDir)
- child.stdout.on('data', (chunk: Buffer) => { stdout.push(chunk) })
- child.stderr.on('data', (chunk: Buffer) => { stderr.push(chunk) })
- let graceTimer: NodeJS.Timeout | undefined
- // Failed spawns use pid -1 so kill remains a no-op.
- const pid = child.pid ?? -1
- const kill = (): void => {
- if (graceTimer !== undefined) return // escalation already in flight
- killGroup(pid, 'SIGTERM')
- graceTimer = setTimeout(() => { killGroup(pid, 'SIGKILL') }, spec.graceMs)
- }
- // The executor owns timeout classification; this layer only reacts to abort.
- const onAbort = (): void => { kill() }
- spec.signal?.addEventListener('abort', onAbort, { once: true })
- // Stdin writes are best-effort; process exit and captured output remain authoritative.
- if (child.stdin !== null) {
- child.stdin.on('error', () => { /* stdin write is best-effort; outcome rides on exit/output. */ })
- child.stdin.end(spec.stdin)
- }
- const done = new Promise<SpawnOutcome>((resolve, reject) => {
- let settled = false
- let pipeDrainTimer: NodeJS.Timeout | undefined
- const settle = (exitCode: number | null, signal: NodeJS.Signals | null): void => {
- if (settled) return
- settled = true
- child.stdout.destroy()
- child.stderr.destroy()
- cleanup()
- resolve({
- exitCode,
- signal,
- stdout: stdout.finalize(),
- stderr: stderr.finalize(),
- })
- }
- child.on('error', (error) => {
- // No meaningful close outcome follows a spawn failure.
- settled = true
- cleanup()
- reject(error)
- })
- child.on('exit', (exitCode, signal) => {
- pipeDrainTimer = setTimeout(() => { settle(exitCode, signal) }, spec.graceMs)
- })
- child.on('close', settle)
- function cleanup(): void {
- if (graceTimer !== undefined) clearTimeout(graceTimer)
- if (pipeDrainTimer !== undefined) clearTimeout(pipeDrainTimer)
- spec.signal?.removeEventListener('abort', onAbort)
- }
- })
- return { pid, stdout, stderr, done, kill }
- }
|