executor.spec.ts 18 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421
  1. import { mkdtempSync } from 'node:fs'
  2. import { tmpdir } from 'node:os'
  3. import { join } from 'node:path'
  4. import { describe, expect, it } from 'vitest'
  5. import { Context } from '@deepseek-ai/cordis'
  6. import { LocalBashExecutor } from '@deepseek-ai/dsh-bash-local'
  7. import SubprocessRuntime from '@deepseek-ai/dsh-subprocess'
  8. import LocalSubprocessRuntime from '@deepseek-ai/dsh-subprocess-local'
  9. import { MAX_TIMER_DELAY_MS } from '@deepseek-ai/dsh-timeout'
  10. import type { ShellProcess } from '@deepseek-ai/dsh-shell'
  11. import type { SubprocessHandle, SubprocessOutputReader, SubprocessSpawnSpec } from '@deepseek-ai/dsh-subprocess'
  12. const spillDir = mkdtempSync(join(tmpdir(), 'dsh-bash-exec-spec-'))
  13. async function setup(config: ConstructorParameters<typeof LocalBashExecutor>[1] = {}) {
  14. const ctx = new Context()
  15. await ctx.plugin(LocalSubprocessRuntime)
  16. ;(ctx.subprocess as LocalSubprocessRuntime).internals = { spillDir }
  17. // A short kill grace via the REAL config path, so escalation tests stay fast.
  18. await ctx.plugin(LocalBashExecutor, { graceMs: 200, ...config })
  19. const bash = ctx.shell as LocalBashExecutor
  20. return { ctx, bash }
  21. }
  22. class RejectingSubprocessRuntime extends SubprocessRuntime {
  23. private readonly reader: SubprocessOutputReader = {
  24. readFrom: () => ({ text: '', lossy: false, nextOffset: 0 }),
  25. }
  26. constructor(ctx: Context, private readonly failure: unknown, private readonly processId = 123) {
  27. super(ctx)
  28. }
  29. override async resolveExecutable(command: string): Promise<string> { return command }
  30. override spawnTerminal(): Promise<never> { throw new Error('bash spawns pipes, never terminals') }
  31. override spawn(_spec: SubprocessSpawnSpec): SubprocessHandle {
  32. return {
  33. pid: this.processId,
  34. stdin: undefined,
  35. stdout: undefined,
  36. stderr: undefined,
  37. collected: { stdout: this.reader, stderr: this.reader },
  38. done: Promise.resolve().then(() => { throw this.failure }),
  39. terminate: () => {},
  40. waitForExit: async () => true,
  41. }
  42. }
  43. }
  44. class ObservingBashExecutor extends LocalBashExecutor {
  45. spawnFailed: boolean | undefined
  46. protected override onProcessDone(
  47. _proc: ShellProcess,
  48. _stderr: string,
  49. spawnFailed: boolean,
  50. ): void {
  51. this.spawnFailed = spawnFailed
  52. }
  53. }
  54. /**
  55. * Poll a handle's consuming readOutput until the ACCUMULATED delta contains
  56. * `expected`; returns the accumulation (reads never re-deliver, so the caller
  57. * gets everything produced up to the match).
  58. */
  59. async function readUntil(proc: ShellProcess, expected: string, timeoutMs = 5_000): Promise<string> {
  60. const deadline = Date.now() + timeoutMs
  61. let all = ''
  62. while (Date.now() < deadline) {
  63. all += proc.readOutput().delta
  64. if (all.includes(expected)) return all
  65. await new Promise(resolve => setTimeout(resolve, 20))
  66. }
  67. throw new Error(`process output did not include ${JSON.stringify(expected)}; accumulated ${JSON.stringify(all)}`)
  68. }
  69. describe('LocalBashExecutor.run', () => {
  70. it('resolves with output and the effective timeout', async () => {
  71. const { bash } = await setup({ timeoutMs: 5_000 })
  72. const result = await bash.run(bash.resolve({ command: 'echo hi' }))
  73. expect(result.exitCode).toBe(0)
  74. expect(result.stdout.text).toBe('hi\n')
  75. expect(result.timeoutMs).toBe(5_000)
  76. })
  77. it('uses config cwd, overridable per call', async () => {
  78. const { bash } = await setup({ cwd: '/tmp' })
  79. const fromConfig = await bash.run(bash.resolve({ command: 'pwd' }))
  80. expect(fromConfig.stdout.text.trim()).toMatch(/\/tmp$/)
  81. const fromCall = await bash.run(bash.resolve({ command: 'pwd', workdir: '/' }))
  82. expect(fromCall.stdout.text.trim()).toBe('/')
  83. })
  84. it('defaults cwd to process.cwd()', async () => {
  85. const { bash } = await setup()
  86. const result = await bash.run(bash.resolve({ command: 'pwd' }))
  87. expect(result.stdout.text.trim()).toBe(process.cwd())
  88. })
  89. it('caps per-call timeouts at maxTimeoutMs', async () => {
  90. const { bash } = await setup({ timeoutMs: 1_000, maxTimeoutMs: 2_000 })
  91. const result = await bash.run(bash.resolve({ command: 'true', timeoutMs: 99_999 }))
  92. expect(result.timeoutMs).toBe(2_000)
  93. })
  94. it('rejects invalid numeric config and timeout overrides', async () => {
  95. await expect(setup({ timeoutMs: Number.NaN })).rejects.toThrow(/timeoutMs/)
  96. await expect(setup({ maxTimeoutMs: 0 })).rejects.toThrow(/maxTimeoutMs/)
  97. await expect(setup({ maxOutputBytes: -1 })).rejects.toThrow(/maxOutputBytes/)
  98. await expect(setup({ maxSpillBytes: 0 })).rejects.toThrow(/maxSpillBytes/)
  99. await expect(setup({ graceMs: 0 })).rejects.toThrow(/graceMs/)
  100. await expect(setup({ graceMs: MAX_TIMER_DELAY_MS + 1 }))
  101. .rejects.toThrow(`graceMs must be no greater than ${MAX_TIMER_DELAY_MS}`)
  102. const { bash } = await setup()
  103. expect(() => bash.resolve({ command: 'true', timeoutMs: Number.NaN })).toThrow(/request\.timeoutMs/)
  104. expect(() => bash.resolve({ command: 'true', timeoutMs: -1 })).toThrow(/request\.timeoutMs/)
  105. expect(() => bash.resolve({ command: 'true', stdoutMaxBytes: Number.NaN })).toThrow(/request\.stdoutMaxBytes/)
  106. expect(() => bash.resolve({ command: 'true', stdoutMaxBytes: -1 })).toThrow(/request\.stdoutMaxBytes/)
  107. })
  108. it('defaults stdoutMaxBytes to maxOutputBytes and lets foreground callers raise stdout only', async () => {
  109. const { bash } = await setup({ maxOutputBytes: 100 })
  110. expect(bash.resolve({ command: 'true' }).stdoutMaxBytes).toBe(100)
  111. const result = await bash.run(bash.resolve({
  112. command: 'printf "%.0sx" $(seq 1 500); printf "%.0se" $(seq 1 500) >&2',
  113. stdoutMaxBytes: 500,
  114. }))
  115. expect(result.stdout.truncated).toBe(false)
  116. expect(result.stdout.text).toBe('x'.repeat(500))
  117. expect(result.stderr.truncated).toBe(true)
  118. expect(result.stderr.text.length).toBeLessThanOrEqual(100)
  119. })
  120. it('per-call timeout takes precedence under the cap and kills on expiry', async () => {
  121. const { bash } = await setup({ timeoutMs: 60_000 })
  122. const result = await bash.run(bash.resolve({ command: 'sleep 60', timeoutMs: 100 }))
  123. expect(result.timedOut).toBe(true)
  124. // Mutually exclusive: a timeout classifies as timedOut, never also aborted.
  125. expect(result.aborted).toBe(false)
  126. expect(result.timeoutMs).toBe(100)
  127. })
  128. it('propagates abort signals', async () => {
  129. const { bash } = await setup()
  130. const controller = new AbortController()
  131. const pending = bash.run(bash.resolve({ command: 'sleep 60', signal: controller.signal }))
  132. setTimeout(() => { controller.abort() }, 50)
  133. const result = await pending
  134. expect(result.aborted).toBe(true)
  135. // Mutually exclusive: an upstream cancel classifies as aborted, never also timedOut.
  136. expect(result.timedOut).toBe(false)
  137. })
  138. it('classifies a self-killed command as neither timed out nor aborted', async () => {
  139. // The command kills itself (SIGTERM) with no timeout and no upstream abort:
  140. // the deadline signal never fires, so both classifications are false — the
  141. // fused-signal classification reports the cause that cut the command short,
  142. // and here nothing the executor owns did.
  143. const { bash } = await setup({ timeoutMs: 60_000 })
  144. const result = await bash.run(bash.resolve({ command: 'kill -TERM $$' }))
  145. expect(result.signal).toBe('SIGTERM')
  146. expect(result.timedOut).toBe(false)
  147. expect(result.aborted).toBe(false)
  148. })
  149. it('rejects on spawn failure (bad workdir)', async () => {
  150. const { bash } = await setup()
  151. await expect(bash.run(bash.resolve({ command: 'true', workdir: '/nonexistent-dsh' }))).rejects.toThrow(/ENOENT/)
  152. })
  153. it('resolve() carries stdin/env/dshEnv onto the spec, and run() threads them to the command', async () => {
  154. const { bash } = await setup()
  155. const spec = bash.resolve({
  156. command: 'cat; echo "[$SEAM_VAR][$DSH_SEAM_VAR]"',
  157. stdin: 'piped\n',
  158. env: { SEAM_VAR: 'env-ok' },
  159. dshEnv: { DSH_SEAM_VAR: 'dsh-ok' },
  160. })
  161. // resolve() keeps the optional input/environment fields verbatim.
  162. expect(spec.stdin).toBe('piped\n')
  163. expect(spec.env).toEqual({ SEAM_VAR: 'env-ok' })
  164. expect(spec.dshEnv).toEqual({ DSH_SEAM_VAR: 'dsh-ok' })
  165. const result = await bash.run(spec)
  166. expect(result.stdout.text).toBe('piped\n[env-ok][dsh-ok]\n')
  167. })
  168. it('resolve() omits stdin/env/dshEnv when the request supplies none', async () => {
  169. const { bash } = await setup()
  170. const spec = bash.resolve({ command: 'true' })
  171. expect('stdin' in spec).toBe(false)
  172. expect('env' in spec).toBe(false)
  173. expect('dshEnv' in spec).toBe(false)
  174. })
  175. })
  176. describe('LocalBashExecutor.start (background process handles)', () => {
  177. it('start returns immediately with a running handle that settles as completed', async () => {
  178. const { bash } = await setup()
  179. const before = Date.now()
  180. const proc = bash.start(bash.resolve({ command: 'sleep 0.2; echo done' }))
  181. expect(Date.now() - before).toBeLessThan(150)
  182. expect(proc.status).toBe('running')
  183. await proc.done
  184. expect(proc.status).toBe('completed')
  185. expect(proc.exitCode).toBe(0)
  186. })
  187. it('threads stdin and extra env into a background process', async () => {
  188. const { bash } = await setup()
  189. const proc = bash.start(bash.resolve({
  190. command: 'cat; echo "[$BG_VAR][$DSH_BG_VAR]"',
  191. stdin: 'bg-stdin\n',
  192. env: { BG_VAR: 'bg-env' },
  193. dshEnv: { DSH_BG_VAR: 'bg-dsh-env' },
  194. }))
  195. const output = await readUntil(proc, '[bg-env][bg-dsh-env]')
  196. expect(output).toContain('bg-stdin')
  197. await proc.done
  198. expect(proc.exitCode).toBe(0)
  199. })
  200. it('readOutput is consuming: increments are never re-delivered, and reads stay valid after exit', async () => {
  201. const { bash } = await setup()
  202. const proc = bash.start(bash.resolve({ command: 'echo first; sleep 1; echo second' }))
  203. const first = await readUntil(proc, 'first\n')
  204. expect(first).toBe('first\n')
  205. await proc.done
  206. // Read-after-exit returns the remaining buffered output — once.
  207. const second = proc.readOutput()
  208. expect(second.delta).toBe('second\n')
  209. expect(second.lossy).toBe(false)
  210. expect(proc.readOutput().delta).toBe('')
  211. })
  212. it('readOutput marks stderr sections', async () => {
  213. const { bash } = await setup()
  214. const proc = bash.start(bash.resolve({ command: 'echo out; echo err >&2' }))
  215. await proc.done
  216. expect(proc.readOutput().delta).toBe('out\n[stderr]\nerr\n')
  217. })
  218. it('readOutput reports stderr-only deltas without a leading newline', async () => {
  219. const { bash } = await setup()
  220. const proc = bash.start(bash.resolve({ command: 'echo err >&2' }))
  221. await proc.done
  222. expect(proc.readOutput().delta).toBe('[stderr]\nerr\n')
  223. })
  224. it('readOutput adds a separator only when stdout lacks a trailing newline', async () => {
  225. const { bash } = await setup()
  226. const proc = bash.start(bash.resolve({ command: 'printf out; echo err >&2' }))
  227. await proc.done
  228. expect(proc.readOutput().delta).toBe('out\n[stderr]\nerr\n')
  229. })
  230. it('readOutput flags lossy reads and reports stdout spill paths', async () => {
  231. const { bash } = await setup({ maxOutputBytes: 100 })
  232. const proc = bash.start(bash.resolve({ command: 'for i in $(seq 1 100); do printf "line-%04d\\n" $i; done' }))
  233. await proc.done
  234. const read = proc.readOutput()
  235. // Window slid past offset 0 → lossy, spill path points at the full stream.
  236. expect(read.lossy).toBe(true)
  237. expect(read.stdoutSpillPath).toBeDefined()
  238. })
  239. it('readOutput reports stderr spill paths', async () => {
  240. const { bash } = await setup({ maxOutputBytes: 100 })
  241. const proc = bash.start(bash.resolve({ command: 'for i in $(seq 1 100); do printf "line-%04d\\n" $i >&2; done' }))
  242. await proc.done
  243. const read = proc.readOutput()
  244. expect(read.lossy).toBe(true)
  245. expect(read.stderrSpillPath).toBeDefined()
  246. expect(read.delta).toContain('[stderr]')
  247. })
  248. it('kill() terminates the process group: true once, false after settlement', async () => {
  249. const { bash } = await setup()
  250. const proc = bash.start(bash.resolve({ command: 'sleep 60' }))
  251. expect(proc.kill()).toBe(true)
  252. await proc.done
  253. expect(proc.status).toBe('killed')
  254. expect(proc.signal).toBe('SIGTERM')
  255. expect(proc.kill()).toBe(false)
  256. })
  257. it('kill() returns false for a naturally completed process', async () => {
  258. const { bash } = await setup()
  259. const proc = bash.start(bash.resolve({ command: 'true' }))
  260. await proc.done
  261. expect(proc.status).toBe('completed')
  262. expect(proc.kill()).toBe(false)
  263. })
  264. it('kill escalation uses the configured graceMs (a TERM-trapping process dies by SIGKILL)', async () => {
  265. const { bash } = await setup() // setup pins graceMs: 200 via config
  266. // The child echoes AFTER arming the trap, so waiting for the marker
  267. // guarantees SIGTERM is already ignored when the kill lands (a fixed sleep
  268. // is load-flaky: a slow spawn would take the SIGTERM before the trap).
  269. const proc = bash.start(bash.resolve({ command: 'trap \'\' TERM; echo armed; sleep 60' }))
  270. await readUntil(proc, 'armed')
  271. proc.kill()
  272. await proc.done
  273. expect(proc.status).toBe('killed')
  274. expect(proc.signal).toBe('SIGKILL')
  275. })
  276. it('a spec.signal abort settles the handle as killed, not completed', async () => {
  277. const { bash } = await setup()
  278. const controller = new AbortController()
  279. const proc = bash.start(bash.resolve({ command: 'sleep 60', signal: controller.signal }))
  280. controller.abort()
  281. await proc.done
  282. expect(proc.status).toBe('killed')
  283. expect(proc.signal).toBe('SIGTERM')
  284. })
  285. it('a self-signal exit settles the handle as killed, not completed', async () => {
  286. const { bash } = await setup()
  287. const proc = bash.start(bash.resolve({ command: 'kill -TERM $$' }))
  288. await proc.done
  289. expect(proc.status).toBe('killed')
  290. expect(proc.exitCode).toBeNull()
  291. expect(proc.signal).toBe('SIGTERM')
  292. })
  293. it('a background spawn failure settles as killed with the error readable on stderr', async () => {
  294. const { bash } = await setup()
  295. const proc = bash.start(bash.resolve({ command: 'true', workdir: '/nonexistent-dsh' }))
  296. // done resolves (never rejects) even though the process never ran.
  297. await expect(proc.done).resolves.toBeUndefined()
  298. expect(proc.status).toBe('killed')
  299. expect(proc.readOutput().delta).toContain('spawn failed:')
  300. })
  301. it('does not label a post-start provider rejection as a spawn failure', async () => {
  302. const ctx = new Context()
  303. const failure = Object.assign(new Error('managed owner became unreadable'), {
  304. code: 'ENOENT',
  305. syscall: 'spawn bash',
  306. path: 'bash',
  307. })
  308. new RejectingSubprocessRuntime(ctx, failure)
  309. await ctx.plugin(ObservingBashExecutor)
  310. const bash = ctx.shell as ObservingBashExecutor
  311. const proc = bash.start(bash.resolve({ command: 'true' }))
  312. await proc.done
  313. expect(proc.readOutput().delta).toContain('subprocess failed:')
  314. expect(proc.readOutput().delta).toBe('')
  315. expect(bash.spawnFailed).toBe(false)
  316. })
  317. it.each([
  318. ['non-object rejection', undefined, 'subprocess failed:', false],
  319. ['non-string syscall', { syscall: 1 }, 'subprocess failed:', false],
  320. ['non-spawn syscall', { syscall: 'kill', path: 'bash' }, 'subprocess failed:', false],
  321. ['matching syscall without path', { syscall: 'spawn bash' }, 'spawn failed:', true],
  322. ])('classifies a pre-start %s from structured evidence', async (_label, failure, note, spawnFailed) => {
  323. const ctx = new Context()
  324. new RejectingSubprocessRuntime(ctx, failure, -1)
  325. await ctx.plugin(ObservingBashExecutor)
  326. const bash = ctx.shell as ObservingBashExecutor
  327. const proc = bash.start(bash.resolve({ command: 'true' }))
  328. await proc.done
  329. expect(proc.readOutput().delta).toContain(note)
  330. expect(bash.spawnFailed).toBe(spawnFailed)
  331. })
  332. })
  333. describe('process lifecycle ownership (the subprocess service, not the executor)', () => {
  334. it('a background process survives executor-fiber disposal and dies with the subprocess service', async () => {
  335. const ctx = new Context()
  336. const managerFiber = await ctx.plugin(LocalSubprocessRuntime)
  337. ;(ctx.subprocess as LocalSubprocessRuntime).internals = { spillDir }
  338. const executorFiber = await ctx.plugin(LocalBashExecutor, { graceMs: 200 })
  339. const bash = ctx.shell as LocalBashExecutor
  340. // The child prints its own pid ($$ = the detached bash group leader) so
  341. // the test can probe liveness through the public read API alone.
  342. const proc = bash.start(bash.resolve({ command: 'echo $$; sleep 60' }))
  343. const pid = Number((await readUntil(proc, '\n')).trim())
  344. expect(Number.isInteger(pid) && pid > 0).toBe(true)
  345. // Executor reload/disposal leaves background work running — the
  346. // handle stays live and readable, mirroring the job runtime's
  347. // registrations-outlive-producer-fibers contract.
  348. await executorFiber.dispose()
  349. expect(proc.status).toBe('running')
  350. expect(() => process.kill(pid, 0)).not.toThrow()
  351. // Service disposal kills the group and AWAITS its exit (no orphans).
  352. await managerFiber.dispose()
  353. expect(() => process.kill(pid, 0)).toThrow()
  354. await proc.done
  355. expect(proc.status).toBe('killed')
  356. })
  357. it('service disposal escalates to SIGKILL for TERM-trapping children and settles handles', async () => {
  358. const ctx = new Context()
  359. const managerFiber = await ctx.plugin(LocalSubprocessRuntime)
  360. ;(ctx.subprocess as LocalSubprocessRuntime).internals = { spillDir }
  361. await ctx.plugin(LocalBashExecutor, { graceMs: 200 })
  362. const bash = ctx.shell as LocalBashExecutor
  363. const finished = bash.start(bash.resolve({ command: 'echo done' }))
  364. await finished.done
  365. expect(finished.status).toBe('completed')
  366. const trapping = bash.start(bash.resolve({ command: 'trap \'\' TERM; echo armed; sleep 60' }))
  367. await readUntil(trapping, 'armed')
  368. await managerFiber.dispose()
  369. // A settled process was untouched; the live one died by escalation.
  370. expect(finished.status).toBe('completed')
  371. await trapping.done
  372. expect(trapping.status).toBe('killed')
  373. expect(trapping.signal).toBe('SIGKILL')
  374. })
  375. })