index.ts 19 KB

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