run.ts 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352
  1. /**
  2. * Process plumbing for the local bash executor: detached process-group spawn,
  3. * tail-keep output with spill files, and SIGTERM→SIGKILL escalation. This layer
  4. * reacts to an abort signal; the executor owns deadlines and classifies causes.
  5. * @module dsh-bash-local/run
  6. */
  7. import { type ChildProcessByStdio, spawn } from 'node:child_process'
  8. import type { Readable, Writable } from 'node:stream'
  9. import { randomBytes } from 'node:crypto'
  10. import { closeSync, mkdtempSync, openSync, writeSync } from 'node:fs'
  11. import { tmpdir } from 'node:os'
  12. import { join } from 'node:path'
  13. import { DSH_ENV_PREFIX } from '@deepseek-ai/dsh-bash'
  14. import type { CollectedOutput, DshEnvironment } from '@deepseek-ai/dsh-bash'
  15. /**
  16. * Model-friendly environment overrides: disable colors, pagers, and
  17. * interactive terminal features that would garble tool output (the same set
  18. * Codex hardcodes; Claude Code achieves it via TERM=dumb).
  19. */
  20. export const ENV_OVERRIDES = {
  21. NO_COLOR: '1',
  22. TERM: 'dumb',
  23. PAGER: 'cat',
  24. GIT_PAGER: 'cat',
  25. } as const
  26. /**
  27. * Credential-shaped env vars are NOT forwarded to commands (the harness's
  28. * own DEEPSEEK_API_KEY must not leak into `env` output, tool results, or
  29. * spill files). Same default pattern as Codex's env policy; a future config
  30. * can whitelist specific vars when a workflow genuinely needs one.
  31. */
  32. export const SENSITIVE_ENV_PATTERN = /KEY|SECRET|TOKEN/i
  33. /**
  34. * Build a child environment from scrubbed ambient values, terminal overrides,
  35. * ordinary caller entries, and a managed `DSH_*` snapshot. Ambient managed
  36. * names are removed; ordinary and managed entries reject the other channel's
  37. * namespace before `dshEnv` merges last.
  38. * @param extra - caller entries; `DSH_*` names are rejected.
  39. * @param dshEnv - managed entries; non-`DSH_*` names are rejected.
  40. * @returns the environment to hand to `spawn` for the child process.
  41. */
  42. export function childEnv(
  43. extra?: Readonly<Record<string, string>>,
  44. dshEnv?: DshEnvironment,
  45. ): NodeJS.ProcessEnv {
  46. const env: NodeJS.ProcessEnv = {}
  47. for (const [key, value] of Object.entries(process.env)) {
  48. if (!SENSITIVE_ENV_PATTERN.test(key) && !key.startsWith(DSH_ENV_PREFIX)) env[key] = value
  49. }
  50. for (const key of Object.keys(extra ?? {})) {
  51. if (key.startsWith(DSH_ENV_PREFIX)) {
  52. throw new Error(`ordinary bash env cannot set reserved variable "${key}"; use dshEnv`)
  53. }
  54. }
  55. for (const key of Object.keys(dshEnv ?? {})) {
  56. if (!key.startsWith(DSH_ENV_PREFIX)) {
  57. throw new Error(`managed bash env cannot set ordinary variable "${key}"; use env`)
  58. }
  59. }
  60. return { ...env, ...ENV_OVERRIDES, ...extra, ...dshEnv }
  61. }
  62. /** What to run and under which limits (resolved — no defaults in here). */
  63. export interface SpawnSpec {
  64. command: string
  65. cwd: string
  66. /** Stdout in-memory cap; overflow spills to disk (tail kept in memory). */
  67. stdoutMaxBytes: number
  68. /** Stderr in-memory cap; overflow spills to disk (tail kept in memory). */
  69. stderrMaxBytes: number
  70. /** Grace period between the SIGTERM and the SIGKILL escalation on a kill. */
  71. graceMs: number
  72. /**
  73. * Abort signal — kills the process group when it fires. The executor owns
  74. * timing: `run()` passes a fused timeout/cancel deadline signal (see
  75. * `@deepseek-ai/dsh-timeout`), `start()` passes the bare upstream signal.
  76. * runBash only listens and kills; it does NOT classify why (the executor
  77. * reads the signal's reason afterward).
  78. */
  79. signal?: AbortSignal | undefined
  80. /**
  81. * Bytes to write to the child's stdin, then close it. Absent (or empty)
  82. * leaves stdin closed/empty. Set by in-process plugins (the hooks bridges);
  83. * the model-facing `dsh-tool-bash` tool does not thread model input here.
  84. */
  85. stdin?: string | undefined
  86. /**
  87. * Ordinary environment entries merged after the credential scrub and
  88. * terminal overrides. `DSH_*` names are rejected and belong in `dshEnv`.
  89. */
  90. env?: Record<string, string> | undefined
  91. /** Harness-owned entries; non-`DSH_*` names are rejected before spawn. */
  92. dshEnv?: DshEnvironment | undefined
  93. }
  94. /**
  95. * Raw outcome of one closed process (before result shaping). Deliberately
  96. * carries NO timeout/cancel classification: runBash kills on abort but does not
  97. * decide why — the executor's `run()`/`start()` reads the deadline signal it
  98. * owns to classify `timedOut`/`aborted` (see the package README).
  99. */
  100. export interface SpawnOutcome {
  101. exitCode: number | null
  102. signal: NodeJS.Signals | null
  103. stdout: CollectedOutput
  104. stderr: CollectedOutput
  105. }
  106. /** Injectable knobs so tests can exercise spill behavior without the OS tmpdir. */
  107. export interface RunInternals {
  108. /** Directory for spill files (defaults to the OS temp dir). */
  109. spillDir?: string
  110. }
  111. /** Default SIGTERM→SIGKILL grace period (the `graceMs` config; matches OpenCode's 3s). */
  112. export const DEFAULT_GRACE_MS = 3_000
  113. let spillCounter = 0
  114. let defaultSpillDir: string | undefined
  115. /**
  116. * The default spill location: a private (0700) per-process directory under
  117. * the OS tmpdir, created lazily. Predictable world-readable paths would let
  118. * other local users read command output or pre-create symlinks.
  119. */
  120. function privateSpillDir(): string {
  121. defaultSpillDir ??= mkdtempSync(join(tmpdir(), 'dsh-bash-'))
  122. return defaultSpillDir
  123. }
  124. /**
  125. * Collects one stream with a bounded in-memory tail. The FULL stream is
  126. * always recoverable: on first overflow a spill file is created and every
  127. * chunk (including those already collected) is appended there.
  128. *
  129. * Tail-keep rationale (pi/OpenCode): errors and final results cluster at the
  130. * end of command output; the spill file covers the head.
  131. */
  132. export class OutputCollector {
  133. private chunks: Buffer[] = []
  134. private bytes = 0
  135. private dropped = false
  136. private spillFd: number | undefined
  137. private spillFile: string | undefined
  138. /** Total bytes ever pushed (not just retained). */
  139. private total = 0
  140. constructor(
  141. private readonly maxBytes: number,
  142. private readonly label: string,
  143. private readonly spillDir: string,
  144. ) {}
  145. /**
  146. * Ingest one stream chunk, counting it toward the whole-stream total. On
  147. * first overflow of the in-memory cap a spill file is opened and every chunk
  148. * (already-collected ones included) is appended there from then on; the
  149. * in-memory tail then drops whole chunks from its head (or the head of a
  150. * single over-cap chunk) until it fits the cap again.
  151. * @param chunk - the raw bytes from one stream 'data' event.
  152. */
  153. push(chunk: Buffer): void {
  154. this.total += chunk.length
  155. const overflows = this.bytes + chunk.length > this.maxBytes
  156. if (overflows || this.spillFd !== undefined) this.spillAll(chunk)
  157. this.chunks.push(chunk)
  158. this.bytes += chunk.length
  159. while (this.bytes > this.maxBytes && this.chunks.length > 1) {
  160. // Drop whole chunks from the head; pipe chunks are small (≤64KiB), so
  161. // the retained tail tracks the cap closely enough for a model-facing
  162. // truncation boundary. (length > 1 was just checked — shift() returns.)
  163. const head = this.chunks.shift() as Buffer
  164. this.bytes -= head.length
  165. this.dropped = true
  166. }
  167. if (this.bytes > this.maxBytes && this.chunks.length === 1) {
  168. // A single chunk larger than the cap: keep its tail.
  169. const only = this.chunks[0] as Buffer
  170. this.chunks[0] = only.subarray(only.length - this.maxBytes)
  171. this.bytes = this.maxBytes
  172. this.dropped = true
  173. }
  174. }
  175. /** Open the spill file lazily and append `chunk` (and any prior chunks once). */
  176. private spillAll(chunk: Buffer): void {
  177. if (this.spillFd === undefined) {
  178. // Random suffix + O_EXCL + no-follow-equivalent ('wx' fails on any
  179. // existing path, symlink or not) + owner-only mode: defeats spill-path
  180. // prediction and symlink planting in shared tmp dirs.
  181. this.spillFile = join(
  182. this.spillDir,
  183. `dsh-bash-${process.pid}-${++spillCounter}-${randomBytes(6).toString('hex')}-${this.label}.log`,
  184. )
  185. this.spillFd = openSync(this.spillFile, 'wx', 0o600)
  186. for (const prior of this.chunks) writeSync(this.spillFd, prior)
  187. }
  188. writeSync(this.spillFd, chunk)
  189. }
  190. /**
  191. * Incremental read in whole-stream byte coordinates: returns everything
  192. * pushed since `fromByte`. When `fromByte` has already slid out of the
  193. * in-memory tail window, the read is `lossy` — it returns the whole
  194. * retained tail and the gap is only recoverable from the spill file.
  195. * @param fromByte - whole-stream offset to resume from (a prior read's `nextOffset`; 0 for the first read).
  196. * @returns the delta text, the offset for the next read, the `lossy` flag, and the spill path when one was created.
  197. */
  198. readFrom(fromByte: number): { text: string; nextOffset: number; lossy: boolean; spillPath?: string } {
  199. const windowStart = this.total - this.bytes
  200. const buffer = Buffer.concat(this.chunks)
  201. const lossy = fromByte < windowStart
  202. const slice = lossy ? buffer : buffer.subarray(fromByte - windowStart)
  203. return {
  204. text: slice.toString('utf8'),
  205. nextOffset: this.total,
  206. lossy,
  207. ...this.spillFile !== undefined ? { spillPath: this.spillFile } : {},
  208. }
  209. }
  210. /**
  211. * Close the spill file (if any) and return the final output. A failed close
  212. * (delayed writeback fault) stops advertising the spill path — the file may
  213. * be missing its tail — but still returns the in-memory result.
  214. * @returns the final collected output: tail text, truncation flag, and the spill path when intact.
  215. */
  216. finalize(): CollectedOutput {
  217. if (this.spillFd !== undefined) {
  218. try {
  219. closeSync(this.spillFd)
  220. } catch {
  221. // A delayed writeback failure makes the spill unreliable; keep finalize
  222. // total but stop advertising that file.
  223. this.spillFile = undefined
  224. }
  225. this.spillFd = undefined
  226. }
  227. return {
  228. text: Buffer.concat(this.chunks).toString('utf8'),
  229. truncated: this.dropped,
  230. ...this.spillFile !== undefined ? { spillPath: this.spillFile } : {},
  231. }
  232. }
  233. }
  234. /**
  235. * Send `sig` to a detached process group. Never throws: delivery races process
  236. * exit and may run in a timer callback, so failures are contained and a
  237. * non-positive pid is a no-op.
  238. * @param pid - the group leader's pid; non-positive means the spawn failed and the call is a no-op.
  239. * @param sig - the signal to deliver to the whole group.
  240. */
  241. export function killGroup(pid: number, sig: NodeJS.Signals): void {
  242. if (pid <= 0) return
  243. try {
  244. process.kill(-pid, sig)
  245. } catch {
  246. // Swallow: see contract above.
  247. }
  248. }
  249. /**
  250. * A live bash child process: the promise resolves when the process closes;
  251. * `kill()` starts the SIGTERM→grace→SIGKILL escalation on its group.
  252. */
  253. export interface RunningBash {
  254. /** Process id (group leader); -1 when the spawn itself failed. */
  255. readonly pid: number
  256. /** stdout/stderr collectors (live — background polling reads incrementally). */
  257. readonly stdout: OutputCollector
  258. readonly stderr: OutputCollector
  259. /** Resolves when the process closes; rejects only for spawn-level failures. */
  260. readonly done: Promise<SpawnOutcome>
  261. /** Begin SIGTERM→grace→SIGKILL on the process group. Idempotent. */
  262. kill(): void
  263. }
  264. /**
  265. * Spawn one isolated `bash -c` process group and collect its output.
  266. * Runtime exits resolve as {@link SpawnOutcome}; only spawn failures reject.
  267. * @param spec - fully resolved command, cwd, limits, and cancellation.
  268. * @param internals - test-only process and spill-directory overrides.
  269. * @returns live process handle and outcome promise.
  270. */
  271. // XXX(stateful-shell): evaluate persistent cwd or PTY sessions when workflows require shell state.
  272. export function runBash(spec: SpawnSpec, internals: RunInternals = {}): RunningBash {
  273. const spillDir = internals.spillDir ?? privateSpillDir()
  274. if (spec.signal?.aborted) {
  275. throw new Error(`aborted before spawn: ${String(spec.signal.reason ?? 'aborted')}`)
  276. }
  277. // Keep absent stdin as /dev/null; literal tuples preserve non-null output types.
  278. const env = childEnv(spec.env, spec.dshEnv)
  279. const child: ChildProcessByStdio<Writable | null, Readable, Readable> = spec.stdin !== undefined
  280. ? spawn('bash', ['-c', spec.command], { cwd: spec.cwd, env, stdio: ['pipe', 'pipe', 'pipe'], detached: true })
  281. : spawn('bash', ['-c', spec.command], { cwd: spec.cwd, env, stdio: ['ignore', 'pipe', 'pipe'], detached: true })
  282. const stdout = new OutputCollector(spec.stdoutMaxBytes, 'stdout', spillDir)
  283. const stderr = new OutputCollector(spec.stderrMaxBytes, 'stderr', spillDir)
  284. child.stdout.on('data', (chunk: Buffer) => { stdout.push(chunk) })
  285. child.stderr.on('data', (chunk: Buffer) => { stderr.push(chunk) })
  286. let graceTimer: NodeJS.Timeout | undefined
  287. // Failed spawns use pid -1 so kill remains a no-op.
  288. const pid = child.pid ?? -1
  289. const kill = (): void => {
  290. if (graceTimer !== undefined) return // escalation already in flight
  291. killGroup(pid, 'SIGTERM')
  292. graceTimer = setTimeout(() => { killGroup(pid, 'SIGKILL') }, spec.graceMs)
  293. }
  294. // The executor owns timeout classification; this layer only reacts to abort.
  295. const onAbort = (): void => { kill() }
  296. spec.signal?.addEventListener('abort', onAbort, { once: true })
  297. // Stdin writes are best-effort; process exit and captured output remain authoritative.
  298. if (child.stdin !== null) {
  299. child.stdin.on('error', () => { /* stdin write is best-effort; outcome rides on exit/output. */ })
  300. child.stdin.end(spec.stdin)
  301. }
  302. const done = new Promise<SpawnOutcome>((resolve, reject) => {
  303. child.on('error', (error) => {
  304. // No meaningful close outcome follows a spawn failure.
  305. cleanup()
  306. reject(error)
  307. })
  308. child.on('close', (exitCode, signal) => {
  309. cleanup()
  310. resolve({
  311. exitCode,
  312. signal,
  313. stdout: stdout.finalize(),
  314. stderr: stderr.finalize(),
  315. })
  316. })
  317. function cleanup(): void {
  318. if (graceTimer !== undefined) clearTimeout(graceTimer)
  319. spec.signal?.removeEventListener('abort', onAbort)
  320. }
  321. })
  322. return { pid, stdout, stderr, done, kill }
  323. }