subagent: implement structured output for in-process backends

The seam vocabulary (SubagentStartRequest.outputSchema, SubagentResult
.structured) existed but no in-process backend honored it — spawn/fork
advertised outputSchema: false. This lands the missing half:

- dsh-tools gains a structured-output JSON Schema subset (json-schema.ts):
  StructuredOutputSchema, assertSupportedOutputSchema (rejects loud outside
  the enforced subset, every violation listed), validateStructuredValue
  (path-qualified issues, total). outputSchema's seam type becomes this raw
  JSON-Schema subset instead of the author-facing SchemaSpec DSL — the schema
  travels verbatim to the model as a forced tool's parameters.
- dsh-subagent-inprocess gains the shared structured runtime: one global
  structured_output capture tool (placeholder parameters) + a prepend:true
  agent/request listener doing FINAL-REQUEST enforcement (strip for plain
  agents, per-run schema for structured children — survives downstream
  request-replacing listeners) + an agent/turn-continuation veto that stops
  a child's turn once captured (no wasted extra model step). Lifetime is
  refcounted by backends (plugin lifetime) AND live runs (start→settle).
- startInProcessRun drives the capture: subset asserted before the child
  exists, instruction appended to the child's system prompt, clean-finish
  nudge loop (structuredNudgeRetries, backend Config, default 1), captured
  value on result.structured; a clean finish without a capture settles
  'error' (never a silent success with a missing field).
- spawn + fork flip outputSchema: true and inject 'tools'.
This commit is contained in:
Tianyi Cui
2026-07-05 11:35:39 +08:00
parent 2bad139ece
commit dafb81be7b
26 changed files with 1402 additions and 65 deletions

View File

