| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382 |
- import { mkdtempSync } from 'node:fs'
- import { tmpdir } from 'node:os'
- import { join } from 'node:path'
- import { describe, expect, it, vi } from 'vitest'
- import { Context } from 'cordis'
- import { LocalBashExecutor } from '@deepseek-ai/dsh-bash-local'
- import { BashTaskId } from '@deepseek-ai/dsh-bash'
- import type { BashTaskRead } from '@deepseek-ai/dsh-bash'
- const spillDir = mkdtempSync(join(tmpdir(), 'dsh-bash-exec-spec-'))
- async function setup(config: ConstructorParameters<typeof LocalBashExecutor>[1] = {}) {
- const ctx = new Context()
- // A short kill grace via the REAL config path, so escalation tests stay fast.
- await ctx.plugin(LocalBashExecutor, { graceMs: 200, ...config })
- const bash = ctx.bash as LocalBashExecutor
- bash.internals = { spillDir }
- return { ctx, bash }
- }
- /** Poll until a pid no longer exists. */
- async function waitGone(pid: number, timeoutMs = 5_000): Promise<void> {
- const deadline = Date.now() + timeoutMs
- while (Date.now() < deadline) {
- try {
- process.kill(pid, 0)
- } catch {
- return
- }
- await new Promise(resolve => setTimeout(resolve, 20))
- }
- throw new Error(`pid ${pid} still alive after ${timeoutMs}ms`)
- }
- async function readUntil(
- bash: LocalBashExecutor,
- id: BashTaskId,
- expected: string,
- timeoutMs = 5_000,
- ): Promise<BashTaskRead> {
- const deadline = Date.now() + timeoutMs
- let last: BashTaskRead | undefined
- let delta = ''
- while (Date.now() < deadline) {
- last = bash.readOutput(id)
- delta += last.delta
- if (delta.includes(expected)) return { ...last, delta }
- await new Promise(resolve => setTimeout(resolve, 20))
- }
- throw new Error(`task ${id} output did not include ${JSON.stringify(expected)}; output was ${JSON.stringify(delta)}, last delta was ${JSON.stringify(last?.delta ?? '')}`)
- }
- describe('LocalBashExecutor.run', () => {
- it('resolves with output and the effective timeout', async () => {
- const { bash } = await setup({ timeoutMs: 5_000 })
- const result = await bash.run(bash.resolve({ command: 'echo hi' }))
- expect(result.exitCode).toBe(0)
- expect(result.stdout.text).toBe('hi\n')
- expect(result.timeoutMs).toBe(5_000)
- })
- it('uses config cwd, overridable per call', async () => {
- const { bash } = await setup({ cwd: '/tmp' })
- const fromConfig = await bash.run(bash.resolve({ command: 'pwd' }))
- expect(fromConfig.stdout.text.trim()).toMatch(/\/tmp$/)
- const fromCall = await bash.run(bash.resolve({ command: 'pwd', workdir: '/' }))
- expect(fromCall.stdout.text.trim()).toBe('/')
- })
- it('defaults cwd to process.cwd()', async () => {
- const { bash } = await setup()
- const result = await bash.run(bash.resolve({ command: 'pwd' }))
- expect(result.stdout.text.trim()).toBe(process.cwd())
- })
- it('caps per-call timeouts at maxTimeoutMs', async () => {
- const { bash } = await setup({ timeoutMs: 1_000, maxTimeoutMs: 2_000 })
- const result = await bash.run(bash.resolve({ command: 'true', timeoutMs: 99_999 }))
- expect(result.timeoutMs).toBe(2_000)
- })
- it('rejects invalid numeric config and timeout overrides', async () => {
- 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({ graceMs: 0 })).rejects.toThrow(/graceMs/)
- const { bash } = await setup()
- expect(() => bash.resolve({ command: 'true', timeoutMs: Number.NaN })).toThrow(/request\.timeoutMs/)
- expect(() => bash.resolve({ command: 'true', timeoutMs: -1 })).toThrow(/request\.timeoutMs/)
- })
- it('kill escalation uses the configured graceMs (a TERM-trapping task dies by SIGKILL)', async () => {
- const { bash } = await setup() // setup pins graceMs: 200 via config
- const task = bash.start(bash.resolve({ command: 'trap \'\' TERM; echo ready; while :; do sleep 60 & wait $!; done' }))
- await readUntil(bash, task.id, 'ready\n')
- bash.kill(task.id)
- await task.done
- expect(task.signal).toBe('SIGKILL')
- })
- it('per-call timeout takes precedence under the cap and kills on expiry', async () => {
- const { bash } = await setup({ timeoutMs: 60_000 })
- const result = await bash.run(bash.resolve({ command: 'sleep 60', timeoutMs: 100 }))
- expect(result.timedOut).toBe(true)
- // Mutually exclusive: a timeout classifies as timedOut, never also aborted.
- expect(result.aborted).toBe(false)
- expect(result.timeoutMs).toBe(100)
- })
- it('propagates abort signals', async () => {
- const { bash } = await setup()
- const controller = new AbortController()
- const pending = bash.run(bash.resolve({ command: 'sleep 60', signal: controller.signal }))
- setTimeout(() => { controller.abort() }, 50)
- const result = await pending
- expect(result.aborted).toBe(true)
- // Mutually exclusive: an upstream cancel classifies as aborted, never also timedOut.
- expect(result.timedOut).toBe(false)
- })
- it('classifies a self-killed command as neither timed out nor aborted', async () => {
- // The command kills itself (SIGTERM) with no timeout and no upstream abort:
- // the deadline signal never fires, so both classifications are false — the
- // fused-signal classification reports the cause that cut the command short,
- // and here nothing the executor owns did.
- const { bash } = await setup({ timeoutMs: 60_000 })
- const result = await bash.run(bash.resolve({ command: 'kill -TERM $$' }))
- expect(result.signal).toBe('SIGTERM')
- expect(result.timedOut).toBe(false)
- expect(result.aborted).toBe(false)
- })
- it('rejects on spawn failure (bad workdir)', async () => {
- const { bash } = await setup()
- await expect(bash.run(bash.resolve({ command: 'true', workdir: '/nonexistent-dsh' }))).rejects.toThrow(/ENOENT/)
- })
- it('resolve() carries stdin/env onto the spec, and run() threads them to the command', async () => {
- const { bash } = await setup()
- const spec = bash.resolve({ command: 'cat; echo "[$DSH_SEAM_VAR]"', stdin: 'piped\n', env: { DSH_SEAM_VAR: 'env-ok' } })
- // resolve() keeps the stdin/env fields verbatim (optional, no default).
- expect(spec.stdin).toBe('piped\n')
- expect(spec.env).toEqual({ DSH_SEAM_VAR: 'env-ok' })
- const result = await bash.run(spec)
- expect(result.stdout.text).toBe('piped\n[env-ok]\n')
- })
- it('resolve() omits stdin/env when the request supplies neither', async () => {
- const { bash } = await setup()
- const spec = bash.resolve({ command: 'true' })
- expect('stdin' in spec).toBe(false)
- expect('env' in spec).toBe(false)
- })
- })
- describe('LocalBashExecutor background tasks', () => {
- it('start returns immediately with a registered running task', async () => {
- const { bash } = await setup()
- const before = Date.now()
- const task = bash.start(bash.resolve({ command: 'sleep 0.2; echo done' }))
- expect(Date.now() - before).toBeLessThan(150)
- expect(task.status).toBe('running')
- expect(bash.get(task.id)).toBe(task)
- expect(bash.list()).toContain(task)
- await task.done
- expect(task.status).toBe('completed')
- expect(task.exitCode).toBe(0)
- })
- it('assigns sequential ids', async () => {
- const { bash } = await setup()
- const first = bash.start(bash.resolve({ command: 'true' }))
- const second = bash.start(bash.resolve({ command: 'true' }))
- expect(first.id).toBe('bash-1')
- expect(second.id).toBe('bash-2')
- await Promise.all([first.done, second.done])
- })
- it('threads stdin and extra env into a background task', async () => {
- const { bash } = await setup()
- const task = bash.start(bash.resolve({
- command: 'cat; echo "[$DSH_BG_VAR]"',
- stdin: 'bg-stdin\n',
- env: { DSH_BG_VAR: 'bg-env' },
- }))
- const read = await readUntil(bash, task.id, '[bg-env]')
- expect(read.delta).toContain('bg-stdin')
- await task.done
- expect(task.exitCode).toBe(0)
- })
- it('readOutput returns increments without re-delivery', async () => {
- const { bash } = await setup()
- const task = bash.start(bash.resolve({ command: 'echo first; sleep 1; echo second' }))
- const first = await readUntil(bash, task.id, 'first\n')
- expect(first.delta).toBe('first\n')
- expect(first.lossy).toBe(false)
- await task.done
- const second = bash.readOutput(task.id)
- expect(second.delta).toBe('second\n')
- const third = bash.readOutput(task.id)
- expect(third.delta).toBe('')
- })
- it('readOutput marks stderr sections', async () => {
- const { bash } = await setup()
- const task = bash.start(bash.resolve({ command: 'echo out; echo err >&2' }))
- await task.done
- const read = bash.readOutput(task.id)
- expect(read.delta).toBe('out\n[stderr]\nerr\n')
- })
- it('readOutput reports stderr-only deltas without a leading newline', async () => {
- const { bash } = await setup()
- const task = bash.start(bash.resolve({ command: 'echo err >&2' }))
- await task.done
- expect(bash.readOutput(task.id).delta).toBe('[stderr]\nerr\n')
- })
- it('readOutput flags lossy reads and reports spill paths', async () => {
- const { bash } = await setup({ maxOutputBytes: 100 })
- const task = bash.start(bash.resolve({ command: 'for i in $(seq 1 100); do printf "line-%04d\\n" $i; done' }))
- await task.done
- const read = bash.readOutput(task.id)
- // Window slid past offset 0 → lossy, spill path points at the full stream.
- expect(read.lossy).toBe(true)
- expect(read.stdoutSpillPath).toBeDefined()
- })
- it('readOutput throws for unknown ids', async () => {
- const { bash } = await setup()
- expect(() => bash.readOutput(BashTaskId('nope'))).toThrow(/unknown bash task "nope"/)
- })
- it('kill terminates the process group and reports status killed', async () => {
- const { bash } = await setup()
- const task = bash.start(bash.resolve({ command: 'sleep 60' }))
- expect(bash.kill(task.id)).toBe(true)
- await task.done
- expect(task.status).toBe('killed')
- expect(task.signal).toBe('SIGTERM')
- })
- it('kill returns false for finished tasks and throws for unknown ids', async () => {
- const { bash } = await setup()
- const task = bash.start(bash.resolve({ command: 'true' }))
- await task.done
- expect(bash.kill(task.id)).toBe(false)
- expect(() => bash.kill(BashTaskId('nope'))).toThrow(/unknown bash task "nope"/)
- })
- it('notifies onTaskDone listeners on completion', async () => {
- const { bash } = await setup()
- const seen: [string, string][] = []
- bash.onTaskDone(task => void seen.push([task.id, task.status]))
- const task = bash.start(bash.resolve({ command: 'true' }))
- await task.done
- expect(seen).toEqual([[task.id, 'completed']])
- })
- it('notifies onTaskDone for killed tasks too', async () => {
- const { bash } = await setup()
- const listener = vi.fn()
- bash.onTaskDone(listener)
- const task = bash.start(bash.resolve({ command: 'sleep 60' }))
- bash.kill(task.id)
- await task.done
- expect(listener).toHaveBeenCalledWith(task)
- expect(task.status).toBe('killed')
- })
- it('marks tasks killed when the background spawn itself fails', async () => {
- const { bash } = await setup()
- const listener = vi.fn()
- bash.onTaskDone(listener)
- const task = bash.start(bash.resolve({ command: 'true', workdir: '/nonexistent-dsh' }))
- await task.done
- expect(task.status).toBe('killed')
- expect(listener).toHaveBeenCalledWith(task)
- expect(bash.readOutput(task.id).delta).toContain('spawn failed')
- })
- it('readOutput adds a separator only when stdout lacks a trailing newline', async () => {
- const { bash } = await setup()
- const task = bash.start(bash.resolve({ command: 'printf out; echo err >&2' }))
- await task.done
- expect(bash.readOutput(task.id).delta).toBe('out\n[stderr]\nerr\n')
- })
- it('readOutput reports stderr spill paths', async () => {
- const { bash } = await setup({ maxOutputBytes: 100 })
- const task = bash.start(bash.resolve({ command: 'for i in $(seq 1 100); do printf "line-%04d\\n" $i >&2; done' }))
- await task.done
- const read = bash.readOutput(task.id)
- expect(read.lossy).toBe(true)
- expect(read.stderrSpillPath).toBeDefined()
- expect(read.delta).toContain('[stderr]')
- })
- it('disposing with already-finished tasks only kills the running ones', async () => {
- const ctx = new Context()
- const fiber = await ctx.plugin(LocalBashExecutor, { graceMs: 200 })
- const bash = ctx.bash as LocalBashExecutor
- bash.internals = { spillDir }
- const finished = bash.start(bash.resolve({ command: 'true' }))
- await finished.done
- const running = bash.start(bash.resolve({ command: 'sleep 60' }))
- await fiber.dispose()
- await running.done
- expect(finished.status).toBe('completed')
- expect(running.signal).toBe('SIGTERM')
- expect(bash.list()).toEqual([])
- })
- it('disposing the executor fiber kills running tasks (no orphans)', async () => {
- const ctx = new Context()
- const fiber = await ctx.plugin(LocalBashExecutor, { graceMs: 200 })
- const bash = ctx.bash as LocalBashExecutor
- bash.internals = { spillDir }
- const listener = vi.fn()
- bash.onTaskDone(listener)
- const task = bash.start(bash.resolve({ command: 'sleep 60' }))
- const running = bash.get(task.id)!
- await new Promise(resolve => setTimeout(resolve, 50))
- // Grab the pid before dispose clears the registry.
- const pid = (running as unknown as { running: { pid: number } }).running.pid
- await fiber.dispose()
- await waitGone(pid)
- expect(bash.list()).toEqual([])
- // Listener silenced by base-class teardown — no late notifications.
- expect(listener).not.toHaveBeenCalled()
- })
- })
- describe('executor cancellation, callback, and disposal contracts', () => {
- it('start honors a pre-aborted or later-aborted AbortSignal', async () => {
- const { bash } = await setup()
- const controller = new AbortController()
- const task = bash.start(bash.resolve({ command: 'sleep 60', signal: controller.signal }))
- controller.abort()
- await task.done
- expect(task.status).toBe('killed')
- expect(task.signal).toBe('SIGTERM')
- })
- it('a throwing onTaskDone listener does not reject task.done or starve later listeners', async () => {
- const { bash } = await setup()
- const errorSpy = vi.spyOn(console, 'error').mockImplementation(() => undefined)
- const second = vi.fn()
- try {
- bash.onTaskDone(() => { throw new Error('listener bug') })
- bash.onTaskDone(second)
- const task = bash.start(bash.resolve({ command: 'true' }))
- await expect(task.done).resolves.toBeUndefined()
- expect(second).toHaveBeenCalledWith(task)
- expect(errorSpy).toHaveBeenCalled()
- } finally {
- errorSpy.mockRestore()
- }
- })
- it('dispose AWAITS a TERM-trapping process (SIGKILL escalation included)', async () => {
- const ctx = new Context()
- const fiber = await ctx.plugin(LocalBashExecutor, { graceMs: 200 })
- const bash = ctx.bash as LocalBashExecutor
- bash.internals = { spillDir }
- const task = bash.start(bash.resolve({ command: 'trap \'\' TERM; sleep 60' }))
- await new Promise(resolve => setTimeout(resolve, 100))
- const pid = (task as unknown as { running: { pid: number } }).running.pid
- await fiber.dispose()
- // Disposal itself waited: the pid must already be gone, no grace left.
- expect(() => process.kill(pid, 0)).toThrow()
- expect(task.status).toBe('killed')
- })
- })
|