Files
deepseek-harness/packages/subagent/subagent-acp/src/run.ts
T

389 lines
18 KiB
TypeScript

/**
* The out-of-process ACP subagent run driver. Spawns a child agent as a
* subprocess, speaks the Agent Client Protocol (ACP) to it over stdio as the
* CLIENT, drives one session to completion, and shapes the result into a
* {@link SubagentResult}. The mirror image of the server-side bridge in
* `@deepseek-ai/dsh-acp` (which is the ACP *agent* side): here we are the ACP
* *client*, so we CALL `initialize`/`newSession`/`prompt`/`cancel` and we
* IMPLEMENT the `Client` callbacks (`sessionUpdate`, `requestPermission`).
*
* One subprocess per run (fresh-process-per-run): `start` spawns, runs exactly
* one ACP session, and `dispose` kills the subprocess and awaits its exit.
* Persistent-process pooling is a future optimization (see the RFC).
*
* TODO(acp-subagent-replay): snapshot-tier coverage of an ACP child is a
* distinct replay shape — each child is its own PROCESS with its own
* single-agent replay (the child boots under `DSH_SNAPSHOT=replay` with its own
* sessions-root + fixture), unlike the in-process per-session keying in
* `dsh-llm-replay`. Deferred to a follow-up; keyless coverage here is via a
* scripted mock ACP server subprocess, and the with-key e2e drives the real
* `acp-agent` example. See the ACP-subagent-backend RFC.
*
* @module @deepseek-ai/dsh-subagent-acp/run
*/
import { spawn } from 'node:child_process'
import { randomUUID } from 'node:crypto'
import { Readable, Writable } from 'node:stream'
import {
ClientSideConnection,
ndJsonStream,
PROTOCOL_VERSION,
type Agent as AcpAgent,
type Client,
type ContentBlock as AcpContentBlock,
type RequestPermissionRequest,
type RequestPermissionResponse,
type SessionNotification,
type StopReason,
} from '@agentclientprotocol/sdk'
import { AgentId } from '@deepseek-ai/dsh-agent'
import type { ContentBlock } from '@deepseek-ai/dsh-llm'
import type { SubagentResult, SubagentRun, SubagentStartRequest, SubagentStopReason } from '@deepseek-ai/dsh-subagent'
import { buildChildEnv, disposeChildProcess, spawnFailure } from '@deepseek-ai/dsh-subagent-subprocess'
/**
* How the client answers a child's `session/request_permission`. The first cut
* does not surface permission prompts to a human, so every request is
* auto-answered by this fixed policy:
*
* - `reject` — decline every prompt (answer `cancelled`). Safe default: a child
* that asks before a side effect does not get to take it.
* - `allow` — approve every prompt by selecting its first `allow_*` option (or,
* if none is offered, `cancelled`). Use when the child is trusted to act.
*/
export type PermissionPolicy = 'allow' | 'reject'
/** Resolved spawn spec for an ACP child process (no defaults — see Config). */
export interface AcpRunSpec {
/** The executable to spawn (the child ACP agent). */
command: string
/** Arguments passed to {@link command}. */
args: string[]
/** Working directory for the child process AND its ACP session `cwd`. */
cwd: string
/** How to auto-answer the child's permission prompts. */
permission: PermissionPolicy
/**
* Extra environment variables to ADD for the child (e.g. the child harness's
* `DEEPSEEK_API_KEY`). Merged on top of the scrubbed ambient env — see
* {@link buildChildEnv}. A value here is forwarded even if its name matches
* the credential-scrub pattern (an explicit opt-in for the child's own creds).
*/
env: Record<string, string>
/**
* Grace period (ms) for the child's EOF-driven quiesce in
* {@link SubagentRun.dispose} — the window to flush persistence and tear down
* its OWN nested subprocesses before the parent escalates to a signal. The
* plugin fills this from its `disposeEofGraceMs` config.
*/
disposeEofGraceMs: number
/**
* Grace period (ms) between `SIGTERM` and the `SIGKILL` escalation in
* {@link SubagentRun.dispose}. The plugin fills this from its
* `disposeGraceMs` config.
*/
disposeGraceMs: number
/**
* Sink for a child-level failure that the run flattened into a stop reason
* (the seam contract forbids `result` rejecting). The driver calls this with
* the original error and the chosen stop reason so the fault is preserved
* rather than silently lost; the provider wires it to `ctx.logger.warn`.
* A throw from the sink itself is contained — it cannot reject `result`.
* Optional — omitted in a unit test that asserts the stop reason directly.
*/
onError?: (error: Error, stopReason: SubagentStopReason) => void
}
/**
* Default grace for the child's EOF-driven quiesce on dispose (the
* `disposeEofGraceMs` config) — the window for it to flush persistence and tear
* down its OWN nested subprocesses (which may run their own `SIGTERM`→`SIGKILL`
* escalation) before the parent escalates to a signal. Deliberately LARGER than
* {@link DEFAULT_DISPOSE_GRACE_MS}: a cooperative child whose teardown is itself
* waiting on a signal-trapping grandchild (e.g. a bash subprocess in its own ~3s
* SIGTERM→SIGKILL grace) plus a final flush needs MORE than a single
* signal-grace of headroom, or the parent's SIGTERM cuts it off exactly as it
* reaches its own SIGKILL+flush. The child is an arbitrary ACP agent, so this is
* a standalone generous default, NOT derived from any child's internals.
*/
export const DEFAULT_DISPOSE_EOF_GRACE_MS = 6_000
/** Default grace between SIGTERM and SIGKILL on dispose (the `disposeGraceMs` config; mirrors the bash executor). */
export const DEFAULT_DISPOSE_GRACE_MS = 3_000
/**
* Map an ACP {@link StopReason} to a harness {@link SubagentStopReason}.
* @param reason - the terminal reason from the child's `session/prompt` response.
* @returns the harness equivalent; `max_turn_requests` and any unknown future
* variant map to `error`, so an unclean stop is never reported as `completed`.
*/
export function acpStopReason(reason: StopReason): SubagentStopReason {
switch (reason) {
case 'end_turn':
return 'completed'
case 'max_tokens':
return 'max-tokens'
case 'refusal':
return 'refusal'
case 'cancelled':
return 'aborted'
// `max_turn_requests` (the child hit its turn-request budget) has no direct
// harness equivalent and means the task did NOT finish cleanly — surface it
// as a generic failure so the consumer maps it to an isError result rather
// than reporting a partial answer as success.
case 'max_turn_requests':
return 'error'
// ACP StopReason is a closed wire union, but a future SDK could add a
// variant; treat an unknown terminal reason as a failure (never silently
// 'completed').
default:
return 'error'
}
}
/**
* Collect the text of an ACP content block (non-text blocks contribute nothing).
* @param content - the content block off a streamed `agent_message_chunk`.
* @returns the block's text, or `''` for a non-text block.
*/
export function acpContentText(content: AcpContentBlock): string {
return content.type === 'text' ? content.text : ''
}
/**
* Translate the harness prompt blocks into ACP prompt blocks (text only).
* @param prompt - the harness prompt; non-text blocks are dropped.
* @returns the ACP text blocks, in order.
*/
export function toAcpPrompt(prompt: ContentBlock[]): AcpContentBlock[] {
const blocks: AcpContentBlock[] = []
for (const block of prompt) {
if (block.type === 'text') blocks.push({ type: 'text', text: block.text })
}
return blocks
}
/** Normalize an unknown thrown value to an Error (the catch binding is `unknown`). */
function toError(value: unknown): Error {
// The catch only sees rejections from the ACP SDK RPCs and the spawn `error`
// event, which are always `Error`s; the `String(value)` arm is a defensive
// fallback for a non-Error throw that the typed surfaces cannot produce.
/* v8 ignore next */
return value instanceof Error ? value : new Error(String(value))
}
/**
* Start an out-of-process ACP child for `request` and return a {@link SubagentRun}.
*
* Spawns the configured command, wraps its stdio in an ACP `ClientSideConnection`,
* and drives one session: `initialize` → `newSession` → `prompt`. The accumulated
* `agent_message_chunk` text is the result output; the prompt's terminal
* `StopReason` maps to the stop reason. `result` never REJECTS on a child-level
* failure after publication resolves with `stopReason: 'error'`. A spawn,
* initialize, new-session, or pre-publication cancellation failure instead
* rejects only after the process has been reaped. `dispose()` requests ACP
* cancellation, then kills and reaps the subprocess.
* @param request - the start request; its signal is the cancellation channel.
* @param spec - the resolved spawn spec: command/args/cwd, env, permission
* policy, dispose graces, and the optional error sink.
* @returns the ready run handle for the child subprocess.
*/
export async function startAcpRun(request: SubagentStartRequest, spec: AcpRunSpec): Promise<SubagentRun> {
const id = AgentId(randomUUID())
if (request.signal.aborted) throw new Error('subagent request was aborted before the ACP child started')
// Spawn the child ACP agent. stdin = ACP request channel, stdout = ACP
// response channel, stderr = INHERIT so the child's diagnostics surface on the
// parent's stderr (no separate capture to drain — we don't fold child stderr
// into the result; the seam reports only output + stop reason).
const child = spawn(spec.command, spec.args, {
cwd: spec.cwd,
env: buildChildEnv(spec.env),
stdio: ['pipe', 'pipe', 'inherit'],
})
// Same-tick capture (the library's contract): a spawn-level failure (e.g.
// ENOENT for a bad command) is an `error` EVENT that would crash the parent
// unheard; the result path races this promise, so a bad command settles
// `error` like any child failure.
const spawnFailed = spawnFailure(child)
// One memoized quiescence transaction is shared by startup rollback and the
// published run's disposer. Once start fulfills, only the holder can invoke
// it; before fulfillment the provider invokes it on every failure path.
let processDisposal: Promise<void> | undefined
const disposeProcess = (): Promise<void> => (processDisposal ??= disposeChildProcess(child, {
disposeEofGraceMs: spec.disposeEofGraceMs,
disposeGraceMs: spec.disposeGraceMs,
}))
// Accumulate the child's streamed assistant text — the SubagentResult output.
const output: string[] = []
// `cancelled` records that the required signal or disposal requested cancel, so a
// run torn down before the prompt resolves settles `aborted` rather than the
// generic error mapping. Held on a mutable object so the async closures that
// set it (the abort listener) and the IIFE that reads it don't fight TS's
// control-flow narrowing of a bare `let` (which would type the catch-time read
// as always-`false`).
const flags = { cancelled: false }
const makeClient = (_agent: AcpAgent): Client => ({
sessionUpdate(params: SessionNotification): Promise<void> {
const update = params.update
if (update.sessionUpdate === 'agent_message_chunk') {
output.push(acpContentText(update.content))
}
// Other updates (thoughts, tool calls, plans) are consumed but not
// surfaced in this cut — the subagent returns only its final answer.
return Promise.resolve()
},
requestPermission(params: RequestPermissionRequest): Promise<RequestPermissionResponse> {
// Auto-answer by the configured policy. `allow` selects the first
// allow-shaped option the child offered; if it offered none (or we
// reject), answer `cancelled` so the child does not proceed.
if (spec.permission === 'allow') {
const allow = params.options.find(o => o.kind === 'allow_once' || o.kind === 'allow_always')
if (allow !== undefined) {
return Promise.resolve({ outcome: { outcome: 'selected', optionId: allow.optionId } })
}
}
return Promise.resolve({ outcome: { outcome: 'cancelled' } })
},
})
const conn = new ClientSideConnection(
makeClient,
ndJsonStream(
Writable.toWeb(child.stdin) as WritableStream<Uint8Array>,
Readable.toWeb(child.stdout) as ReadableStream<Uint8Array>,
),
)
let sessionId: string | undefined
// Resolves when a cancel is requested, so `result` can settle `aborted` even
// if the child never cooperates with `session/cancel` (it ignores the notify,
// or the prompt wedges). The result path races this against the ACP drive: the
// FIRST to settle wins, so signal/dispose cancellation always honors the contract (`result`
// settles `aborted`) without waiting on a non-cooperative child. `dispose`
// still kills the process and reaps it; this only unblocks `result`. The
// executor runs synchronously, so `signalCancelSettled` is assigned before the
// Promise constructor returns (the `!` asserts the definite assignment).
let signalCancelSettled!: () => void
const cancelSettled = new Promise<void>((resolve) => { signalCancelSettled = resolve })
const requestCancel = (): void => {
if (flags.cancelled) return
flags.cancelled = true
signalCancelSettled()
// Best-effort: tell the child to cancel the in-flight turn. Swallows a
// rejection — the session may not exist yet, or the pipe may be gone; the
// dispose path kills the process regardless. If the session has NOT been
// created yet (cancel raced ahead of `newSession`), the `cancelled` flag
// alone carries it: the result path re-checks the flag after each await and
// settles `aborted` without running the prompt. The `.catch` swallow is
// defensive for a narrow transport race (child gone mid-send) — v8-ignored
// because dispose kills the process regardless, so it can't be hit in tests.
/* v8 ignore next */
if (sessionId !== undefined) void conn.cancel({ sessionId }).catch(() => { /* child gone / no session */ })
}
const onAbort = (): void => { requestCancel() }
request.signal.addEventListener('abort', onAbort, { once: true })
// The accumulated child text as harness ContentBlocks (empty array when the
// child streamed nothing). Read at every return so a partial answer survives
// a later cancel/error.
const collectOutput = (): ContentBlock[] => {
const text = output.join('')
return text.length > 0 ? [{ type: 'text', text }] : []
}
// Establish the remote session before publishing a handle. Any failure owns
// the still-private process and therefore reaps it before rejecting.
try {
await Promise.race([
(async (): Promise<void> => {
await conn.initialize({
protocolVersion: PROTOCOL_VERSION,
// Advertise NO optional client capabilities (no fs, no terminal): the
// child self-serves in its own process.
clientCapabilities: {},
})
const session = await conn.newSession({ cwd: spec.cwd, mcpServers: [] })
sessionId = session.sessionId
if (flags.cancelled) throw new Error('subagent cancelled before the ACP session started')
})(),
spawnFailed.then((err): never => { throw err }),
cancelSettled.then((): never => { throw new Error('subagent cancelled before the ACP session started') }),
])
} catch (error: unknown) {
request.signal.removeEventListener('abort', onAbort)
await disposeProcess()
if (flags.cancelled) throw new Error('subagent request was aborted before the ACP child started')
throw toError(error)
}
const result: Promise<SubagentResult> = (async (): Promise<SubagentResult> => {
try {
// Race two post-publication outcomes, first to settle wins:
// - prompt: the normal remote turn;
// - cancelSettled: a cancel was requested — settle `aborted` immediately
// rather than waiting on a child that may ignore `session/cancel` or
// wedge the prompt (`result` settles `aborted`). After `newSession`
// succeeds, transport/process failure rejects the in-flight prompt RPC.
const prompt = async (): Promise<SubagentResult> => {
// The startup phase cannot fulfill without assigning the session id.
const promptResult = await conn.prompt({ sessionId: sessionId as string, prompt: toAcpPrompt(request.prompt) })
return { output: collectOutput(), stopReason: acpStopReason(promptResult.stopReason) }
}
return await Promise.race([
prompt(),
cancelSettled.then((): SubagentResult => ({ output: collectOutput(), stopReason: 'aborted' })),
])
} catch (error: unknown) {
// A deterministic cancellation resolves `cancelSettled` before its
// best-effort ACP cancel can reject the prompt. This fallback is only for
// a process/pipe rejection already queued when the abort event fires; its
// first-outcome ordering cannot be forced without a timing-dependent test.
/* v8 ignore next */
if (flags.cancelled) return { output: collectOutput(), stopReason: 'aborted' }
// The seam contract: result resolves (never rejects) on a child-level
// failure. Startup failures were already rejected before publication;
// every rejection here is a prompt transport/RPC failure.
// Flatten to `error` and surface the original via onError so a real fault
// is preserved rather than silently lost.
try {
spec.onError?.(toError(error), 'error')
} catch {
// Swallows only the caller-supplied sink's OWN throw: an unguarded
// sink exception would reject `result` and break the contract above.
// The child-level failure being reported still settles as `error`.
}
return { output: collectOutput(), stopReason: 'error' }
} finally {
request.signal.removeEventListener('abort', onAbort)
}
})()
let disposal: Promise<void> | undefined
return {
id,
result,
dispose(): Promise<void> {
if (disposal !== undefined) return disposal
request.signal.removeEventListener('abort', onAbort)
requestCancel()
// Quiescent teardown via the shared ladder (stdin EOF → SIGTERM →
// SIGKILL, awaiting the actual exit). For THIS child the EOF tier is the
// one that matters: our acp-agent has NO SIGTERM handler in a normal
// session — it tears down via the server bridge's connection-close path
// (conn.closed → per-agent dispose → final session/flush), driven by the
// stdin EOF, NOT by a signal — and a prompt response can resolve from a
// turn/end BEFORE that post-turn flush lands, so the child still has
// durable work owed when dispose runs (hence the wide EOF grace; see
// DEFAULT_DISPOSE_EOF_GRACE_MS).
disposal = disposeProcess()
return disposal
},
}
}