fix: address codex review round 1
- Strict steer now rejects the two windows where an acknowledged message would be silently dropped: the closed-turn durability-flush window (status still running, loop strands drained steering) and a committed structured capture (terminal turn-stop discards late steering). Seam JSDoc, catalog doc, README, and the Agent Note bilingual pair state the tightened contract; new keyless tests pin both rejections. - Continuable background delegation now fails loud when the advertised send_message tool is not registered, instead of starting a durable child the model cannot continue. The acp-agent example already loads the control tool; the tool-catalog boot recipe is unaffected because capability wording is harvested at mount.
This commit is contained in:
@@ -30,7 +30,7 @@ The required request signal covers both startup and the live run. Before publica
|
||||
|
||||
After fulfillment, the caller owns the run. Provider-plugin unload does not revoke it. `dispose()` removes the live abort listener, records cancellation, and delegates to the returned `AgentHandle.dispose()`, whose memoized quiescence transaction stops the loop, removes the agent and session, and unwinds scoped registrations. Cancellation owns every non-completed in-flight outcome and reports `aborted`; an already-completed turn remains completed.
|
||||
|
||||
Runs expose the strict `steer` capability: a synchronous `AgentStatus.running` check and `Agent.steer()` call share one frame, so delivery joins the observed turn or throws. The Agent-level idle fallback (queue and start a new turn) is deliberately not reachable through the run — that would start an untracked turn after the run's result was read.
|
||||
Runs expose the strict `steer` capability: the synchronous checks and the `Agent.steer()` call share one frame, so delivery joins the observed turn or throws. Delivery requires `AgentStatus.running`, an open turn in the child log (status stays `running` through a closed turn's durability flush, where the loop would strand the message), and no committed structured capture (whose terminal stop makes the loop discard late steering). The Agent-level idle fallback (queue and start a new turn) is deliberately not reachable through the run — that would start an untracked turn after the run's result was read.
|
||||
|
||||
## Spawn and fork inputs
|
||||
|
||||
|
||||
@@ -259,13 +259,30 @@ function driveTurn(
|
||||
return handle.dispose()
|
||||
},
|
||||
steer(content: ContentBlock[]): void {
|
||||
// Strict live delivery: the synchronous running check and Agent.steer()
|
||||
// Strict live delivery: the synchronous checks and the Agent.steer()
|
||||
// call share one frame, so delivery joins the observed turn or throws.
|
||||
// Agent.steer()'s own idle fallback would instead QUEUE the message and
|
||||
// start a new, untracked turn after this run's result was read.
|
||||
if (child.status !== 'running') {
|
||||
throw new Error(`subagent child "${childId}" is not running; the message was not delivered`)
|
||||
}
|
||||
// The status stays `running` through the closed turn's durability flush,
|
||||
// and the loop DISCARDS terminal-stopped steering drained after turn
|
||||
// close instead of recording it. Requiring an open turn keeps
|
||||
// acknowledged delivery honest.
|
||||
const lastBoundary = child.session.events.findLast(
|
||||
event => event.type === 'turn/start' || event.type === 'turn/end',
|
||||
)
|
||||
if (lastBoundary?.type !== 'turn/start') {
|
||||
throw new Error(`subagent child "${childId}" turn has already closed; the message was not delivered`)
|
||||
}
|
||||
// A committed structured capture makes the pending `agent/turn-stop`
|
||||
// checkpoint terminal, and the loop then discards late steering. The
|
||||
// capture is synchronously observable, so reject rather than
|
||||
// acknowledge a message the run is about to drop.
|
||||
if (structured?.captured() !== undefined) {
|
||||
throw new Error(`subagent child "${childId}" already reported its structured result; the message was not delivered`)
|
||||
}
|
||||
child.steer(createUserMessage({ content, source: { kind: 'user' } }))
|
||||
},
|
||||
}
|
||||
|
||||
@@ -120,6 +120,33 @@ describe('in-process structured output', () => {
|
||||
await run.dispose()
|
||||
})
|
||||
|
||||
it('strict steer rejects delivery once the structured result is captured', async () => {
|
||||
// Hold the capture's tool result open so the child is observably running
|
||||
// with a committed capture: the pending agent/turn-stop checkpoint is
|
||||
// terminal, and the loop would DISCARD a steering message, so an
|
||||
// acknowledged delivery here would be a lie.
|
||||
let releaseResult: (() => void) | undefined
|
||||
const { ctx, parent } = await setup([
|
||||
toolCallResponse('c1', STRUCTURED_OUTPUT_TOOL, { answer: 7 }),
|
||||
])
|
||||
ctx.on('agent/post-step', (agent) => {
|
||||
if (agent.session.header.parentSession === undefined || releaseResult !== undefined) return
|
||||
return new Promise<void>((resolve) => { releaseResult = resolve })
|
||||
})
|
||||
const run = await ctx.subagents.start('spawn', structuredRequest(parent))
|
||||
await new Promise<void>((resolve) => {
|
||||
const timer = setInterval(() => {
|
||||
if (releaseResult !== undefined) { clearInterval(timer); resolve() }
|
||||
}, 5)
|
||||
})
|
||||
expect(() => { run.steer!([{ type: 'text', text: 'one more thing' }]) })
|
||||
.toThrow(/already reported its structured result; the message was not delivered/)
|
||||
releaseResult!()
|
||||
const result = await run.result
|
||||
expect(result.structured).toEqual({ answer: 7 })
|
||||
await run.dispose()
|
||||
})
|
||||
|
||||
it('denies tool calls that FOLLOW the capture in the same response — terminal means terminal', async () => {
|
||||
// One model response carrying structured_output FIRST and a side-effecting
|
||||
// call after it: the continuation veto only fires at step end, so without
|
||||
|
||||
@@ -245,4 +245,45 @@ describe('startInProcessRun', () => {
|
||||
expect(ctx.agents.list()).toHaveLength(beforeAgents)
|
||||
expect(ctx.sessions.list()).toHaveLength(beforeSessions)
|
||||
})
|
||||
|
||||
it('strict steer rejects a settled child instead of queueing an untracked turn', async () => {
|
||||
const { ctx, parent } = await setup([textResponse('done')])
|
||||
const run = await startInProcessRun(request(parent), {})
|
||||
await run.result
|
||||
// The child is idle after its turn: Agent.steer() would silently QUEUE.
|
||||
expect(() => { run.steer!([{ type: 'text', text: 'late' }]) })
|
||||
.toThrow(/not running; the message was not delivered/)
|
||||
const child = ctx.agents.get(run.id)!
|
||||
expect(child.session.events.some(event => event.type === 'steering/message')).toBe(false)
|
||||
await run.dispose()
|
||||
})
|
||||
|
||||
it('strict steer rejects the closed-turn flush window where the loop discards steering', async () => {
|
||||
// Hold the turn-end durability flush open: the turn has closed in the log
|
||||
// and status is still `running`, exactly the window where the loop would
|
||||
// discard a drained steering message instead of recording it.
|
||||
const { ctx, parent } = await setup([textResponse('quick')])
|
||||
let releaseFlush: (() => void) | undefined
|
||||
ctx.on('session/flush', (session) => {
|
||||
if (session.header.parentSession === undefined || releaseFlush !== undefined) return
|
||||
const lastEnd = session.events.findLast(event => event.type === 'turn/end')
|
||||
if (lastEnd === undefined) return
|
||||
return new Promise<void>((resolve) => { releaseFlush = resolve })
|
||||
})
|
||||
const run = await startInProcessRun(request(parent), {})
|
||||
const child = ctx.agents.get(run.id)!
|
||||
// Wait until the child's turn has closed while the flush keeps it running.
|
||||
await new Promise<void>((resolve) => {
|
||||
const timer = setInterval(() => {
|
||||
if (releaseFlush !== undefined) { clearInterval(timer); resolve() }
|
||||
}, 5)
|
||||
})
|
||||
expect(child.status).toBe('running')
|
||||
expect(() => { run.steer!([{ type: 'text', text: 'into the void' }]) })
|
||||
.toThrow(/turn has already closed; the message was not delivered/)
|
||||
releaseFlush!()
|
||||
await run.result
|
||||
expect(child.session.events.some(event => event.type === 'steering/message')).toBe(false)
|
||||
await run.dispose()
|
||||
})
|
||||
})
|
||||
|
||||
@@ -219,10 +219,11 @@ export interface SubagentRun {
|
||||
/**
|
||||
* OPTIONAL (strict live-steering capability): deliver additional content to
|
||||
* the actively running child turn. STRICT means delivery joins the observed
|
||||
* turn or fails — the implementation must synchronously require the child to
|
||||
* be running with no asynchronous boundary before delivery, and must not
|
||||
* fall back to a queue path that could start a new, untracked turn after
|
||||
* this run has settled. Throws when the child is not running. A run
|
||||
* turn or fails — the implementation must synchronously verify, with no
|
||||
* asynchronous boundary before delivery, that the child is running and its
|
||||
* turn can still record the message, and must not fall back to a queue path
|
||||
* that could start a new, untracked turn or silently drop the message after
|
||||
* this run has settled. Throws when delivery cannot join the turn. A run
|
||||
* represents one disposable activation, so it has no cold-resume operation;
|
||||
* resuming a settled child goes through {@link SubagentProvider.resume}.
|
||||
*/
|
||||
|
||||
@@ -285,6 +285,13 @@ export function apply(ctx: Context, config: Config): void {
|
||||
if (control === undefined) {
|
||||
throw new Error('continuable background subagents unavailable: load @deepseek-ai/dsh-subagent-control and @deepseek-ai/dsh-tool-tasks')
|
||||
}
|
||||
// The schema above tells the model to follow up with
|
||||
// `send_message`; starting a durable child the model cannot
|
||||
// continue would make that advertisement false. Sibling load order
|
||||
// is undetermined at mount, so the check lives at the operation.
|
||||
if (ctx.tools.get('send_message') === undefined) {
|
||||
throw new Error('continuable background subagents unavailable: load @deepseek-ai/dsh-tool-subagent-control (the advertised send_message tool is not registered)')
|
||||
}
|
||||
// The control service owns the durable child id, descriptor
|
||||
// snapshot, Task registration, and settle-then-dispose ordering; a
|
||||
// synchronous validation failure rejects the call with no Task.
|
||||
|
||||
@@ -17,6 +17,7 @@ import type { SubagentStartRequest } from '@deepseek-ai/dsh-subagent'
|
||||
import LocalTaskService from '@deepseek-ai/dsh-tasks-local'
|
||||
import SubagentControlService from '@deepseek-ai/dsh-subagent-control'
|
||||
import * as SubagentSpawn from '@deepseek-ai/dsh-subagent-spawn'
|
||||
import * as ToolSubagentControl from '@deepseek-ai/dsh-tool-subagent-control'
|
||||
import * as ToolTasks from '@deepseek-ai/dsh-tool-tasks'
|
||||
import { MockAdapter, textResponse } from '../../../core/agent-loop/tests/mock-adapter.ts'
|
||||
import * as mock from './scripted-provider.ts'
|
||||
@@ -825,7 +826,7 @@ describe('dsh-tool-subagent continuable background mode', () => {
|
||||
})
|
||||
|
||||
/** Boot the real continuable stack: loop, persistence, spawn, tasks, control. */
|
||||
async function continuableSetup() {
|
||||
async function continuableSetup(options: { controlTool?: boolean } = {}) {
|
||||
const ctx = new Context()
|
||||
await mountAgentLoopTestDependencies(ctx)
|
||||
const root = mkdtempSync(path.join(tmpdir(), 'dsh-tool-subagent-continuable-'))
|
||||
@@ -837,6 +838,7 @@ describe('dsh-tool-subagent continuable background mode', () => {
|
||||
await ctx.plugin(LocalTaskService)
|
||||
await ctx.plugin(ToolTasks, {})
|
||||
await ctx.plugin(SubagentControlService)
|
||||
if (options.controlTool !== false) await ctx.plugin(ToolSubagentControl)
|
||||
await ctx.plugin(tool, { provider: 'spawn' })
|
||||
ctx.llm.registerAdapter(['mock'], new MockAdapter([
|
||||
textResponse('continuable answer'),
|
||||
@@ -886,6 +888,21 @@ describe('dsh-tool-subagent continuable background mode', () => {
|
||||
expect(result.isError).toBe(true)
|
||||
expect(text(result)).toContain('load @deepseek-ai/dsh-subagent-control')
|
||||
})
|
||||
|
||||
it('fails loud when the advertised send_message tool is not registered', async () => {
|
||||
// The schema tells the model to follow up with send_message; starting a
|
||||
// durable child the model cannot continue would make that false.
|
||||
const { ctx, parent } = await continuableSetup({ controlTool: false })
|
||||
const result = await callSubagent(
|
||||
ctx,
|
||||
{ description: 'd', prompt: 'p', run_in_background: true },
|
||||
{ agent: parent },
|
||||
)
|
||||
expect(result.isError).toBe(true)
|
||||
expect(text(result)).toContain('load @deepseek-ai/dsh-tool-subagent-control')
|
||||
// Nothing was started: no Task exists for the parent.
|
||||
expect(ctx.tasks.list(parent)).toEqual([])
|
||||
})
|
||||
})
|
||||
|
||||
describe('background preflight failure (no orphaned child, by construction)', () => {
|
||||
|
||||
Reference in New Issue
Block a user