@@ -12,12 +12,13 @@ The seam this rides on: `CreateAgentOptions.seed` (added on `dsh-agent`, threade
## Capabilities
`{ outputSchema: false, depthLimit: true, toolFilter: false }` — identical to spawn (the depth/model/output behavior is the shared driver's).
`{ outputSchema: true, depthLimit: true, toolFilter: false }` — identical to spawn (the depth/model/structured-output behavior is the shared driver's).
## Config
| Key | Meaning |
|---|---|
| `providerName` | Registry name on `ctx.subagents` (default `fork`). |
| `structuredNudgeRetries` | How many times a structured run re-prompts a child that finished cleanly without calling `structured_output` (default 1). |
See [`dsh-subagent-spawn`](../subagent-spawn/README.md) for the run lifecycle, model inheritance, and depth tracking — all shared.

View File

@@ -25,19 +25,25 @@ import z from 'schemastery'
import type { SessionEvent } from '@deepseek-ai/dsh-session'
import type { Agent } from '@deepseek-ai/dsh-agent'
import type { SubagentCapabilities, SubagentProvider, SubagentStartRequest } from '@deepseek-ai/dsh-subagent'
import { startInProcessRun } from '@deepseek-ai/dsh-subagent-inprocess'
import { acquireStructuredRuntime, startInProcessRun } from '@deepseek-ai/dsh-subagent-inprocess'
export const name = 'subagent-fork'
export const inject = ['subagents', 'agents']
export const inject = ['subagents', 'agents', 'tools']
/** Config: the registry name to register the provider under. */
/** Config: the registry name to register the provider under, plus structured-run tuning. */
export interface Config {
/** Provider name on `ctx.subagents` (default `fork`). */
providerName: string
/**
* How many times a structured run re-prompts a child that finished cleanly
* without calling `structured_output` before giving up (default 1).
*/
structuredNudgeRetries: number
}
export const Config: z<Config> = z.object({
providerName: z.string().default('fork'),
structuredNudgeRetries: z.natural().default(1),
})
/**
@@ -57,18 +63,24 @@ export function completedTurnPrefix(parent: Agent): SessionEvent[] {
}
/**
* The fork provider. Supports `depthLimit`; NOT `outputSchema`/`toolFilter` this
* cut (the service rejects a request needing either before `start` runs).
* The fork provider. Supports `depthLimit` and `outputSchema` (via the shared
* in-process structured runtime); NOT `toolFilter` this cut (the service
* rejects a request needing it before `start` runs).
*/
class ForkProvider implements SubagentProvider {
readonly capabilities: SubagentCapabilities = { outputSchema: false, depthLimit: true, toolFilter: false }
readonly capabilities: SubagentCapabilities = { outputSchema: true, depthLimit: true, toolFilter: false }
constructor(readonly name: string, private readonly ctx: Context) {}
constructor(
readonly name: string,
private readonly ctx: Context,
private readonly structuredNudgeRetries: number,
) {}
start(request: SubagentStartRequest) {
const seed = completedTurnPrefix(request.parent)
return startInProcessRun(this.ctx, request, {
providerName: this.name,
structuredNudgeRetries: this.structuredNudgeRetries,
// Only pass a seed when there's a completed turn to inherit; an empty seed
// is equivalent to a fresh child, so omit it to keep the session unseeded.
...seed.length > 0 ? { seed } : {},
@@ -77,5 +89,12 @@ class ForkProvider implements SubagentProvider {
}
export function apply(ctx: Context, config: Config): void {
ctx.subagents.registerProvider(new ForkProvider(config.providerName, ctx))
// Hold the structured runtime for the plugin's lifetime (see the spawn
// backend — same two-level lifetime: backends for availability, runs for
// mid-run survival across a backend unload).
ctx.effect(() => {
const acquisition = acquireStructuredRuntime(ctx)
return () => { acquisition.release() }
}, 'subagent-fork structured runtime')
ctx.subagents.registerProvider(new ForkProvider(config.providerName, ctx, config.structuredNudgeRetries))
}

View File

@@ -30,8 +30,8 @@ async function setup(script: Script) {
await ctx.plugin(Invariants)
await ctx.plugin(AgentLoop, { agents: [] })
await ctx.plugin(SubagentService)
await ctx.plugin(Spawn, { providerName: 'spawn' })
await ctx.plugin(fork, { providerName: 'fork' })
await ctx.plugin(Spawn, { providerName: 'spawn', structuredNudgeRetries: 1 })
await ctx.plugin(fork, { providerName: 'fork', structuredNudgeRetries: 1 })
ctx.llm.registerAdapter(['mock'], new MockAdapter(script))
const parent = ctx.agentLoop.create(AgentId('parent'), { model: 'mock' })
return { ctx, parent }

View File

@@ -37,7 +37,7 @@ async function setup(script: Script) {
await ctx.plugin(Invariants)
await ctx.plugin(AgentLoop, { agents: [] })
await ctx.plugin(SubagentService)
await ctx.plugin(fork, { providerName: 'fork' })
await ctx.plugin(fork, { providerName: 'fork', structuredNudgeRetries: 1 })
ctx.llm.registerAdapter(['mock'], new MockAdapter(script))
const parent = ctx.agentLoop.create(AgentId('parent'), { model: 'mock' })
return { ctx, parent }
@@ -161,16 +161,20 @@ describe('dsh-subagent-fork', () => {
await run.dispose()
})
it('advertises depthLimit but not outputSchema/toolFilter', async () => {
it('advertises depthLimit and outputSchema but not toolFilter', async () => {
const { ctx } = await setup([])
expect(ctx.subagents.getProvider('fork')!.capabilities).toEqual({ outputSchema: false, depthLimit: true, toolFilter: false })
expect(ctx.subagents.getProvider('fork')!.capabilities).toEqual({ outputSchema: true, depthLimit: true, toolFilter: false })
})
it('unregisters the provider when its fiber is disposed (HMR safety)', async () => {
const ctx = new Context()
await ctx.plugin(SubagentService)
await ctx.plugin(AgentRegistry)
const fiber = await ctx.plugin(fork, { providerName: 'fork' })
// The backend injects 'tools' for the structured runtime, so the registry
// (and its systemPrompt dependency) must be live for the fiber to activate.
await ctx.plugin(SystemPrompt)
await ctx.plugin(ToolRegistry)
const fiber = await ctx.plugin(fork, { providerName: 'fork', structuredNudgeRetries: 1 })
expect(ctx.subagents.list()).toEqual(['fork'])
await fiber.dispose()
expect(ctx.subagents.list()).toEqual([])
@@ -179,12 +183,12 @@ describe('dsh-subagent-fork', () => {
it('has the namespace-plugin export shape (no stray default)', () => {
expect('default' in fork).toBe(false)
expect(fork.name).toBe('subagent-fork')
expect(fork.inject).toEqual(['subagents', 'agents'])
expect(fork.inject).toEqual(['subagents', 'agents', 'tools'])
const loader = Object.create(Loader.prototype) as Loader
const unwrapped = loader.unwrapExports(fork) as Record<string, unknown>
expect(unwrapped).toBe(fork)
expect(unwrapped.name).toBe('subagent-fork')
expect(unwrapped.inject).toEqual(['subagents', 'agents'])
expect(unwrapped.inject).toEqual(['subagents', 'agents', 'tools'])
expect(typeof unwrapped.apply).toBe('function')
})
})

View File

@@ -8,16 +8,27 @@ The shared **in-process subagent run driver**. A pure library (no provider, no r
Runs a child as a child [`Agent`](../../core/agent) on the same cordis context (`ctx.agents`):
1. computes child depth = `depthOf(parent) + 1`; if `request.maxDepth` is set and exceeded, throws `SubagentDepthError` (the `depthLimit` capability);
2. creates a child via `ctx.agents.create` with a fresh `AgentId`/`SessionId`, the parent's `cwd` + `parentSession` lineage, the optional `options.seed` (fork's completed-turn prefix; omitted for a fresh child), and `agentOptions` (the child inherits the **parent's model** by default — a child with no model can't run — overridable via `request.agentOptions.model`; the system prompt is NOT inherited);
3. drives the one-shot: `child.send(prompt)` then `await child.whenIdle()` (ordering matters — `send` enqueues synchronously, so `whenIdle` observes the queued work and resolves on the child's `running → idle` transition, never before the turn starts);
4. reads the result, scoped to the child's OWN events (everything at or after `seedLength`, so a seeded child that produced no message of its own never returns the seeded parent's last message): the last `assistant/message` content (deep-cloned — the log is frozen) and the last `turn/end.reason` mapped to a `SubagentStopReason`.
1. computes child depth = `depthOf(parent) + 1`; if `request.maxDepth` is set and exceeded, throws `SubagentDepthError` (the `depthLimit` capability); a `request.outputSchema` is asserted against the supported subset (`assertSupportedOutputSchema` from [dsh-tools](../../core/tools/README.md)) before any child exists;
2. creates a child via `ctx.agents.create` with a fresh `AgentId`/`SessionId`, the parent's `cwd` + `parentSession` lineage, the optional `options.seed` (fork's completed-turn prefix; omitted for a fresh child), and `agentOptions` (the child inherits the **parent's model** by default — a child with no model can't run — overridable via `request.agentOptions.model`; the system prompt is NOT inherited; a structured run appends the `structured_output` instruction after the caller's prompt);
3. drives the one-shot: `child.send(prompt)` then `await child.whenIdle()` (ordering matters — `send` enqueues synchronously, so `whenIdle` observes the queued work and resolves on the child's `running → idle` transition, never before the turn starts); a structured child that finished a turn CLEANLY without calling `structured_output` is re-prompted (a nudge — a fresh turn) up to `options.structuredNudgeRetries` times;
4. reads the result, scoped to the child's OWN events (everything at or after `seedLength`, so a seeded child that produced no message of its own never returns the seeded parent's last message): the last `assistant/message` content (deep-cloned — the log is frozen) and the last `turn/end.reason` mapped to a `SubagentStopReason`. A structured run surfaces the captured value as `result.structured`; a structured child that finished cleanly WITHOUT ever capturing settles `error` (a clean finish without the demanded result is a failure, not a success with a missing field).
`dispose()` delegates to `AgentHandle.dispose()` (stop loop → await quiescence → remove session); `cancel()` cancels the child's in-flight turn. A cancel landing before any `turn/end` (the pre-turn window) still settles `aborted`, honoring the cancel contract rather than the generic no-turn `error`.
### `InProcessRunOptions`
`{ providerName: string; seed?: SessionEvent[] }` — the per-backend inputs: the provider name (for error context) and the optional child-session seed.
`{ providerName: string; seed?: SessionEvent[]; structuredNudgeRetries: number }` — the per-backend inputs: the provider name (for error context), the optional child-session seed, and the structured-run nudge budget (REQUIRED, resolved from the backend's validated Config — the driver never fills it with a hidden default).
### Structured output: `acquireStructuredRuntime(ctx): StructuredAcquisition`
The mechanism behind `outputSchema` for in-process children. One globally registered `structured_output` capture tool (its registered parameters are a placeholder) plus two listeners, registered once per root context and shared by every holder:
- an `agent/request` waterfall listener registered `prepend: true` that post-processes `await next()` — **final-request enforcement**: the request that hits the wire never carries `structured_output` for an agent without a structured run, and always carries the run's OWN schema (as the tool's `parameters`) for one that has it. Per-agent shaping lives here because the tool registry and prompt assembly are context-global while schemas differ per concurrent child; cooperative mutate-then-`next()` would not survive a downstream listener returning a replacement request.
- an `agent/turn-continuation` listener that stops a child's turn once its output is captured, so a successful capture doesn't buy a wasted extra model step.
The capture tool validates each call against the run's schema (`validateStructuredValue`) — violations become an `INVALID_ARGS` isError result the model retries in-turn; a valid call records the value.
Lifetime is refcounted with two kinds of holder: each backend acquires for its plugin lifetime (`apply`), and each structured RUN holds its own acquisition from start to settle — so unregistration can never precede a live run's settle, and the runtime disposes only when the last backend AND the last run are gone. `release()` is idempotent per acquisition.
### `depthOf(agent): number`

View File

@@ -26,6 +26,7 @@
"@deepseek-ai/dsh-llm": "^0.0.1",
"@deepseek-ai/dsh-session": "^0.0.1",
"@deepseek-ai/dsh-subagent": "^0.0.1",
"@deepseek-ai/dsh-tools": "^0.0.1",
"cordis": "^4.0.0-rc.6"
},
"devDependencies": {

View File

@@ -18,7 +18,22 @@ import type { Context } from 'cordis'
import { AgentId, type Agent, type AgentHandle, type AgentOptions } from '@deepseek-ai/dsh-agent'
import { SessionId, type SessionEvent, type TurnEndReason } from '@deepseek-ai/dsh-session'
import type { ContentBlock } from '@deepseek-ai/dsh-llm'
import { assertSupportedOutputSchema } from '@deepseek-ai/dsh-tools'
import type { SubagentResult, SubagentRun, SubagentStartRequest, SubagentStopReason } from '@deepseek-ai/dsh-subagent'
import {
acquireStructuredRuntime,
STRUCTURED_OUTPUT_INSTRUCTION,
STRUCTURED_OUTPUT_NUDGE,
type StructuredAcquisition,
} from './structured.ts'
export {
acquireStructuredRuntime,
STRUCTURED_OUTPUT_TOOL,
STRUCTURED_OUTPUT_INSTRUCTION,
STRUCTURED_OUTPUT_NUDGE,
type StructuredAcquisition,
} from './structured.ts'
declare module '@deepseek-ai/dsh-agent' {
interface AgentOptions {
@@ -76,6 +91,13 @@ export interface InProcessRunOptions {
* parent's log (FORK), or `undefined` for a fresh child (SPAWN).
*/
readonly seed?: SessionEvent[]
/**
* How many times a structured run re-prompts a child that finished a turn
* cleanly WITHOUT calling `structured_output` (see the structured module).
* REQUIRED, resolved from the backend's validated Config — per the explicit-
* defaulting rule, the driver never fills it with a hidden fallback.
*/
readonly structuredNudgeRetries: number
}
/**
@@ -98,6 +120,10 @@ export function startInProcessRun(
if (request.maxDepth !== undefined && childDepth > request.maxDepth) {
throw new SubagentDepthError(childDepth, request.maxDepth)
}
// Assert the schema subset BEFORE any child exists (the service has already
// capability-gated; this rejects a schema outside the enforced subset loud).
const schema = request.outputSchema
if (schema !== undefined) assertSupportedOutputSchema(schema)
const childId = AgentId(randomUUID())
// The child's OWN events begin after the seed (fork seeds the parent's
@@ -109,13 +135,24 @@ export function startInProcessRun(
// Inherit the parent's model by default (a child with no model cannot run);
// an explicit `request.agentOptions.model` overrides it. The parent's
// systemPrompt is NOT inherited — a fresh child is a clean specialist unless
// the caller supplies one.
// the caller supplies one. A structured run appends the structured_output
// instruction after whatever prompt the caller supplied.
const callerPrompt = request.agentOptions?.systemPrompt
const systemPrompt = schema === undefined
? callerPrompt
: [callerPrompt, STRUCTURED_OUTPUT_INSTRUCTION].filter(text => text !== undefined && text.length > 0).join('\n\n')
const agentOptions: AgentOptions = {
...request.parent.options.model !== undefined ? { model: request.parent.options.model } : {},
...request.agentOptions,
...systemPrompt !== undefined ? { systemPrompt } : {},
subagentDepth: childDepth,
}
// The structured runtime is held for the WHOLE run (acquired before the child
// exists, released when the result settles), so a backend hot-reload mid-run
// cannot unregister the capture tool out from under this live child.
const structured: StructuredAcquisition | undefined = schema !== undefined ? acquireStructuredRuntime(ctx) : undefined
const handle: AgentHandle = ctx.agents.create({
agentId: childId,
sessionId: SessionId(randomUUID()),
@@ -130,6 +167,7 @@ export function startInProcessRun(
agentOptions,
})
const child = handle.agent
if (structured && schema !== undefined) structured.attach(child, schema)
// Bridge the request's abort signal to the child (the consumer also bridges
// its own exec.signal, but a backend-level bridge keeps the contract local).
@@ -154,9 +192,30 @@ export function startInProcessRun(
if (request.signal?.aborted) return { output: [], stopReason: 'aborted' }
child.send(request.prompt)
await child.whenIdle()
return readResult(child, seedLength, cancelled)
if (structured) {
// Nudge loop: a child that finished a turn CLEANLY without calling
// structured_output gets re-prompted, up to the backend-configured
// retry count. An errored/aborted turn is not nudged — its failure is
// the honest result. (This also covers a cancel: a cancelled turn ends
// `aborted`, and a pre-turn cancel leaves no `turn/end` at all, so
// neither reads `completed`.)
let nudges = options.structuredNudgeRetries
while (
structured.captured(child) === undefined && nudges > 0
&& lastOwnTurnEnd(child, seedLength)?.data.reason.kind === 'completed'
) {
nudges -= 1
child.send([{ type: 'text', text: STRUCTURED_OUTPUT_NUDGE }])
await child.whenIdle()
}
}
return readResult(child, seedLength, cancelled, structured ? { captured: structured.captured(child) } : undefined)
} finally {
request.signal?.removeEventListener('abort', onAbort)
if (structured) {
structured.detach(child)
structured.release()
}
}
})()
@@ -173,6 +232,12 @@ export function startInProcessRun(
}
}
/** The child's OWN last `turn/end` event (events at or after `seedLength`), if any. */
function lastOwnTurnEnd(child: Agent, seedLength: number): SessionEvent<'turn/end'> | undefined {
return child.session.events.slice(seedLength)
.findLast((e): e is SessionEvent<'turn/end'> => e.type === 'turn/end')
}
/**
* Read a settled child's terminal result from its session log, scoped to the
* child's OWN events (everything at or after `seedLength` — fork seeds the
@@ -184,12 +249,29 @@ export function startInProcessRun(
* logged (a cancel landed in the pre-turn window, before any turn ran), the
* run settles `aborted` per the {@link SubagentRun.cancel} contract rather than
* the generic no-turn `error`.
*
* A structured run (`structured` present) additionally reports the captured
* value on {@link SubagentResult.structured}. A structured child that finished
* CLEANLY without ever capturing (the nudges ran out) settles `error` — a clean
* finish without the demanded structured result is a failure, not a success
* with a missing field; a non-`completed` reason keeps its own honest mapping.
*/
function readResult(child: Agent, seedLength: number, cancelled: boolean): SubagentResult {
function readResult(
child: Agent,
seedLength: number,
cancelled: boolean,
structured?: { captured?: { value: unknown } | undefined },
): SubagentResult {
const own = child.session.events.slice(seedLength)
const lastMessage = own.findLast((e): e is SessionEvent<'assistant/message'> => e.type === 'assistant/message')
const lastEnd = own.findLast((e): e is SessionEvent<'turn/end'> => e.type === 'turn/end')
const output: ContentBlock[] = lastMessage ? structuredClone(lastMessage.data.content) : []
if (lastEnd === undefined && cancelled) return { output, stopReason: 'aborted' }
return { output, stopReason: toStopReason(lastEnd?.data.reason) }
const stopReason: SubagentStopReason = lastEnd === undefined && cancelled
? 'aborted'
: toStopReason(lastEnd?.data.reason)
if (structured) {
if (structured.captured) return { output, structured: structured.captured.value, stopReason }
if (stopReason === 'completed') return { output, stopReason: 'error' }
}
return { output, stopReason }
}

View File

@@ -0,0 +1,193 @@
/**
* Structured-output support for the in-process subagent backends: the mechanism
* behind `SubagentStartRequest.outputSchema` for children that run as agents on
* the same context.
*
* The model-facing surface is one globally registered `structured_output` tool
* whose REGISTERED parameters are a placeholder — the real schema is per run.
* Because the tool registry and prompt assembly are context-global while
* schemas differ per child (two concurrent structured runs may carry different
* schemas), per-agent shaping happens on the `agent/request` waterfall with a
* `prepend: true` listener that post-processes `await next()` — FINAL-REQUEST
* enforcement: whatever downstream listeners mutated or replaced, the request
* that hits the wire never carries `structured_output` for an agent without a
* structured run, and always carries the run's OWN schema for one that has it.
* (Cooperative mutate-then-`next()` would not survive a downstream listener
* returning a replacement request — see the waterfall composition caveat in
* docs/architecture.md.)
*
* A companion `agent/turn-continuation` listener stops a child's turn once its
* output is captured — without it, the loop's default "had tool calls ⇒
* continue" buys a wasted extra model step per structured child.
*
* Lifetime is refcounted with two kinds of holder: each backend acquires for
* its plugin lifetime (so the tool exists before any run), and each structured
* RUN acquires from start to settle (so a backend hot-reload mid-run cannot
* unregister the capture tool out from under a live child). Registrations are
* effects on the ROOT context — their natural upper bound is app teardown — and
* the refcount disposes them when the last holder releases.
*
* @module @deepseek-ai/dsh-subagent-inprocess/structured
*/
import type { Context } from 'cordis'
import type { Agent } from '@deepseek-ai/dsh-agent'
import type { ContentBlock, GenerateOptions, ToolSchema } from '@deepseek-ai/dsh-llm'
import type { ContinuationDecision } from '@deepseek-ai/dsh-agent'
import type { ToolExecution } from '@deepseek-ai/dsh-tools'
import { ToolArgsError, validateStructuredValue, type StructuredOutputSchema } from '@deepseek-ai/dsh-tools'
/** The model-facing tool name a structured child must call to finish. */
export const STRUCTURED_OUTPUT_TOOL = 'structured_output'
/** The per-child instruction appended to a structured child's system prompt. */
export const STRUCTURED_OUTPUT_INSTRUCTION
= 'When you have your final answer, you MUST report it by calling the '
+ `\`${STRUCTURED_OUTPUT_TOOL}\` tool with arguments matching its parameter schema exactly. `
+ 'Do not finish with a plain text answer: only the tool call counts as your result.'
/** The nudge sent when a structured child finishes cleanly without calling the tool. */
export const STRUCTURED_OUTPUT_NUDGE
= `You finished without calling \`${STRUCTURED_OUTPUT_TOOL}\`. `
+ `Call \`${STRUCTURED_OUTPUT_TOOL}\` now with your final result matching its parameter schema.`
/** One structured run's state: the schema to enforce and the captured value, once recorded. */
interface RunState {
readonly schema: StructuredOutputSchema
captured?: { value: unknown }
}
/** The per-root-context runtime: run states plus the shared registrations. */
interface StructuredRuntime {
refs: number
readonly states: WeakMap<Agent, RunState>
readonly disposers: (() => void)[]
}
/** One root context ⇒ one runtime (multi-app test isolation). */
const runtimes = new WeakMap<Context, StructuredRuntime>()
/**
* One holder's handle on the shared structured runtime. `release()` is
* idempotent per acquisition; the runtime's registrations are disposed when the
* LAST holder (backend plugin or live run) releases.
*/
export interface StructuredAcquisition {
/** Enforce `schema` on `agent`'s requests and start capturing its `structured_output` call. */
attach(agent: Agent, schema: StructuredOutputSchema): void
/** The captured value, once the child called the tool with valid arguments. */
captured(agent: Agent): { value: unknown } | undefined
/** Stop enforcing/capturing for `agent` (WeakMap-backed; safe to call twice). */
detach(agent: Agent): void
/** Drop this holder's reference (idempotent); the last release unregisters everything. */
release(): void
}
/**
* Acquire the per-root-context structured runtime, registering the capture tool
* and the two waterfall listeners on the FIRST acquisition. See the module doc
* for the enforcement and lifetime design.
* @param ctx - any context of the app; the runtime keys off `ctx.root`.
* @returns this holder's handle (attach/captured/detach + idempotent release).
*/
export function acquireStructuredRuntime(ctx: Context): StructuredAcquisition {
const root: Context = ctx.root
let runtime = runtimes.get(root)
if (!runtime) {
runtime = { refs: 0, states: new WeakMap(), disposers: [] }
runtimes.set(root, runtime)
registerRuntime(root, runtime)
}
runtime.refs += 1
let released = false
return {
attach(agent: Agent, schema: StructuredOutputSchema): void {
runtime.states.set(agent, { schema })
},
captured(agent: Agent): { value: unknown } | undefined {
return runtime.states.get(agent)?.captured
},
detach(agent: Agent): void {
runtime.states.delete(agent)
},
release(): void {
if (released) return
released = true
runtime.refs -= 1
if (runtime.refs > 0) return
runtimes.delete(root)
for (const dispose of runtime.disposers.splice(0)) dispose()
},
}
}
/** Register the capture tool + the two listeners on the root context (first acquire). */
function registerRuntime(root: Context, runtime: StructuredRuntime): void {
// The registered parameters are a PLACEHOLDER: the request listener below
// swaps in the run's real schema per child, and strips the tool entirely for
// every agent without a structured run — so this shape is never model-visible.
runtime.disposers.push(root.tools.register({
name: STRUCTURED_OUTPUT_TOOL,
description:
'Report your final structured result. Call this exactly once, when your answer is complete; '
+ 'the arguments must match this tool\'s parameter schema exactly.',
parameters: { type: 'object', properties: {} },
execute(args: unknown, exec: ToolExecution): Promise<ContentBlock[]> {
const state = exec.agent ? runtime.states.get(exec.agent) : undefined
if (!state) {
// Reachable only if a non-structured agent somehow calls the tool (the
// request listener strips it, so the model never sees it) — fail loud
// rather than capture into nowhere.
throw new Error(`${STRUCTURED_OUTPUT_TOOL} is only available to subagents started with an output schema`)
}
const violations = validateStructuredValue(state.schema, args)
// ToolArgsError → isError result with INVALID_ARGS: the model retries
// within the same turn, exactly like a schema-validated defineTool call.
if (violations.length > 0) throw new ToolArgsError(violations)
state.captured = { value: args }
return Promise.resolve([{ type: 'text', text: 'Structured output recorded.' }])
},
}))
// FINAL-REQUEST enforcement (prepend: true = first registered = OUTERMOST
// wrapper): post-process whatever the downstream listeners and the core
// produced, so a downstream listener returning a replacement request cannot
// leak the tool to other agents or erase the child's schema.
runtime.disposers.push(root.on('agent/request', async function (
this: unknown, agent: Agent, _turn: number, _step: number, _options: GenerateOptions, next: () => Promise<GenerateOptions>,
): Promise<GenerateOptions> {
const final = await next()
const state = runtime.states.get(agent)
if (state) {
const schemaEntry: ToolSchema = {
name: STRUCTURED_OUTPUT_TOOL,
description:
'Report your final structured result. Call this exactly once, when your answer is complete; '
+ 'the arguments must match this tool\'s parameter schema exactly.',
// ToolSchema.parameters is the wire-level JSON Schema object; the
// asserted subset type is structurally exactly that.
parameters: state.schema as unknown as Record<string, unknown>,
}
final.tools = [...(final.tools ?? []).filter(tool => tool.name !== STRUCTURED_OUTPUT_TOOL), schemaEntry]
return final
}
// No structured run: strip the placeholder if present; leave an absent
// tools field absent (an adapter may treat `tools: []` and no tools
// differently on the wire).
if (final.tools?.some(tool => tool.name === STRUCTURED_OUTPUT_TOOL)) {
final.tools = final.tools.filter(tool => tool.name !== STRUCTURED_OUTPUT_TOOL)
}
return final
}, { prepend: true }))
// Stop a structured child's turn once its output is captured: the default
// "had tool calls ⇒ continue" would otherwise buy a wasted extra model step
// after every successful capture.
runtime.disposers.push(root.on('agent/turn-continuation', function (
this: unknown, agent: Agent, _turn: number, _decision: ContinuationDecision, next: () => Promise<ContinuationDecision>,
): Promise<ContinuationDecision> {
if (runtime.states.get(agent)?.captured) return Promise.resolve({ action: 'stop' })
return next()
}))
}

View File

@@ -0,0 +1,393 @@
import { describe, expect, it } from 'vitest'
import { Context } from 'cordis'
import LlmService, { type GenerateOptions } from '@deepseek-ai/dsh-llm'
import SessionStore from '@deepseek-ai/dsh-session'
import SystemPrompt from '@deepseek-ai/dsh-system-prompt'
import ToolRegistry from '@deepseek-ai/dsh-tools'
import AgentRegistry, { AgentId } from '@deepseek-ai/dsh-agent'
import AgentLoop from '@deepseek-ai/dsh-agent-loop'
import * as Invariants from '@deepseek-ai/dsh-invariants'
import SubagentService, { type SubagentStartRequest } from '@deepseek-ai/dsh-subagent'
import type { StructuredOutputSchema } from '@deepseek-ai/dsh-tools'
import { MockAdapter, textResponse, toolCallResponse } from '../../../core/agent-loop/tests/mock-adapter.ts'
import * as spawn from '../../subagent-spawn/src/index.ts'
import * as fork from '../../subagent-fork/src/index.ts'
import {
acquireStructuredRuntime,
STRUCTURED_OUTPUT_INSTRUCTION,
STRUCTURED_OUTPUT_TOOL,
} from '../src/structured.ts'
type Script = ConstructorParameters<typeof MockAdapter>[0]
const SCHEMA: StructuredOutputSchema = {
type: 'object',
properties: { answer: { type: 'number' }, note: { type: 'string' } },
required: ['answer'],
}
/**
* Real loop + scripted mock model + the REAL spawn backend (which acquires the
* structured runtime at apply, exactly as shipped). The mock model script
* drives the child's structured_output calls.
*/
async function setup(script: Script, options?: { nudges?: number; withFork?: boolean }) {
const ctx = new Context()
const adapter = new MockAdapter(script)
await ctx.plugin(LlmService)
await ctx.plugin(SessionStore)
await ctx.plugin(SystemPrompt)
await ctx.plugin(ToolRegistry)
await ctx.plugin(AgentRegistry)
await ctx.plugin(Invariants)
await ctx.plugin(AgentLoop, { agents: [] })
await ctx.plugin(SubagentService)
const fiber = await ctx.plugin(spawn, { providerName: 'spawn', structuredNudgeRetries: options?.nudges ?? 1 })
const forkFiber = options?.withFork
? await ctx.plugin(fork, { providerName: 'fork', structuredNudgeRetries: options?.nudges ?? 1 })
: undefined
ctx.llm.registerAdapter(['mock'], adapter)
const parent = ctx.agentLoop.create(AgentId('parent'), { model: 'mock' })
return { ctx, parent, adapter, fiber, forkFiber }
}
function structuredRequest(parent: SubagentStartRequest['parent'], extra?: Partial<SubagentStartRequest>): SubagentStartRequest {
return { prompt: [{ type: 'text', text: 'produce the answer' }], parent, outputSchema: SCHEMA, ...extra }
}
/** The tool names of one recorded model request. */
function toolNames(request: GenerateOptions): string[] {
return (request.tools ?? []).map(tool => tool.name)
}
describe('in-process structured output', () => {
it('captures a valid structured_output call and surfaces result.structured', async () => {
const { ctx, parent } = await setup([
toolCallResponse('c1', STRUCTURED_OUTPUT_TOOL, { answer: 42, note: 'done' }),
])
const run = ctx.subagents.start('spawn', structuredRequest(parent))
const result = await run.result
expect(result.stopReason).toBe('completed')
expect(result.structured).toEqual({ answer: 42, note: 'done' })
await run.dispose()
})
it('stops the turn after a successful capture — no extra model step is spent', async () => {
const { ctx, parent, adapter } = await setup([
toolCallResponse('c1', STRUCTURED_OUTPUT_TOOL, { answer: 1 }),
textResponse('MUST NOT BE CONSUMED'),
])
const run = ctx.subagents.start('spawn', structuredRequest(parent))
await run.result
// Default continuation would run a second step after the tool call; the
// structured runtime's turn-continuation veto stops the turn instead.
expect(adapter.requests.length).toBe(1)
await run.dispose()
})
it('an invalid call gets an INVALID_ARGS isError result and the model retries in-turn', async () => {
const { ctx, parent } = await setup([
toolCallResponse('c1', STRUCTURED_OUTPUT_TOOL, { answer: 'not-a-number' }),
toolCallResponse('c2', STRUCTURED_OUTPUT_TOOL, { answer: 7 }),
])
const run = ctx.subagents.start('spawn', structuredRequest(parent))
const result = await run.result
expect(result.structured).toEqual({ answer: 7 })
expect(result.stopReason).toBe('completed')
// The child's log carries the isError tool/result for the invalid call.
const child = ctx.agents.get(run.id)!
const results = child.session.events.filter(e => e.type === 'tool/result')
expect(results.length).toBe(2)
expect((results[0]!.data as { isError?: boolean }).isError).toBe(true)
await run.dispose()
})
it('nudges a child that finished cleanly without calling the tool, then captures', async () => {
const { ctx, parent } = await setup([
textResponse('here is my answer in prose'),
toolCallResponse('c1', STRUCTURED_OUTPUT_TOOL, { answer: 3 }),
])
const run = ctx.subagents.start('spawn', structuredRequest(parent))
const result = await run.result
expect(result.structured).toEqual({ answer: 3 })
expect(result.stopReason).toBe('completed')
// The nudge is a real user-visible message in the child's log.
const child = ctx.agents.get(run.id)!
const users = child.session.events.filter(e => e.type === 'user/message')
expect(users.length).toBe(2)
await run.dispose()
})
it('settles error when the nudges run out without a capture', async () => {
const { ctx, parent, adapter } = await setup([
textResponse('prose only'),
textResponse('still prose'),
], { nudges: 1 })
const run = ctx.subagents.start('spawn', structuredRequest(parent))
const result = await run.result
expect(result.stopReason).toBe('error')
expect(result.structured).toBeUndefined()
expect(adapter.requests.length).toBe(2)
await run.dispose()
})
it('zero nudge retries fails immediately after the first clean prose finish', async () => {
const { ctx, parent, adapter } = await setup([textResponse('prose')], { nudges: 0 })
const run = ctx.subagents.start('spawn', structuredRequest(parent))
const result = await run.result
expect(result.stopReason).toBe('error')
expect(adapter.requests.length).toBe(1)
await run.dispose()
})
it('a child that errored is NOT nudged (its failure is the honest result)', async () => {
// Script exhaustion on the first call → the child turn errors.
const { ctx, parent, adapter } = await setup([], { nudges: 3 })
const run = ctx.subagents.start('spawn', structuredRequest(parent))
const result = await run.result
expect(result.stopReason).toBe('error')
expect(adapter.requests.length).toBe(1)
await run.dispose()
})
it('rejects a schema outside the subset loud, before any child exists', async () => {
const { ctx, parent } = await setup([])
expect(() => ctx.subagents.start('spawn', structuredRequest(parent, {
outputSchema: { type: 'object', oneOf: [] } as unknown as StructuredOutputSchema,
}))).toThrow(/unsupported output schema/)
expect(ctx.agents.get(AgentId('parent'))).toBeDefined()
})
it('appends the structured instruction to the child system prompt (caller prompt preserved)', async () => {
const { ctx, parent } = await setup([toolCallResponse('c1', STRUCTURED_OUTPUT_TOOL, { answer: 1 })])
const run = ctx.subagents.start('spawn', structuredRequest(parent, {
agentOptions: { systemPrompt: 'You are a counter.' },
}))
await run.result
const child = ctx.agents.get(run.id)!
expect(child.options.systemPrompt).toBe(`You are a counter.\n\n${STRUCTURED_OUTPUT_INSTRUCTION}`)
await run.dispose()
})
it('a structured child WITHOUT a caller prompt gets exactly the instruction', async () => {
const { ctx, parent } = await setup([toolCallResponse('c1', STRUCTURED_OUTPUT_TOOL, { answer: 1 })])
const run = ctx.subagents.start('spawn', structuredRequest(parent))
await run.result
const child = ctx.agents.get(run.id)!
expect(child.options.systemPrompt).toBe(STRUCTURED_OUTPUT_INSTRUCTION)
await run.dispose()
})
describe('final-request enforcement (the prepend agent/request listener)', () => {
it('a structured child sees structured_output with ITS schema; a plain agent never sees the tool', async () => {
const { ctx, parent, adapter } = await setup([
// Parent turn (a plain agent): must NOT see the tool.
textResponse('parent answer'),
// Child turn: must see it, with the run's schema.
toolCallResponse('c1', STRUCTURED_OUTPUT_TOOL, { answer: 42 }),
])
parent.send([{ type: 'text', text: 'hello' }])
await parent.whenIdle()
expect(toolNames(adapter.requests[0]!)).not.toContain(STRUCTURED_OUTPUT_TOOL)
const run = ctx.subagents.start('spawn', structuredRequest(parent))
await run.result
const childRequest = adapter.requests[1]!
expect(toolNames(childRequest)).toContain(STRUCTURED_OUTPUT_TOOL)
const entry = childRequest.tools!.find(tool => tool.name === STRUCTURED_OUTPUT_TOOL)!
expect(entry.parameters).toEqual(SCHEMA)
await run.dispose()
})
it('two concurrent structured children each see their OWN schema', async () => {
const otherSchema: StructuredOutputSchema = {
type: 'object',
properties: { verdict: { type: 'string', enum: ['real', 'bogus'] } },
required: ['verdict'],
}
const { ctx, parent, adapter } = await setup([
(options: GenerateOptions) => {
// Answer with whatever schema this child was given — proves each
// request carried the right one regardless of scheduling order.
const entry = options.tools!.find(tool => tool.name === STRUCTURED_OUTPUT_TOOL)!
const args = 'verdict' in (entry.parameters.properties as Record<string, unknown>)
? { verdict: 'real' }
: { answer: 1 }
return toolCallResponse('c1', STRUCTURED_OUTPUT_TOOL, args)
},
(options: GenerateOptions) => {
const entry = options.tools!.find(tool => tool.name === STRUCTURED_OUTPUT_TOOL)!
const args = 'verdict' in (entry.parameters.properties as Record<string, unknown>)
? { verdict: 'real' }
: { answer: 1 }
return toolCallResponse('c2', STRUCTURED_OUTPUT_TOOL, args)
},
])
const runA = ctx.subagents.start('spawn', structuredRequest(parent))
const runB = ctx.subagents.start('spawn', structuredRequest(parent, { outputSchema: otherSchema }))
const [a, b] = await Promise.all([runA.result, runB.result])
expect(a.structured).toEqual({ answer: 1 })
expect(b.structured).toEqual({ verdict: 'real' })
const schemas = adapter.requests.map(request =>
request.tools!.find(tool => tool.name === STRUCTURED_OUTPUT_TOOL)!.parameters)
expect(schemas).toContainEqual(SCHEMA)
expect(schemas).toContainEqual(otherSchema)
await runA.dispose()
await runB.dispose()
})
it('wins against a downstream listener that REPLACES the request object', async () => {
const { ctx, parent, adapter } = await setup([
toolCallResponse('c1', STRUCTURED_OUTPUT_TOOL, { answer: 5 }),
])
// A downstream (non-prepend) listener that returns a brand-new request —
// the composition caveat that erases cooperative mutations. Registered
// AFTER the runtime's prepend listener, so it runs INSIDE it.
ctx.on('agent/request', async (_agent, _turn, _step, _options, next) => {
const replaced = await next()
return { ...replaced, tools: [...(replaced.tools ?? [])] }
})
const run = ctx.subagents.start('spawn', structuredRequest(parent))
const result = await run.result
expect(result.structured).toEqual({ answer: 5 })
const entry = adapter.requests[0]!.tools!.find(tool => tool.name === STRUCTURED_OUTPUT_TOOL)
expect(entry).toBeDefined()
expect(entry!.parameters).toEqual(SCHEMA)
await run.dispose()
})
it('a non-structured agent request keeps tools ABSENT when it had none (no tools: [] materialized)', async () => {
const { parent, adapter } = await setup([
// The registry contributes the placeholder via prompt assembly, so
// tools is an array in the raw request — but after stripping the
// placeholder (its ONLY entry), the field must not be re-added as a
// different shape.
textResponse('plain'),
])
parent.send([{ type: 'text', text: 'q' }])
await parent.whenIdle()
const request = adapter.requests[0]!
expect(toolNames(request)).not.toContain(STRUCTURED_OUTPUT_TOOL)
await new Promise(resolve => setTimeout(resolve, 0))
})
it('handles a request with NO tools field at all, for plain and structured agents alike', async () => {
// Drive the waterfall directly with a toolless request — the enforcement
// listener must tolerate `tools: undefined` on both branches: leave it
// absent for a plain agent, and create the array for a structured child.
const { ctx, parent } = await setup([])
const bare: GenerateOptions = { model: 'mock', messages: [] }
const plain = await ctx.waterfall('agent/request', parent, 1, 1, bare, () => Promise.resolve(bare))
expect(plain.tools).toBeUndefined()
const acquisition = acquireStructuredRuntime(ctx)
acquisition.attach(parent, SCHEMA)
const bare2: GenerateOptions = { model: 'mock', messages: [] }
const shaped = await ctx.waterfall('agent/request', parent, 1, 1, bare2, () => Promise.resolve(bare2))
expect(shaped.tools!.map(tool => tool.name)).toEqual([STRUCTURED_OUTPUT_TOOL])
acquisition.detach(parent)
acquisition.release()
})
})
describe('runtime lifetime (refcount: backends + live runs)', () => {
it('registers the capture tool while a backend is loaded and unregisters when the last unloads', async () => {
const { ctx, fiber, forkFiber } = await setup([], { withFork: true })
expect(ctx.tools.get(STRUCTURED_OUTPUT_TOOL)).toBeDefined()
await fiber.dispose()
// fork still holds a reference.
expect(ctx.tools.get(STRUCTURED_OUTPUT_TOOL)).toBeDefined()
await forkFiber!.dispose()
expect(ctx.tools.get(STRUCTURED_OUTPUT_TOOL)).toBeUndefined()
})
it('a live run-level acquisition keeps the runtime registered after EVERY backend unloads', async () => {
// Simulates the run-holder half of the two-level lifetime: a structured
// run acquires at start and releases at settle, so registration ordering
// is settle-then-unregister even if all backends unload first. (A real
// in-process child dies WITH its backend's fiber — the acquisition's
// observable job is this ordering, which a manual holder pins directly.)
const { ctx, fiber, forkFiber } = await setup([], { withFork: true })
const runHolder = acquireStructuredRuntime(ctx)
await fiber.dispose()
await forkFiber!.dispose()
expect(ctx.tools.get(STRUCTURED_OUTPUT_TOOL)).toBeDefined()
runHolder.release()
expect(ctx.tools.get(STRUCTURED_OUTPUT_TOOL)).toBeUndefined()
})
it('a structured run releases its acquisition when it settles (backend unload mid-run)', async () => {
const { ctx, parent, fiber } = await setup(['hang'])
const run = ctx.subagents.start('spawn', structuredRequest(parent))
// Let the child's step start streaming, then unload the backend. The
// backend owns the child agent, so the unload tears the child down and
// the run settles — releasing its own acquisition on the way out.
await new Promise(resolve => setTimeout(resolve, 30))
await fiber.dispose()
const result = await run.result
expect(result.stopReason).toBe('error')
// Both holders (backend + run) released — nothing keeps the runtime now.
expect(ctx.tools.get(STRUCTURED_OUTPUT_TOOL)).toBeUndefined()
await run.dispose()
})
it('fork children capture structured output through the same runtime', async () => {
const { ctx, parent } = await setup([
toolCallResponse('c1', STRUCTURED_OUTPUT_TOOL, { answer: 9 }),
], { withFork: true })
const run = ctx.subagents.start('fork', structuredRequest(parent))
const result = await run.result
expect(result.structured).toEqual({ answer: 9 })
await run.dispose()
})
it('acquisition release is idempotent (double release cannot underflow the refcount)', async () => {
const ctx = new Context()
await ctx.plugin(SystemPrompt)
await ctx.plugin(ToolRegistry)
const first = acquireStructuredRuntime(ctx)
const second = acquireStructuredRuntime(ctx)
first.release()
first.release()
// The second holder still keeps the tool registered.
expect(ctx.tools.get(STRUCTURED_OUTPUT_TOOL)).toBeDefined()
second.release()
expect(ctx.tools.get(STRUCTURED_OUTPUT_TOOL)).toBeUndefined()
})
it('attach/captured/detach manage per-agent state through the acquisition surface', async () => {
const { ctx, parent } = await setup([])
const acquisition = acquireStructuredRuntime(ctx)
expect(acquisition.captured(parent)).toBeUndefined()
acquisition.attach(parent, SCHEMA)
expect(acquisition.captured(parent)).toBeUndefined()
acquisition.detach(parent)
acquisition.detach(parent)
acquisition.release()
// The backend still holds its own reference from setup().
expect(ctx.tools.get(STRUCTURED_OUTPUT_TOOL)).toBeDefined()
})
})
it('a direct structured_output call from an agent WITHOUT a structured run is an isError', async () => {
const { ctx, parent } = await setup([])
const result = await ctx.tools.execute({
callId: 'x' as never,
name: STRUCTURED_OUTPUT_TOOL,
arguments: { answer: 1 },
agent: parent,
})
expect(result.isError).toBe(true)
expect(result.content[0]).toMatchObject({ type: 'text' })
})
it('a structured_output call with NO calling agent at all is an isError', async () => {
const { ctx } = await setup([])
const result = await ctx.tools.execute({
callId: 'x' as never,
name: STRUCTURED_OUTPUT_TOOL,
arguments: { answer: 1 },
})
expect(result.isError).toBe(true)
})
})

View File

@@ -51,7 +51,7 @@ describe('depthOf', () => {
describe('startInProcessRun', () => {
it('drives a fresh child (no seed) to completion and returns its output', async () => {
const { ctx, parent } = await setup([textResponse('driver child answer')])
const run = startInProcessRun(ctx, { prompt: [{ type: 'text', text: 'do X' }], parent }, { providerName: 'spawn' })
const run = startInProcessRun(ctx, { prompt: [{ type: 'text', text: 'do X' }], parent }, { providerName: 'spawn', structuredNudgeRetries: 1 })
const result = await run.result
expect(result.stopReason).toBe('completed')
expect(text(result.output)).toBe('driver child answer')
@@ -61,7 +61,7 @@ describe('startInProcessRun', () => {
it('throws SubagentDepthError when the child would exceed maxDepth', async () => {
const { ctx, parent } = await setup([])
expect(() => startInProcessRun(ctx, { prompt: [{ type: 'text', text: 'p' }], parent, maxDepth: 0 }, { providerName: 'spawn' }))
expect(() => startInProcessRun(ctx, { prompt: [{ type: 'text', text: 'p' }], parent, maxDepth: 0 }, { providerName: 'spawn', structuredNudgeRetries: 1 }))
.toThrow(SubagentDepthError)
})
@@ -73,7 +73,7 @@ describe('startInProcessRun', () => {
parent.send([{ type: 'text', text: 'parent q' }])
await parent.whenIdle()
const seed = parent.session.events.slice()
const run = startInProcessRun(ctx, { prompt: [{ type: 'text', text: 'child q' }], parent }, { providerName: 'fork', seed })
const run = startInProcessRun(ctx, { prompt: [{ type: 'text', text: 'child q' }], parent }, { providerName: 'fork', structuredNudgeRetries: 1, seed })
const result = await run.result
expect(result.stopReason).toBe('completed')
expect(text(result.output)).toBe('seeded child reply')

View File

@@ -25,6 +25,9 @@
},
{
"path": "../subagent"
},
{
"path": "../../core/tools"
}
]
}

View File

@@ -6,14 +6,15 @@ The run mechanics live in the shared [`@deepseek-ai/dsh-subagent-inprocess`](../
## What it does
`start(request)` delegates to `startInProcessRun(ctx, request, { providerName })` with no seed: a fresh child agent with the parent's `cwd`/`parentSession` lineage and (by default) the parent's model. See the [driver README](../subagent-inprocess/README.md) for the full lifecycle (depth check, one-shot drive, result read, dispose).
`start(request)` delegates to `startInProcessRun(ctx, request, { providerName, structuredNudgeRetries })` with no seed: a fresh child agent with the parent's `cwd`/`parentSession` lineage and (by default) the parent's model. See the [driver README](../subagent-inprocess/README.md) for the full lifecycle (depth check, one-shot drive, result read, dispose).
## Capabilities
`{ outputSchema: false, depthLimit: true, toolFilter: false }`. It constructs the child, so it enforces a recursion cap; structured output and tool-scoping are deferred (the service rejects a request needing either before `start` runs).
`{ outputSchema: true, depthLimit: true, toolFilter: false }`. It constructs the child, so it enforces a recursion cap, and it supports structured output via the driver's shared [structured runtime](../subagent-inprocess/README.md) (the backend acquires it for its plugin lifetime; each structured run holds its own acquisition until it settles). Tool-scoping is deferred (the service rejects a request needing it before `start` runs).
## Config
| Key | Meaning |
|---|---|
| `providerName` | Registry name on `ctx.subagents` (default `spawn`). |
| `structuredNudgeRetries` | How many times a structured run re-prompts a child that finished cleanly without calling `structured_output` (default 1). |

View File

@@ -9,6 +9,11 @@
* ({@link startInProcessRun}); this backend just passes NO seed (a fresh
* child). The fork backend is an independent peer over the same driver.
*
* Structured output (`outputSchema`) is supported via the driver's shared
* structured runtime: the backend acquires it for its plugin lifetime (so the
* capture tool and request-shaping listeners exist before any run), and each
* structured run holds its own acquisition until it settles.
*
* Plugin export shape: named `name`/`inject`/`Config`/`apply`, NO default.
*
* @module @deepseek-ai/dsh-subagent-spawn
@@ -17,38 +22,61 @@
import type { Context } from 'cordis'
import z from 'schemastery'
import type { SubagentCapabilities, SubagentProvider, SubagentStartRequest } from '@deepseek-ai/dsh-subagent'
import { startInProcessRun } from '@deepseek-ai/dsh-subagent-inprocess'
import { acquireStructuredRuntime, startInProcessRun } from '@deepseek-ai/dsh-subagent-inprocess'
export const name = 'subagent-spawn'
export const inject = ['subagents', 'agents']
export const inject = ['subagents', 'agents', 'tools']
/** Config: the registry name to register the provider under. */
/** Config: the registry name to register the provider under, plus structured-run tuning. */
export interface Config {
/** Provider name on `ctx.subagents` (default `spawn`). */
providerName: string
/**
* How many times a structured run re-prompts a child that finished cleanly
* without calling `structured_output` before giving up (default 1).
*/
structuredNudgeRetries: number
}
export const Config: z<Config> = z.object({
providerName: z.string().default('spawn'),
structuredNudgeRetries: z.natural().default(1),
})
/**
* The spawn provider. Supports `depthLimit` (it constructs the child, so it can
* enforce a recursion cap) but NOT `outputSchema` or `toolFilter` in this cut —
* a request that needs either is rejected by the service before `start` runs.
* enforce a recursion cap) and `outputSchema` (via the shared in-process
* structured runtime); NOT `toolFilter` in this cut — a request that needs it
* is rejected by the service before `start` runs.
*/
class SpawnProvider implements SubagentProvider {
readonly capabilities: SubagentCapabilities = { outputSchema: false, depthLimit: true, toolFilter: false }
readonly capabilities: SubagentCapabilities = { outputSchema: true, depthLimit: true, toolFilter: false }
constructor(readonly name: string, private readonly ctx: Context) {}
constructor(
readonly name: string,
private readonly ctx: Context,
private readonly structuredNudgeRetries: number,
) {}
start(request: SubagentStartRequest) {
// Fresh child: no seed. The shared driver mints ids, stamps cwd/lineage/
// depth, drives the one-shot, and maps the result.
return startInProcessRun(this.ctx, request, { providerName: this.name })
// depth, drives the one-shot (including the structured capture/nudge loop
// when the request carries an outputSchema), and maps the result.
return startInProcessRun(this.ctx, request, {
providerName: this.name,
structuredNudgeRetries: this.structuredNudgeRetries,
})
}
}
export function apply(ctx: Context, config: Config): void {
ctx.subagents.registerProvider(new SpawnProvider(config.providerName, ctx))
// Hold the structured runtime for the plugin's lifetime, so the capture tool
// and its request-shaping listeners are registered before the first
// structured run and torn down when the last backend unloads (live runs hold
// their own acquisitions, so an unload mid-run cannot strand a child).
ctx.effect(() => {
const acquisition = acquireStructuredRuntime(ctx)
return () => { acquisition.release() }
}, 'subagent-spawn structured runtime')
ctx.subagents.registerProvider(new SpawnProvider(config.providerName, ctx, config.structuredNudgeRetries))
}

View File

@@ -31,7 +31,7 @@ export async function spawnHarness(workdir: string): Promise<Context> {
await ctx.plugin(LocalBashExecutor, { cwd: workdir, timeoutMs: 30_000 })
await ctx.plugin(ToolBash)
await ctx.plugin(SubagentService)
await ctx.plugin(Spawn, { providerName: 'spawn' })
await ctx.plugin(Spawn, { providerName: 'spawn', structuredNudgeRetries: 1 })
// The model-facing subagent tool, bound to the spawn backend.
await ctx.plugin(ToolSubagent, { provider: 'spawn' })
return ctx

View File

@@ -34,7 +34,7 @@ async function setup(script: Script) {
await ctx.plugin(Invariants)
await ctx.plugin(AgentLoop, { agents: [] })
await ctx.plugin(SubagentService)
await ctx.plugin(spawn, { providerName: 'spawn' })
await ctx.plugin(spawn, { providerName: 'spawn', structuredNudgeRetries: 1 })
ctx.llm.registerAdapter(['mock'], adapter)
const parent = ctx.agentLoop.create(AgentId('parent'), { model: 'mock' })
return { ctx, parent, adapter }
@@ -241,17 +241,21 @@ describe('dsh-subagent-spawn', () => {
await parentHandle.dispose()
})
it('advertises depthLimit but not outputSchema/toolFilter', async () => {
it('advertises depthLimit and outputSchema but not toolFilter', async () => {
const { ctx } = await setup([])
const provider = ctx.subagents.getProvider('spawn')!
expect(provider.capabilities).toEqual({ outputSchema: false, depthLimit: true, toolFilter: false })
expect(provider.capabilities).toEqual({ outputSchema: true, depthLimit: true, toolFilter: false })
})
it('unregisters the provider when its fiber is disposed (HMR safety)', async () => {
const ctx = new Context()
await ctx.plugin(SubagentService)
await ctx.plugin(AgentRegistry)
const fiber = await ctx.plugin(spawn, { providerName: 'spawn' })
// The backend injects 'tools' for the structured runtime, so the registry
// (and its systemPrompt dependency) must be live for the fiber to activate.
await ctx.plugin(SystemPrompt)
await ctx.plugin(ToolRegistry)
const fiber = await ctx.plugin(spawn, { providerName: 'spawn', structuredNudgeRetries: 1 })
expect(ctx.subagents.list()).toEqual(['spawn'])
await fiber.dispose()
expect(ctx.subagents.list()).toEqual([])
@@ -260,12 +264,12 @@ describe('dsh-subagent-spawn', () => {
it('has the namespace-plugin export shape (no stray default)', () => {
expect('default' in spawn).toBe(false)
expect(spawn.name).toBe('subagent-spawn')
expect(spawn.inject).toEqual(['subagents', 'agents'])
expect(spawn.inject).toEqual(['subagents', 'agents', 'tools'])
const loader = Object.create(Loader.prototype) as Loader
const unwrapped = loader.unwrapExports(spawn) as Record<string, unknown>
expect(unwrapped).toBe(spawn)
expect(unwrapped.name).toBe('subagent-spawn')
expect(unwrapped.inject).toEqual(['subagents', 'agents'])
expect(unwrapped.inject).toEqual(['subagents', 'agents', 'tools'])
expect(typeof unwrapped.apply).toBe('function')
})
})

View File

@@ -8,7 +8,7 @@
import type { Agent, AgentId, AgentOptions } from '@deepseek-ai/dsh-agent'
import type { ContentBlock } from '@deepseek-ai/dsh-llm'
import type { SchemaSpec } from '@deepseek-ai/dsh-tools'
import type { StructuredOutputSchema } from '@deepseek-ai/dsh-tools'
/**
* Which START-TIME features a provider supports. Checked by the service
@@ -56,12 +56,16 @@ export interface SubagentStartRequest {
/** Per-child agent options (model, system prompt). */
agentOptions?: AgentOptions
/**
* Optional structured-output schema. When set AND the provider's
* {@link SubagentCapabilities.outputSchema} is `true`, the child's final
* answer is shaped to this schema and surfaced as {@link SubagentResult.structured}.
* Optional structured-output schema — an object-rooted JSON Schema within the
* enforced subset (see `assertSupportedOutputSchema` in dsh-tools; a schema
* outside the subset is rejected loud at start). When set AND the provider's
* {@link SubagentCapabilities.outputSchema} is `true`, the child is driven to
* report a value matching this schema, surfaced as
* {@link SubagentResult.structured}. The schema must be plain host-realm JSON
* data — a caller holding foreign-realm data materializes it first.
* Requesting it against a provider that lacks the capability is rejected at start.
*/
outputSchema?: SchemaSpec
outputSchema?: StructuredOutputSchema
/**
* Optional recursion cap (max delegation depth below this child). Requires
* {@link SubagentCapabilities.depthLimit}; rejected at start otherwise.

View File

@@ -124,7 +124,7 @@ describe('SubagentService', () => {
describe('start-time capability validation (fail loud, before any child)', () => {
it.each([
{ field: 'outputSchema', request: baseRequest({ outputSchema: { x: { type: 'string' } } }) },
{ field: 'outputSchema', request: baseRequest({ outputSchema: { type: 'object', properties: { x: { type: 'string' } } } }) },
{ field: 'maxDepth', request: baseRequest({ maxDepth: 2 }) },
{ field: 'toolFilter', request: baseRequest({ toolFilter: { deny: ['bash'] } }) },
])('rejects $field against a provider that lacks the capability — before start() runs', ({ request }) => {
@@ -149,7 +149,7 @@ describe('SubagentService', () => {
await ctx.plugin(SubagentService)
const provider = new StubProvider('strong', ALL_CAPS)
ctx.subagents.registerProvider(provider)
ctx.subagents.start('strong', baseRequest({ outputSchema: { x: { type: 'string' } }, maxDepth: 1 }))
ctx.subagents.start('strong', baseRequest({ outputSchema: { type: 'object', properties: { x: { type: 'string' } } }, maxDepth: 1 }))
expect(provider.startCount).toBe(1)
})
})