run.ts 18 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418
  1. /**
  2. * Process plumbing for the local bash executor: spawn, output collection
  3. * with tail-keep + spill-to-disk truncation, and process-group kill with
  4. * SIGTERM→SIGKILL escalation.
  5. *
  6. * Everything here is deliberately free of Cordis concepts so it can be unit
  7. * tested in isolation; `LocalBashExecutor` owns lifecycle and configuration.
  8. *
  9. * Design notes (surveyed against Claude Code, OpenCode, Codex, and pi — see
  10. * the package README): spawn-per-call with `detached: true` so the child
  11. * leads its own process group; kills target the group (`kill(-pid)`) so
  12. * pipelines and subshells die with the parent. SIGTERM first, SIGKILL after a
  13. * grace period (OpenCode's escalation; Codex/pi jump straight to SIGKILL).
  14. *
  15. * @module dsh-bash-local/run
  16. */
  17. import { type ChildProcessByStdio, spawn } from 'node:child_process'
  18. import type { Readable, Writable } from 'node:stream'
  19. import { randomBytes } from 'node:crypto'
  20. import { closeSync, mkdtempSync, openSync, writeSync } from 'node:fs'
  21. import { tmpdir } from 'node:os'
  22. import { join } from 'node:path'
  23. import type { CollectedOutput } from '@deepseek-ai/dsh-bash'
  24. /**
  25. * Model-friendly environment overrides: disable colors, pagers, and
  26. * interactive terminal features that would garble tool output (the same set
  27. * Codex hardcodes; Claude Code achieves it via TERM=dumb).
  28. */
  29. export const ENV_OVERRIDES = {
  30. NO_COLOR: '1',
  31. TERM: 'dumb',
  32. PAGER: 'cat',
  33. GIT_PAGER: 'cat',
  34. } as const
  35. /**
  36. * Credential-shaped env vars are NOT forwarded to commands (the harness's
  37. * own DEEPSEEK_API_KEY must not leak into `env` output, tool results, or
  38. * spill files). Same default pattern as Codex's env policy; a future config
  39. * can whitelist specific vars when a workflow genuinely needs one.
  40. */
  41. export const SENSITIVE_ENV_PATTERN = /KEY|SECRET|TOKEN/i
  42. /**
  43. * `process.env` minus credential-shaped vars, plus the model-friendly
  44. * overrides, plus any caller-supplied `extra` entries.
  45. *
  46. * Layering matters: the scrub drops `process.env` credentials, then
  47. * `ENV_OVERRIDES` forces the model-friendly terminal vars, then `extra` is
  48. * merged LAST so an explicit caller entry wins even when its name matches the
  49. * scrub pattern (the scrub is the control that stops the HARNESS's ambient
  50. * credentials leaking into a spawned command; a caller that explicitly sets a
  51. * var named a value it already holds, not that ambient secret). `extra` is set
  52. * by in-process plugins (the hooks bridges), not the model — `dsh-tool-bash`
  53. * builds its request from named fields only and does not forward model input
  54. * here (see its README, § "The tool builds its request from named args only").
  55. * @param extra - caller-supplied entries merged last; an explicit entry wins even against the scrub and the overrides.
  56. * @returns the environment to hand to `spawn` for the child process.
  57. */
  58. export function childEnv(extra?: Record<string, string>): NodeJS.ProcessEnv {
  59. const env: NodeJS.ProcessEnv = {}
  60. for (const [key, value] of Object.entries(process.env)) {
  61. if (!SENSITIVE_ENV_PATTERN.test(key)) env[key] = value
  62. }
  63. return { ...env, ...ENV_OVERRIDES, ...extra }
  64. }
  65. /** What to run and under which limits (resolved — no defaults in here). */
  66. export interface SpawnSpec {
  67. command: string
  68. cwd: string
  69. /** Kill the process group after this many milliseconds. 0 = no timeout. */
  70. timeoutMs: number
  71. /** Per-stream in-memory cap; overflow spills to disk (tail kept in memory). */
  72. maxOutputBytes: number
  73. /** Grace period between the SIGTERM and the SIGKILL escalation on a kill. */
  74. graceMs: number
  75. /** Abort signal — kills the process group when fired. */
  76. signal?: AbortSignal | undefined
  77. /**
  78. * Bytes to write to the child's stdin, then close it. Absent (or empty)
  79. * leaves stdin closed/empty. Set by in-process plugins (the hooks bridges);
  80. * the model-facing `dsh-tool-bash` tool does not thread model input here.
  81. */
  82. stdin?: string | undefined
  83. /**
  84. * Extra environment entries, merged onto the scrubbed env AFTER the
  85. * credential scrub and the model-friendly overrides (so an explicit entry
  86. * wins). Set by in-process plugins; the model-facing tool does not forward
  87. * model input here.
  88. */
  89. env?: Record<string, string> | undefined
  90. }
  91. /** Raw outcome of one closed process (before result shaping). */
  92. export interface SpawnOutcome {
  93. exitCode: number | null
  94. signal: NodeJS.Signals | null
  95. timedOut: boolean
  96. aborted: boolean
  97. stdout: CollectedOutput
  98. stderr: CollectedOutput
  99. }
  100. /** Injectable knobs so tests can exercise spill behavior without the OS tmpdir. */
  101. export interface RunInternals {
  102. /** Directory for spill files (defaults to the OS temp dir). */
  103. spillDir?: string
  104. }
  105. /** Default SIGTERM→SIGKILL grace period (the `graceMs` config; matches OpenCode's 3s). */
  106. export const DEFAULT_GRACE_MS = 3_000
  107. let spillCounter = 0
  108. let defaultSpillDir: string | undefined
  109. /**
  110. * The default spill location: a private (0700) per-process directory under
  111. * the OS tmpdir, created lazily. Predictable world-readable paths would let
  112. * other local users read command output or pre-create symlinks.
  113. */
  114. function privateSpillDir(): string {
  115. defaultSpillDir ??= mkdtempSync(join(tmpdir(), 'dsh-bash-'))
  116. return defaultSpillDir
  117. }
  118. /**
  119. * Collects one stream with a bounded in-memory tail. The FULL stream is
  120. * always recoverable: on first overflow a spill file is created and every
  121. * chunk (including those already collected) is appended there.
  122. *
  123. * Tail-keep rationale (pi/OpenCode): errors and final results cluster at the
  124. * end of command output; the spill file covers the head.
  125. */
  126. export class OutputCollector {
  127. private chunks: Buffer[] = []
  128. private bytes = 0
  129. private dropped = false
  130. private spillFd: number | undefined
  131. private spillFile: string | undefined
  132. /** Total bytes ever pushed (not just retained). */
  133. private total = 0
  134. constructor(
  135. private readonly maxBytes: number,
  136. private readonly label: string,
  137. private readonly spillDir: string,
  138. ) {}
  139. /**
  140. * Ingest one stream chunk, counting it toward the whole-stream total. On
  141. * first overflow of the in-memory cap a spill file is opened and every chunk
  142. * (already-collected ones included) is appended there from then on; the
  143. * in-memory tail then drops whole chunks from its head (or the head of a
  144. * single over-cap chunk) until it fits the cap again.
  145. * @param chunk - the raw bytes from one stream 'data' event.
  146. */
  147. push(chunk: Buffer): void {
  148. this.total += chunk.length
  149. const overflows = this.bytes + chunk.length > this.maxBytes
  150. if (overflows || this.spillFd !== undefined) this.spillAll(chunk)
  151. this.chunks.push(chunk)
  152. this.bytes += chunk.length
  153. while (this.bytes > this.maxBytes && this.chunks.length > 1) {
  154. // Drop whole chunks from the head; pipe chunks are small (≤64KiB), so
  155. // the retained tail tracks the cap closely enough for a model-facing
  156. // truncation boundary. (length > 1 was just checked — shift() returns.)
  157. const head = this.chunks.shift() as Buffer
  158. this.bytes -= head.length
  159. this.dropped = true
  160. }
  161. if (this.bytes > this.maxBytes && this.chunks.length === 1) {
  162. // A single chunk larger than the cap: keep its tail.
  163. const only = this.chunks[0] as Buffer
  164. this.chunks[0] = only.subarray(only.length - this.maxBytes)
  165. this.bytes = this.maxBytes
  166. this.dropped = true
  167. }
  168. }
  169. /** Open the spill file lazily and append `chunk` (and any prior chunks once). */
  170. private spillAll(chunk: Buffer): void {
  171. if (this.spillFd === undefined) {
  172. // Random suffix + O_EXCL + no-follow-equivalent ('wx' fails on any
  173. // existing path, symlink or not) + owner-only mode: defeats spill-path
  174. // prediction and symlink planting in shared tmp dirs.
  175. this.spillFile = join(
  176. this.spillDir,
  177. `dsh-bash-${process.pid}-${++spillCounter}-${randomBytes(6).toString('hex')}-${this.label}.log`,
  178. )
  179. this.spillFd = openSync(this.spillFile, 'wx', 0o600)
  180. for (const prior of this.chunks) writeSync(this.spillFd, prior)
  181. }
  182. writeSync(this.spillFd, chunk)
  183. }
  184. // TODO(snapshot-scope): `snapshot()` has one internal caller (`finalize()` at
  185. // the bottom of this file) and `totalBytes` is read only by a test. The live
  186. // background-poll path goes through `readFrom()`, so inline snapshot() into
  187. // finalize() and drop or privatize the totalBytes getter.
  188. /**
  189. * Read the collected tail without finalizing (the final-result snapshot).
  190. * @returns the retained tail text, the truncation flag, and the spill path when one was created.
  191. */
  192. snapshot(): CollectedOutput {
  193. return {
  194. text: Buffer.concat(this.chunks).toString('utf8'),
  195. truncated: this.dropped,
  196. ...this.spillFile !== undefined ? { spillPath: this.spillFile } : {},
  197. }
  198. }
  199. /** Total bytes ever pushed (including bytes dropped from memory). */
  200. get totalBytes(): number {
  201. return this.total
  202. }
  203. /**
  204. * Incremental read in whole-stream byte coordinates: returns everything
  205. * pushed since `fromByte`. When `fromByte` has already slid out of the
  206. * in-memory tail window, the read is `lossy` — it returns the whole
  207. * retained tail and the gap is only recoverable from the spill file.
  208. * @param fromByte - whole-stream offset to resume from (a prior read's `nextOffset`; 0 for the first read).
  209. * @returns the delta text, the offset for the next read, the `lossy` flag, and the spill path when one was created.
  210. */
  211. readFrom(fromByte: number): { text: string; nextOffset: number; lossy: boolean; spillPath?: string } {
  212. const windowStart = this.total - this.bytes
  213. const buffer = Buffer.concat(this.chunks)
  214. const lossy = fromByte < windowStart
  215. const slice = lossy ? buffer : buffer.subarray(fromByte - windowStart)
  216. return {
  217. text: slice.toString('utf8'),
  218. nextOffset: this.total,
  219. lossy,
  220. ...this.spillFile !== undefined ? { spillPath: this.spillFile } : {},
  221. }
  222. }
  223. /**
  224. * Close the spill file (if any) and return the final output. A failed close
  225. * (delayed writeback fault) stops advertising the spill path — the file may
  226. * be missing its tail — but still returns the in-memory result.
  227. * @returns the final collected output: tail text, truncation flag, and the spill path when intact.
  228. */
  229. finalize(): CollectedOutput {
  230. if (this.spillFd !== undefined) {
  231. try {
  232. closeSync(this.spillFd)
  233. } catch {
  234. // close can surface delayed writeback failures (for example EIO/ENOSPC)
  235. // after writeSync appeared to succeed. Keep finalize total so runBash's
  236. // close handler still resolves, but stop advertising a spill file that
  237. // may be missing its tail.
  238. this.spillFile = undefined
  239. }
  240. this.spillFd = undefined
  241. }
  242. return this.snapshot()
  243. }
  244. }
  245. /**
  246. * Send `sig` to the process GROUP led by `pid` (requires the child to have
  247. * been spawned with `detached: true`). NEVER throws: kills race process exit
  248. * by design (ESRCH), and the other failure modes (EPERM from setuid
  249. * children, …) fire inside timer callbacks where a throw would crash the
  250. * host process — a kill that cannot be delivered is reported by the process
  251. * NOT dying, which callers already handle via escalation/timeouts. No-op for
  252. * non-positive pids (spawn never started a process).
  253. * @param pid - the group leader's pid; non-positive means the spawn failed and the call is a no-op.
  254. * @param sig - the signal to deliver to the whole group.
  255. */
  256. export function killGroup(pid: number, sig: NodeJS.Signals): void {
  257. if (pid <= 0) return
  258. try {
  259. process.kill(-pid, sig)
  260. } catch {
  261. // Swallow: see contract above.
  262. }
  263. }
  264. /**
  265. * A live bash child process: the promise resolves when the process closes;
  266. * `kill()` starts the SIGTERM→grace→SIGKILL escalation on its group.
  267. */
  268. export interface RunningBash {
  269. /** Process id (group leader); -1 when the spawn itself failed. */
  270. readonly pid: number
  271. /** stdout/stderr collectors (live — background polling reads incrementally). */
  272. readonly stdout: OutputCollector
  273. readonly stderr: OutputCollector
  274. /** Resolves when the process closes; rejects only for spawn-level failures. */
  275. readonly done: Promise<SpawnOutcome>
  276. /** Begin SIGTERM→grace→SIGKILL on the process group. Idempotent. */
  277. kill(): void
  278. }
  279. /**
  280. * Spawn `bash -c <command>` in its own process group and collect output.
  281. *
  282. * Outcome semantics: the returned promise REJECTS only for spawn-level
  283. * failures (bad cwd → ENOENT, missing binary, pre-aborted signal); every
  284. * runtime outcome — nonzero exit, timeout kill, abort kill, signal death —
  285. * RESOLVES with a {@link SpawnOutcome} describing what happened, so callers
  286. * shape one consistent report for the model.
  287. *
  288. * XXX(stateful-shell): per the agent-tool survey there are two proven
  289. * stateful designs worth revisiting — Claude Code persists ONLY cwd between
  290. * calls (captures `pwd -P` after each command), and Codex keeps whole PTY
  291. * exec sessions addressable via session ids + stdin writes. We deliberately
  292. * spawn a fresh non-login `bash -c` per call for determinism (no rc files,
  293. * no inherited shell state); revisit when real workflows demand it.
  294. * @param spec - the fully-resolved run (command, cwd, limits); no defaulting happens here.
  295. * @param internals - test-only knobs; omitted fields fall back to the private per-process spill dir.
  296. * @returns the live handle: pid, the two live collectors, the outcome promise, and `kill()`.
  297. */
  298. export function runBash(spec: SpawnSpec, internals: RunInternals = {}): RunningBash {
  299. const spillDir = internals.spillDir ?? privateSpillDir()
  300. if (spec.signal?.aborted) {
  301. throw new Error(`aborted before spawn: ${String(spec.signal.reason ?? 'aborted')}`)
  302. }
  303. // stdin is a pipe ONLY when the caller supplied bytes; with none it is `ignore`
  304. // (fd 0 → /dev/null) — the exact pre-seam default. This matters: a spawn pipe
  305. // and /dev/null are NOT observationally identical (node's pipe is an AF_UNIX
  306. // socket, so a command that probes stdin's type — `test -c /dev/stdin`, `stat
  307. // /proc/self/fd/0` — sees a char device vs a socket), so the no-stdin path
  308. // (every model-driven call) must keep /dev/null rather than regress to a socket.
  309. // Two LITERAL `stdio` tuples (not one variable tuple): only a literal lets the
  310. // typed `spawn` overload infer non-null stdout/stderr, which the
  311. // `ChildProcessByStdio` annotation captures (stdin `Writable | null`; stdout/
  312. // stderr the non-null `Readable` the collectors attach to without a cast).
  313. const env = childEnv(spec.env)
  314. const child: ChildProcessByStdio<Writable | null, Readable, Readable> = spec.stdin !== undefined
  315. ? spawn('bash', ['-c', spec.command], { cwd: spec.cwd, env, stdio: ['pipe', 'pipe', 'pipe'], detached: true })
  316. : spawn('bash', ['-c', spec.command], { cwd: spec.cwd, env, stdio: ['ignore', 'pipe', 'pipe'], detached: true })
  317. const stdout = new OutputCollector(spec.maxOutputBytes, 'stdout', spillDir)
  318. const stderr = new OutputCollector(spec.maxOutputBytes, 'stderr', spillDir)
  319. child.stdout.on('data', (chunk: Buffer) => { stdout.push(chunk) })
  320. child.stderr.on('data', (chunk: Buffer) => { stderr.push(chunk) })
  321. let timedOut = false
  322. let aborted = false
  323. let killTimer: NodeJS.Timeout | undefined
  324. let graceTimer: NodeJS.Timeout | undefined
  325. // pid is undefined when the spawn itself fails (bad cwd, missing binary);
  326. // the 'error' handler rejects `done` and kills become no-ops via pid -1.
  327. const pid = child.pid ?? -1
  328. const kill = (): void => {
  329. if (graceTimer !== undefined) return // escalation already in flight
  330. killGroup(pid, 'SIGTERM')
  331. graceTimer = setTimeout(() => { killGroup(pid, 'SIGKILL') }, spec.graceMs)
  332. }
  333. if (spec.timeoutMs > 0) {
  334. killTimer = setTimeout(() => {
  335. timedOut = true
  336. kill()
  337. }, spec.timeoutMs)
  338. }
  339. const onAbort = (): void => {
  340. aborted = true
  341. kill()
  342. }
  343. spec.signal?.addEventListener('abort', onAbort, { once: true })
  344. // Write stdin and close it, but ONLY when the caller supplied bytes — with no
  345. // stdin, fd 0 is `ignore` (/dev/null) and `child.stdin` is null. The error
  346. // handler must exist whenever we write: an unhandled 'error' on the stream
  347. // would throw and crash the host. We swallow the error rather than reject
  348. // `done`, and that is correct for ANY stdin-write error, not just the common
  349. // one — the stdin write is BEST-EFFORT, while the command's authoritative
  350. // outcome is its exit code + captured output, which the `close` handler reports
  351. // regardless of whether the write landed. The expected case is EPIPE (the child
  352. // exited without reading, so closing our end of a still-full pipe fails); a rare
  353. // non-EPIPE pipe fault means the command ran with incomplete stdin, and it
  354. // surfaces that itself through its own exit/output (e.g. a hook that gets
  355. // truncated JSON errors out) — rejecting here would instead discard that real
  356. // output and turn it into an opaque infrastructure error, which is worse.
  357. if (child.stdin !== null) {
  358. child.stdin.on('error', () => { /* stdin write is best-effort; outcome rides on exit/output. */ })
  359. child.stdin.end(spec.stdin)
  360. }
  361. const done = new Promise<SpawnOutcome>((resolve, reject) => {
  362. child.on('error', (error) => {
  363. // Spawn-level failure (ENOENT cwd, EACCES, …): no close event with
  364. // meaningful output follows; clean up and reject.
  365. cleanup()
  366. reject(error)
  367. })
  368. child.on('close', (exitCode, signal) => {
  369. cleanup()
  370. resolve({
  371. exitCode,
  372. signal,
  373. timedOut,
  374. aborted,
  375. stdout: stdout.finalize(),
  376. stderr: stderr.finalize(),
  377. })
  378. })
  379. function cleanup(): void {
  380. if (killTimer !== undefined) clearTimeout(killTimer)
  381. if (graceTimer !== undefined) clearTimeout(graceTimer)
  382. spec.signal?.removeEventListener('abort', onAbort)
  383. }
  384. })
  385. return { pid, stdout, stderr, done, kill }
  386. }