diff --git a/packages/subagent/subagent-inprocess/src/structured.ts b/packages/subagent/subagent-inprocess/src/structured.ts index 370f7a538d..a557e44785 100644 --- a/packages/subagent/subagent-inprocess/src/structured.ts +++ b/packages/subagent/subagent-inprocess/src/structured.ts @@ -36,14 +36,20 @@ * closes the within-step window the continuation veto cannot: a * `tools/pre-execute` deny for any call arriving after the agent's capture, so * a response that lists `structured_output` before further tool calls cannot - * run side effects after the final answer was accepted. + * run side effects after the final answer was accepted. A fourth, + * `tools/post-execute`, is the capture COMMIT: the tool body only stages the + * validated value, and it becomes the run's captured result only when the + * final post-execute decision accepts the call — a blocking hook downstream + * yields `isError` in the log, and the run must not report success for it. * - * 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. + * Lifetime is refcounted by structured RUNS: each acquires from start to + * settle, so the registrations exist exactly while at least one structured + * child is live — a plain deployment that never passes `outputSchema` carries + * no always-on global state, and a backend hot-reload mid-run cannot + * unregister the capture tool out from under a live child (the run holds its + * own acquisition). Registrations land on the ROOT context and the refcount + * disposes them when the last run settles; the next structured run + * re-registers them. * * @module @deepseek-ai/dsh-subagent-inprocess/structured */ @@ -53,7 +59,7 @@ import type { Agent } from '@deepseek-ai/dsh-agent' import type { ContentBlock, ToolSchema } from '@deepseek-ai/dsh-llm' import type { ContinuationDecision } from '@deepseek-ai/dsh-agent' import type { AssembleContext, PromptAssembly } from '@deepseek-ai/dsh-system-prompt' -import type { PreToolDecision, ToolExecution } from '@deepseek-ai/dsh-tools' +import type { PostToolDecision, PreToolDecision, ToolExecution, ToolExecutionResult } 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. */ @@ -75,6 +81,15 @@ export const STRUCTURED_OUTPUT_INSTRUCTION /** One structured run's state: the schema to enforce and the captured value, once recorded. */ interface RunState { readonly schema: StructuredOutputSchema + /** + * A validated value awaiting the post-execute verdict on ITS OWN call. Set + * by the capture tool's body, promoted to {@link RunState.captured} only + * when the final `tools/post-execute` decision accepts the call — a + * downstream block turns the logged result into `isError`, and a value + * committed at body time would let the run report success for a call the + * model saw fail. + */ + pending?: { value: unknown } captured?: { value: unknown } } @@ -106,7 +121,7 @@ export interface StructuredAcquisition { /** * Acquire the per-root-context structured runtime, registering the capture tool - * and the two waterfall listeners on the FIRST acquisition. See the module doc + * and the runtime's 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). @@ -179,7 +194,9 @@ function registerRuntime(root: Context, runtime: StructuredRuntime): void { // 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 } + // Two-phase commit: the body only STAGES the value; the post-execute + // listener below promotes it once the final decision accepts the call. + state.pending = { value: args } return Promise.resolve([{ type: 'text', text: 'Structured output recorded.' }]) }, }) @@ -243,6 +260,33 @@ function registerRuntime(root: Context, runtime: StructuredRuntime): void { return next() }, { prepend: true })) + // The capture COMMIT: promote the staged value only when the final + // post-execute decision accepts the call. The capture tool's body cannot + // decide — `tools/post-execute` runs after it, and a blocking listener (a + // PostToolUse hook) turns the logged result into `isError` feedback; a value + // committed at body time would make readResult report `structured` success + // for a call whose result the model and session log saw fail. `prepend: + // true` = outermost at registration time, so `await next()` returns the + // COMPOSED downstream decision — the same final verdict the registry maps + // onto the result. (A later-registered outer listener that blocks without + // delegating skips this commit entirely: the staged value is dropped and the + // run errors — failure-safe in the same direction.) The staging slot clears + // on every path, including a rejecting downstream listener. + runtime.disposers.push(root.on('tools/post-execute', async function ( + this: unknown, exec: ToolExecution, _result: ToolExecutionResult, next: () => Promise, + ): Promise { + const state = exec.agent ? runtime.states.get(exec.agent) : undefined + if (!state || exec.name !== STRUCTURED_OUTPUT_TOOL || state.pending === undefined) return next() + const pending = state.pending + try { + const decision = await next() + if (decision.kind === 'accept') state.captured = pending + return decision + } finally { + delete state.pending + } + }, { prepend: true })) + // Terminal means terminal WITHIN the step, not only at its end: the // turn-continuation veto above runs after every call in the current model // response has executed, so a response that puts `structured_output` before diff --git a/packages/subagent/subagent-inprocess/tests/structured.spec.ts b/packages/subagent/subagent-inprocess/tests/structured.spec.ts index ba693cfecd..0de036a0a0 100644 --- a/packages/subagent/subagent-inprocess/tests/structured.spec.ts +++ b/packages/subagent/subagent-inprocess/tests/structured.spec.ts @@ -11,8 +11,7 @@ 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 '@deepseek-ai/dsh-subagent-spawn' -import * as fork from '@deepseek-ai/dsh-subagent-fork' +import { startInProcessRun } from '../src/index.ts' import { acquireStructuredRuntime, STRUCTURED_OUTPUT_INSTRUCTION, @@ -28,11 +27,14 @@ const SCHEMA: StructuredOutputSchema = { } /** - * 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. + * Real loop + scripted mock model + an INLINE spawn-shaped provider over the + * shared driver. The concrete backend plugins are deliberately NOT loaded — + * they would devDep-cycle this package (spawn/fork already depend on the + * driver), and the runtime under test is the driver's; plugin-level structured + * coverage lives in the spawn/fork specs. The mock model script drives the + * child's structured_output calls. */ -async function setup(script: Script, options?: { withFork?: boolean }) { +async function setup(script: Script) { const ctx = new Context() const adapter = new MockAdapter(script) await ctx.plugin(LlmService) @@ -43,13 +45,15 @@ async function setup(script: Script, options?: { withFork?: boolean }) { await ctx.plugin(Invariants) await ctx.plugin(AgentLoop, { agents: [] }) await ctx.plugin(SubagentService) - const fiber = await ctx.plugin(spawn, { providerName: 'spawn' }) - const forkFiber = options?.withFork - ? await ctx.plugin(fork, { providerName: 'fork' }) - : undefined + const disposeProvider = ctx.subagents.registerProvider({ + name: 'spawn', + capabilities: { outputSchema: true, depthLimit: true, toolFilter: false }, + inheritsParentContext: false, + start: (request: SubagentStartRequest) => startInProcessRun(ctx, request, { providerName: 'spawn' }), + }) ctx.llm.registerAdapter(['mock'], adapter) const parent = ctx.agentLoop.create(AgentId('parent'), { model: 'mock' }) - return { ctx, parent, adapter, fiber, forkFiber } + return { ctx, parent, adapter, disposeProvider } } function structuredRequest(parent: SubagentStartRequest['parent'], extra?: Partial): SubagentStartRequest { @@ -263,6 +267,62 @@ describe('in-process structured output', () => { expect(ctx.agents.get(AgentId('parent'))).toBeDefined() }) + it('a schema carrying non-JSON values fails as OutputSchemaError, never as a raw clone error', async () => { + const { ctx, parent } = await setup([]) + // Assertion runs BEFORE the defensive structuredClone: a function-valued + // annotation must surface as the subset violation it is, not escape as + // structuredClone's DataCloneError. + expect(() => ctx.subagents.start('spawn', structuredRequest(parent, { + outputSchema: { type: 'object', default: () => {} } as unknown as StructuredOutputSchema, + }))).toThrow(/unsupported output schema.*annotation must be JSON data/) + }) + + it('a post-execute BLOCK on the capture call denies the capture: log and result agree on failure', async () => { + const { ctx, parent, adapter } = await setup([ + toolCallResponse('c1', STRUCTURED_OUTPUT_TOOL, { answer: 7 }), + textResponse('continues after the blocked capture'), + ]) + // A PostToolUse-style hook, registered AFTER the runtime (so the runtime's + // prepend commit listener stays outermost and composes this verdict). + ctx.on('tools/post-execute', (exec, _result, next) => { + if (exec.name === STRUCTURED_OUTPUT_TOOL) { + return Promise.resolve({ kind: 'block' as const, feedback: [{ type: 'text' as const, text: 'capture rejected by hook' }] }) + } + return next() + }) + const run = ctx.subagents.start('spawn', structuredRequest(parent)) + const result = await run.result + // No capture was committed: the run reports the schema shortfall... + expect(result.structured).toBeUndefined() + expect(result.stopReason).toBe('error') + // ...the logged tool result is the blocked isError with the feedback... + const child = ctx.agents.get(run.id)! + const results = child.session.events.filter(e => e.type === 'tool/result') + expect((results[0]!.data as { isError?: boolean }).isError).toBe(true) + expect(JSON.stringify((results[0]!.data as { content: unknown }).content)).toContain('capture rejected by hook') + // ...and the turn CONTINUED past the blocked call (no captured veto): + // the model got to react to the failure with a second step. + expect(adapter.requests.length).toBe(2) + await run.dispose() + }) + + it('a post-execute accept-with-replacement still commits the capture', async () => { + const { ctx, parent } = await setup([ + toolCallResponse('c1', STRUCTURED_OUTPUT_TOOL, { answer: 8 }), + ]) + ctx.on('tools/post-execute', (exec, _result, next) => { + if (exec.name === STRUCTURED_OUTPUT_TOOL) { + return Promise.resolve({ kind: 'accept' as const, content: [{ type: 'text' as const, text: 'recorded (rewritten)' }] }) + } + return next() + }) + const run = ctx.subagents.start('spawn', structuredRequest(parent)) + const result = await run.result + expect(result.stopReason).toBe('completed') + expect(result.structured).toEqual({ answer: 8 }) + await run.dispose() + }) + it('appends the structured instruction to the child REQUEST\'s system text (base prompt preserved)', async () => { const { ctx, parent, adapter } = await setup([toolCallResponse('c1', STRUCTURED_OUTPUT_TOOL, { answer: 1 })]) // A context-wide section stands in for the deployment persona: the @@ -298,6 +358,21 @@ describe('in-process structured output', () => { }) describe('final-request enforcement (the prepend agent/request listener)', () => { + it('a plain agent assembling while the runtime is LIVE gets the placeholder stripped', async () => { + // Run-scoped acquisition means a plain deployment never registers the + // tool at all; the strip branch exists for the CONCURRENT case — a plain + // agent taking a turn while some structured child holds the runtime open. + const { ctx, parent, adapter } = await setup([textResponse('parent answer')]) + const hold = acquireStructuredRuntime(ctx) + parent.send([{ type: 'text', text: 'hello' }]) + await parent.whenIdle() + // The placeholder IS in the registry during this turn; the assembly the + // loop rendered must not carry it for an agent without a structured run. + expect(ctx.tools.get(STRUCTURED_OUTPUT_TOOL)).toBeDefined() + expect(toolNames(adapter.requests[0]!)).not.toContain(STRUCTURED_OUTPUT_TOOL) + hold.release() + }) + 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. @@ -396,10 +471,13 @@ describe('in-process structured output', () => { // and shape a structured agent's assembly on the same path the loop // renders and logs as the request header. const { ctx, parent } = await setup([]) + const acquisition = acquireStructuredRuntime(ctx) + // Bare assemble WHILE the runtime is live: the no-agent branch must + // strip the registered placeholder (before the acquisition there is + // nothing to strip — run-scoped registration). const bare = await ctx.systemPrompt.assemble({}) expect(bare.tools.map(tool => tool.name)).not.toContain(STRUCTURED_OUTPUT_TOOL) - const acquisition = acquireStructuredRuntime(ctx) acquisition.attach(parent, SCHEMA) const shaped = await ctx.systemPrompt.assemble({ agent: parent }) expect(shaped.tools.map(tool => tool.name)).toContain(STRUCTURED_OUTPUT_TOOL) @@ -412,57 +490,37 @@ describe('in-process structured output', () => { }) }) - 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 () => { + describe('runtime lifetime (refcount: live structured runs)', () => { + it('the runtime exists exactly while structured runs are live: nothing before, nothing after', async () => { const { ctx, parent } = await setup([ - toolCallResponse('c1', STRUCTURED_OUTPUT_TOOL, { answer: 9 }), - ], { withFork: true }) - const run = ctx.subagents.start('fork', structuredRequest(parent)) + toolCallResponse('c1', STRUCTURED_OUTPUT_TOOL, { answer: 4 }), + ]) + // No always-on global state: a context that has run no structured child + // carries no capture tool. + expect(ctx.tools.get(STRUCTURED_OUTPUT_TOOL)).toBeUndefined() + const run = ctx.subagents.start('spawn', structuredRequest(parent)) const result = await run.result - expect(result.structured).toEqual({ answer: 9 }) + // The capture succeeded — the registrations existed while the run lived. + expect(result.structured).toEqual({ answer: 4 }) + // The run's settle released the last acquisition. + expect(ctx.tools.get(STRUCTURED_OUTPUT_TOOL)).toBeUndefined() await run.dispose() }) + it('concurrent structured runs share one runtime; the last settle disposes it', async () => { + const { ctx, parent } = await setup([ + toolCallResponse('c1', STRUCTURED_OUTPUT_TOOL, { answer: 1 }), + toolCallResponse('c2', STRUCTURED_OUTPUT_TOOL, { answer: 2 }), + ]) + const first = ctx.subagents.start('spawn', structuredRequest(parent)) + const second = ctx.subagents.start('spawn', structuredRequest(parent)) + const [a, b] = await Promise.all([first.result, second.result]) + expect([a.structured, b.structured].sort()).toEqual([{ answer: 1 }, { answer: 2 }].sort()) + expect(ctx.tools.get(STRUCTURED_OUTPUT_TOOL)).toBeUndefined() + await first.dispose() + await second.dispose() + }) + it('acquisition release is idempotent (double release cannot underflow the refcount)', async () => { const ctx = new Context() await ctx.plugin(SystemPrompt) @@ -513,13 +571,16 @@ describe('in-process structured output', () => { 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() + // That manual acquisition was the ONLY holder - release disposes. + expect(ctx.tools.get(STRUCTURED_OUTPUT_TOOL)).toBeUndefined() }) }) it('a direct structured_output call from an agent WITHOUT a structured run is an isError', async () => { const { ctx, parent } = await setup([]) + // Hold the runtime open (run-scoped: nothing is registered otherwise) so + // the call reaches the capture tool's own fail-loud guard, not UNKNOWN_TOOL. + const hold = acquireStructuredRuntime(ctx) const result = await ctx.tools.execute({ callId: 'x' as never, name: STRUCTURED_OUTPUT_TOOL, @@ -527,16 +588,19 @@ describe('in-process structured output', () => { agent: parent, }) expect(result.isError).toBe(true) - expect(result.content[0]).toMatchObject({ type: 'text' }) + expect(JSON.stringify(result.content)).toContain('only available to subagents') + hold.release() }) it('a structured_output call with NO calling agent at all is an isError', async () => { const { ctx } = await setup([]) + const hold = acquireStructuredRuntime(ctx) const result = await ctx.tools.execute({ callId: 'x' as never, name: STRUCTURED_OUTPUT_TOOL, arguments: { answer: 1 }, }) expect(result.isError).toBe(true) + hold.release() }) })