run.ts 15 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399
  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, unlinkSync, 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. /** Per-stream spill-file cap; larger streams retain only their in-memory tail. */
  71. maxSpillBytes: number
  72. /** Grace period for kill escalation and for inherited pipes after shell exit. */
  73. graceMs: number
  74. /**
  75. * Abort signal — kills the process group when it fires. The executor owns
  76. * timing: `run()` passes a fused timeout/cancel deadline signal (see
  77. * `@deepseek-ai/dsh-timeout`), `start()` passes the bare upstream signal.
  78. * runBash only listens and kills; it does NOT classify why (the executor
  79. * reads the signal's reason afterward).
  80. */
  81. signal?: AbortSignal | undefined
  82. /**
  83. * Bytes to write to the child's stdin, then close it. Absent (or empty)
  84. * leaves stdin closed/empty. Set by in-process plugins (the hooks bridges);
  85. * the model-facing `dsh-tool-bash` tool does not thread model input here.
  86. */
  87. stdin?: string | undefined
  88. /**
  89. * Ordinary environment entries merged after the credential scrub and
  90. * terminal overrides. `DSH_*` names are rejected and belong in `dshEnv`.
  91. */
  92. env?: Record<string, string> | undefined
  93. /** Harness-owned entries; non-`DSH_*` names are rejected before spawn. */
  94. dshEnv?: DshEnvironment | undefined
  95. }
  96. /**
  97. * Raw outcome of one closed process (before result shaping). Deliberately
  98. * carries NO timeout/cancel classification: runBash kills on abort but does not
  99. * decide why — the executor's `run()`/`start()` reads the deadline signal it
  100. * owns to classify `timedOut`/`aborted` (see the package README).
  101. */
  102. export interface SpawnOutcome {
  103. exitCode: number | null
  104. signal: NodeJS.Signals | null
  105. stdout: CollectedOutput
  106. stderr: CollectedOutput
  107. }
  108. /** Injectable knobs so tests can exercise spill behavior without the OS tmpdir. */
  109. export interface RunInternals {
  110. /** Directory for spill files (defaults to the OS temp dir). */
  111. spillDir?: string
  112. }
  113. /** Default SIGTERM→SIGKILL grace period (the `graceMs` config; matches OpenCode's 3s). */
  114. export const DEFAULT_GRACE_MS = 3_000
  115. /** Default per-stream spill cap (the `maxSpillBytes` config). */
  116. export const DEFAULT_MAX_SPILL_BYTES = 64 * 1024 * 1024
  117. let spillCounter = 0
  118. let defaultSpillDir: string | undefined
  119. /**
  120. * The default spill location: a private (0700) per-process directory under
  121. * the OS tmpdir, created lazily. Predictable world-readable paths would let
  122. * other local users read command output or pre-create symlinks.
  123. */
  124. function privateSpillDir(): string {
  125. defaultSpillDir ??= mkdtempSync(join(tmpdir(), 'dsh-bash-'))
  126. return defaultSpillDir
  127. }
  128. /**
  129. * Collects one stream with a bounded in-memory tail. On first overflow a
  130. * spill file is created and every chunk (including those already collected)
  131. * is appended there while the full stream remains within `maxSpillBytes`.
  132. *
  133. * Tail-keep rationale (pi/OpenCode): errors and final results cluster at the
  134. * end of command output; the spill file covers the head.
  135. */
  136. export class OutputCollector {
  137. private chunks: Buffer[] = []
  138. private bytes = 0
  139. private dropped = false
  140. private spillFd: number | undefined
  141. private spillFile: string | undefined
  142. private spillDisabled = false
  143. /** Total bytes ever pushed (not just retained). */
  144. private total = 0
  145. constructor(
  146. private readonly maxBytes: number,
  147. private readonly maxSpillBytes: number,
  148. private readonly label: string,
  149. private readonly spillDir: string,
  150. ) {}
  151. /**
  152. * Ingest one stream chunk, counting it toward the whole-stream total. On
  153. * first overflow of the in-memory cap a spill file is opened and every chunk
  154. * (already-collected ones included) is appended there from then on; the
  155. * in-memory tail then drops whole chunks from its head (or the head of a
  156. * single over-cap chunk) until it fits the cap again.
  157. * @param chunk - the raw bytes from one stream 'data' event.
  158. */
  159. push(chunk: Buffer): void {
  160. this.total += chunk.length
  161. const overflows = this.bytes + chunk.length > this.maxBytes
  162. if (!this.spillDisabled && (overflows || this.spillFd !== undefined)) this.spillAll(chunk)
  163. this.chunks.push(chunk)
  164. this.bytes += chunk.length
  165. while (this.bytes > this.maxBytes && this.chunks.length > 1) {
  166. // Drop whole chunks from the head; pipe chunks are small (≤64KiB), so
  167. // the retained tail tracks the cap closely enough for a model-facing
  168. // truncation boundary. (length > 1 was just checked — shift() returns.)
  169. const head = this.chunks.shift() as Buffer
  170. this.bytes -= head.length
  171. this.dropped = true
  172. }
  173. if (this.bytes > this.maxBytes && this.chunks.length === 1) {
  174. // A single chunk larger than the cap: keep its tail.
  175. const only = this.chunks[0] as Buffer
  176. this.chunks[0] = only.subarray(only.length - this.maxBytes)
  177. this.bytes = this.maxBytes
  178. this.dropped = true
  179. }
  180. }
  181. /** Open the spill file lazily and append `chunk` (and any prior chunks once). */
  182. private spillAll(chunk: Buffer): void {
  183. if (this.total > this.maxSpillBytes) {
  184. this.discardSpill()
  185. return
  186. }
  187. if (this.spillFd === undefined) {
  188. // Random suffix + O_EXCL + no-follow-equivalent ('wx' fails on any
  189. // existing path, symlink or not) + owner-only mode: defeats spill-path
  190. // prediction and symlink planting in shared tmp dirs.
  191. this.spillFile = join(
  192. this.spillDir,
  193. `dsh-bash-${process.pid}-${++spillCounter}-${randomBytes(6).toString('hex')}-${this.label}.log`,
  194. )
  195. this.spillFd = openSync(this.spillFile, 'wx', 0o600)
  196. for (const prior of this.chunks) writeSync(this.spillFd, prior)
  197. }
  198. writeSync(this.spillFd, chunk)
  199. }
  200. /** Stop spilling and remove the file once it can no longer hold the complete stream. */
  201. private discardSpill(): void {
  202. const fd = this.spillFd
  203. const file = this.spillFile
  204. this.spillFd = undefined
  205. this.spillFile = undefined
  206. this.spillDisabled = true
  207. if (fd !== undefined) {
  208. try {
  209. closeSync(fd)
  210. } catch {
  211. // Retain the descriptor so finalize can retry the failed close.
  212. this.spillFd = fd
  213. }
  214. }
  215. if (file !== undefined) {
  216. try {
  217. unlinkSync(file)
  218. } catch {
  219. // A failed unlink leaves at most maxSpillBytes behind, never an unbounded file.
  220. }
  221. }
  222. }
  223. /**
  224. * Incremental read in whole-stream byte coordinates: returns everything
  225. * pushed since `fromByte`. When `fromByte` has already slid out of the
  226. * in-memory tail window, the read is `lossy` — it returns the whole
  227. * retained tail and the gap is only recoverable from the spill file.
  228. * @param fromByte - whole-stream offset to resume from (a prior read's `nextOffset`; 0 for the first read).
  229. * @returns the delta text, the offset for the next read, the `lossy` flag, and the spill path when one was created.
  230. */
  231. readFrom(fromByte: number): { text: string; nextOffset: number; lossy: boolean; spillPath?: string } {
  232. const windowStart = this.total - this.bytes
  233. const buffer = Buffer.concat(this.chunks)
  234. const lossy = fromByte < windowStart
  235. const slice = lossy ? buffer : buffer.subarray(fromByte - windowStart)
  236. return {
  237. text: slice.toString('utf8'),
  238. nextOffset: this.total,
  239. lossy,
  240. ...this.spillFile !== undefined ? { spillPath: this.spillFile } : {},
  241. }
  242. }
  243. /**
  244. * Close the spill file (if any) and return the final output. A failed close
  245. * (delayed writeback fault) stops advertising the spill path — the file may
  246. * be missing its tail — but still returns the in-memory result.
  247. * @returns the final collected output: tail text, truncation flag, and the spill path when intact.
  248. */
  249. finalize(): CollectedOutput {
  250. if (this.spillFd !== undefined) {
  251. try {
  252. closeSync(this.spillFd)
  253. } catch {
  254. // A delayed writeback failure makes the spill unreliable; keep finalize
  255. // total but stop advertising that file.
  256. this.spillFile = undefined
  257. }
  258. this.spillFd = undefined
  259. }
  260. return {
  261. text: Buffer.concat(this.chunks).toString('utf8'),
  262. truncated: this.dropped,
  263. ...this.spillFile !== undefined ? { spillPath: this.spillFile } : {},
  264. }
  265. }
  266. }
  267. /**
  268. * Send `sig` to a detached process group. Never throws: delivery races process
  269. * exit and may run in a timer callback, so failures are contained and a
  270. * non-positive pid is a no-op.
  271. * @param pid - the group leader's pid; non-positive means the spawn failed and the call is a no-op.
  272. * @param sig - the signal to deliver to the whole group.
  273. */
  274. export function killGroup(pid: number, sig: NodeJS.Signals): void {
  275. if (pid <= 0) return
  276. try {
  277. process.kill(-pid, sig)
  278. } catch {
  279. // Swallow: see contract above.
  280. }
  281. }
  282. /**
  283. * A live bash child process: the promise resolves when the process closes;
  284. * `kill()` starts the SIGTERM→grace→SIGKILL escalation on its group.
  285. */
  286. export interface RunningBash {
  287. /** Process id (group leader); -1 when the spawn itself failed. */
  288. readonly pid: number
  289. /** stdout/stderr collectors (live — background polling reads incrementally). */
  290. readonly stdout: OutputCollector
  291. readonly stderr: OutputCollector
  292. /** Resolves when the process closes; rejects only for spawn-level failures. */
  293. readonly done: Promise<SpawnOutcome>
  294. /** Begin SIGTERM→grace→SIGKILL on the process group. Idempotent. */
  295. kill(): void
  296. }
  297. /**
  298. * Spawn one isolated `bash -c` process group and collect its output.
  299. * Runtime exits resolve as {@link SpawnOutcome}; only spawn failures reject.
  300. * @param spec - fully resolved command, cwd, limits, and cancellation.
  301. * @param internals - test-only process and spill-directory overrides.
  302. * @returns live process handle and outcome promise.
  303. */
  304. // XXX(stateful-shell): evaluate persistent cwd or PTY sessions when workflows require shell state.
  305. export function runBash(spec: SpawnSpec, internals: RunInternals = {}): RunningBash {
  306. const spillDir = internals.spillDir ?? privateSpillDir()
  307. if (spec.signal?.aborted) {
  308. throw new Error(`aborted before spawn: ${String(spec.signal.reason ?? 'aborted')}`)
  309. }
  310. // Keep absent stdin as /dev/null; literal tuples preserve non-null output types.
  311. const env = childEnv(spec.env, spec.dshEnv)
  312. const child: ChildProcessByStdio<Writable | null, Readable, Readable> = spec.stdin !== undefined
  313. ? spawn('bash', ['-c', spec.command], { cwd: spec.cwd, env, stdio: ['pipe', 'pipe', 'pipe'], detached: true })
  314. : spawn('bash', ['-c', spec.command], { cwd: spec.cwd, env, stdio: ['ignore', 'pipe', 'pipe'], detached: true })
  315. const stdout = new OutputCollector(spec.stdoutMaxBytes, spec.maxSpillBytes, 'stdout', spillDir)
  316. const stderr = new OutputCollector(spec.stderrMaxBytes, spec.maxSpillBytes, 'stderr', spillDir)
  317. child.stdout.on('data', (chunk: Buffer) => { stdout.push(chunk) })
  318. child.stderr.on('data', (chunk: Buffer) => { stderr.push(chunk) })
  319. let graceTimer: NodeJS.Timeout | undefined
  320. // Failed spawns use pid -1 so kill remains a no-op.
  321. const pid = child.pid ?? -1
  322. const kill = (): void => {
  323. if (graceTimer !== undefined) return // escalation already in flight
  324. killGroup(pid, 'SIGTERM')
  325. graceTimer = setTimeout(() => { killGroup(pid, 'SIGKILL') }, spec.graceMs)
  326. }
  327. // The executor owns timeout classification; this layer only reacts to abort.
  328. const onAbort = (): void => { kill() }
  329. spec.signal?.addEventListener('abort', onAbort, { once: true })
  330. // Stdin writes are best-effort; process exit and captured output remain authoritative.
  331. if (child.stdin !== null) {
  332. child.stdin.on('error', () => { /* stdin write is best-effort; outcome rides on exit/output. */ })
  333. child.stdin.end(spec.stdin)
  334. }
  335. const done = new Promise<SpawnOutcome>((resolve, reject) => {
  336. let settled = false
  337. let pipeDrainTimer: NodeJS.Timeout | undefined
  338. const settle = (exitCode: number | null, signal: NodeJS.Signals | null): void => {
  339. if (settled) return
  340. settled = true
  341. child.stdout.destroy()
  342. child.stderr.destroy()
  343. cleanup()
  344. resolve({
  345. exitCode,
  346. signal,
  347. stdout: stdout.finalize(),
  348. stderr: stderr.finalize(),
  349. })
  350. }
  351. child.on('error', (error) => {
  352. // No meaningful close outcome follows a spawn failure.
  353. settled = true
  354. cleanup()
  355. reject(error)
  356. })
  357. child.on('exit', (exitCode, signal) => {
  358. pipeDrainTimer = setTimeout(() => { settle(exitCode, signal) }, spec.graceMs)
  359. })
  360. child.on('close', settle)
  361. function cleanup(): void {
  362. if (graceTimer !== undefined) clearTimeout(graceTimer)
  363. if (pipeDrainTimer !== undefined) clearTimeout(pipeDrainTimer)
  364. spec.signal?.removeEventListener('abort', onAbort)
  365. }
  366. })
  367. return { pid, stdout, stderr, done, kill }
  368. }