Add the ACP subagent backend: out-of-process delegation (PR3)
The first OUT-OF-PROCESS subagent backend, proving the seam generalizes past the in-process backends. @deepseek-ai/dsh-subagent-acp runs each child agent in a spawned subprocess, driven over the Agent Client Protocol as the CLIENT — the direction-inverted twin of the dsh-acp server bridge. Point the configured command at the acp-agent example and the harness talks to its own process. - Fresh process per run: start spawns, runs one ACP session (initialize → newSession → prompt), dispose kills the subprocess and awaits its exit. - Minimal client stub: advertises no fs/terminal; accumulates agent_message_chunk text as the result output; auto-answers session/request_permission by a configured policy (reject default / allow). No start-time capabilities (an out-of-process child can't enforce the parent's depth/tool-filter); ignores request.parent; injects only `subagents`. - StopReason mapping (end_turn→completed, cancelled→aborted, …); result resolves error/aborted on a child failure, never rejects (seam contract). - Security: credential-shaped ambient env vars are scrubbed; the child's own key is forwarded only via explicit config.env. A spawn-level error (ENOENT) is captured and raced against the ACP drive so a bad command settles error rather than crashing the parent. Testing designed at every tier: keyless integration drives a scripted mock ACP server subprocess (cancellation incl. the pre-newSession race and a torn-pipe-after-cancel, permission auto-answer, non-message updates, spawn failure, HMR, export shape) at 100% coverage; a with-key e2e drives the REAL acp-agent example process (PONG + real file write, verified on disk) — the harness driving itself. Snapshot coverage of an ACP child is deferred as TODO(acp-subagent-replay) (each child is its own process with its own replay). Stayed on @agentclientprotocol/sdk 0.25.1: the proposed 0.28.x bump only deprecates the stable ClientSideConnection/AgentSideConnection API this layer uses (33 sites incl. the server bridge), turning no-deprecated red across code this PR shouldn't rewrite — that fluent-API migration is its own follow-up. The backend needs nothing 0.28.x adds. This completes the subagent seam stack (PR1 interface → PR2 in-process → PR2.5 snapshot infra → PR3 ACP); the seam RFC moves to implemented/, amended.
This commit is contained in:
90
packages/subagent/subagent-acp/src/index.ts
Normal file
90
packages/subagent/subagent-acp/src/index.ts
Normal file
@@ -0,0 +1,90 @@
|
||||
/**
|
||||
* The out-of-process ACP subagent backend: registers a {@link SubagentProvider}
|
||||
* on `ctx.subagents` that runs each child agent in a SPAWNED SUBPROCESS, driven
|
||||
* over the Agent Client Protocol (ACP) as the client. The parent process is the
|
||||
* ACP client; the child is any ACP agent (point the configured command at the
|
||||
* `acp-agent` example to "talk to our own process").
|
||||
*
|
||||
* Unlike the in-process backends (`-spawn`/`-fork`), the child does NOT share
|
||||
* this cordis context — it is a separate process with its own session, model
|
||||
* client, and tools. So this backend injects only `subagents` (no `agents`),
|
||||
* advertises NO start-time capabilities (an out-of-process child cannot enforce
|
||||
* the parent's depth/tool-filter), and ignores `request.parent`.
|
||||
*
|
||||
* Plugin export shape: named `name`/`inject`/`Config`/`apply`, NO default
|
||||
* export (the cordis Loader's `unwrapExports` does `exports.default ?? exports`,
|
||||
* so a stray default would drop the namespace — see docs/postmortem/0001).
|
||||
*
|
||||
* @module @deepseek-ai/dsh-subagent-acp
|
||||
*/
|
||||
|
||||
import type { Context } from 'cordis'
|
||||
import z from 'schemastery'
|
||||
import type { SubagentCapabilities, SubagentProvider, SubagentStartRequest } from '@deepseek-ai/dsh-subagent'
|
||||
import { type AcpRunSpec, type PermissionPolicy, startAcpRun } from './run.ts'
|
||||
|
||||
export const name = 'subagent-acp'
|
||||
export const inject = ['subagents']
|
||||
|
||||
/** Config: how to spawn and drive the child ACP agent process. */
|
||||
export interface Config {
|
||||
/** Provider name on `ctx.subagents` (default `acp`). */
|
||||
providerName: string
|
||||
/** The executable to spawn for each run (the child ACP agent). */
|
||||
command: string
|
||||
/** Arguments passed to {@link command}. */
|
||||
args: string[]
|
||||
/**
|
||||
* Working directory for the child process and its ACP session. Defaults to
|
||||
* the parent process's cwd when omitted.
|
||||
*/
|
||||
cwd?: string
|
||||
/**
|
||||
* How to auto-answer the child's `session/request_permission` prompts:
|
||||
* `reject` (default — decline every prompt) or `allow` (approve via the first
|
||||
* allow-shaped option). The first cut surfaces no prompt to a human.
|
||||
*/
|
||||
permission: PermissionPolicy
|
||||
/**
|
||||
* Extra environment variables for the child process — e.g. the child
|
||||
* harness's own `DEEPSEEK_API_KEY`. Forwarded on top of a credential-scrubbed
|
||||
* copy of the parent env, so an explicit key here reaches the child while
|
||||
* ambient secrets do not leak implicitly.
|
||||
*/
|
||||
env: Record<string, string>
|
||||
}
|
||||
|
||||
export const Config: z<Config> = z.object({
|
||||
providerName: z.string().default('acp'),
|
||||
command: z.string().required(),
|
||||
args: z.array(z.string()).default([]),
|
||||
cwd: z.string(),
|
||||
permission: z.union(['allow', 'reject'] as const).default('reject'),
|
||||
env: z.dict(z.string()).default({}),
|
||||
})
|
||||
|
||||
/**
|
||||
* The ACP provider. Advertises NO start-time capabilities: an out-of-process
|
||||
* child cannot honor `outputSchema`/`maxDepth`/`toolFilter` (the service rejects
|
||||
* a request needing any of them before `start` runs).
|
||||
*/
|
||||
class AcpProvider implements SubagentProvider {
|
||||
readonly capabilities: SubagentCapabilities = { outputSchema: false, depthLimit: false, toolFilter: false }
|
||||
|
||||
constructor(readonly name: string, private readonly config: Config) {}
|
||||
|
||||
start(request: SubagentStartRequest) {
|
||||
const spec: AcpRunSpec = {
|
||||
command: this.config.command,
|
||||
args: this.config.args,
|
||||
cwd: this.config.cwd ?? process.cwd(),
|
||||
permission: this.config.permission,
|
||||
env: this.config.env,
|
||||
}
|
||||
return startAcpRun(request, spec)
|
||||
}
|
||||
}
|
||||
|
||||
export function apply(ctx: Context, config: Config): void {
|
||||
ctx.subagents.registerProvider(new AcpProvider(config.providerName, config))
|
||||
}
|
||||
292
packages/subagent/subagent-acp/src/run.ts
Normal file
292
packages/subagent/subagent-acp/src/run.ts
Normal file
@@ -0,0 +1,292 @@
|
||||
/**
|
||||
* 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, type ChildProcess } 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'
|
||||
|
||||
/**
|
||||
* 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>
|
||||
}
|
||||
|
||||
/**
|
||||
* Credential-shaped ambient env vars are NOT forwarded to the child by default
|
||||
* (the parent harness's own `DEEPSEEK_API_KEY`/secrets must not leak into a
|
||||
* spawned process implicitly). Same pattern as the bash executor. The child
|
||||
* agent needs its OWN credentials to reach a model — those are supplied
|
||||
* explicitly via {@link AcpRunSpec.env}, which is layered on top AFTER the
|
||||
* scrub, so an intended `DEEPSEEK_API_KEY` survives while an incidental
|
||||
* `AWS_SECRET_ACCESS_KEY` does not.
|
||||
*/
|
||||
export const SENSITIVE_ENV_PATTERN = /KEY|SECRET|TOKEN/i
|
||||
|
||||
/** The ambient env minus credential-shaped vars, plus the spec's explicit env. */
|
||||
export function buildChildEnv(extra: Record<string, string>): NodeJS.ProcessEnv {
|
||||
const env: NodeJS.ProcessEnv = {}
|
||||
for (const [key, value] of Object.entries(process.env)) {
|
||||
if (!SENSITIVE_ENV_PATTERN.test(key)) env[key] = value
|
||||
}
|
||||
return { ...env, ...extra }
|
||||
}
|
||||
|
||||
/** Map an ACP {@link StopReason} to a harness {@link SubagentStopReason}. */
|
||||
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). */
|
||||
export function acpContentText(content: AcpContentBlock): string {
|
||||
return content.type === 'text' ? content.text : ''
|
||||
}
|
||||
|
||||
/** Translate the harness prompt blocks into ACP prompt blocks (text only). */
|
||||
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
|
||||
}
|
||||
|
||||
/** Resolve once the child process exits (any code/signal); immediate if gone. */
|
||||
function waitForExit(child: ChildProcess): Promise<void> {
|
||||
if (child.exitCode !== null || child.signalCode !== null) return Promise.resolve()
|
||||
return new Promise<void>(resolve => child.once('exit', () => { resolve() }))
|
||||
}
|
||||
|
||||
/**
|
||||
* 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 (a spawn/transport/RPC error resolves with `stopReason: 'error'`), per
|
||||
* the seam contract. `cancel()` sends `session/cancel`; `dispose()` kills the
|
||||
* subprocess and awaits its exit (quiescent teardown).
|
||||
*/
|
||||
export function startAcpRun(request: SubagentStartRequest, spec: AcpRunSpec): SubagentRun {
|
||||
const id = AgentId(randomUUID())
|
||||
|
||||
// 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'],
|
||||
})
|
||||
// A spawn-level failure (e.g. ENOENT for a bad command) is emitted as an
|
||||
// `error` event, NOT a thrown exception — without a listener Node treats it as
|
||||
// an unhandled error and crashes the parent. Capture it into a promise the
|
||||
// result path races, so a bad command settles `error` like any child failure.
|
||||
const spawnFailed = new Promise<Error>((resolve) => {
|
||||
child.once('error', (err) => { resolve(err) })
|
||||
})
|
||||
|
||||
// Accumulate the child's streamed assistant text — the SubagentResult output.
|
||||
const output: string[] = []
|
||||
// `cancelled` records that a cancel was requested (signal or cancel()), so a
|
||||
// run torn down before the prompt resolves settles `aborted` rather than the
|
||||
// generic error mapping.
|
||||
let 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
|
||||
const requestCancel = (): void => {
|
||||
cancelled = true
|
||||
// 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 })
|
||||
|
||||
const result: Promise<SubagentResult> = (async (): Promise<SubagentResult> => {
|
||||
// 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 }] : []
|
||||
}
|
||||
try {
|
||||
// An already-aborted request never runs the child.
|
||||
if (request.signal?.aborted) {
|
||||
cancelled = true
|
||||
return { output: [], stopReason: 'aborted' }
|
||||
}
|
||||
// Race the ACP drive against a spawn failure: a bad command never speaks
|
||||
// ACP, so `initialize` would hang forever — the spawn `error` event is the
|
||||
// only signal, and a rejected race settles the run `error` via the catch.
|
||||
const driveAcp = async (): Promise<SubagentResult> => {
|
||||
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
|
||||
// A cancel that raced ahead of `newSession` set `cancelled` but could not
|
||||
// send `session/cancel` (no session id yet). Honor it here: settle
|
||||
// `aborted` without ever issuing the prompt, rather than running the child
|
||||
// to completion and ignoring the cancel.
|
||||
if (cancelled) return { output: collectOutput(), stopReason: 'aborted' }
|
||||
const promptResult = await conn.prompt({ sessionId, prompt: toAcpPrompt(request.prompt) })
|
||||
return { output: collectOutput(), stopReason: acpStopReason(promptResult.stopReason) }
|
||||
}
|
||||
return await Promise.race([
|
||||
driveAcp(),
|
||||
spawnFailed.then((err): SubagentResult => { throw err }),
|
||||
])
|
||||
} catch {
|
||||
// The seam contract: result resolves (never rejects) on a child-level
|
||||
// failure. A spawn/transport/RPC error becomes an error/aborted result —
|
||||
// `aborted` if a cancel was requested (the failure is the cancellation
|
||||
// surfacing as a torn pipe / rejected RPC), else a genuine `error`.
|
||||
return { output: collectOutput(), stopReason: cancelled ? 'aborted' : 'error' }
|
||||
}
|
||||
})()
|
||||
|
||||
return {
|
||||
id,
|
||||
result,
|
||||
cancel(_reason?: string): void {
|
||||
requestCancel()
|
||||
},
|
||||
async dispose(): Promise<void> {
|
||||
request.signal?.removeEventListener('abort', onAbort)
|
||||
// Kill the subprocess and AWAIT its exit (quiescent teardown — dispose
|
||||
// must reach quiescence, not merely request it). SIGTERM first; the child
|
||||
// is our own short-lived ACP agent, so a graceful term is enough. Guard
|
||||
// the kill: the process may already be gone.
|
||||
if (child.exitCode === null && child.signalCode === null) {
|
||||
child.kill('SIGTERM')
|
||||
}
|
||||
await waitForExit(child)
|
||||
},
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user