fix(bash-local): bound inherited pipe drain and spill files

This commit is contained in:
Yichen Jiang
2026-07-17 17:21:10 +08:00
parent 44e126e306
commit a6915745e0
6 changed files with 151 additions and 27 deletions

View File

@@ -1,6 +1,6 @@
# @deepseek-ai/dsh-bash-local
Local-subprocess implementation of the `@deepseek-ai/dsh-bash` executor seam: `LocalBashExecutor` spawns `bash -c <command>` per call in its own process group, collects bounded output with full-stream spill files, and escalates kills SIGTERM→SIGKILL across the whole group.
Local-subprocess implementation of the `@deepseek-ai/dsh-bash` executor seam: `LocalBashExecutor` spawns `bash -c <command>` per call in its own process group, collects bounded output with size-limited full-stream spill files, and escalates kills SIGTERM→SIGKILL across the whole group.
The package root exports the default and named `LocalBashExecutor` plugin plus its `Config`; subprocess plumbing stays internal to the implementation package.
@@ -14,7 +14,8 @@ The package root exports the default and named `LocalBashExecutor` plugin plus i
timeoutMs: 120000 # default foreground timeout
maxTimeoutMs: 600000 # cap for per-call overrides
maxOutputBytes: 64000 # per-stream in-memory cap; overflow spills to disk
graceMs: 3000 # SIGTERM→SIGKILL escalation grace on kills
maxSpillBytes: 67108864 # per-stream full-output spill cap
graceMs: 3000 # kill escalation and post-exit pipe-drain grace
```
## Behavior (and where it came from)
@@ -22,8 +23,8 @@ The package root exports the default and named `LocalBashExecutor` plugin plus i
Design surveyed against the bash tools of Claude Code, OpenCode, Codex, and pi; the notable choices:
- **Spawn per call, no shell state** — every call is a fresh non-login `bash -c` (deterministic; no rc files). All four surveyed tools spawn per call. `XXX(stateful-shell)` in `src/run.ts` records the two proven stateful designs (Claude Code's cwd-only persistence; Codex's PTY exec sessions) for when real workflows demand them.
- **Process-group kills with escalation** — children are spawned `detached` (own process group); kills send SIGTERM to the group, then SIGKILL after the `graceMs` grace (default 3s — OpenCode's escalation; pipelines and subshells die with the parent). ESRCH is tolerated; daemons that re-parent away from the group can still survive — same caveat as the surveyed tools.
- **Tail-keep truncation + spill files** — output beyond `maxOutputBytes` keeps the in-memory TAIL (errors/results cluster at the end — pi/OpenCode rationale) while the FULL stream is appended to a temp file whose path is reported when available. If the final spill close reports a delayed writeback failure, the executor still returns the tail but withholds the path rather than advertising a possibly incomplete file.
- **Process-group kills with escalation** — children are spawned `detached` (own process group); kills send SIGTERM to the group, then SIGKILL after the `graceMs` grace (default 3s — OpenCode's escalation; pipelines and subshells die with the parent). After the main shell exits, inherited stdout/stderr pipes receive the same bounded drain grace so a surviving descendant cannot hold the command open indefinitely. ESRCH is tolerated; daemons that re-parent away from the group can still survive — same caveat as the surveyed tools.
- **Tail-keep truncation + bounded spill files** — output beyond `maxOutputBytes` keeps the in-memory TAIL (errors/results cluster at the end — pi/OpenCode rationale) while the FULL stream is appended to a temp file whose path is reported when available. A stream larger than `maxSpillBytes` discards its now-incomplete spill and returns only the marked truncated tail. If the final spill close reports a delayed writeback failure, the executor likewise withholds the path rather than advertising an incomplete file.
- **Model-friendly env + credential scrub** — `process.env` minus credential-shaped vars (`*KEY*`/`*SECRET*`/`*TOKEN*`), then `NO_COLOR=1 TERM=dumb PAGER=cat GIT_PAGER=cat` (Codex's hardcoded set) so pagers and ANSI color don't garble results. This scrub is the security control that keeps the harness's *ambient* credentials out of a spawned command. A spec's `env` is merged LAST (after the scrub), so a caller's explicit entry — a value it already holds — wins even on a credential-shaped name. The spec's `stdin`, when supplied, is written to the child and closed; with none supplied, fd 0 is `/dev/null` — the exact pre-seam default, so a command that probes stdin's file type is unaffected. Both `env`/`stdin` are set by in-process plugins (the hooks bridges); the model-facing tool doesn't expose them. See [the bash-stdin-env RFC](../../../docs/rfc/implemented/architecture/2026-06-30-bash-stdin-env-trusted-plugin-surface.md).
- **Background processes** — `start()` returns a live `BashProcess` handle immediately, no timeout applies (Claude Code detaches timeouts when backgrounding), the handle's `readOutput()` is incremental with whole-stream byte offsets, and disposal kills every running process and awaits its exit. Everything task-shaped (ids, ownership, polling, notices) lives in the generic [`ctx.tasks` runtime](../../tasks/tasks/README.md), which the tool layer registers the handle with — this executor never sees a session or a registry.
@@ -37,6 +38,6 @@ Indirectly, through `dsh-tool-bash`, which renders this executor's bounded stdou
- **No persistent shell or PTY** — every call starts a fresh non-login `bash -c`; cwd-only persistence and interactive terminal sessions remain deferred until a real workflow requires them.
- **POSIX-only** — the `bash` binary, detached process groups, group kills, and SIGTERM→SIGKILL escalation are hardcoded; Windows is unsupported.
- **The credential scrub is a name heuristic** — `*KEY*`/`*SECRET*`/`*TOKEN*` only; differently-named secrets (e.g. `*PASSWORD*`) pass through, and a whitelist for over-scrubbed vars is noted future work.
- **Spill files are never deleted** — full-output recovery files (and the private per-process spill dir) accumulate under the OS tmpdir until something external cleans them.
- **Completed spill files are not deleted** — bounded full-output recovery files (and the private per-process spill dir) accumulate under the OS tmpdir until something external cleans them; oversize incomplete spills are deleted immediately.
The raw process handling lives in `src/run.ts`; `src/index.ts` is the service wiring.

View File

@@ -10,7 +10,7 @@ import z from 'schemastery'
import { BashExecutor } from '@deepseek-ai/dsh-bash'
import type { BashExecRequest, BashExecSpec, BashProcess, BashProcessRead, BashRunResult } from '@deepseek-ai/dsh-bash'
import { clampTimeout, deadline, timeoutOf } from '@deepseek-ai/dsh-timeout'
import { DEFAULT_GRACE_MS, runBash } from './run.ts'
import { DEFAULT_GRACE_MS, DEFAULT_MAX_SPILL_BYTES, runBash } from './run.ts'
import type { RunInternals, RunningBash } from './run.ts'
/** Plugin config (all optional — `static Config` supplies the defaults). */
@@ -23,7 +23,9 @@ export interface Config {
maxTimeoutMs?: number
/** Per-stream in-memory output cap; overflow spills to a temp file. */
maxOutputBytes?: number
/** Grace period between the SIGTERM and the SIGKILL escalation on a kill. */
/** 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
}
@@ -46,6 +48,7 @@ export class LocalBashExecutor extends BashExecutor {
timeoutMs: z.number().default(120_000),
maxTimeoutMs: z.number().default(600_000),
maxOutputBytes: z.number().default(64_000),
maxSpillBytes: z.number().default(DEFAULT_MAX_SPILL_BYTES),
graceMs: z.number().default(DEFAULT_GRACE_MS),
})
@@ -64,6 +67,7 @@ export class LocalBashExecutor extends BashExecutor {
assertPositiveFinite('timeoutMs', this.config.timeoutMs)
assertPositiveFinite('maxTimeoutMs', this.config.maxTimeoutMs)
assertPositiveFinite('maxOutputBytes', this.config.maxOutputBytes)
assertPositiveFinite('maxSpillBytes', this.config.maxSpillBytes)
assertPositiveFinite('graceMs', this.config.graceMs)
ctx.effect(() => async () => {
// Await closure so even a TERM-trapping child cannot outlive the fiber.
@@ -112,6 +116,7 @@ export class LocalBashExecutor extends BashExecutor {
command: spec.command,
cwd: spec.workdir,
maxOutputBytes: this.config.maxOutputBytes,
maxSpillBytes: this.config.maxSpillBytes,
graceMs: this.config.graceMs,
signal: d.signal,
stdin: spec.stdin,
@@ -129,6 +134,7 @@ export class LocalBashExecutor extends BashExecutor {
command: spec.command,
cwd: spec.workdir,
maxOutputBytes: this.config.maxOutputBytes,
maxSpillBytes: this.config.maxSpillBytes,
graceMs: this.config.graceMs,
signal: spec.signal,
stdin: spec.stdin,

View File

@@ -8,7 +8,7 @@
import { type ChildProcessByStdio, spawn } from 'node:child_process'
import type { Readable, Writable } from 'node:stream'
import { randomBytes } from 'node:crypto'
import { closeSync, mkdtempSync, openSync, writeSync } from 'node:fs'
import { closeSync, mkdtempSync, openSync, unlinkSync, writeSync } from 'node:fs'
import { tmpdir } from 'node:os'
import { join } from 'node:path'
import type { CollectedOutput } from '@deepseek-ai/dsh-bash'
@@ -54,7 +54,9 @@ export interface SpawnSpec {
cwd: string
/** Per-stream in-memory cap; overflow spills to disk (tail kept in memory). */
maxOutputBytes: number
/** Grace period between the SIGTERM and the SIGKILL escalation on a kill. */
/** 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
@@ -101,6 +103,9 @@ export interface RunInternals {
/** 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
@@ -115,9 +120,9 @@ function privateSpillDir(): string {
}
/**
* Collects one stream with a bounded in-memory tail. The FULL stream is
* always recoverable: on first overflow a spill file is created and every
* chunk (including those already collected) is appended there.
* 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.
@@ -128,11 +133,13 @@ export class OutputCollector {
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,
) {}
@@ -148,7 +155,7 @@ export class OutputCollector {
push(chunk: Buffer): void {
this.total += chunk.length
const overflows = this.bytes + chunk.length > this.maxBytes
if (overflows || this.spillFd !== undefined) this.spillAll(chunk)
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) {
@@ -170,6 +177,10 @@ export class OutputCollector {
/** 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
@@ -184,6 +195,30 @@ export class OutputCollector {
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
@@ -283,8 +318,8 @@ export function runBash(spec: SpawnSpec, internals: RunInternals = {}): RunningB
? 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.maxOutputBytes, 'stdout', spillDir)
const stderr = new OutputCollector(spec.maxOutputBytes, 'stderr', spillDir)
const stdout = new OutputCollector(spec.maxOutputBytes, spec.maxSpillBytes, 'stdout', spillDir)
const stderr = new OutputCollector(spec.maxOutputBytes, spec.maxSpillBytes, 'stderr', spillDir)
child.stdout.on('data', (chunk: Buffer) => { stdout.push(chunk) })
child.stderr.on('data', (chunk: Buffer) => { stderr.push(chunk) })
@@ -310,12 +345,13 @@ export function runBash(spec: SpawnSpec, internals: RunInternals = {}): RunningB
}
const done = new Promise<SpawnOutcome>((resolve, reject) => {
child.on('error', (error) => {
// No meaningful close outcome follows a spawn failure.
cleanup()
reject(error)
})
child.on('close', (exitCode, signal) => {
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,
@@ -323,9 +359,20 @@ export function runBash(spec: SpawnSpec, internals: RunInternals = {}): RunningB
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)
}
})

View File

@@ -66,6 +66,7 @@ describe('LocalBashExecutor.run', () => {
await expect(setup({ timeoutMs: Number.NaN })).rejects.toThrow(/timeoutMs/)
await expect(setup({ maxTimeoutMs: 0 })).rejects.toThrow(/maxTimeoutMs/)
await expect(setup({ maxOutputBytes: -1 })).rejects.toThrow(/maxOutputBytes/)
await expect(setup({ maxSpillBytes: 0 })).rejects.toThrow(/maxSpillBytes/)
await expect(setup({ graceMs: 0 })).rejects.toThrow(/graceMs/)
const { bash } = await setup()

View File

@@ -1,11 +1,14 @@
import { mkdtempSync, readFileSync, statSync } from 'node:fs'
import { mkdtempSync, readFileSync, statSync, unlinkSync } from 'node:fs'
import { tmpdir } from 'node:os'
import { dirname, join } from 'node:path'
import { describe, expect, it, vi } from 'vitest'
import { killGroup, OutputCollector, runBash } from '../src/run.ts'
import type { RunningBash } from '../src/run.ts'
const { failNextClose } = vi.hoisted(() => ({ failNextClose: { value: false } }))
const { failNextClose, failNextUnlink } = vi.hoisted(() => ({
failNextClose: { value: false },
failNextUnlink: { value: false },
}))
vi.mock('node:fs', async (importOriginal) => {
const actual = await importOriginal<typeof import('node:fs')>()
return {
@@ -17,6 +20,13 @@ vi.mock('node:fs', async (importOriginal) => {
}
actual.closeSync(fd)
},
unlinkSync(path: Parameters<typeof actual.unlinkSync>[0]): void {
if (failNextUnlink.value) {
failNextUnlink.value = false
throw Object.assign(new Error('simulated EIO on unlink'), { code: 'EIO' })
}
actual.unlinkSync(path)
},
}
})
@@ -27,6 +37,7 @@ function spec(command: string, overrides: Partial<Parameters<typeof runBash>[0]>
command,
cwd: process.cwd(),
maxOutputBytes: 64_000,
maxSpillBytes: 64 * 1024 * 1024,
graceMs: 3_000,
...overrides,
}
@@ -171,6 +182,22 @@ describe('runBash', () => {
const result = await running.done
expect(result.signal).toBe('SIGTERM')
})
it('bounds inherited-pipe draining after the shell exits', async () => {
const pidFile = join(spillDir, `pipe-holder-${Date.now()}.pid`)
const started = Date.now()
const running = runBash(spec(`sleep 60 & echo $! > ${pidFile}; echo shell-done`, { graceMs: 100 }))
const descendant = await waitForPidFile(pidFile)
try {
const result = await running.done
expect(Date.now() - started).toBeLessThan(1_000)
expect(result.exitCode).toBe(0)
expect(result.stdout.text).toBe('shell-done\n')
} finally {
process.kill(descendant, 'SIGKILL')
await waitGone(descendant)
}
})
})
describe('stdin and extra env (set by in-process plugins)', () => {
@@ -266,7 +293,7 @@ describe('output truncation and spill', () => {
describe('OutputCollector', () => {
it('keeps the tail of a single oversized chunk', () => {
const collector = new OutputCollector(10, 'test', spillDir)
const collector = new OutputCollector(10, 100, 'test', spillDir)
collector.push(Buffer.from('0123456789abcdef'))
const out = collector.finalize()
expect(out.text).toBe('6789abcdef')
@@ -275,7 +302,7 @@ describe('OutputCollector', () => {
})
it('readFrom returns increments and flags lossy reads', () => {
const collector = new OutputCollector(10, 'test', spillDir)
const collector = new OutputCollector(10, 100, 'test', spillDir)
collector.push(Buffer.from('aaaaa'))
const first = collector.readFrom(0)
expect(first.text).toBe('aaaaa')
@@ -296,7 +323,7 @@ describe('OutputCollector', () => {
})
it('contains close failures and drops the spill path', () => {
const collector = new OutputCollector(4, 'closefail', spillDir)
const collector = new OutputCollector(4, 100, 'closefail', spillDir)
collector.push(Buffer.from('aaaa'))
collector.push(Buffer.from('bbbb'))
expect(collector.readFrom(0).spillPath).toBeDefined()
@@ -310,6 +337,46 @@ describe('OutputCollector', () => {
expect(out!.truncated).toBe(true)
expect(out!.spillPath).toBeUndefined()
})
it('discards a spill that exceeds its configured cap', () => {
const collector = new OutputCollector(4, 8, 'bounded', spillDir)
collector.push(Buffer.from('aaaa'))
collector.push(Buffer.from('bbbb'))
const spillPath = collector.readFrom(0).spillPath!
expect(readFileSync(spillPath, 'utf8')).toBe('aaaabbbb')
collector.push(Buffer.from('c'))
collector.push(Buffer.from('dddd'))
const out = collector.finalize()
expect(out.text).toBe('dddd')
expect(out.truncated).toBe(true)
expect(out.spillPath).toBeUndefined()
expect(() => readFileSync(spillPath)).toThrow()
})
it('does not create a spill when the first overflowing chunk exceeds the cap', () => {
const collector = new OutputCollector(4, 4, 'no-spill', spillDir)
collector.push(Buffer.from('abcdefgh'))
const out = collector.finalize()
expect(out.text).toBe('efgh')
expect(out.truncated).toBe(true)
expect(out.spillPath).toBeUndefined()
})
it('contains cleanup failures while disabling an oversize spill', () => {
const collector = new OutputCollector(4, 8, 'cleanup-fail', spillDir)
collector.push(Buffer.from('aaaa'))
collector.push(Buffer.from('bbbb'))
const spillPath = collector.readFrom(0).spillPath!
failNextClose.value = true
failNextUnlink.value = true
expect(() => { collector.push(Buffer.from('c')) }).not.toThrow()
expect(failNextClose.value).toBe(false)
expect(failNextUnlink.value).toBe(false)
expect(collector.finalize().spillPath).toBeUndefined()
unlinkSync(spillPath)
})
})
describe('killGroup', () => {