fix: honor the teardown order on owner unload; make the structured commit unconditional
Adversarial-review findings (own reviewer agent), each verified and pinned: B1: agents.register() returned a wrapper lambda, so the factory composite's yield could not identity-nest it — on OWNER unload the unregistration (and agent/disposed) disposed as a concurrent sibling, firing mid-drain while the final turn was still closing (pre-existing on master; this branch's docs re-assert the order, so it must be true). register() now returns the EXACT cordis effect disposer (the Scope.rawDispose move); the composite nests it and owner unload runs stop/drain -> unregister -> detach -> scope like every other path. Regression test pins turn-end before disposed before detach on owner unload. B2: the structured two-phase commit could promote a stale stage when a later capture call REUSED the orphaned stage's call id with a body that never staged (denied downstream, or invalid args throwing pre-stage). The runtime's pre-execute listener now clears any stale stage unconditionally when a new capture call enters the pipeline — only a call's own body can stage for its commit; the call-id mismatch guard becomes a defensive second layer. Repro test: blocked capture then same-id invalid call. C1: an explicit empty toolFilter config now fails at plugin LOAD (the check is self-contained) instead of killing every delegation at child setup. C2: Scope.dispose/ScopeHost.dispose @returns state the single-shot repeat-call semantics honestly.
This commit is contained in:
@@ -170,6 +170,40 @@ describe('agent scope lifecycle', () => {
|
|||||||
expect(heard).toEqual(['a1:2'])
|
expect(heard).toEqual(['a1:2'])
|
||||||
})
|
})
|
||||||
|
|
||||||
|
it('owner unload honors the documented teardown order: unregistration AFTER the drain, before detach', async () => {
|
||||||
|
const ctx = await harness()
|
||||||
|
let handle!: ReturnType<typeof ctx.agents.create>
|
||||||
|
const owner = await ctx.plugin(Object.assign((inner: Context) => {
|
||||||
|
handle = inner.agents.create({ agentId: AgentId('o1'), sessionId: SessionId('o1-s'), agentOptions: { model: 'mock' } })
|
||||||
|
}, { inject: ['agents'] }))
|
||||||
|
const { agent } = handle
|
||||||
|
|
||||||
|
const order: string[] = []
|
||||||
|
ctx.on('session/event', (_s, event) => {
|
||||||
|
if (event.type === 'turn/end') order.push('turn-end')
|
||||||
|
})
|
||||||
|
ctx.on('agent/disposed', () => {
|
||||||
|
order.push(`disposed(listed=${ctx.agents.get(AgentId('o1')) !== undefined})`)
|
||||||
|
order.push(`session-still-stored=${ctx.sessions.get(SessionId('o1-s')) !== undefined}`)
|
||||||
|
})
|
||||||
|
|
||||||
|
// Open a turn so the drain has real work: the loop must finish it BEFORE
|
||||||
|
// the registry entry goes away (the agent/disposed contract: "its fiber
|
||||||
|
// and any in-flight turn have been torn down"). Wait for the turn to be
|
||||||
|
// OPEN in the log — a dispose landing in the pre-step window would drop
|
||||||
|
// the queued prompt without ever opening a turn.
|
||||||
|
const turnOpen = new Promise<void>((resolve) => {
|
||||||
|
const off = ctx.on('session/event', (_s, event) => {
|
||||||
|
if (event.type === 'turn/start') { off(); resolve() }
|
||||||
|
})
|
||||||
|
})
|
||||||
|
agent.send(text('work'))
|
||||||
|
await turnOpen
|
||||||
|
await owner.dispose()
|
||||||
|
expect(order).toEqual(['turn-end', 'disposed(listed=false)', 'session-still-stored=true'])
|
||||||
|
expect(ctx.sessions.get(SessionId('o1-s'))).toBeUndefined()
|
||||||
|
})
|
||||||
|
|
||||||
it('handle.dispose() during owner unload still awaits true quiescence (shared boundary)', async () => {
|
it('handle.dispose() during owner unload still awaits true quiescence (shared boundary)', async () => {
|
||||||
const ctx = await harness()
|
const ctx = await harness()
|
||||||
let handle!: ReturnType<typeof ctx.agents.create>
|
let handle!: ReturnType<typeof ctx.agents.create>
|
||||||
|
|||||||
@@ -208,9 +208,16 @@ export class AgentRegistry extends Service {
|
|||||||
* (calling through `agent.ctx` scopes EFFECTS; dispatch scoping always
|
* (calling through `agent.ctx` scopes EFFECTS; dispatch scoping always
|
||||||
* requires passing the carrier). Returns the disposer.
|
* requires passing the carrier). Returns the disposer.
|
||||||
* @param agent - the already-constructed agent to record in the store.
|
* @param agent - the already-constructed agent to record in the store.
|
||||||
* @returns the disposer that removes the agent and emits `agent/disposed`.
|
* @returns the EXACT Cordis effect disposer (single-shot; a repeat call
|
||||||
|
* returns undefined without awaiting an in-flight teardown). Exact
|
||||||
|
* identity is load-bearing: a composite (generator) effect that owns a
|
||||||
|
* teardown ORDER — the agent factory's lifecycle chain — must yield THIS
|
||||||
|
* function so Cordis nests the unregistration at that yield position;
|
||||||
|
* yielding a wrapper would leave it disposing as a concurrent sibling on
|
||||||
|
* owner unload, unregistering the agent (and emitting `agent/disposed`)
|
||||||
|
* while its final turn is still draining.
|
||||||
*/
|
*/
|
||||||
register(agent: Agent): () => void {
|
register(agent: Agent): () => Promise<void> | void {
|
||||||
const dispose = this.ctx.effect(function* (this: AgentRegistry) {
|
const dispose = this.ctx.effect(function* (this: AgentRegistry) {
|
||||||
if (this.store.has(agent.id)) {
|
if (this.store.has(agent.id)) {
|
||||||
throw new Error(`agent "${agent.id}" is already registered`)
|
throw new Error(`agent "${agent.id}" is already registered`)
|
||||||
@@ -242,9 +249,7 @@ export class AgentRegistry extends Service {
|
|||||||
}
|
}
|
||||||
this.ctx.emit(scopeTarget(agent, agent), 'agent/created', agent)
|
this.ctx.emit(scopeTarget(agent, agent), 'agent/created', agent)
|
||||||
}.bind(this), 'agents.register()')
|
}.bind(this), 'agents.register()')
|
||||||
// ctx.effect's disposer returns Promise<void>; our disposer API is
|
return dispose
|
||||||
// synchronous fire-and-forget — discard the (always-resolved) promise.
|
|
||||||
return () => void dispose()
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|||||||
@@ -83,7 +83,12 @@ export interface Scope {
|
|||||||
* returns undefined the second time; this wrapper Promise-normalizes it).
|
* returns undefined the second time; this wrapper Promise-normalizes it).
|
||||||
* After disposal the scoped context is inert — a further registration
|
* After disposal the scoped context is inert — a further registration
|
||||||
* through it throws Cordis's INACTIVE_EFFECT.
|
* through it throws Cordis's INACTIVE_EFFECT.
|
||||||
* @returns resolves when every registration's disposer has settled.
|
* @returns for the call that initiates teardown: resolves when every
|
||||||
|
* registration's disposer has settled. A repeat/racing call resolves
|
||||||
|
* immediately WITHOUT awaiting the in-flight teardown (the underlying
|
||||||
|
* Cordis disposer is single-shot) — a caller needing a shared quiescence
|
||||||
|
* boundary across racing disposers keeps its own completion promise (the
|
||||||
|
* agent factory's pattern).
|
||||||
*/
|
*/
|
||||||
dispose(): Promise<void>
|
dispose(): Promise<void>
|
||||||
}
|
}
|
||||||
@@ -250,7 +255,8 @@ export interface ScopeHost {
|
|||||||
mint(key: ScopeKey): Scope
|
mint(key: ScopeKey): Scope
|
||||||
/**
|
/**
|
||||||
* Dispose the host fiber and with it every scope minted through it.
|
* Dispose the host fiber and with it every scope minted through it.
|
||||||
* @returns resolves when all collected disposers have settled.
|
* @returns resolves when all collected disposers have settled (first call;
|
||||||
|
* a repeat call resolves immediately — single-shot, like Scope.dispose).
|
||||||
*/
|
*/
|
||||||
dispose(): Promise<void>
|
dispose(): Promise<void>
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -159,6 +159,13 @@ export function attachStructuredRuntime(childCtx: Context, schema: StructuredOut
|
|||||||
reason: `structured output already recorded: the run is complete, so \`${exec.name}\` is not executed`,
|
reason: `structured output already recorded: the run is complete, so \`${exec.name}\` is not executed`,
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
// A NEW capture call invalidates any stale stage UNCONDITIONALLY, before
|
||||||
|
// dispatch: only THIS call's own body may stage for this call's commit.
|
||||||
|
// Without this, a stale entry orphaned by an outer short-circuited chain
|
||||||
|
// could be promoted by a later call REUSING the same call id whose body
|
||||||
|
// never staged (pre-execute-denied downstream, or invalid args throwing
|
||||||
|
// before the stage) — reporting success for a value the model saw fail.
|
||||||
|
if (exec.name === STRUCTURED_OUTPUT_TOOL) pending = undefined
|
||||||
return next()
|
return next()
|
||||||
}, { prepend: true })
|
}, { prepend: true })
|
||||||
|
|
||||||
@@ -171,13 +178,14 @@ export function attachStructuredRuntime(childCtx: Context, schema: StructuredOut
|
|||||||
this: unknown, exec: ToolExecution, _result: ToolExecutionResult, next: () => Promise<PostToolDecision>,
|
this: unknown, exec: ToolExecution, _result: ToolExecutionResult, next: () => Promise<PostToolDecision>,
|
||||||
): Promise<PostToolDecision> {
|
): Promise<PostToolDecision> {
|
||||||
if (exec.name !== STRUCTURED_OUTPUT_TOOL || pending === undefined) return next()
|
if (exec.name !== STRUCTURED_OUTPUT_TOOL || pending === undefined) return next()
|
||||||
|
/* v8 ignore start -- defensive second layer: the pre-execute clear above
|
||||||
|
* already drops every stale stage before a new capture call dispatches,
|
||||||
|
* so a call-id mismatch cannot be reached through the tool pipeline */
|
||||||
if (pending.callId !== exec.callId) {
|
if (pending.callId !== exec.callId) {
|
||||||
// A stale stage from a different call: an outer listener short-circuited
|
|
||||||
// that call's post-execute chain past this commit, so its verdict never
|
|
||||||
// reached us and the value must never be promoted — drop it.
|
|
||||||
pending = undefined
|
pending = undefined
|
||||||
return next()
|
return next()
|
||||||
}
|
}
|
||||||
|
/* v8 ignore stop */
|
||||||
const staged = pending
|
const staged = pending
|
||||||
try {
|
try {
|
||||||
const decision = await next()
|
const decision = await next()
|
||||||
|
|||||||
@@ -530,4 +530,42 @@ describe('in-process structured output', () => {
|
|||||||
expect(valid.isError).toBeFalsy()
|
expect(valid.isError).toBeFalsy()
|
||||||
await run.dispose()
|
await run.dispose()
|
||||||
})
|
})
|
||||||
|
|
||||||
|
it('a later capture call REUSING a stale stage\'s call id never promotes it (unconditional commit safety)', async () => {
|
||||||
|
const { ctx, parent } = await setup([
|
||||||
|
toolCallResponse('c1', STRUCTURED_OUTPUT_TOOL, { answer: 1 }),
|
||||||
|
])
|
||||||
|
const run = ctx.subagents.start('spawn', structuredRequest(parent))
|
||||||
|
const child = ctx.agents.get(run.id)!
|
||||||
|
// Orphan a stage: an outer short-circuiting post-execute BLOCK on the
|
||||||
|
// first capture (its chain never reaches the commit listener).
|
||||||
|
let blocks = 1
|
||||||
|
ctx.on('tools/post-execute', (exec, _result, next) => {
|
||||||
|
if (exec.name === STRUCTURED_OUTPUT_TOOL && blocks > 0) {
|
||||||
|
blocks -= 1
|
||||||
|
return Promise.resolve({ kind: 'block' as const, feedback: [{ type: 'text' as const, text: 'rejected' }] })
|
||||||
|
}
|
||||||
|
return next()
|
||||||
|
}, { prepend: true })
|
||||||
|
await run.result
|
||||||
|
// A SECOND capture call with the SAME call id whose body never stages
|
||||||
|
// (invalid args throw before the stage): the stale value must not ride
|
||||||
|
// its acceptance.
|
||||||
|
const reused = await ctx.tools.execute({
|
||||||
|
callId: 'c1' as never,
|
||||||
|
name: STRUCTURED_OUTPUT_TOOL,
|
||||||
|
arguments: { answer: 'not-a-number' },
|
||||||
|
agent: child,
|
||||||
|
})
|
||||||
|
expect(reused.isError).toBe(true)
|
||||||
|
// Nothing was ever committed: a fresh valid call is still required.
|
||||||
|
const valid = await ctx.tools.execute({
|
||||||
|
callId: 'c1' as never,
|
||||||
|
name: STRUCTURED_OUTPUT_TOOL,
|
||||||
|
arguments: { answer: 5 },
|
||||||
|
agent: child,
|
||||||
|
})
|
||||||
|
expect(valid.isError).toBeFalsy()
|
||||||
|
await run.dispose()
|
||||||
|
})
|
||||||
})
|
})
|
||||||
|
|||||||
@@ -184,6 +184,12 @@ export function providerWording(inherits: boolean): { description: string; promp
|
|||||||
}
|
}
|
||||||
|
|
||||||
export function apply(ctx: Context, config: Config): void {
|
export function apply(ctx: Context, config: Config): void {
|
||||||
|
// Misconfiguration fails loud AT LOAD (the check is self-contained): an
|
||||||
|
// explicit `toolFilter: {}` would otherwise pass the capability gate and
|
||||||
|
// kill every delegation later, in the child-setup `restrict({})` throw.
|
||||||
|
if (config.toolFilter !== undefined && config.toolFilter.allow === undefined && config.toolFilter.deny === undefined) {
|
||||||
|
throw new Error('tool-subagent: `toolFilter` is configured but names neither `allow` nor `deny` — remove the key or fill the filter')
|
||||||
|
}
|
||||||
// The tool MIRRORS its provider's lifecycle instead of assuming load order:
|
// The tool MIRRORS its provider's lifecycle instead of assuming load order:
|
||||||
// the cordis Loader starts sibling entries concurrently, so "backend listed
|
// the cordis Loader starts sibling entries concurrently, so "backend listed
|
||||||
// first in cordis.yml" does not guarantee "provider registered first", and
|
// first in cordis.yml" does not guarantee "provider registered first", and
|
||||||
|
|||||||
@@ -510,4 +510,19 @@ describe('dsh-tool-subagent', () => {
|
|||||||
expect(seen?.toolFilter).toEqual({ deny: ['subagent'] })
|
expect(seen?.toolFilter).toEqual({ deny: ['subagent'] })
|
||||||
expect(seen?.toolFilter).not.toHaveProperty('allow')
|
expect(seen?.toolFilter).not.toHaveProperty('allow')
|
||||||
})
|
})
|
||||||
|
|
||||||
|
it('an explicit empty toolFilter fails at plugin load, not at first delegation', async () => {
|
||||||
|
const ctx = new Context()
|
||||||
|
await ctx.plugin(SystemPrompt)
|
||||||
|
await ctx.plugin(ToolRegistry)
|
||||||
|
await ctx.plugin(SubagentService)
|
||||||
|
ctx.subagents.registerProvider({
|
||||||
|
name: 'p',
|
||||||
|
capabilities: { outputSchema: false, depthLimit: false, toolFilter: true, persona: false },
|
||||||
|
inheritsParentContext: false,
|
||||||
|
start: () => { throw new Error('unreachable') },
|
||||||
|
})
|
||||||
|
const fiber = ctx.plugin(tool, { provider: 'p', toolFilter: {} })
|
||||||
|
await expect(fiber).rejects.toThrow(/names neither `allow` nor `deny`/)
|
||||||
|
})
|
||||||
})
|
})
|
||||||
|
|||||||
Reference in New Issue
Block a user