fix review finding: the capture commits only on the final post-execute accept
The cross-seam blocker: structured_output recorded its value in the tool BODY, before tools/post-execute could block the call — a PostToolUse hook's block turned the logged result into isError while readResult still returned structured success and the continuation veto ended the turn. Two-phase commit: the body validates and STAGES (RunState.pending); a fourth runtime listener on tools/post-execute — prepend, so await next() returns the composed final decision — promotes the stage to captured only on an accepted call, and clears it on every path. A block now yields a consistent pair: the model and log see the isError feedback, the run settles error with no structured value, and the turn continues so the model can react. Regressions: block denies the capture end-to-end; accept-with-replacement still commits.
This commit is contained in:
@@ -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<PostToolDecision>,
|
||||
): Promise<PostToolDecision> {
|
||||
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
|
||||
|
||||
@@ -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>): 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()
|
||||
})
|
||||
})
|
||||
|
||||
Reference in New Issue
Block a user