stdio-chat.ts 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403
  1. /**
  2. * The stdio app's readline UI: reads lines from stdin → `agent.send()`/
  3. * `steer()`, and renders the durable transcript to stdout. A UI is "just a
  4. * plugin" — it consumes the `session/event` feed (the assistant token stream,
  5. * turn/step boundaries, tool activity, todos) plus a few `agent/*` control
  6. * events (`agent/status`, `agent/created`/`agent/disposed`) and the `agents`
  7. * service. Dimmed chain-of-thought rendering plus robust piped-stdin EOF→idle
  8. * exit handling, configured via {@link Config}.
  9. *
  10. * An internal module of the stdio app, not a package of its own: the app's
  11. * front-door cluster always includes this UI, and nothing else composes it.
  12. * The export shape stays named `name`/`inject`/`Config`/`apply` — the plugin
  13. * contract the app's `ctx.plugin(uiStdio, …)` mount consumes.
  14. *
  15. * @module @deepseek-ai/dsh-stdio-agent/stdio-chat
  16. */
  17. import { createInterface } from 'node:readline'
  18. import type { Readable, Writable } from 'node:stream'
  19. import type { Context } from 'cordis'
  20. import z from 'schemastery'
  21. import { AgentId } from '@deepseek-ai/dsh-agent'
  22. import {
  23. UserInteractionError,
  24. type AskUserQuestionAnswer,
  25. type AskUserQuestionAnswerItem,
  26. type AskUserQuestionItem,
  27. type AskUserQuestionOption,
  28. type AskUserQuestionRequest,
  29. } from '@deepseek-ai/dsh-user-interaction'
  30. export const name = 'ui-stdio'
  31. export const inject = ['agents', 'userInteraction']
  32. /** Serializable plugin configuration (cordis-native, schemastery). */
  33. export interface Config {
  34. /** Banner printed once on start, before the first `> ` prompt. */
  35. welcome?: string
  36. // TODO(fixed-stdio-agent): this app-internal plugin is mounted only for the
  37. // precreated `main` agent; remove configurability and its config-only test.
  38. /** Id of the agent stdin drives (`send`/`steer`) and whose status gates the EOF exit; rendering is global. Defaults to `'main'`. */
  39. agent?: string
  40. }
  41. export const Config: z<Config> = z.object({
  42. welcome: z.string().default('ready.'),
  43. agent: z.string().default('main'),
  44. })
  45. /**
  46. * Process-I/O seam — the side-effecting handles the plugin would otherwise
  47. * reach for as globals. Defaulted to the real `process` streams in
  48. * {@link apply}; injected by tests so the EOF, render, and disposal branches
  49. * are exercised without hijacking globals. Deliberately NOT part of the
  50. * serializable {@link Config} (streams/functions don't belong in YAML config).
  51. */
  52. export interface StdioRuntime {
  53. /** Line source (default `process.stdin`). */
  54. input: Readable
  55. /** Render sink (default `process.stdout`). */
  56. output: Writable
  57. /** Process-exit hook (default `process.exit`); called once on stdin EOF. */
  58. exit: (code: number) => void
  59. }
  60. function isTTYPair(input: Readable, output: Writable): boolean {
  61. return Boolean((input as { isTTY?: boolean }).isTTY && (output as { isTTY?: boolean }).isTTY)
  62. }
  63. interface PendingQuestion {
  64. request: AskUserQuestionRequest
  65. questionIndex: number
  66. answers: AskUserQuestionAnswerItem[]
  67. resolve(answer: AskUserQuestionAnswer): void
  68. reject(error: unknown): void
  69. onAbort: () => void
  70. }
  71. type OptionSelection =
  72. | { kind: 'selected'; options: AskUserQuestionOption[] }
  73. | { kind: 'custom' }
  74. | { kind: 'invalid' }
  75. /**
  76. * The plugin body, parameterized over its I/O runtime. `apply` is the thin
  77. * production wrapper that binds the real `process` streams; tests call this
  78. * directly with fakes. Returns nothing — all registration is via `ctx.on`/
  79. * `ctx.effect`, so fiber disposal tears every listener and the readline
  80. * interface down.
  81. * @param ctx - the context supplying the `agents` service and the event feeds.
  82. * @param config - the plugin config; defaults are re-applied here for direct
  83. * callers that bypass Loader validation.
  84. * @param runtime - the process-I/O seam (line source, render sink, exit hook).
  85. */
  86. export function createStdioChat(ctx: Context, config: Config, runtime: StdioRuntime): void {
  87. // Default here too (not just via schemastery's `.default()`): this helper is
  88. // exported and called directly by tests / programmatic consumers that bypass
  89. // Loader validation, so it must be self-contained rather than trusting the
  90. // cast — `config.welcome as string` would otherwise be `undefined` on `{}`.
  91. const welcome = config.welcome ?? 'ready.'
  92. const agentId = AgentId(config.agent ?? 'main')
  93. const { input, output, exit } = runtime
  94. // Render label lookup: the `turn/start` session event carries only the turn
  95. // number, so to print the short agent id (`[main turn 1]`) we map the
  96. // session's id to its agent's id. The session id is not reliably the agent id
  97. // (a session can be created with an explicit/client-supplied id), so build the
  98. // map from `agent/created` rather than parsing the id string. Seed from the
  99. // registry's current agents first: an agent registered before this plugin
  100. // installed (e.g. the pre-created `main` agent, or any agent surviving an HMR
  101. // reload of just this fiber) already fired its `agent/created`, so the live
  102. // listener alone would miss it and its turns would fall back to the raw
  103. // session id.
  104. const labelBySession = new Map<string, string>()
  105. for (const agent of ctx.agents.list()) labelBySession.set(agent.session.header.id, agent.id)
  106. ctx.on('agent/created', (agent) => { labelBySession.set(agent.session.header.id, agent.id) })
  107. ctx.on('agent/disposed', (agent) => { labelBySession.delete(agent.session.header.id) })
  108. // Transcript rendering off the durable `session/event` feed — the assistant
  109. // token stream, turn/step boundaries, tool activity, and todos all come from
  110. // the one canonical stream (no agent/* mirrors). A single listener over the
  111. // append order keeps `inReasoning` transitions deterministic across chunk and
  112. // boundary events.
  113. let inReasoning = false
  114. ctx.on('session/event', (session, event) => {
  115. if (event.type === 'assistant/chunk') {
  116. const { chunk } = event.data
  117. if (chunk.type === 'reasoning-delta') {
  118. // Dim the chain-of-thought so the final answer stands out.
  119. if (!inReasoning) output.write('\x1B[2m')
  120. inReasoning = true
  121. output.write(chunk.text)
  122. } else if (chunk.type === 'text-delta') {
  123. if (inReasoning) output.write('\x1B[0m\n')
  124. inReasoning = false
  125. output.write(chunk.text)
  126. }
  127. } else if (event.type === 'turn/start') {
  128. const label = labelBySession.get(session.header.id) ?? session.header.id
  129. output.write(`\n[${label} turn ${event.data.turn}] `)
  130. } else if (event.type === 'turn/end') {
  131. if (inReasoning) output.write('\x1B[0m')
  132. inReasoning = false
  133. output.write('\n> ')
  134. } else if (event.type === 'tool/call') {
  135. const { name: toolName, arguments: args } = event.data
  136. if (inReasoning) output.write('\x1B[0m')
  137. inReasoning = false
  138. output.write(`\n [tool call] ${toolName}(${args})`)
  139. } else if (event.type === 'tool/result') {
  140. const { content } = event.data
  141. const text = content.filter(block => block.type === 'text').map(block => block.text).join('')
  142. output.write(`\n [tool result] ${text}\n `)
  143. } else if (event.type === 'todo/write') {
  144. if (inReasoning) output.write('\x1B[0m')
  145. inReasoning = false
  146. const glyph = (status: string): string =>
  147. status === 'completed' ? '[x]' : status === 'in_progress' ? '[~]' : '[ ]'
  148. const lines = event.data.todos.map(todo => ` ${glyph(todo.status)} ${todo.content}`).join('\n')
  149. output.write(`\n [todos]\n${lines}\n `)
  150. }
  151. })
  152. ctx.effect(() => {
  153. const reader = createInterface({ input, output, terminal: isTTYPair(input, output) })
  154. // Piped-input exit, once stdin reaches EOF:
  155. // - If no line ever submitted work (empty stdin, blank-only lines), exit
  156. // immediately — no turn will ever start, so there is nothing to wait
  157. // for. (Gating on an observed 'running' here would hang forever.)
  158. // - If work WAS submitted, exit the next time the agent settles to idle
  159. // AFTER having run. Two subtleties this handles: the loop batches
  160. // several queued messages into ONE turn (one idle), so we don't count
  161. // sends; and agent.send() does NOT synchronously flip status to
  162. // 'running', so requiring an observed 'running' first (`sawRunning`)
  163. // avoids exiting in the gap before the turn starts and dropping work.
  164. let stdinClosed = false
  165. let disposed = false
  166. let submittedWork = false
  167. let sawRunning = false
  168. let exitTimer: ReturnType<typeof setTimeout> | undefined
  169. let activeQuestion: PendingQuestion | undefined
  170. const questionQueue: PendingQuestion[] = []
  171. const maybeExit = (): void => {
  172. if (disposed || !stdinClosed) return
  173. // No work submitted: nothing will ever run, exit straight away.
  174. // Work submitted: wait until a turn has run and the agent is idle.
  175. if (submittedWork) {
  176. if (!sawRunning) return
  177. const agent = ctx.agents.get(agentId)
  178. if (agent && agent.status !== 'idle') return // a turn is still running
  179. }
  180. // Let any final output flush, then exit. The handle is tracked so the
  181. // disposer can cancel it — a dispose within the flush window must not let
  182. // the process exit out from under HMR. Re-entrant `maybeExit` calls (e.g.
  183. // repeated idle signals) coalesce onto the one pending timer.
  184. if (exitTimer !== undefined) {
  185. return // exit already scheduled — coalesce re-entrant calls
  186. }
  187. exitTimer = setTimeout(() => { exit(0) }, 200)
  188. }
  189. const disposeStatusListener = ctx.on('agent/status', (subject, status) => {
  190. if (subject.id !== agentId) return
  191. if (status === 'running') sawRunning = true
  192. if (status === 'idle') maybeExit()
  193. })
  194. const activeQuestionItem = (pending: PendingQuestion): AskUserQuestionItem =>
  195. pending.request.questions[pending.questionIndex] as AskUserQuestionItem
  196. const renderQuestion = (pending: PendingQuestion): void => {
  197. const question = activeQuestionItem(pending)
  198. const options = question.options ?? []
  199. output.write('\n')
  200. output.write(question.header ? `[${question.header}] ${question.question}\n` : `${question.question}\n`)
  201. options.forEach((option, index) => {
  202. output.write(` ${index + 1}. ${option.label}\n`)
  203. if (option.description) output.write(` ${option.description}\n`)
  204. })
  205. output.write('> ')
  206. }
  207. const removeAbortListener = (pending: PendingQuestion): void => {
  208. pending.request.signal?.removeEventListener('abort', pending.onAbort)
  209. }
  210. const startNextQuestion = (): void => {
  211. if (activeQuestion !== undefined) return
  212. const pending = questionQueue.shift()
  213. if (pending === undefined) return
  214. // The queue never contains an aborted pending ask: the seam rejects an
  215. // already-aborted request synchronously, and queued asks attach their
  216. // abort listener before enqueueing.
  217. activeQuestion = pending
  218. renderQuestion(pending)
  219. }
  220. const disposeQuestion = (pending: PendingQuestion): void => {
  221. removeAbortListener(pending)
  222. pending.reject(new UserInteractionError('ask_user_question was interrupted before the user answered', 'ASK_ABORTED'))
  223. }
  224. const disposePendingQuestions = (): void => {
  225. if (activeQuestion !== undefined) {
  226. disposeQuestion(activeQuestion)
  227. activeQuestion = undefined
  228. }
  229. for (const pending of questionQueue.splice(0)) {
  230. disposeQuestion(pending)
  231. }
  232. }
  233. const finishQuestion = (pending: PendingQuestion): void => {
  234. activeQuestion = undefined
  235. removeAbortListener(pending)
  236. pending.resolve({ answers: pending.answers })
  237. output.write('\n')
  238. startNextQuestion()
  239. }
  240. const answerCurrentQuestion = (pending: PendingQuestion, answer: AskUserQuestionAnswerItem): void => {
  241. pending.answers.push(answer)
  242. pending.questionIndex += 1
  243. if (pending.questionIndex >= pending.request.questions.length) {
  244. finishQuestion(pending)
  245. return
  246. }
  247. renderQuestion(pending)
  248. }
  249. const selectedOptions = (text: string, options: AskUserQuestionOption[], multiSelect: boolean): OptionSelection => {
  250. if (text === '') return { kind: 'invalid' }
  251. if (!multiSelect) {
  252. if (!/^\d+$/.test(text)) return { kind: 'custom' }
  253. const selected = options[Number(text) - 1]
  254. return selected === undefined ? { kind: 'invalid' } : { kind: 'selected', options: [selected] }
  255. }
  256. const indices = text.split(/[,\s]+/).filter(Boolean)
  257. if (indices.length === 0) return { kind: 'invalid' }
  258. if (indices.some(part => !/^\d+$/.test(part))) return { kind: 'custom' }
  259. const uniqueIndices = [...new Set(indices)]
  260. const selected = uniqueIndices.map(part => options[Number(part) - 1])
  261. return selected.some(option => option === undefined)
  262. ? { kind: 'invalid' }
  263. : { kind: 'selected', options: selected as AskUserQuestionOption[] }
  264. }
  265. const answerQuestion = (line: string): void => {
  266. const pending = activeQuestion as PendingQuestion
  267. const question = activeQuestionItem(pending)
  268. const text = line.trim()
  269. const options = question.options ?? []
  270. const selection = options.length > 0
  271. ? selectedOptions(text, options, question.multiSelect ?? false)
  272. : { kind: text === '' ? 'invalid' : 'custom' } as OptionSelection
  273. if (selection.kind === 'selected') {
  274. answerCurrentQuestion(pending, { id: question.id, selected: selection.options.map(option => option.label) })
  275. return
  276. }
  277. if (selection.kind === 'custom' && text !== '') {
  278. answerCurrentQuestion(pending, { id: question.id, selected: [], custom: text })
  279. return
  280. }
  281. output.write(options.length > 0
  282. ? 'Please enter one of the option numbers'
  283. + (question.multiSelect ? ' (comma or space separated)' : '')
  284. + ' or a custom answer'
  285. + '.\n> '
  286. : 'Please enter an answer.\n> ')
  287. }
  288. const disposeUserInteractionProvider = ctx.userInteraction.registerProvider({
  289. ask(request) {
  290. if (disposed || stdinClosed) {
  291. return Promise.reject(
  292. new UserInteractionError('ask_user_question cannot be answered because stdin is closed', 'ASK_ABORTED'),
  293. )
  294. }
  295. return new Promise<AskUserQuestionAnswer>((resolve, reject) => {
  296. const pending: PendingQuestion = {
  297. request,
  298. questionIndex: 0,
  299. answers: [],
  300. resolve,
  301. reject,
  302. onAbort: () => {
  303. if (activeQuestion === pending) {
  304. activeQuestion = undefined
  305. disposeQuestion(pending)
  306. startNextQuestion()
  307. return
  308. }
  309. // If it is not active, this listener can only fire while the ask
  310. // remains queued; settled asks remove the listener first.
  311. questionQueue.splice(questionQueue.indexOf(pending), 1)
  312. disposeQuestion(pending)
  313. },
  314. }
  315. request.signal?.addEventListener('abort', pending.onAbort, { once: true })
  316. questionQueue.push(pending)
  317. startNextQuestion()
  318. })
  319. },
  320. })
  321. reader.on('line', (line) => {
  322. if (activeQuestion !== undefined) {
  323. answerQuestion(line)
  324. return
  325. }
  326. const text = line.trim()
  327. if (!text) return
  328. const agent = ctx.agents.get(agentId)
  329. if (!agent) {
  330. ctx.logger.error('ui-stdio: agent "%s" is not running', agentId)
  331. return
  332. }
  333. submittedWork = true
  334. if (agent.status === 'running') {
  335. agent.steer([{ type: 'text', text }])
  336. } else {
  337. agent.send([{ type: 'text', text }])
  338. }
  339. })
  340. reader.on('close', () => {
  341. // Fires for BOTH stdin EOF and plugin disposal (reader.close() below);
  342. // `disposed` guards teardown so HMR/dispose never exits the process.
  343. stdinClosed = true
  344. if (!disposed) disposePendingQuestions()
  345. maybeExit()
  346. })
  347. output.write(`${welcome}\n> `)
  348. return () => {
  349. disposed = true
  350. if (exitTimer !== undefined) clearTimeout(exitTimer)
  351. disposePendingQuestions()
  352. disposeUserInteractionProvider()
  353. disposeStatusListener()
  354. reader.close()
  355. }
  356. }, 'ui-stdio')
  357. }
  358. /**
  359. * Cordis entry point. Binds the real `process` streams and delegates to
  360. * {@link createStdioChat}; the indirection keeps the side-effecting handles out
  361. * of the testable core, which is why the unit suite drives `createStdioChat`
  362. * directly. This thin wrapper is exercised end-to-end by the keyless
  363. * Loader-path e2e smoke in `examples/echo-agent` (the real product entry).
  364. */
  365. /* v8 ignore start -- production stdio wiring; testable core is createStdioChat() (covered), exercised e2e by echo-agent keyless smoke */
  366. export function apply(ctx: Context, config: Config): void {
  367. createStdioChat(ctx, config, {
  368. input: process.stdin,
  369. output: process.stdout,
  370. exit: code => process.exit(code),
  371. })
  372. }
  373. /* v8 ignore stop */