cli.spec.ts 20 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474
  1. import { readdir, mkdtemp } from 'node:fs/promises'
  2. import { tmpdir } from 'node:os'
  3. import { join, resolve } from 'node:path'
  4. import { Context } from 'cordis'
  5. import type { Agent } from '@deepseek-ai/dsh-agent'
  6. import { CallId, LlmAdapter, type GenerateOptions, type StreamChunk, type TokenUsage } from '@deepseek-ai/dsh-llm'
  7. import { SessionId, type SessionEvent, type TurnEndReason } from '@deepseek-ai/dsh-session'
  8. import { afterEach, describe, expect, it } from 'vitest'
  9. import * as cliDemo from '../src/index.ts'
  10. import {
  11. executeCli,
  12. formatTurnFailure,
  13. parseCliArgs,
  14. runOneShot,
  15. type CliResult,
  16. } from '../src/cli.ts'
  17. type ScriptEntry = readonly StreamChunk[] | 'hang'
  18. class ScriptedAdapter extends LlmAdapter {
  19. readonly requests: GenerateOptions[] = []
  20. private cursor = 0
  21. constructor(private readonly script: readonly ScriptEntry[]) {
  22. super()
  23. }
  24. async * stream(options: GenerateOptions): AsyncIterable<StreamChunk> {
  25. this.requests.push(options)
  26. const entry = this.script[this.cursor++]
  27. if (entry === undefined) throw new Error('script exhausted')
  28. if (entry === 'hang') {
  29. yield { type: 'block-start', index: 0, blockType: 'text' }
  30. yield { type: 'text-delta', index: 0, text: 'partial' }
  31. await new Promise<void>((_resolve, reject) => {
  32. if (options.signal?.aborted === true) {
  33. reject(new Error('aborted'))
  34. return
  35. }
  36. options.signal?.addEventListener('abort', () => { reject(new Error('aborted')) }, { once: true })
  37. })
  38. return
  39. }
  40. for (const chunk of entry) yield chunk
  41. }
  42. }
  43. function textResponse(text: string, usage?: TokenUsage, finish: 'stop' | 'max-tokens' = 'stop'): StreamChunk[] {
  44. return [
  45. { type: 'block-start', index: 0, blockType: 'text' },
  46. { type: 'text-delta', index: 0, text },
  47. { type: 'block-end', index: 0, block: { type: 'text', text } },
  48. ...usage === undefined ? [] : [{ type: 'usage', usage } as const],
  49. { type: 'finish', reason: { kind: finish } },
  50. ]
  51. }
  52. function toolResponse(usage: TokenUsage): StreamChunk[] {
  53. const id = CallId('cli-call')
  54. const args = JSON.stringify({ text: 'round trip' })
  55. return [
  56. { type: 'block-start', index: 0, blockType: 'text' },
  57. { type: 'text-delta', index: 0, text: 'working' },
  58. { type: 'block-end', index: 0, block: { type: 'text', text: 'working' } },
  59. { type: 'block-start', index: 1, blockType: 'tool-call' },
  60. { type: 'tool-call-delta', index: 1, id, name: 'echo', argumentsDelta: args },
  61. { type: 'block-end', index: 1, block: { type: 'tool-call', id, name: 'echo', arguments: args } },
  62. { type: 'usage', usage },
  63. { type: 'finish', reason: { kind: 'tool-calls' } },
  64. ]
  65. }
  66. function reasoningResponse(text: string): StreamChunk[] {
  67. return [
  68. { type: 'block-start', index: 0, blockType: 'reasoning' },
  69. { type: 'reasoning-delta', index: 0, text },
  70. { type: 'block-end', index: 0, block: { type: 'reasoning', text } },
  71. { type: 'finish', reason: { kind: 'stop' } },
  72. ]
  73. }
  74. interface Harness {
  75. readonly ctx: Context
  76. readonly agent: Agent
  77. readonly persistenceRoot: string
  78. }
  79. const liveContexts: Context[] = []
  80. async function harness(script: readonly ScriptEntry[]): Promise<Harness> {
  81. const root = await mkdtemp(join(tmpdir(), 'dsh-cli-runner-'))
  82. const skillHome = await mkdtemp(join(tmpdir(), 'dsh-cli-runner-skills-'))
  83. const ctx = new Context()
  84. liveContexts.push(ctx)
  85. await ctx.plugin(cliDemo, {
  86. provider: 'mock',
  87. model: 'mock',
  88. persistenceRoot: root,
  89. skills: { local: { dshHome: join(skillHome, '.dsh'), agentsHome: join(skillHome, '.agents') } },
  90. workspaceContext: false,
  91. })
  92. await new Promise(resolve => setTimeout(resolve, 80))
  93. ctx.llm.registerAdapter(['mock'], new ScriptedAdapter(script))
  94. ctx.tools.register({
  95. name: 'echo',
  96. description: 'Echo text.',
  97. parameters: { text: { type: 'string', required: true } },
  98. execute: async args => [{ type: 'text', text: `ECHO: ${(args as { text: string }).text}` }],
  99. })
  100. const [agent] = ctx.agents.roots()
  101. if (agent === undefined) throw new Error('test main agent missing')
  102. return { ctx, agent, persistenceRoot: root }
  103. }
  104. async function invoke(
  105. ctx: Context,
  106. args: readonly string[],
  107. options: { signal?: AbortSignal; failStdout?: boolean; failDispose?: boolean } = {},
  108. ): Promise<{ code: number; stdout: string; stderr: string }> {
  109. let stdout = ''
  110. let stderr = ''
  111. const code = await executeCli(args, {
  112. cwd: '/tmp/cli-cwd',
  113. ...options.signal === undefined ? {} : { signal: options.signal },
  114. boot: async () => ctx,
  115. loadEnv: () => {},
  116. writeStdout: (chunk) => {
  117. if (options.failStdout === true) throw new Error('stdout closed')
  118. stdout += chunk
  119. },
  120. writeStderr: (chunk) => { stderr += chunk },
  121. ...options.failDispose === true
  122. ? { dispose: async (target: Context) => {
  123. await target.fiber.dispose()
  124. throw new Error('dispose exploded')
  125. } }
  126. : {},
  127. })
  128. return { code, stdout, stderr }
  129. }
  130. afterEach(async () => {
  131. await Promise.all(liveContexts.splice(0).map(ctx => ctx.fiber.dispose()))
  132. })
  133. describe('parseCliArgs', () => {
  134. it('parses defaults, explicit options, spaces, and an option-like task after --', () => {
  135. expect(parseCliArgs(['task with spaces'])).toEqual({
  136. kind: 'run', configPath: './cordis.yml', outputFormat: 'text', task: 'task with spaces',
  137. })
  138. expect(parseCliArgs(['--config', 'custom.yml', '--output-format', 'stream-json', 'do it'])).toEqual({
  139. kind: 'run', configPath: 'custom.yml', outputFormat: 'stream-json', task: 'do it',
  140. })
  141. expect(parseCliArgs(['--', '-task'])).toMatchObject({ task: '-task' })
  142. expect(parseCliArgs(['--help', 'ignored'])).toEqual({ kind: 'help' })
  143. })
  144. it('rejects missing, blank, extra, invalid-format, and unsupported flags', () => {
  145. expect(() => parseCliArgs([])).toThrow('received 0')
  146. expect(() => parseCliArgs([' '])).toThrow('must not be blank')
  147. expect(() => parseCliArgs(['one', 'two'])).toThrow('received 2')
  148. expect(() => parseCliArgs(['--output-format', 'xml', 'task'])).toThrow('unsupported output format')
  149. expect(() => parseCliArgs(['-p', 'task'])).toThrow('Unknown option')
  150. })
  151. })
  152. describe('runOneShot and executeCli', () => {
  153. it('prints help and argument diagnostics without booting or contaminating stdout', async () => {
  154. let booted = false
  155. let stdout = ''
  156. let stderr = ''
  157. const runtime = {
  158. boot: async (): Promise<Context> => { booted = true; throw new Error('unexpected') },
  159. writeStdout: (chunk: string): void => { stdout += chunk },
  160. writeStderr: (chunk: string): void => { stderr += chunk },
  161. }
  162. expect(await executeCli(['--help'], runtime)).toBe(0)
  163. expect(stdout).toContain('Usage: dsh-cli-demo')
  164. stdout = ''
  165. expect(await executeCli([], runtime)).toBe(1)
  166. expect(stdout).toBe('')
  167. expect(stderr).toContain('received 0')
  168. expect(booted).toBe(false)
  169. })
  170. it('leaves stdout empty for environment and boot failures and resolves the default config', async () => {
  171. let bootPath = ''
  172. let stderr = ''
  173. const code = await executeCli(['task'], {
  174. cwd: '/tmp/cli-work',
  175. loadEnv: (_name, _dir, warn) => { warn('env warning\n') },
  176. boot: async (_name, path) => { bootPath = path; throw 'boot exploded' },
  177. writeStdout: () => { throw new Error('stdout must stay empty') },
  178. writeStderr: (chunk) => { stderr += chunk },
  179. })
  180. expect(code).toBe(1)
  181. expect(bootPath).toBe(resolve('/tmp/cli-work/cordis.yml'))
  182. expect(stderr).toContain('env warning')
  183. expect(stderr).toContain('boot exploded')
  184. })
  185. it('contains a thrown value whose inspection and coercion both fail', async () => {
  186. const hostile = new Proxy({}, {
  187. getPrototypeOf: () => { throw new Error('prototype trap escaped') },
  188. get: (target, key, receiver) => {
  189. if (key === Symbol.toPrimitive) throw new Error('coercion escaped')
  190. return Reflect.get(target, key, receiver) as unknown
  191. },
  192. })
  193. let stdout = ''
  194. let stderr = ''
  195. const code = await executeCli(['task'], {
  196. boot: async () => { throw hostile },
  197. loadEnv: () => {},
  198. writeStdout: (chunk) => { stdout += chunk },
  199. writeStderr: (chunk) => { stderr += chunk },
  200. })
  201. expect(code).toBe(1)
  202. expect(stdout).toBe('')
  203. expect(stderr).toBe('dsh-cli-demo: [unrenderable thrown value]\n')
  204. })
  205. it('interrupts Loader boot and contains every late boot outcome', async () => {
  206. const abort = new AbortController()
  207. const lateContext = new Context()
  208. liveContexts.push(lateContext)
  209. const boot = Promise.withResolvers<Context>()
  210. const disposed = Promise.withResolvers<undefined>()
  211. let disposeCalls = 0
  212. let stderr = ''
  213. const running = executeCli(['task'], {
  214. signal: abort.signal,
  215. boot: () => boot.promise,
  216. loadEnv: () => {},
  217. writeStdout: () => {},
  218. writeStderr: (chunk) => { stderr += chunk },
  219. dispose: async (ctx) => {
  220. disposeCalls += 1
  221. await ctx.fiber.dispose()
  222. disposed.resolve(undefined)
  223. },
  224. })
  225. abort.abort('received SIGTERM')
  226. await expect(running).resolves.toBe(1)
  227. expect(stderr).toContain('received SIGTERM')
  228. expect(disposeCalls).toBe(0)
  229. boot.resolve(lateContext)
  230. await disposed.promise
  231. expect(disposeCalls).toBe(1)
  232. const rejectedBoot = Promise.withResolvers<Context>()
  233. const rejectedAbort = new AbortController()
  234. const rejected = executeCli(['task'], {
  235. signal: rejectedAbort.signal,
  236. boot: () => rejectedBoot.promise,
  237. loadEnv: () => {},
  238. writeStdout: () => {},
  239. writeStderr: () => {},
  240. })
  241. rejectedAbort.abort('stop rejected boot')
  242. await expect(rejected).resolves.toBe(1)
  243. rejectedBoot.reject(new Error('late boot rejection'))
  244. await Promise.resolve()
  245. let ordinaryBootStderr = ''
  246. const ordinaryBootFailure = await executeCli(['task'], {
  247. signal: new AbortController().signal,
  248. boot: async () => { throw new Error('ordinary boot failure') },
  249. loadEnv: () => {},
  250. writeStdout: () => {},
  251. writeStderr: (chunk) => { ordinaryBootStderr += chunk },
  252. })
  253. expect(ordinaryBootFailure).toBe(1)
  254. expect(ordinaryBootStderr).toContain('ordinary boot failure')
  255. const failedCleanupBoot = Promise.withResolvers<Context>()
  256. const failedCleanupAbort = new AbortController()
  257. const cleanupFailure = Promise.withResolvers<undefined>()
  258. const failedCleanupContext = new Context()
  259. liveContexts.push(failedCleanupContext)
  260. const failedCleanup = executeCli(['task'], {
  261. signal: failedCleanupAbort.signal,
  262. boot: () => failedCleanupBoot.promise,
  263. loadEnv: () => {},
  264. writeStdout: () => {},
  265. writeStderr: (chunk) => {
  266. if (chunk.includes('dispose after interrupted boot failed: late cleanup')) cleanupFailure.resolve(undefined)
  267. },
  268. dispose: async (ctx) => {
  269. await ctx.fiber.dispose()
  270. throw new Error('late cleanup')
  271. },
  272. })
  273. failedCleanupAbort.abort('stop failed cleanup boot')
  274. await expect(failedCleanup).resolves.toBe(1)
  275. failedCleanupBoot.resolve(failedCleanupContext)
  276. await cleanupFailure.promise
  277. })
  278. it('renders text, flushes a persisted fresh session, and disposes the context', async () => {
  279. const { ctx, agent, persistenceRoot } = await harness([textResponse('final answer')])
  280. const output = await invoke(ctx, ['task'])
  281. expect(output).toEqual({ code: 0, stdout: 'final answer\n', stderr: '' })
  282. expect(agent.status).toBe('disposed')
  283. const files = await readdir(persistenceRoot, { recursive: true })
  284. expect(files.some(file => file.endsWith('.jsonl'))).toBe(true)
  285. })
  286. it('sums usage across tool steps and selects the last text-bearing assistant message', async () => {
  287. const first = { inputTokens: 10, outputTokens: 3, cacheReadTokens: 2, cacheWriteTokens: 1 }
  288. const second = { inputTokens: 7, outputTokens: 5, cacheReadTokens: 4, reasoningTokens: 6 }
  289. const { ctx } = await harness([toolResponse(first), textResponse('done', second)])
  290. const output = await invoke(ctx, ['--output-format', 'json', 'task'])
  291. const result = JSON.parse(output.stdout) as CliResult
  292. expect(output.code).toBe(0)
  293. expect(result).toMatchObject({ type: 'result', success: true, turn: 1, result: 'done', reason: { kind: 'completed' } })
  294. expect(result.usage).toEqual({
  295. inputTokens: 17,
  296. outputTokens: 8,
  297. cacheReadTokens: 6,
  298. cacheWriteTokens: 1,
  299. reasoningTokens: 6,
  300. })
  301. })
  302. it('keeps the prior text when a later assistant message has no text blocks', async () => {
  303. const { ctx } = await harness([
  304. toolResponse({ inputTokens: 1, outputTokens: 1 }),
  305. reasoningResponse('reasoning only'),
  306. ])
  307. const result = await runOneShot(ctx, { task: 'task' })
  308. expect(result.result).toBe('working')
  309. })
  310. it('streams only the correlated main message turn and then the result envelope', async () => {
  311. const { ctx, agent } = await harness([textResponse('streamed')])
  312. const other = ctx.sessions.create(SessionId('unrelated'))
  313. let injected = false
  314. ctx.on('agent/queued', (subject) => {
  315. if (subject !== agent || injected) return
  316. injected = true
  317. agent.inject([{ type: 'text', text: 'startup injection' }], { source: { kind: 'plugin', plugin: 'test' } })
  318. other.append('turn/start', { turn: 1, trigger: { kind: 'injection', source: { kind: 'plugin', plugin: 'test' } } })
  319. other.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  320. })
  321. const output = await invoke(ctx, ['--output-format', 'stream-json', 'task'])
  322. const lines = output.stdout.trimEnd().split('\n').map(line => JSON.parse(line) as Record<string, unknown>)
  323. const events = lines.slice(0, -1).map(line => line['event'] as SessionEvent)
  324. expect(lines.at(-1)).toMatchObject({ type: 'result', success: true, turn: 2, result: 'streamed' })
  325. expect(events[0]).toMatchObject({ type: 'turn/start', data: { turn: 2, trigger: { kind: 'message' } } })
  326. expect(events.at(-1)).toMatchObject({ type: 'turn/end', data: { turn: 2 } })
  327. expect(lines.slice(0, -1).every(line => line['sessionId'] === agent.session.id)).toBe(true)
  328. expect(events.some(event => event.type === 'context/message')).toBe(false)
  329. })
  330. it('emits partial data and a diagnostic for non-completed turns', async () => {
  331. const { ctx } = await harness([textResponse('partial', { inputTokens: 2, outputTokens: 3 }, 'max-tokens')])
  332. const output = await invoke(ctx, ['--output-format', 'json', 'task'])
  333. expect(JSON.parse(output.stdout)).toMatchObject({ success: false, result: 'partial', reason: { kind: 'max-tokens' } })
  334. expect(output.code).toBe(1)
  335. expect(output.stderr).toContain('output-token limit')
  336. })
  337. it('cancels an active turn, emits its durable aborted result, and disposes', async () => {
  338. const { ctx, agent } = await harness(['hang'])
  339. const abort = new AbortController()
  340. let started!: () => void
  341. const running = new Promise<void>((resolveStarted) => { started = resolveStarted })
  342. ctx.on('session/event', (session, event) => {
  343. if (session === agent.session && event.type === 'assistant/chunk') started()
  344. })
  345. const outcome = invoke(ctx, ['--output-format', 'json', 'task'], { signal: abort.signal })
  346. await running
  347. abort.abort('received SIGINT')
  348. const output = await outcome
  349. expect(JSON.parse(output.stdout)).toMatchObject({ success: false, reason: { kind: 'aborted', reason: 'received SIGINT' } })
  350. expect(output.code).toBe(1)
  351. expect(output.stderr).toContain('was aborted: received SIGINT')
  352. expect(agent.status).toBe('disposed')
  353. })
  354. it('contains stream-writer failures, cancels, flushes, and returns the output error', async () => {
  355. const { ctx, agent } = await harness(['hang'])
  356. await expect(runOneShot(ctx, {
  357. task: 'task',
  358. onEvent: () => { throw new Error('stream sink failed') },
  359. })).rejects.toThrow('stream sink failed')
  360. expect(agent.status).toBe('idle')
  361. })
  362. it('handles cancellation before submission, a missing main agent, and final-output failure', async () => {
  363. const early = await harness([textResponse('unused')])
  364. const fakeSignal = {
  365. aborted: true,
  366. reason: undefined,
  367. } as unknown as AbortSignal
  368. await expect(runOneShot(early.ctx, { task: 'task', signal: fakeSignal })).rejects.toThrow('interrupted')
  369. const preBootAbort = new AbortController()
  370. preBootAbort.abort('before boot completed')
  371. const preBoot = await invoke(early.ctx, ['task'], { signal: preBootAbort.signal })
  372. expect(preBoot).toMatchObject({ code: 1, stdout: '' })
  373. expect(preBoot.stderr).toContain('before boot completed')
  374. const empty = new Context()
  375. liveContexts.push(empty)
  376. await expect(runOneShot(empty, { task: 'task' })).rejects.toThrow('exactly one top-level agent')
  377. const final = await harness([textResponse('answer')])
  378. const output = await invoke(final.ctx, ['task'], { failStdout: true })
  379. expect(output.code).toBe(1)
  380. expect(output.stdout).toBe('')
  381. expect(output.stderr).toContain('stdout closed')
  382. expect(final.agent.status).toBe('disposed')
  383. const disposal = await harness([textResponse('answer')])
  384. const disposalOutput = await invoke(disposal.ctx, ['task'], { failDispose: true })
  385. expect(disposalOutput).toMatchObject({ code: 1, stdout: 'answer\n' })
  386. expect(disposalOutput.stderr).toContain('dispose exploded')
  387. })
  388. it('reports disposal failure alongside an earlier run failure', async () => {
  389. const ctx = new Context()
  390. liveContexts.push(ctx)
  391. const output = await invoke(ctx, ['task'], { failDispose: true })
  392. expect(output).toEqual({
  393. code: 1,
  394. stdout: '',
  395. stderr: 'dsh-cli-demo: config must create exactly one top-level agent, found 0\n'
  396. + 'dsh-cli-demo: dispose failed: dispose exploded\n',
  397. })
  398. })
  399. it('cancels startup work and queued work before the correlated turn begins', async () => {
  400. const startup = await harness(['hang'])
  401. let started!: () => void
  402. const running = new Promise<void>((resolveStarted) => { started = resolveStarted })
  403. startup.ctx.on('session/event', (session, event) => {
  404. if (session === startup.agent.session && event.type === 'assistant/chunk') started()
  405. })
  406. startup.agent.send([{ type: 'text', text: 'first' }])
  407. await running
  408. const startupAbort = new AbortController()
  409. const waiting = runOneShot(startup.ctx, { task: 'second', signal: startupAbort.signal })
  410. startupAbort.abort('cancel startup')
  411. await expect(waiting).rejects.toThrow('cancel startup')
  412. await startup.agent.whenIdle()
  413. const queued = await harness([textResponse('unused')])
  414. const queuedAbort = new AbortController()
  415. queued.ctx.on('agent/queued', (agent) => {
  416. if (agent === queued.agent) queuedAbort.abort('cancel queued')
  417. })
  418. await expect(runOneShot(queued.ctx, { task: 'task', signal: queuedAbort.signal })).rejects.toThrow('cancel queued')
  419. await queued.agent.whenIdle()
  420. })
  421. })
  422. describe('formatTurnFailure', () => {
  423. it('diagnoses every durable reason and preserves merge-extensible unknowns', () => {
  424. const cases: [TurnEndReason, string][] = [
  425. [{ kind: 'completed' }, 'completed'],
  426. [{ kind: 'aborted' }, 'was aborted'],
  427. [{ kind: 'aborted', reason: 'stop' }, 'was aborted: stop'],
  428. [{ kind: 'error', step: 2, message: 'bad' }, 'failed at step 2: bad'],
  429. [{ kind: 'disposed' }, 'was disposed'],
  430. [{ kind: 'max-tokens' }, 'output-token limit'],
  431. [{ kind: 'rejected', reason: 'policy' }, 'was rejected: policy'],
  432. [{ kind: 'interrupted' }, 'persistence recovery'],
  433. ]
  434. for (const [reason, expected] of cases) expect(formatTurnFailure(reason)).toContain(expected)
  435. expect(formatTurnFailure({ kind: 'extension' } as unknown as TurnEndReason)).toContain('extension')
  436. })
  437. })