cli.ts 18 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470
  1. /**
  2. * Command parser and one-turn driver for `dsh-cli-demo`. The executable wrapper
  3. * owns process signals; this module owns output, durability, and cleanup.
  4. * @module @deepseek-ai/dsh-cli-demo/cli
  5. */
  6. import { parseArgs } from 'node:util'
  7. import type { Context } from 'cordis'
  8. import type { Agent } from '@deepseek-ai/dsh-agent'
  9. import { createUserMessage, type TokenUsage } from '@deepseek-ai/dsh-llm'
  10. import type { SessionEvent, TurnEndReason } from '@deepseek-ai/dsh-session'
  11. import { boot, loadEnv, resolveConfigPath } from '@deepseek-ai/dsh-app-boot'
  12. const CLI_NAME = 'dsh-cli-demo'
  13. const DEFAULT_CONFIG_PATH = './cordis.yml'
  14. const OUTPUT_FORMATS = ['text', 'json', 'stream-json'] as const
  15. const USAGE = `Usage: ${CLI_NAME} [--config path] [--output-format text|json|stream-json] (-p <task> | <task>)\n`
  16. /** Supported CLI output encodings. */
  17. export type OutputFormat = typeof OUTPUT_FORMATS[number]
  18. /** Parsed command: help exits before boot; run carries one validated task. */
  19. export type CliCommand =
  20. | { readonly kind: 'help' }
  21. | {
  22. readonly kind: 'run'
  23. readonly configPath: string
  24. readonly outputFormat: OutputFormat
  25. readonly task: string
  26. }
  27. /** DSH-native final record emitted by JSON modes. */
  28. export interface CliResult {
  29. readonly type: 'result'
  30. readonly success: boolean
  31. readonly sessionId: string
  32. readonly turn: number
  33. readonly result: string
  34. readonly reason: TurnEndReason
  35. readonly usage?: TokenUsage
  36. }
  37. /** Options for one turn against the configured top-level agent. */
  38. export interface OneShotOptions {
  39. /** Exactly one nonblank user task. */
  40. readonly task: string
  41. /** Optional signal that cancels the selected agent. */
  42. readonly signal?: AbortSignal
  43. /** Synchronous task-turn observer; a throw cancels the agent and fails the run after flush. */
  44. readonly onEvent?: (sessionId: string, event: SessionEvent) => void
  45. }
  46. /** Injectable process boundaries used by {@link executeCli}. */
  47. export interface CliRuntime {
  48. /** Process cwd for config resolution and `.env` loading. */
  49. readonly cwd?: string
  50. /** Cancellation signal, normally aborted by SIGINT or SIGTERM. */
  51. readonly signal?: AbortSignal
  52. /** Loader boot boundary. */
  53. readonly boot?: (name: string, absoluteConfigPath: string) => Promise<Context>
  54. /** Optional `.env` loader boundary. */
  55. readonly loadEnv?: (name: string, dir: string, warn: (line: string) => void) => void
  56. /** Stdout sink; throws are treated as output failures. */
  57. readonly writeStdout?: (chunk: string) => unknown
  58. /** Stderr diagnostic sink. */
  59. readonly writeStderr?: (chunk: string) => unknown
  60. /** Context disposal boundary. */
  61. readonly dispose?: (ctx: Context) => Promise<void>
  62. }
  63. interface ParsedArguments {
  64. readonly values: {
  65. readonly config?: string
  66. readonly 'output-format'?: string
  67. readonly help?: boolean
  68. readonly prompt?: string
  69. }
  70. readonly positionals: string[]
  71. }
  72. class CliArgumentError extends Error {
  73. constructor(message: string) {
  74. super(message)
  75. this.name = 'CliArgumentError'
  76. }
  77. }
  78. class CliInterruptedError extends Error {
  79. constructor(reason: string) {
  80. super(reason)
  81. this.name = 'CliInterruptedError'
  82. }
  83. }
  84. /** Render an arbitrary value without trusting its type traps or string coercion. */
  85. function renderUnknown(value: unknown): string {
  86. try {
  87. return String(value)
  88. } catch {
  89. return '[unrenderable thrown value]'
  90. }
  91. }
  92. /** Normalize an arbitrary thrown value without letting inspection escape containment. */
  93. function toError(error: unknown): Error {
  94. try {
  95. if (error instanceof Error) return error
  96. } catch {
  97. // A hostile proxy may throw during instanceof; use the total renderer below.
  98. }
  99. return new Error(renderUnknown(error))
  100. }
  101. function interruptionReason(signal: AbortSignal): string {
  102. return signal.reason === undefined ? 'interrupted' : renderUnknown(signal.reason)
  103. }
  104. /**
  105. * Parse the bin arguments and enforce the one-positional-task contract.
  106. * @param args - arguments after the executable name.
  107. * @returns a help or run command.
  108. * @throws {@link CliArgumentError} for unknown flags, invalid formats, or task cardinality.
  109. */
  110. export function parseCliArgs(args: readonly string[]): CliCommand {
  111. let parsed: ParsedArguments
  112. try {
  113. parsed = parseArgs({
  114. args: [...args],
  115. options: {
  116. config: { type: 'string' },
  117. 'output-format': { type: 'string' },
  118. help: { type: 'boolean' },
  119. prompt: { type: 'string', short: 'p' },
  120. },
  121. allowPositionals: true,
  122. strict: true,
  123. })
  124. } catch (error: unknown) {
  125. throw new CliArgumentError(toError(error).message)
  126. }
  127. if (parsed.values.help === true) return { kind: 'help' }
  128. const prompt = parsed.values.prompt
  129. if (prompt !== undefined && parsed.positionals.length > 0) {
  130. throw new CliArgumentError('-p/--prompt and a positional task are mutually exclusive')
  131. }
  132. if (prompt === undefined && parsed.positionals.length !== 1) {
  133. throw new CliArgumentError(`expected exactly one positional task or -p, received ${parsed.positionals.length} positional(s)`)
  134. }
  135. // Cardinality was checked above, so the fallback index zero exists.
  136. // oxlint-disable-next-line typescript/no-non-null-assertion
  137. const task = prompt ?? parsed.positionals[0]!
  138. if (task.trim().length === 0) throw new CliArgumentError('task must not be blank')
  139. const requestedFormat = parsed.values['output-format'] ?? 'text'
  140. if (!OUTPUT_FORMATS.some(format => format === requestedFormat)) {
  141. throw new CliArgumentError(`unsupported output format ${JSON.stringify(requestedFormat)}`)
  142. }
  143. return {
  144. kind: 'run',
  145. configPath: parsed.values.config ?? DEFAULT_CONFIG_PATH,
  146. outputFormat: requestedFormat as OutputFormat,
  147. task,
  148. }
  149. }
  150. function addUsage(total: TokenUsage | undefined, step: TokenUsage): TokenUsage {
  151. const next: TokenUsage = {
  152. inputTokens: (total?.inputTokens ?? 0) + step.inputTokens,
  153. outputTokens: (total?.outputTokens ?? 0) + step.outputTokens,
  154. }
  155. for (const key of ['cacheReadTokens', 'cacheWriteTokens', 'reasoningTokens'] as const) {
  156. if (total?.[key] !== undefined || step[key] !== undefined) next[key] = (total?.[key] ?? 0) + (step[key] ?? 0)
  157. }
  158. return next
  159. }
  160. function assistantText(event: Extract<SessionEvent, { type: 'assistant/message' }>): string | undefined {
  161. const blocks = event.data.message.content.filter(block => block.type === 'text')
  162. return blocks.length === 0 ? undefined : blocks.map(block => block.text).join('')
  163. }
  164. /** Wait for startup quiescence while making pre-run cancellation terminal. */
  165. async function waitForStartupIdle(agent: Agent, signal?: AbortSignal): Promise<void> {
  166. if (signal === undefined) {
  167. await agent.whenIdle()
  168. return
  169. }
  170. if (signal.aborted) {
  171. agent.cancel({ kind: 'user' })
  172. throw new CliInterruptedError(interruptionReason(signal))
  173. }
  174. await new Promise<void>((resolve, reject) => {
  175. const onAbort = (): void => {
  176. agent.cancel({ kind: 'user' })
  177. reject(new CliInterruptedError(interruptionReason(signal)))
  178. }
  179. signal.addEventListener('abort', onAbort, { once: true })
  180. void agent.whenIdle().then(resolve, reject).finally(() => {
  181. signal.removeEventListener('abort', onAbort)
  182. })
  183. })
  184. }
  185. /**
  186. * Run one message-triggered turn on the configured top-level agent, aggregate its
  187. * final text and model usage, wait for idle plus an explicit persistence flush,
  188. * and return its durable ending. Only the selected agent's task turn reaches
  189. * `onEvent`; startup injections and unrelated sessions are ignored. The context
  190. * must contain exactly one top-level agent. Signal abort cancels that agent; an
  191. * abort before the correlated task turn rejects. An observer throw cancels the
  192. * turn and is rethrown after the agent reaches idle and the session flushes.
  193. * @param ctx - settled Loader root containing one agent plus `ctx.sessions`.
  194. * @param options - task, optional cancellation, and optional stream observer.
  195. * @returns the DSH-native result envelope after durable quiescence.
  196. */
  197. export async function runOneShot(ctx: Context, options: OneShotOptions): Promise<CliResult> {
  198. const agents = ctx.get('agents')?.roots() ?? []
  199. const [agent] = agents
  200. if (agent === undefined || agents.length !== 1) {
  201. throw new Error(`config must create exactly one top-level agent, found ${agents.length}`)
  202. }
  203. await waitForStartupIdle(agent, options.signal)
  204. let targetTurn: number | undefined
  205. let reason: TurnEndReason | undefined
  206. let result = ''
  207. const usageByStep = new Map<string, TokenUsage>()
  208. let outputError: Error | undefined
  209. let resolveTurn!: () => void
  210. let rejectTurn!: (error: Error) => void
  211. let firstTurnEnded = false
  212. const turnEnded = new Promise<void>((resolve, reject) => {
  213. resolveTurn = resolve
  214. rejectTurn = reject
  215. })
  216. const settleResolved = (): void => {
  217. if (firstTurnEnded) return
  218. firstTurnEnded = true
  219. resolveTurn()
  220. }
  221. const settleRejected = (error: Error): void => {
  222. // The once-registered abort listener is the only rejecter, and a settled
  223. // prompt makes targetTurn defined so onAbort skips rejection entirely;
  224. // kept for symmetry with settleResolved.
  225. /* v8 ignore next -- unreachable second settlement, see above */
  226. if (firstTurnEnded) return
  227. firstTurnEnded = true
  228. rejectTurn(error)
  229. }
  230. const observe = (sessionId: string, event: SessionEvent): void => {
  231. if (outputError !== undefined || options.onEvent === undefined) return
  232. try {
  233. options.onEvent(sessionId, event)
  234. } catch (error: unknown) {
  235. outputError = toError(error)
  236. agent.cancel({ kind: 'user' })
  237. }
  238. }
  239. const disposeListener = ctx.on('session/event', (session, event) => {
  240. if (session !== agent.session) return
  241. if (targetTurn === undefined) {
  242. if (event.type !== 'turn/start' || event.data.trigger.kind !== 'message') return
  243. targetTurn = event.data.turn
  244. } else if (event.type === 'turn/start' && event.data.trigger.kind === 'retry'
  245. && reason?.kind === 'error') {
  246. targetTurn = event.data.turn
  247. reason = undefined
  248. }
  249. observe(session.id, event)
  250. if (event.type === 'assistant/chunk'
  251. && event.data.turn === targetTurn
  252. && event.data.chunk.type === 'usage') {
  253. usageByStep.set(`${event.data.turn}/${event.data.step}`, event.data.chunk.usage)
  254. }
  255. if (event.type === 'assistant/message' && event.data.turn === targetTurn) {
  256. result = assistantText(event) ?? result
  257. if (event.data.usage !== undefined) {
  258. usageByStep.set(`${event.data.turn}/${event.data.step}`, event.data.usage)
  259. }
  260. }
  261. if (event.type === 'turn/end' && event.data.turn === targetTurn) {
  262. reason = event.data.reason
  263. settleResolved()
  264. }
  265. })
  266. const signal = options.signal
  267. let onAbort: (() => void) | undefined
  268. if (signal !== undefined) {
  269. onAbort = (): void => {
  270. agent.cancel({ kind: 'user' })
  271. if (targetTurn === undefined) settleRejected(new CliInterruptedError(interruptionReason(signal)))
  272. }
  273. signal.addEventListener('abort', onAbort, { once: true })
  274. /* v8 ignore next -- closes the race between startup-idle completion and listener registration */
  275. if (signal.aborted) onAbort()
  276. }
  277. try {
  278. /* v8 ignore next -- skips send only when cancellation wins the listener-registration race above */
  279. if (!firstTurnEnded) { // oxlint-disable-line typescript/no-unnecessary-condition
  280. agent.followup(createUserMessage({ content: [{ type: 'text', text: options.task }], source: { kind: 'user' } }))
  281. }
  282. await turnEnded
  283. } finally {
  284. if (onAbort !== undefined) signal?.removeEventListener('abort', onAbort)
  285. await agent.whenIdle()
  286. disposeListener()
  287. }
  288. /* v8 ignore next 3 -- turnEnded resolves only from the matching branch that assigns both values */
  289. if (targetTurn === undefined || reason === undefined) {
  290. throw new Error('task ended without a correlated turn/end event')
  291. }
  292. await ctx.sessions.flush(agent.session)
  293. if (outputError !== undefined) throw outputError
  294. const usage = [...usageByStep.values()].reduce<TokenUsage | undefined>(addUsage, undefined)
  295. return {
  296. type: 'result',
  297. success: reason.kind === 'completed',
  298. sessionId: agent.session.id,
  299. turn: targetTurn,
  300. result,
  301. reason,
  302. ...usage === undefined ? {} : { usage },
  303. }
  304. }
  305. function renderResult(outputFormat: OutputFormat, result: CliResult): string {
  306. return outputFormat === 'text' ? `${result.result}\n` : `${JSON.stringify(result)}\n`
  307. }
  308. /**
  309. * Race Loader boot with cancellation without abandoning a context that becomes
  310. * available after the caller has been released. Waiting for that late context
  311. * would recreate the signal hang, so its disposal and diagnostics run detached.
  312. */
  313. async function bootInterruptibly(
  314. start: () => Promise<Context>,
  315. signal: AbortSignal | undefined,
  316. disposeLateContext: (ctx: Context) => Promise<void>,
  317. reportLateDisposalFailure: (error: unknown) => void,
  318. ): Promise<Context> {
  319. if (signal === undefined) return await start()
  320. if (signal.aborted) throw new CliInterruptedError(interruptionReason(signal))
  321. let onAbort!: () => void
  322. const interruptedBoot = new Promise<never>((_resolve, reject) => {
  323. onAbort = (): void => {
  324. reject(new CliInterruptedError(interruptionReason(signal)))
  325. }
  326. signal.addEventListener('abort', onAbort, { once: true })
  327. /* v8 ignore next -- closes registration against a non-standard synchronously mutating signal */
  328. if (signal.aborted) onAbort()
  329. })
  330. const booting = Promise.resolve().then(start)
  331. try {
  332. return await Promise.race([booting, interruptedBoot])
  333. } catch (error: unknown) {
  334. // The awaited race permits the signal to change after the preflight check.
  335. // oxlint-disable-next-line typescript/no-unnecessary-condition
  336. if (signal.aborted) {
  337. void booting.then(
  338. async (lateContext) => {
  339. try {
  340. await disposeLateContext(lateContext)
  341. } catch (error: unknown) {
  342. reportLateDisposalFailure(error)
  343. }
  344. },
  345. () => {},
  346. )
  347. }
  348. throw error
  349. } finally {
  350. signal.removeEventListener('abort', onAbort)
  351. }
  352. }
  353. /**
  354. * Render a non-completed turn reason for stderr.
  355. * @param reason - durable turn ending to describe.
  356. * @returns a concise diagnostic fragment.
  357. */
  358. export function formatTurnFailure(reason: TurnEndReason): string {
  359. switch (reason.kind) {
  360. case 'completed': return 'completed'
  361. case 'aborted': return 'was aborted'
  362. case 'error': return `failed at step ${reason.step}: ${'failure' in reason ? reason.failure.message : reason.message}`
  363. case 'disposed': return 'was disposed'
  364. case 'max-tokens': return 'reached the model output-token limit'
  365. case 'interrupted': return 'was interrupted during persistence recovery'
  366. default: return `ended with ${JSON.stringify(reason)}`
  367. }
  368. }
  369. /**
  370. * Execute one CLI invocation. Argument and boot failures never write stdout;
  371. * context disposal is awaited before return, and its failure does not replace
  372. * an earlier diagnostic.
  373. * @param args - arguments after the executable name.
  374. * @param runtime - optional injected process boundaries for tests and embedding.
  375. * @returns the ordinary process exit code; the thin bin overrides it for Unix signals.
  376. */
  377. export async function executeCli(args: readonly string[], runtime: CliRuntime = {}): Promise<number> {
  378. /* v8 ignore next -- default process sinks are exercised by the built-bin smoke */
  379. const writeStdout = runtime.writeStdout ?? (chunk => process.stdout.write(chunk))
  380. /* v8 ignore next -- default process sinks are exercised by the built-bin smoke */
  381. const writeStderr = runtime.writeStderr ?? (chunk => process.stderr.write(chunk))
  382. let command: CliCommand
  383. try {
  384. command = parseCliArgs(args)
  385. } catch (error: unknown) {
  386. writeStderr(`${CLI_NAME}: ${toError(error).message}\n${USAGE}`)
  387. return 1
  388. }
  389. if (command.kind === 'help') {
  390. writeStdout(USAGE)
  391. return 0
  392. }
  393. /* v8 ignore next -- default process cwd is exercised by the built-bin smoke */
  394. const cwd = runtime.cwd ?? process.cwd()
  395. /* v8 ignore next -- default env/boot boundaries are exercised by the Loader and built-bin smokes */
  396. const loadEnvironment = runtime.loadEnv ?? loadEnv
  397. /* v8 ignore next -- default env/boot boundaries are exercised by the Loader and built-bin smokes */
  398. const bootContext = runtime.boot ?? boot
  399. /* v8 ignore next -- default disposal is exercised by the built-bin smoke */
  400. const disposeContext = runtime.dispose ?? (target => target.fiber.dispose())
  401. let ctx: Context | undefined
  402. let exitCode = 1
  403. let diagnostic: string | undefined
  404. try {
  405. loadEnvironment(CLI_NAME, cwd, line => writeStderr(line))
  406. ctx = await bootInterruptibly(
  407. () => bootContext(CLI_NAME, resolveConfigPath(command.configPath, undefined, cwd)),
  408. runtime.signal,
  409. disposeContext,
  410. error => writeStderr(`${CLI_NAME}: dispose after interrupted boot failed: ${toError(error).message}\n`),
  411. )
  412. const result = await runOneShot(ctx, {
  413. task: command.task,
  414. ...runtime.signal === undefined ? {} : { signal: runtime.signal },
  415. ...command.outputFormat === 'stream-json'
  416. ? { onEvent: (sessionId: string, event: SessionEvent) => {
  417. writeStdout(`${JSON.stringify({ type: 'session_event', sessionId, event })}\n`)
  418. } }
  419. : {},
  420. })
  421. writeStdout(renderResult(command.outputFormat, result))
  422. exitCode = result.success ? 0 : 1
  423. if (!result.success) diagnostic = `${CLI_NAME}: turn ${result.turn} ${formatTurnFailure(result.reason)}\n`
  424. } catch (error: unknown) {
  425. diagnostic = `${CLI_NAME}: ${toError(error).message}\n`
  426. } finally {
  427. if (ctx !== undefined) {
  428. try {
  429. await disposeContext(ctx)
  430. } catch (error: unknown) {
  431. diagnostic = `${diagnostic ?? ''}${CLI_NAME}: dispose failed: ${toError(error).message}\n`
  432. exitCode = 1
  433. }
  434. }
  435. }
  436. if (diagnostic !== undefined) writeStderr(diagnostic)
  437. return exitCode
  438. }