index.ts 8.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222
  1. /**
  2. * @deepseek-ai/dsh-headless — one-shot direct Agent driver. The bundle patch
  3. * rides over dsh-base without Host, HTTP, or browser plugins; this runner
  4. * creates one Agent through the core registry, drives the task to quiescence,
  5. * streams provider reasoning to stderr, flushes its Session, prints the final
  6. * assistant text to stdout, and exits.
  7. *
  8. * @module @deepseek-ai/dsh-headless
  9. */
  10. import { randomUUID } from 'node:crypto'
  11. import type { Context } from '@deepseek-ai/cordis'
  12. import z from '@deepseek-ai/schemastery'
  13. import { brandString } from '@deepseek-ai/dsh-brand'
  14. import { installModelSelection } from '@deepseek-ai/dsh-agent'
  15. import type { Agent, ModelSelectionRef } from '@deepseek-ai/dsh-agent'
  16. import type {} from '@deepseek-ai/dsh-agent-default-model'
  17. import { createUserMessage } from '@deepseek-ai/dsh-llm'
  18. import { assertNever } from '@deepseek-ai/dsh-util-values'
  19. import type { SessionEvent, SessionId } from '@deepseek-ai/dsh-session'
  20. // Empty type imports carry the loader Context merge for the settlement await
  21. // and the cmdline Context merge for the appExit host value.
  22. import type {} from '@deepseek-ai/cordis-plugin-loader'
  23. import type {} from '@deepseek-ai/dsh-cmdline'
  24. /** Stable Cordis plugin name. */
  25. export const name = 'headless-runner'
  26. /** Core services required before the one-shot turn can start. */
  27. export const inject = ['agentDefaultModel', 'agents', 'sessions']
  28. /** Plugin config: the task resolved from this app's injected provider service. */
  29. export interface Config {
  30. /** The prompt text for the single run. */
  31. task: string
  32. }
  33. export const Config: z<Config> = z.object({
  34. task: z.string().required(),
  35. })
  36. /** Outcome of one owned run interval. */
  37. interface RunOutcome {
  38. text: string
  39. reason: SessionEvent<'turn/end'>['data']['reason'] | undefined
  40. }
  41. /** Process-facing effects of one run: output streams plus the launcher's bounded exit request. */
  42. interface HeadlessIo {
  43. stdout: { write(chunk: string): unknown }
  44. stderr: { write(chunk: string): unknown }
  45. /** Request process exit with `code` after the tree disposes. */
  46. exit(code: number): void
  47. }
  48. /** The process streams the runner writes to; tests substitute captures. */
  49. export const internals: { stdout: HeadlessIo['stdout']; stderr: HeadlessIo['stderr'] } = {
  50. stdout: process.stdout,
  51. stderr: process.stderr,
  52. }
  53. /** Aggregate the last assistant text and turn outcome in one owned interval. */
  54. function summarize(events: readonly SessionEvent[], firstSeq: number): RunOutcome {
  55. let started = false
  56. let text = ''
  57. let reason: SessionEvent<'turn/end'>['data']['reason'] | undefined
  58. for (const event of events) {
  59. if (event.seq < firstSeq) continue
  60. if (event.type === 'turn/start') {
  61. started = true
  62. continue
  63. }
  64. if (!started) continue
  65. if (event.type === 'assistant/message') {
  66. const joined = event.data.message.content
  67. .filter(block => block.type === 'text')
  68. .map(block => block.text)
  69. .join('')
  70. if (joined !== '') text = joined
  71. }
  72. if (event.type === 'turn/end') reason = event.data.reason
  73. }
  74. return { text, reason }
  75. }
  76. /**
  77. * Project provider-reported reasoning from one owned run to stderr as it is
  78. * appended, while keeping final outcome derivation on the durable log.
  79. * @param ctx - plugin context carrying the Session event feed.
  80. * @param agent - the exact Agent whose reasoning belongs to this invocation.
  81. * @param stderr - progress output sink.
  82. * @returns a disposer that also terminates an unterminated reasoning line.
  83. */
  84. function streamReasoning(
  85. ctx: Context,
  86. agent: Agent,
  87. stderr: HeadlessIo['stderr'],
  88. ): () => void {
  89. let started = false
  90. let open = false
  91. let endsWithNewline = true
  92. const close = (): void => {
  93. if (!open) return
  94. if (!endsWithNewline) stderr.write('\n')
  95. open = false
  96. endsWithNewline = true
  97. }
  98. const dispose = ctx.on('session/event', (session, event) => {
  99. if (session !== agent.session) return
  100. if (event.type === 'turn/start') {
  101. close()
  102. started = true
  103. return
  104. }
  105. if (!started || event.type !== 'assistant/chunk') return
  106. const chunk = event.data.chunk
  107. switch (chunk.type) {
  108. case 'reasoning-delta':
  109. if (chunk.text === '') return
  110. if (!open) {
  111. stderr.write('dsh: reasoning:\n')
  112. open = true
  113. }
  114. stderr.write(chunk.text)
  115. endsWithNewline = chunk.text.endsWith('\n')
  116. return
  117. case 'block-start':
  118. if (chunk.blockType !== 'reasoning') close()
  119. return
  120. case 'block-end':
  121. if (chunk.block.type !== 'reasoning') close()
  122. return
  123. case 'usage':
  124. return
  125. case 'text-delta':
  126. case 'tool-call-delta':
  127. case 'finish':
  128. close()
  129. return
  130. /* v8 ignore next -- closed-union exhaustiveness guard */
  131. default:
  132. return assertNever(chunk, 'headless reasoning stream')
  133. }
  134. })
  135. return () => {
  136. dispose()
  137. close()
  138. }
  139. }
  140. /** Report an unexpected direct-driver failure and request a failing exit. */
  141. function fail(io: HeadlessIo, error: unknown): void {
  142. io.stderr.write(`dsh: ${error instanceof Error ? error.message : String(error)}\n`)
  143. io.exit(1)
  144. }
  145. /**
  146. * Run one task through a freshly created Agent and request process exit.
  147. * @param ctx - plugin context carrying the Agent, default model, Session, and launcher IO services.
  148. * @param task - one-shot task text.
  149. * @param io - process-facing effects.
  150. */
  151. async function run(ctx: Context, task: string, io: HeadlessIo): Promise<void> {
  152. // Loader siblings mount concurrently. Await the complete application before
  153. // creating an Agent so its scoped tools and adapters are not half-composed.
  154. await ctx.get('loader')?.await()
  155. const agents = ctx.get('agents')
  156. const defaultModel = ctx.get('agentDefaultModel')
  157. const sessions = ctx.get('sessions')
  158. // Early process shutdown can dispose the tree while settlement is pending.
  159. if (agents === undefined || defaultModel === undefined || sessions === undefined) return
  160. const selection = defaultModel.currentSelection()
  161. // This bundle composes no preset roster, so the model-facing rows sit in the
  162. // host plane and the agent reads them from the global layer. A deployment
  163. // that DOES configure one has to join it here first
  164. // (@deepseek-ai/dsh-agent-presets README, "Composing a child agent").
  165. const { agent } = await agents.create({
  166. sessionId: brandString<SessionId>(`session-${randomUUID()}`),
  167. meta: { cwd: process.cwd() },
  168. agentOptions: { provider: selection.provider, model: selection.model },
  169. setup: (agentCtx) => {
  170. const selected: ModelSelectionRef = { current: selection, assembled: undefined }
  171. installModelSelection(agentCtx, selected)
  172. },
  173. })
  174. await agent.whenIdle()
  175. const firstSeq = agent.session.seq
  176. const stopReasoning = streamReasoning(ctx, agent, io.stderr)
  177. try {
  178. agent.followup(createUserMessage({
  179. content: [{ type: 'text', text: task }],
  180. source: { kind: 'user' },
  181. }))
  182. await agent.whenIdle()
  183. } finally {
  184. stopReasoning()
  185. }
  186. await sessions.flush(agent.session)
  187. const outcome = summarize(agent.session.events, firstSeq)
  188. io.stdout.write(outcome.text + '\n')
  189. if (outcome.reason?.kind === 'error') {
  190. io.stderr.write(`dsh: ${outcome.reason.error.code}: ${outcome.reason.error.message}\n`)
  191. }
  192. io.exit(outcome.reason?.kind === 'completed' ? 0 : 1)
  193. }
  194. /**
  195. * Mount the one-shot direct driver.
  196. * @param ctx - plugin context carrying core services and the launcher-provided exit request.
  197. * @param config - validated task config.
  198. */
  199. export function apply(ctx: Context, config: Config): void {
  200. // Read through the global service store, not the property proxy: appExit is
  201. // an optional host value, never an injected dependency.
  202. const exit = ctx.get('appExit')
  203. if (exit === undefined) {
  204. throw new Error('headless-runner: the launcher must provide ctx.appExit before the tree mounts')
  205. }
  206. const io: HeadlessIo = { stdout: internals.stdout, stderr: internals.stderr, exit }
  207. void run(ctx, config.task, io).catch((error: unknown) => { fail(io, error) })
  208. }