run.ts 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314
  1. /**
  2. * Fresh-process ACP subagent client. Drives one child session and owns cancellation and
  3. * quiescent disposal.
  4. *
  5. * TODO(acp-subagent-replay): add snapshot-tier coverage with a separate replay fixture and
  6. * sessions root inside each child process. Current keyless coverage uses a scripted ACP child;
  7. * with-key coverage drives the real ACP example.
  8. * @module @deepseek-ai/dsh-subagent-acp/run
  9. */
  10. import { spawn } from 'node:child_process'
  11. import { randomUUID } from 'node:crypto'
  12. import { Readable, Writable } from 'node:stream'
  13. import {
  14. ClientSideConnection,
  15. ndJsonStream,
  16. PROTOCOL_VERSION,
  17. type Agent as AcpAgent,
  18. type Client,
  19. type ContentBlock as AcpContentBlock,
  20. type RequestPermissionRequest,
  21. type RequestPermissionResponse,
  22. type SessionNotification,
  23. type StopReason,
  24. } from '@agentclientprotocol/sdk'
  25. import type { ContentBlock } from '@deepseek-ai/dsh-llm'
  26. import { SessionId } from '@deepseek-ai/dsh-session'
  27. import type { SubagentResult, SubagentRun, SubagentStartRequest, SubagentStopReason } from '@deepseek-ai/dsh-subagent'
  28. import { buildChildEnv, disposeChildProcess, spawnFailure } from '@deepseek-ai/dsh-subagent-subprocess'
  29. /** Fixed response to child permission requests: reject by default, or select the first allow option. */
  30. export type PermissionPolicy = 'allow' | 'reject'
  31. /** Resolved spawn spec for an ACP child process (no defaults — see Config). */
  32. export interface AcpRunSpec {
  33. /** The executable to spawn (the child ACP agent). */
  34. command: string
  35. /** Arguments passed to {@link command}. */
  36. args: string[]
  37. /**
  38. * Absolute working directory for the child process AND its ACP session
  39. * `cwd`. The provider resolves it before this spec exists: config override,
  40. * else the delegating parent session's workspace.
  41. */
  42. cwd: string
  43. /** How to auto-answer the child's permission prompts. */
  44. permission: PermissionPolicy
  45. /**
  46. * Extra environment variables to ADD for the child (e.g. the child harness's
  47. * `DEEPSEEK_API_KEY`). Merged on top of the scrubbed ambient env — see
  48. * {@link buildChildEnv}. A value here is forwarded even if its name matches
  49. * the credential-scrub pattern (an explicit opt-in for the child's own creds).
  50. */
  51. env: Record<string, string>
  52. /**
  53. * Grace period (ms) for the child's EOF-driven quiesce in
  54. * {@link SubagentRun.dispose} — the window to flush persistence and tear down
  55. * its OWN nested subprocesses before the parent escalates to a signal. The
  56. * plugin fills this from its `disposeEofGraceMs` config.
  57. */
  58. disposeEofGraceMs: number
  59. /**
  60. * Termination confirmation window (ms) in {@link SubagentRun.dispose}; POSIX applies it after
  61. * `SIGTERM` and `SIGKILL`, while Windows applies it after direct forced termination. The plugin
  62. * fills this from its `disposeGraceMs` config.
  63. */
  64. disposeGraceMs: number
  65. /**
  66. * Sink for a child-level failure that the run flattened into a stop reason
  67. * (the seam contract forbids `result` rejecting). The driver calls this with
  68. * the original error and the chosen stop reason so the fault is preserved
  69. * rather than silently lost; the provider wires it to `ctx.logger.warn`.
  70. * A throw from the sink itself is contained — it cannot reject `result`.
  71. * Optional — omitted in a unit test that asserts the stop reason directly.
  72. */
  73. onError?: (error: Error, stopReason: SubagentStopReason) => void
  74. }
  75. /** EOF grace for child flush and nested-process teardown; wider than the signal grace below. */
  76. export const DEFAULT_DISPOSE_EOF_GRACE_MS = 6_000
  77. /** Default POSIX grace between SIGTERM and SIGKILL on dispose (the `disposeGraceMs` config). */
  78. export const DEFAULT_DISPOSE_GRACE_MS = 3_000
  79. /**
  80. * Map an ACP {@link StopReason} to a harness {@link SubagentStopReason}.
  81. * @param reason - the terminal reason from the child's `session/prompt` response.
  82. * @returns the harness equivalent; `max_turn_requests` and any unknown future
  83. * variant map to `error`, so an unclean stop is never reported as `completed`.
  84. */
  85. export function acpStopReason(reason: StopReason): SubagentStopReason {
  86. switch (reason) {
  87. case 'end_turn':
  88. return 'completed'
  89. case 'max_tokens':
  90. return 'max-tokens'
  91. case 'refusal':
  92. return 'refusal'
  93. case 'cancelled':
  94. return 'aborted'
  95. // `max_turn_requests` (the child hit its turn-request budget) has no direct
  96. // harness equivalent and means the task did NOT finish cleanly — surface it
  97. // as a generic failure so the consumer maps it to an isError result rather
  98. // than reporting a partial answer as success.
  99. case 'max_turn_requests':
  100. return 'error'
  101. // ACP StopReason is a closed wire union, but a future SDK could add a
  102. // variant; treat an unknown terminal reason as a failure (never silently
  103. // 'completed').
  104. default:
  105. return 'error'
  106. }
  107. }
  108. /**
  109. * Collect the text of an ACP content block (non-text blocks contribute nothing).
  110. * @param content - the content block off a streamed `agent_message_chunk`.
  111. * @returns the block's text, or `''` for a non-text block.
  112. */
  113. export function acpContentText(content: AcpContentBlock): string {
  114. return content.type === 'text' ? content.text : ''
  115. }
  116. /**
  117. * Translate the harness prompt blocks into ACP prompt blocks (text only).
  118. * @param prompt - the harness prompt; non-text blocks are dropped.
  119. * @returns the ACP text blocks, in order.
  120. */
  121. export function toAcpPrompt(prompt: ContentBlock[]): AcpContentBlock[] {
  122. const blocks: AcpContentBlock[] = []
  123. for (const block of prompt) {
  124. if (block.type === 'text') blocks.push({ type: 'text', text: block.text })
  125. }
  126. return blocks
  127. }
  128. /** Normalize an unknown thrown value to an Error (the catch binding is `unknown`). */
  129. function toError(value: unknown): Error {
  130. // The catch only sees rejections from the ACP SDK RPCs and the spawn `error`
  131. // event, which are always `Error`s; the `String(value)` arm is a defensive
  132. // fallback for a non-Error throw that the typed surfaces cannot produce.
  133. /* v8 ignore next */
  134. return value instanceof Error ? value : new Error(String(value))
  135. }
  136. /**
  137. * Start and publish one ACP child after initialization and session creation.
  138. * Child failures resolve through the run result; startup failures reject after
  139. * process reap. Disposal cancels, kills, and reaps the child.
  140. * @param request - the start request; its signal is the cancellation channel.
  141. * @param spec - the resolved spawn spec: command/args/cwd, env, permission
  142. * policy, dispose graces, and the optional error sink.
  143. * @returns the ready run handle for the child subprocess.
  144. */
  145. export async function startAcpRun(request: SubagentStartRequest, spec: AcpRunSpec): Promise<SubagentRun> {
  146. if (request.signal.aborted) throw new Error('subagent request was aborted before the ACP child started')
  147. // ACP session ids are unique only within the child server. The lifecycle id
  148. // is minted in the parent namespace so fresh processes cannot collide with
  149. // each other or with a local agent that happens to use the same session id.
  150. const id = SessionId(randomUUID())
  151. // Keep diagnostics on parent stderr; only ACP output contributes to the result.
  152. const child = spawn(spec.command, spec.args, {
  153. cwd: spec.cwd,
  154. env: buildChildEnv(spec.env),
  155. stdio: ['pipe', 'pipe', 'inherit'],
  156. })
  157. // Capture the child-process error event immediately.
  158. const spawnFailed = spawnFailure(child)
  159. // Startup rollback and the published handle share one process teardown.
  160. let processDisposal: Promise<void> | undefined
  161. const disposeProcess = (): Promise<void> => (processDisposal ??= disposeChildProcess(child, {
  162. disposeEofGraceMs: spec.disposeEofGraceMs,
  163. disposeGraceMs: spec.disposeGraceMs,
  164. }))
  165. // Accumulate the child's streamed assistant text — the SubagentResult output.
  166. const output: string[] = []
  167. // Shared mutable state keeps cancellation visible across async closures.
  168. const flags = { cancelled: false }
  169. const makeClient = (_agent: AcpAgent): Client => ({
  170. sessionUpdate(params: SessionNotification): Promise<void> {
  171. const update = params.update
  172. if (update.sessionUpdate === 'agent_message_chunk') {
  173. output.push(acpContentText(update.content))
  174. }
  175. // Other updates (thoughts, tool calls, plans) are consumed but not
  176. // surfaced in this cut — the subagent returns only its final answer.
  177. return Promise.resolve()
  178. },
  179. requestPermission(params: RequestPermissionRequest): Promise<RequestPermissionResponse> {
  180. // Auto-answer by the configured policy. `allow` selects the first
  181. // allow-shaped option the child offered; if it offered none (or we
  182. // reject), answer `cancelled` so the child does not proceed.
  183. if (spec.permission === 'allow') {
  184. const allow = params.options.find(o => o.kind === 'allow_once' || o.kind === 'allow_always')
  185. if (allow !== undefined) {
  186. return Promise.resolve({ outcome: { outcome: 'selected', optionId: allow.optionId } })
  187. }
  188. }
  189. return Promise.resolve({ outcome: { outcome: 'cancelled' } })
  190. },
  191. })
  192. const conn = new ClientSideConnection(
  193. makeClient,
  194. ndJsonStream(
  195. Writable.toWeb(child.stdin) as WritableStream<Uint8Array>,
  196. Readable.toWeb(child.stdout) as ReadableStream<Uint8Array>,
  197. ),
  198. )
  199. let sessionId: string | undefined
  200. // Cancellation settles the result without waiting for a cooperative child.
  201. let signalCancelSettled!: () => void
  202. const cancelSettled = new Promise<void>((resolve) => { signalCancelSettled = resolve })
  203. const requestCancel = (): void => {
  204. if (flags.cancelled) return
  205. flags.cancelled = true
  206. signalCancelSettled()
  207. // Best-effort ACP cancel; process teardown remains authoritative.
  208. /* v8 ignore next */
  209. if (sessionId !== undefined) void conn.cancel({ sessionId }).catch(() => { /* child gone / no session */ })
  210. }
  211. const onAbort = (): void => { requestCancel() }
  212. request.signal.addEventListener('abort', onAbort, { once: true })
  213. // The accumulated child text as harness ContentBlocks (empty array when the
  214. // child streamed nothing). Read at every return so a partial answer survives
  215. // a later cancel/error.
  216. const collectOutput = (): ContentBlock[] => {
  217. const text = output.join('')
  218. return text.length > 0 ? [{ type: 'text', text }] : []
  219. }
  220. // Establish the remote session before publishing a handle. Any failure owns
  221. // the still-private process and therefore reaps it before rejecting.
  222. try {
  223. await Promise.race([
  224. (async (): Promise<void> => {
  225. await conn.initialize({
  226. protocolVersion: PROTOCOL_VERSION,
  227. // Advertise NO optional client capabilities (no fs, no terminal): the
  228. // child self-serves in its own process.
  229. clientCapabilities: {},
  230. })
  231. const session = await conn.newSession({ cwd: spec.cwd, mcpServers: [] })
  232. const returnedSessionId: unknown = Reflect.get(session, 'sessionId')
  233. if (typeof returnedSessionId !== 'string') throw new Error('ACP child published without a session id')
  234. sessionId = returnedSessionId
  235. if (flags.cancelled) throw new Error('subagent cancelled before the ACP session started')
  236. })(),
  237. spawnFailed.then((err): never => { throw err }),
  238. cancelSettled.then((): never => { throw new Error('subagent cancelled before the ACP session started') }),
  239. ])
  240. } catch (error: unknown) {
  241. request.signal.removeEventListener('abort', onAbort)
  242. await disposeProcess()
  243. if (flags.cancelled) throw new Error('subagent request was aborted before the ACP child started')
  244. throw toError(error)
  245. }
  246. // The startup transaction validates the returned id before it can fulfill.
  247. // This assertion carries that cross-closure invariant into TypeScript.
  248. /* v8 ignore next */
  249. if (sessionId === undefined) throw new Error('unreachable: ACP startup fulfilled without a session id')
  250. const remoteSessionId = sessionId
  251. const result: Promise<SubagentResult> = (async (): Promise<SubagentResult> => {
  252. try {
  253. // Race the remote turn against local cancellation.
  254. const prompt = async (): Promise<SubagentResult> => {
  255. // The startup phase cannot fulfill without assigning the session id.
  256. const promptResult = await conn.prompt({ sessionId: remoteSessionId, prompt: toAcpPrompt(request.prompt) })
  257. return { output: collectOutput(), stopReason: acpStopReason(promptResult.stopReason) }
  258. }
  259. return await Promise.race([
  260. prompt(),
  261. cancelSettled.then((): SubagentResult => ({ output: collectOutput(), stopReason: 'aborted' })),
  262. ])
  263. } catch (error: unknown) {
  264. // Cover a process rejection already queued when cancellation arrives.
  265. /* v8 ignore next */
  266. if (flags.cancelled) return { output: collectOutput(), stopReason: 'aborted' }
  267. // Flatten post-publication transport failures while preserving diagnostics.
  268. try {
  269. spec.onError?.(toError(error), 'error')
  270. } catch {
  271. // The diagnostic sink cannot reject the run result.
  272. }
  273. return { output: collectOutput(), stopReason: 'error' }
  274. } finally {
  275. request.signal.removeEventListener('abort', onAbort)
  276. }
  277. })()
  278. let disposal: Promise<void> | undefined
  279. return {
  280. id,
  281. localAgent: undefined,
  282. result,
  283. dispose(): Promise<void> {
  284. if (disposal !== undefined) return disposal
  285. request.signal.removeEventListener('abort', onAbort)
  286. requestCancel()
  287. // The shared platform-aware ladder awaits exit. ACP normally quiesces from
  288. // stdin EOF, including the final flush, so this backend uses a wider EOF
  289. // grace before process termination escalates.
  290. disposal = disposeProcess()
  291. return disposal
  292. },
  293. }
  294. }