|
|
@@ -15,7 +15,7 @@ import type { GenerateOptions, MessageId, StreamChunk } from '@deepseek-ai/dsh-l
|
|
|
import { CallId, createUserMessage, LlmAdapter } from '@deepseek-ai/dsh-llm'
|
|
|
import { defineTool } from '@deepseek-ai/dsh-tools'
|
|
|
import InvariantService from '@deepseek-ai/dsh-invariants'
|
|
|
-import { MockAdapter, textResponse, toolCallResponse } from '../../../core/agent-loop/tests/mock-adapter.ts'
|
|
|
+import { MockAdapter, maxTokensResponse, textResponse, toolCallResponse } from '../../../core/agent-loop/tests/mock-adapter.ts'
|
|
|
import SubagentService, {
|
|
|
SubagentError,
|
|
|
SUBAGENT_DESCRIPTOR_VERSION,
|
|
|
@@ -142,6 +142,18 @@ async function waitNoActivation(ctx: Context, childId: SessionId): Promise<void>
|
|
|
}, { timeout: 5_000 })
|
|
|
}
|
|
|
|
|
|
+/**
|
|
|
+ * Keep the top-level test parent out of a scripted model corpus. Every child
|
|
|
+ * settlement wakes its parent, so a suite that scripts only child responses
|
|
|
+ * would otherwise spend them on the parent's own turns.
|
|
|
+ */
|
|
|
+function parkParent(ctx: Context, parent: Agent): void {
|
|
|
+ ctx.on('agent/pre-step', async ({ agent: subject }, next) => {
|
|
|
+ if (subject !== parent) return next()
|
|
|
+ return { kind: 'reject' as const }
|
|
|
+ })
|
|
|
+}
|
|
|
+
|
|
|
/** Observe calls at the Agent cancellation boundary without a production event. */
|
|
|
function observeCancel(agent: Agent, callback: () => void): void {
|
|
|
const cancel = agent.cancel.bind(agent)
|
|
|
@@ -531,7 +543,10 @@ describe('SubagentService.followup residency routing', () => {
|
|
|
await waitNoActivation(ctx, grandchild.childId)
|
|
|
await waitNoActivation(ctx, started.childId)
|
|
|
const loaded = await ctx.sessionPersistence.load(started.childId)
|
|
|
- expect(userTexts(loaded.events)).toEqual(['child task', 'while waiting'])
|
|
|
+ // This child is itself a parent, so its grandchild's settlement notice is
|
|
|
+ // an ordinary later user message in its log.
|
|
|
+ expect(userTexts(loaded.events).slice(0, 2)).toEqual(['child task', 'while waiting'])
|
|
|
+ expect(userTexts(loaded.events).slice(2).join('\n')).toContain('finished and will do no further work')
|
|
|
})
|
|
|
|
|
|
it('rejects a parent that is not the durable direct parent', async () => {
|
|
|
@@ -1182,6 +1197,7 @@ describe('continuable review regressions', () => {
|
|
|
|
|
|
it('reports this epoch\'s own output, captured while the child was still live', async () => {
|
|
|
const { ctx, parent } = await setup([textResponse('first answer'), textResponse('second answer')])
|
|
|
+ parkParent(ctx, parent)
|
|
|
const ends: SubagentRunEndInfo[] = []
|
|
|
ctx.on('subagent/end', (info) => { ends.push(info) })
|
|
|
|
|
|
@@ -1254,9 +1270,10 @@ describe('continuable review regressions', () => {
|
|
|
await waitNoActivation(ctx, started.childId)
|
|
|
|
|
|
await vi.waitFor(() => { expect(ends).toHaveLength(1) })
|
|
|
- // Reading the whole session would resurrect 'first answer' here.
|
|
|
+ // Reading the whole session would resurrect 'first answer' here. The
|
|
|
+ // rejection discarded the claimed follow-up, so the epoch reads as refused.
|
|
|
expect(ends[0]!.lastAssistantMessage).toBeUndefined()
|
|
|
- expect(ends[0]!.stopReason).toBe('completed')
|
|
|
+ expect(ends[0]!.stopReason).toBe('refusal')
|
|
|
})
|
|
|
|
|
|
it('reports handle-disposal failure on the terminal edge', async () => {
|
|
|
@@ -1439,11 +1456,13 @@ describe('continuable review regressions', () => {
|
|
|
await waitNoActivation(ctx, started.childId)
|
|
|
})
|
|
|
|
|
|
- it('reports completed when no ordinary turn closed', async () => {
|
|
|
+ it('reports a prompt a pre-step rejection discarded as refusal', async () => {
|
|
|
const { ctx, parent } = await setup([])
|
|
|
+ parkParent(ctx, parent)
|
|
|
const ends: SubagentRunEndInfo[] = []
|
|
|
ctx.on('subagent/end', (info) => { ends.push(info) })
|
|
|
- // Block admission so the child's only turn never opens.
|
|
|
+ // A UserPromptSubmit deny or a policy plugin: the child claims its prompt,
|
|
|
+ // the rejection discards it, and no step ever runs.
|
|
|
ctx.on('agent/pre-step', async ({ agent: subject }, next) => {
|
|
|
if (subject === parent) return next()
|
|
|
return { kind: 'reject' }
|
|
|
@@ -1452,8 +1471,10 @@ describe('continuable review regressions', () => {
|
|
|
const started = await ctx.subagents.startContinuable(startSpec(parent))
|
|
|
await waitNoActivation(ctx, started.childId)
|
|
|
|
|
|
+ // The parent would otherwise believe a vetoed delivery was done and never
|
|
|
+ // resend it — the one failure the settlement promise says cannot happen.
|
|
|
await vi.waitFor(() => { expect(ends).toHaveLength(1) })
|
|
|
- expect(ends[0]!.stopReason).toBe('completed')
|
|
|
+ expect(ends[0]!.stopReason).toBe('refusal')
|
|
|
})
|
|
|
|
|
|
it('retains the Activation while an accepted message is still in the inbox', async () => {
|
|
|
@@ -1482,15 +1503,538 @@ describe('continuable review regressions', () => {
|
|
|
expect(ctx.agents.get(started.childId)).toBe(child)
|
|
|
releaseFirst.resolve(undefined)
|
|
|
await waitNoActivation(ctx, started.childId)
|
|
|
- expect(adapter.requests).toHaveLength(2)
|
|
|
+ // Two child turns; the third request is the parent's own turn on the
|
|
|
+ // settlement notice.
|
|
|
+ expect(adapter.requests.filter(request => request.sessionId === started.childId)).toHaveLength(2)
|
|
|
const loaded = await ctx.sessionPersistence.load(started.childId)
|
|
|
expect(hasUserText(loaded.events, 'queued')).toBe(true)
|
|
|
})
|
|
|
})
|
|
|
|
|
|
+/** Every settlement notice this agent received, in order, as flat text. */
|
|
|
+function settlementNotices(agent: Agent): { sender: string; text: string; summary: string }[] {
|
|
|
+ const logged = agent.session.events.flatMap(event => event.type === 'user/message' ? [event.data] : [])
|
|
|
+ return [...logged, ...agent.inbox.nextStep, ...agent.inbox.nextTurn].flatMap((message) => {
|
|
|
+ if (message.source.kind !== 'subagent-settled') return []
|
|
|
+ return [{
|
|
|
+ sender: message.source.senderSessionId,
|
|
|
+ summary: message.source.summary,
|
|
|
+ text: message.content.flatMap(block => block.type === 'text' ? [block.text] : []).join('\n'),
|
|
|
+ }]
|
|
|
+ })
|
|
|
+}
|
|
|
+
|
|
|
+describe('continuable settlement delivery', () => {
|
|
|
+ it('tells the parent what the child finished with, without being asked', async () => {
|
|
|
+ const { ctx, parent } = await setup([textResponse('the answer'), textResponse('parent ack')])
|
|
|
+ const started = await ctx.subagents.startContinuable(startSpec(parent))
|
|
|
+ await waitNoActivation(ctx, started.childId)
|
|
|
+
|
|
|
+ await vi.waitFor(() => { expect(settlementNotices(parent)).toHaveLength(1) })
|
|
|
+ const notice = settlementNotices(parent)[0]!
|
|
|
+ expect(notice.sender).toBe(started.childId)
|
|
|
+ expect(notice.text).toBe(
|
|
|
+ `Background subagent ${started.childId} finished and will do no further work unless you send it more.`
|
|
|
+ + '\nIts closing message:\nthe answer',
|
|
|
+ )
|
|
|
+ // The collapsed row states the outcome without the child's content.
|
|
|
+ expect(notice.summary).toBe(
|
|
|
+ `Background subagent ${started.childId} finished and will do no further work unless you send it more.`,
|
|
|
+ )
|
|
|
+ })
|
|
|
+
|
|
|
+ it('delivers even when the child already reported for itself', async () => {
|
|
|
+ const { ctx, parent } = await setup([textResponse('the answer'), textResponse('parent ack')])
|
|
|
+ const started = await ctx.subagents.startContinuable(startSpec(parent))
|
|
|
+ const child = await vi.waitFor(() => {
|
|
|
+ const live = ctx.agents.get(started.childId)
|
|
|
+ expect(live).toBeDefined()
|
|
|
+ return live!
|
|
|
+ })
|
|
|
+ await ctx.subagents.reportFrom(child, message('an explicit report'), {
|
|
|
+ delivery: 'quiet',
|
|
|
+ signal: testSignal,
|
|
|
+ })
|
|
|
+ await waitNoActivation(ctx, started.childId)
|
|
|
+
|
|
|
+ // The contract is unconditional precisely so the parent-side tool
|
|
|
+ // description can promise it; bookkeeping "did it report?" would make the
|
|
|
+ // promise conditional on a channel this manager does not own.
|
|
|
+ await vi.waitFor(() => { expect(settlementNotices(parent)).toHaveLength(1) })
|
|
|
+ })
|
|
|
+
|
|
|
+ it('delivers the terminal reason when the child never had a chance to report', async () => {
|
|
|
+ const { ctx, parent } = await setup([maxTokensResponse('half an ans'), textResponse('parent ack')])
|
|
|
+ const started = await ctx.subagents.startContinuable(startSpec(parent))
|
|
|
+ await waitNoActivation(ctx, started.childId)
|
|
|
+
|
|
|
+ await vi.waitFor(() => { expect(settlementNotices(parent)).toHaveLength(1) })
|
|
|
+ expect(settlementNotices(parent)[0]!.text).toBe(
|
|
|
+ `Background subagent ${started.childId} ran out of room before it finished.`
|
|
|
+ + '\nIts closing message:\nhalf an ans',
|
|
|
+ )
|
|
|
+ })
|
|
|
+
|
|
|
+ it('tells the parent a policy-rejected delivery was declined, not finished', async () => {
|
|
|
+ const { ctx, parent } = await setup([textResponse('parent ack')])
|
|
|
+ // A pre-step rejection on the child — a UserPromptSubmit deny, a policy
|
|
|
+ // plugin — discards the claimed prompt without running it.
|
|
|
+ ctx.on('agent/pre-step', async ({ agent: subject }, next) => {
|
|
|
+ if (subject === parent) return next()
|
|
|
+ return { kind: 'reject' }
|
|
|
+ })
|
|
|
+
|
|
|
+ const started = await ctx.subagents.startContinuable(startSpec(parent))
|
|
|
+ await waitNoActivation(ctx, started.childId)
|
|
|
+
|
|
|
+ await vi.waitFor(() => { expect(settlementNotices(parent)).toHaveLength(1) })
|
|
|
+ expect(settlementNotices(parent)[0]!.text).toBe(
|
|
|
+ `Background subagent ${started.childId} declined the task.`
|
|
|
+ + '\nIt left no closing message.',
|
|
|
+ )
|
|
|
+ })
|
|
|
+
|
|
|
+ it('reports a turn that failed before reaching its first step', async () => {
|
|
|
+ const releaseFirst = Promise.withResolvers<undefined>()
|
|
|
+ const adapter = new GatedAdapter([
|
|
|
+ { chunks: textResponse('the answer'), gate: releaseFirst.promise },
|
|
|
+ { chunks: textResponse('parent ack') },
|
|
|
+ ])
|
|
|
+ const { ctx, parent } = await setupWith(adapter)
|
|
|
+ // The shipped durability checkpoint (`dsh-session-checkpoint-policy`) is
|
|
|
+ // fail-closed at the step boundary, so a rejected write ends the turn after
|
|
|
+ // it claimed its messages and before it entered a step.
|
|
|
+ ctx.on('agent/pre-step', async ({ agent: subject, turn }, next) => {
|
|
|
+ if (subject.session.header.parentSession === undefined || turn < 2) return next()
|
|
|
+ throw new Error('ENOSPC: no space left on device')
|
|
|
+ })
|
|
|
+
|
|
|
+ const started = await ctx.subagents.startContinuable(startSpec(parent))
|
|
|
+ await followup(ctx, parent, started.childId, message('second task'))
|
|
|
+ releaseFirst.resolve(undefined)
|
|
|
+ await waitNoActivation(ctx, started.childId)
|
|
|
+
|
|
|
+ await vi.waitFor(() => { expect(settlementNotices(parent)).toHaveLength(1) })
|
|
|
+ // The parent must not be told the child finished: the delivery it is still
|
|
|
+ // waiting on was claimed out of the inbox and then swallowed by the failure.
|
|
|
+ const child = await ctx.sessionPersistence.load(started.childId)
|
|
|
+ expect(hasUserText(child.events, 'second task')).toBe(false)
|
|
|
+ expect(settlementNotices(parent)[0]!.text).toBe(
|
|
|
+ `Background subagent ${started.childId} failed before it finished.`
|
|
|
+ + '\nIts closing message:\nthe answer',
|
|
|
+ )
|
|
|
+ })
|
|
|
+
|
|
|
+ it('reports accepted work cut short before its first step as stopped', async () => {
|
|
|
+ const releaseFirst = Promise.withResolvers<undefined>()
|
|
|
+ const releaseCheckpoint = Promise.withResolvers<undefined>()
|
|
|
+ const atCheckpoint = Promise.withResolvers<undefined>()
|
|
|
+ const adapter = new GatedAdapter([{ chunks: textResponse('the answer'), gate: releaseFirst.promise }])
|
|
|
+ const { ctx, parent } = await setupWith(adapter)
|
|
|
+ // A step-boundary participant — the shipped durability checkpoint, a hook,
|
|
|
+ // prompt assembly — holding the child's second turn open before its first
|
|
|
+ // step, which is where teardown cancellation then catches it.
|
|
|
+ ctx.on('agent/pre-step', async ({ agent: subject, turn }, next) => {
|
|
|
+ if (subject.session.header.parentSession === undefined || turn < 2) return next()
|
|
|
+ atCheckpoint.resolve(undefined)
|
|
|
+ await releaseCheckpoint.promise
|
|
|
+ return next()
|
|
|
+ })
|
|
|
+
|
|
|
+ const started = await ctx.subagents.startContinuable(startSpec(parent))
|
|
|
+ // Queued while turn 1 still runs, so turn 2 opens and claims it without a
|
|
|
+ // second model call: the Activation is mid-turn when the drain cancels it.
|
|
|
+ await followup(ctx, parent, started.childId, message('second task'))
|
|
|
+ releaseFirst.resolve(undefined)
|
|
|
+ await atCheckpoint.promise
|
|
|
+ const drained = drainManager(ctx)
|
|
|
+ releaseCheckpoint.resolve(undefined)
|
|
|
+ await drained
|
|
|
+
|
|
|
+ // Turn 2 leaves a balanced no-step `aborted` end, so the log alone would
|
|
|
+ // answer with turn 1's clean completion and tell the parent its still-unrun
|
|
|
+ // task had finished.
|
|
|
+ await vi.waitFor(() => { expect(settlementNotices(parent)).toHaveLength(1) })
|
|
|
+ expect(settlementNotices(parent)[0]!.text).toBe(
|
|
|
+ `Background subagent ${started.childId} was stopped before it finished.`
|
|
|
+ + '\nIts closing message:\nthe answer',
|
|
|
+ )
|
|
|
+ })
|
|
|
+
|
|
|
+ it('reports a child stopped before it ever reached the model as stopped', async () => {
|
|
|
+ const releaseCheckpoint = Promise.withResolvers<undefined>()
|
|
|
+ const atCheckpoint = Promise.withResolvers<undefined>()
|
|
|
+ const { ctx, parent } = await setupWith(new GatedAdapter([]))
|
|
|
+ ctx.on('agent/pre-step', async ({ agent: subject }, next) => {
|
|
|
+ if (subject.session.header.parentSession === undefined) return next()
|
|
|
+ atCheckpoint.resolve(undefined)
|
|
|
+ await releaseCheckpoint.promise
|
|
|
+ return next()
|
|
|
+ })
|
|
|
+
|
|
|
+ const started = await ctx.subagents.startContinuable(startSpec(parent))
|
|
|
+ await atCheckpoint.promise
|
|
|
+ const drained = drainManager(ctx)
|
|
|
+ releaseCheckpoint.resolve(undefined)
|
|
|
+ await drained
|
|
|
+
|
|
|
+ // This epoch closed no stepped turn at all, which on its own reads as "had
|
|
|
+ // nothing to report"; only the interruption distinguishes it from a child
|
|
|
+ // that genuinely finished with no output.
|
|
|
+ await vi.waitFor(() => { expect(settlementNotices(parent)).toHaveLength(1) })
|
|
|
+ expect(settlementNotices(parent)[0]!.text).toBe(
|
|
|
+ `Background subagent ${started.childId} was stopped before it finished.`
|
|
|
+ + '\nIt left no closing message.',
|
|
|
+ )
|
|
|
+ })
|
|
|
+
|
|
|
+ it('reports a child an ancestor interrupted before its first step as stopped', async () => {
|
|
|
+ const atCheckpoint = Promise.withResolvers<undefined>()
|
|
|
+ const releaseCheckpoint = Promise.withResolvers<undefined>()
|
|
|
+ const adapter = new GatedAdapter([{ chunks: textResponse('parent ack') }])
|
|
|
+ const { ctx, parent } = await setupWith(adapter)
|
|
|
+ ctx.on('agent/pre-step', async ({ agent: subject }, next) => {
|
|
|
+ if (subject.session.header.parentSession === undefined) return next()
|
|
|
+ atCheckpoint.resolve(undefined)
|
|
|
+ await releaseCheckpoint.promise
|
|
|
+ return next()
|
|
|
+ })
|
|
|
+
|
|
|
+ const started = await ctx.subagents.startContinuable(startSpec(parent))
|
|
|
+ await atCheckpoint.promise
|
|
|
+ // The shipped interrupt path: nothing about it runs inside this manager, so
|
|
|
+ // no pre-cancel sample could see it — the child's own log has to say so.
|
|
|
+ ctx.subagents.interrupt(started.childId, { kind: 'ancestor', agent: parent })
|
|
|
+ releaseCheckpoint.resolve(undefined)
|
|
|
+ await waitNoActivation(ctx, started.childId)
|
|
|
+
|
|
|
+ await vi.waitFor(() => { expect(settlementNotices(parent)).toHaveLength(1) })
|
|
|
+ expect(settlementNotices(parent)[0]!.text).toBe(
|
|
|
+ `Background subagent ${started.childId} was stopped before it finished.`
|
|
|
+ + '\nIt left no closing message.',
|
|
|
+ )
|
|
|
+ })
|
|
|
+
|
|
|
+ it('reports accepted work cancelled before any turn could open as stopped', async () => {
|
|
|
+ const releaseChild = Promise.withResolvers<undefined>()
|
|
|
+ const releaseGrandchild = Promise.withResolvers<undefined>()
|
|
|
+ const releaseMaintenance = Promise.withResolvers<undefined>()
|
|
|
+ const adapter = new GatedAdapter([
|
|
|
+ { chunks: textResponse('the answer'), gate: releaseChild.promise },
|
|
|
+ { chunks: textResponse('grandchild'), gate: releaseGrandchild.promise },
|
|
|
+ ])
|
|
|
+ const { ctx, parent } = await setupWith(adapter)
|
|
|
+ const started = await ctx.subagents.startContinuable(startSpec(parent))
|
|
|
+ const child = await vi.waitFor(() => {
|
|
|
+ const live = ctx.agents.get(started.childId)
|
|
|
+ expect(live).toBeDefined()
|
|
|
+ return live!
|
|
|
+ })
|
|
|
+ // A descendant keeps the child resident once its own turn closes, so the
|
|
|
+ // maintenance phase below is reachable without racing settlement.
|
|
|
+ const grandchild = await ctx.subagents.startContinuable(startSpec(child))
|
|
|
+ await vi.waitFor(() => { expect(ctx.agents.get(grandchild.childId)).toBeDefined() })
|
|
|
+ releaseChild.resolve(undefined)
|
|
|
+ await vi.waitFor(() => { expect(child.status).toBe('idle') })
|
|
|
+
|
|
|
+ // Context maintenance folds into `idle` and defers waking work, so this
|
|
|
+ // delivery is accepted with no turn to claim it.
|
|
|
+ const maintaining = child.runMaintenance(async () => { await releaseMaintenance.promise })
|
|
|
+ await followup(ctx, parent, started.childId, message('never runs'))
|
|
|
+ const drained = drainManager(ctx)
|
|
|
+ releaseMaintenance.resolve(undefined)
|
|
|
+ releaseGrandchild.resolve(undefined)
|
|
|
+ await maintaining
|
|
|
+ await drained
|
|
|
+
|
|
|
+ // Turn 1 closed cleanly and no later turn opened, so the cancelled queue is
|
|
|
+ // the only record that this epoch was cut short.
|
|
|
+ expect(hasUserText(child.session.events, 'never runs')).toBe(false)
|
|
|
+ await vi.waitFor(() => { expect(settlementNotices(parent)).toHaveLength(1) })
|
|
|
+ expect(settlementNotices(parent)[0]!.text).toBe(
|
|
|
+ `Background subagent ${started.childId} was stopped before it finished.`
|
|
|
+ + '\nIts closing message:\nthe answer',
|
|
|
+ )
|
|
|
+ })
|
|
|
+
|
|
|
+ it('withholds an outcome the harness could not durably release', async () => {
|
|
|
+ const { ctx, parent } = await setup([textResponse('the answer'), textResponse('parent ack')])
|
|
|
+ const started = await ctx.subagents.startContinuable(startSpec(parent))
|
|
|
+ const manager = (ctx.subagents as unknown as {
|
|
|
+ continuations: { activations: Map<SessionId, { handle: { dispose(): Promise<void> } }> }
|
|
|
+ }).continuations
|
|
|
+ const activation = await vi.waitFor(() => {
|
|
|
+ const live = manager.activations.get(started.childId)
|
|
|
+ expect(live).toBeDefined()
|
|
|
+ return live!
|
|
|
+ })
|
|
|
+ const dispose = activation.handle.dispose.bind(activation.handle)
|
|
|
+ activation.handle.dispose = async () => {
|
|
|
+ await dispose()
|
|
|
+ throw new Error('scope unwind failed')
|
|
|
+ }
|
|
|
+
|
|
|
+ await waitNoActivation(ctx, started.childId)
|
|
|
+ await vi.waitFor(() => { expect(settlementNotices(parent)).toHaveLength(1) })
|
|
|
+ expect(settlementNotices(parent)[0]!.text).toBe(
|
|
|
+ `Background subagent ${started.childId} failed before it finished.\nIt left no closing message.`,
|
|
|
+ )
|
|
|
+ })
|
|
|
+
|
|
|
+ it('gives an idle parent one ordinary turn on the notice', async () => {
|
|
|
+ const { ctx, parent, adapter } = await setup([textResponse('the answer'), textResponse('parent ack')])
|
|
|
+ const turnStarts: number[] = []
|
|
|
+ ctx.on('session/event', (session, event) => {
|
|
|
+ if (session.id === parent.id && event.type === 'turn/start') turnStarts.push(event.data.turn)
|
|
|
+ })
|
|
|
+
|
|
|
+ const started = await ctx.subagents.startContinuable(startSpec(parent))
|
|
|
+ await waitNoActivation(ctx, started.childId)
|
|
|
+ await vi.waitFor(() => {
|
|
|
+ expect(adapter.requests.filter(request => request.sessionId === parent.id)).toHaveLength(1)
|
|
|
+ })
|
|
|
+ expect(turnStarts).toEqual([1])
|
|
|
+ })
|
|
|
+
|
|
|
+ it('batches simultaneous notices into one step of a busy parent', async () => {
|
|
|
+ const releaseChildren = Promise.withResolvers<undefined>()
|
|
|
+ const releaseParent = Promise.withResolvers<undefined>()
|
|
|
+ const adapter = new GatedAdapter([
|
|
|
+ { chunks: textResponse('parent works'), gate: releaseParent.promise },
|
|
|
+ { chunks: textResponse('first child'), gate: releaseChildren.promise },
|
|
|
+ { chunks: textResponse('second child'), gate: releaseChildren.promise },
|
|
|
+ { chunks: textResponse('parent reacts') },
|
|
|
+ ])
|
|
|
+ const { ctx, parent } = await setupWith(adapter)
|
|
|
+ // Open a parent turn first, so both notices arrive while it is running.
|
|
|
+ parent.followup(createUserMessage({ content: message('start working'), source: { kind: 'user' } }))
|
|
|
+ await vi.waitFor(() => { expect(parent.status).toBe('running') })
|
|
|
+
|
|
|
+ const first = await ctx.subagents.startContinuable(startSpec(parent))
|
|
|
+ const second = await ctx.subagents.startContinuable(startSpec(parent))
|
|
|
+ releaseChildren.resolve(undefined)
|
|
|
+ await waitNoActivation(ctx, first.childId)
|
|
|
+ await waitNoActivation(ctx, second.childId)
|
|
|
+
|
|
|
+ // Both notices are waiting for the same step boundary, not two turns.
|
|
|
+ expect(parent.inbox.nextStep).toHaveLength(2)
|
|
|
+ expect(parent.inbox.nextTurn).toHaveLength(0)
|
|
|
+ const turnStarts: number[] = []
|
|
|
+ ctx.on('session/event', (session, event) => {
|
|
|
+ if (session.id === parent.id && event.type === 'turn/start') turnStarts.push(event.data.turn)
|
|
|
+ })
|
|
|
+ releaseParent.resolve(undefined)
|
|
|
+ await vi.waitFor(() => { expect(settlementNotices(parent)).toHaveLength(2) })
|
|
|
+ expect(turnStarts).toEqual([])
|
|
|
+ // Both children released together, so which settles first is not ordered.
|
|
|
+ expect(new Set(settlementNotices(parent).map(entry => entry.sender)))
|
|
|
+ .toEqual(new Set([first.childId, second.childId]))
|
|
|
+ })
|
|
|
+
|
|
|
+ it('holds a maintaining parent live until it can read the notice', async () => {
|
|
|
+ const releaseFirst = Promise.withResolvers<undefined>()
|
|
|
+ const releaseSecond = Promise.withResolvers<undefined>()
|
|
|
+ const adapter = new GatedAdapter([
|
|
|
+ { chunks: textResponse('outer') },
|
|
|
+ { chunks: textResponse('first inner'), gate: releaseFirst.promise },
|
|
|
+ { chunks: textResponse('second inner'), gate: releaseSecond.promise },
|
|
|
+ { chunks: textResponse('outer reacts') },
|
|
|
+ { chunks: textResponse('root reacts') },
|
|
|
+ ])
|
|
|
+ const { ctx, parent } = await setupWith(adapter)
|
|
|
+ const outer = await ctx.subagents.startContinuable(startSpec(parent))
|
|
|
+ const middle = await vi.waitFor(() => {
|
|
|
+ const live = ctx.agents.get(outer.childId)
|
|
|
+ expect(live).toBeDefined()
|
|
|
+ return live!
|
|
|
+ })
|
|
|
+ const first = await ctx.subagents.startContinuable(startSpec(middle))
|
|
|
+ const second = await ctx.subagents.startContinuable(startSpec(middle))
|
|
|
+ await vi.waitFor(() => { expect(middle.status).toBe('idle') })
|
|
|
+
|
|
|
+ // `Agent.status` folds maintenance into `idle`, and a waking send behind it
|
|
|
+ // only arms a deferred wake. The first child's release moves the middle
|
|
|
+ // Activation's settlement watcher onto its quiescence race; the second one
|
|
|
+ // then arrives at exactly the point where an unaccounted delivery would be
|
|
|
+ // judged quiet, settled, and cancelled — clearing the inbox it sits in.
|
|
|
+ const maintaining = Promise.withResolvers<undefined>()
|
|
|
+ const maintenance = middle.runMaintenance(async () => { await maintaining.promise })
|
|
|
+ releaseFirst.resolve(undefined)
|
|
|
+ await waitNoActivation(ctx, first.childId)
|
|
|
+ releaseSecond.resolve(undefined)
|
|
|
+ await waitNoActivation(ctx, second.childId)
|
|
|
+ expect(ctx.agents.get(outer.childId)).toBe(middle)
|
|
|
+
|
|
|
+ maintaining.resolve(undefined)
|
|
|
+ await maintenance
|
|
|
+ await vi.waitFor(() => { expect(settlementNotices(middle)).toHaveLength(2) })
|
|
|
+ expect(settlementNotices(middle).map(entry => entry.sender))
|
|
|
+ .toEqual([first.childId, second.childId])
|
|
|
+ await waitNoActivation(ctx, outer.childId)
|
|
|
+ })
|
|
|
+
|
|
|
+ it('delivers before releasing the ownership that lets the parent settle', async () => {
|
|
|
+ const releaseChild = Promise.withResolvers<undefined>()
|
|
|
+ const adapter = new GatedAdapter([
|
|
|
+ { chunks: textResponse('outer') },
|
|
|
+ { chunks: textResponse('inner'), gate: releaseChild.promise },
|
|
|
+ { chunks: textResponse('outer reacts') },
|
|
|
+ ])
|
|
|
+ const { ctx, parent } = await setupWith(adapter)
|
|
|
+ const outer = await ctx.subagents.startContinuable(startSpec(parent))
|
|
|
+ const middle = await vi.waitFor(() => {
|
|
|
+ const live = ctx.agents.get(outer.childId)
|
|
|
+ expect(live).toBeDefined()
|
|
|
+ return live!
|
|
|
+ })
|
|
|
+ const inner = await ctx.subagents.startContinuable(startSpec(middle))
|
|
|
+ await vi.waitFor(() => { expect(middle.status).toBe('idle') })
|
|
|
+
|
|
|
+ const manager = (ctx.subagents as unknown as {
|
|
|
+ continuations: { activations: Map<SessionId, { ownedChildren: Set<SessionId> }> }
|
|
|
+ }).continuations
|
|
|
+ let ownedAtDelivery: SessionId[] | undefined
|
|
|
+ ctx.on('agent/inbox/inserted', ({ agent, message }) => {
|
|
|
+ if (agent !== middle || message.source.kind !== 'subagent-settled') return
|
|
|
+ ownedAtDelivery = [...manager.activations.get(middle.id)!.ownedChildren]
|
|
|
+ })
|
|
|
+
|
|
|
+ releaseChild.resolve(undefined)
|
|
|
+ await waitNoActivation(ctx, inner.childId)
|
|
|
+ // Still owned at delivery: the parent is structurally unable to settle in
|
|
|
+ // the window the notice crosses, rather than winning a race against it.
|
|
|
+ expect(ownedAtDelivery).toEqual([inner.childId])
|
|
|
+ await waitNoActivation(ctx, outer.childId)
|
|
|
+ })
|
|
|
+
|
|
|
+ it('does not wake a parent whose own teardown already began', async () => {
|
|
|
+ const hold = Promise.withResolvers<undefined>()
|
|
|
+ const adapter = new GatedAdapter([{ chunks: textResponse('interrupted'), gate: hold.promise }])
|
|
|
+ const { ctx, parent } = await setupWith(adapter)
|
|
|
+ const started = await ctx.subagents.startContinuable(startSpec(parent))
|
|
|
+ await vi.waitFor(() => { expect(ctx.agents.get(started.childId)).toBeDefined() })
|
|
|
+
|
|
|
+ const drained = drainManager(ctx)
|
|
|
+ hold.resolve(undefined)
|
|
|
+ await drained
|
|
|
+
|
|
|
+ // Delivered and durably logged, but no turn: waking a parent the host is
|
|
|
+ // about to dispose spends a model request nothing reads. What happens to the
|
|
|
+ // message when that parent is disposed next is pinned by the test below.
|
|
|
+ expect(settlementNotices(parent)).toHaveLength(1)
|
|
|
+ expect(settlementNotices(parent)[0]!.text).toBe(
|
|
|
+ `Background subagent ${started.childId} was stopped before it finished.`
|
|
|
+ + '\nIt left no closing message.',
|
|
|
+ )
|
|
|
+ expect(parent.session.events.some(event => event.type === 'agent/inbox/spliced')).toBe(true)
|
|
|
+ expect(parent.session.events.some(event => event.type === 'turn/start')).toBe(false)
|
|
|
+ expect(parent.status).toBe('idle')
|
|
|
+ })
|
|
|
+
|
|
|
+ it('does not wake a parent below a scoped teardown root', async () => {
|
|
|
+ const hold = Promise.withResolvers<undefined>()
|
|
|
+ const adapter = new GatedAdapter([{ chunks: textResponse('interrupted'), gate: hold.promise }])
|
|
|
+ const { ctx, parent } = await setupWith(adapter)
|
|
|
+ const started = await ctx.subagents.startContinuable(startSpec(parent))
|
|
|
+ await vi.waitFor(() => { expect(ctx.agents.get(started.childId)).toBeDefined() })
|
|
|
+
|
|
|
+ const drained = ctx.subagents.drainContinuableDescendants([parent])
|
|
|
+ hold.resolve(undefined)
|
|
|
+ await drained
|
|
|
+
|
|
|
+ expect(settlementNotices(parent)).toHaveLength(1)
|
|
|
+ expect(parent.session.events.some(event => event.type === 'turn/start')).toBe(false)
|
|
|
+ })
|
|
|
+
|
|
|
+ it('records but cannot deliver a teardown notice once the parent is disposed too', async () => {
|
|
|
+ const hold = Promise.withResolvers<undefined>()
|
|
|
+ const adapter = new GatedAdapter([{ chunks: textResponse('interrupted'), gate: hold.promise }])
|
|
|
+ const { ctx } = await setupWith(adapter)
|
|
|
+ const parentId = SessionId('closing-parent')
|
|
|
+ const host = await ctx.agents.create({
|
|
|
+ sessionId: parentId,
|
|
|
+ agentOptions: { provider: 'mock', model: 'mock' },
|
|
|
+ })
|
|
|
+ const started = await ctx.subagents.startContinuable(startSpec(host.agent))
|
|
|
+ await vi.waitFor(() => { expect(ctx.agents.get(started.childId)).toBeDefined() })
|
|
|
+
|
|
|
+ const drained = ctx.subagents.drainContinuableDescendants([host.agent])
|
|
|
+ hold.resolve(undefined)
|
|
|
+ await drained
|
|
|
+ expect(settlementNotices(host.agent)).toHaveLength(1)
|
|
|
+
|
|
|
+ // Disposal is a `keepInbox: false` cancel, so it durably cancels the notice
|
|
|
+ // it never claimed. Teardown delivery therefore reaches a parent that is
|
|
|
+ // still resident — a resumed one reads the log, not a pending message — and
|
|
|
+ // no wording anywhere may promise otherwise.
|
|
|
+ await host.dispose()
|
|
|
+ const resumed = await ctx.agents.resume({
|
|
|
+ resumeSessionId: parentId,
|
|
|
+ agentOptions: { provider: 'mock', model: 'mock' },
|
|
|
+ })
|
|
|
+ expect(settlementNotices(resumed.agent)).toEqual([])
|
|
|
+ await resumed.dispose()
|
|
|
+ // The account is still in the durable log: delivered, then cancelled unread.
|
|
|
+ const persisted = await ctx.sessionPersistence.load(parentId)
|
|
|
+ expect(persisted.events.flatMap(event => event.type === 'agent/inbox/spliced'
|
|
|
+ ? [{ inserted: event.data.inserted.length, removed: event.data.removedCount ?? 0 }]
|
|
|
+ : [])).toEqual([{ inserted: 1, removed: 0 }, { inserted: 0, removed: 1 }])
|
|
|
+ })
|
|
|
+
|
|
|
+ it('drops the notice without disturbing teardown when the parent is gone', async () => {
|
|
|
+ const releaseChild = Promise.withResolvers<undefined>()
|
|
|
+ const adapter = new GatedAdapter([{ chunks: textResponse('answer'), gate: releaseChild.promise }])
|
|
|
+ const { ctx } = await setupWith(adapter)
|
|
|
+ const warnings: string[] = []
|
|
|
+ ctx.logger.warn = (text: string) => { warnings.push(text) }
|
|
|
+ const host = await ctx.agents.create({
|
|
|
+ sessionId: SessionId('disposable-parent'),
|
|
|
+ agentOptions: { provider: 'mock', model: 'mock' },
|
|
|
+ })
|
|
|
+ const started = await ctx.subagents.startContinuable(startSpec(host.agent))
|
|
|
+ const ends: SubagentRunEndInfo[] = []
|
|
|
+ ctx.on('subagent/end', (info) => { ends.push(info) })
|
|
|
+
|
|
|
+ releaseChild.resolve(undefined)
|
|
|
+ await host.dispose()
|
|
|
+ await waitNoActivation(ctx, started.childId)
|
|
|
+ await vi.waitFor(() => { expect(ends).toHaveLength(1) })
|
|
|
+ expect(warnings).toEqual([])
|
|
|
+ })
|
|
|
+
|
|
|
+ it('logs a rejected notice instead of failing the child\'s teardown', async () => {
|
|
|
+ const { ctx, parent } = await setup([textResponse('the answer')])
|
|
|
+ const warnings: string[] = []
|
|
|
+ ctx.logger.warn = (text: string) => { warnings.push(text) }
|
|
|
+ vi.spyOn(parent, 'followup').mockImplementation(() => {
|
|
|
+ throw new Error('parent closed during delivery')
|
|
|
+ })
|
|
|
+ const ends: SubagentRunEndInfo[] = []
|
|
|
+ ctx.on('subagent/end', (info) => { ends.push(info) })
|
|
|
+
|
|
|
+ const started = await ctx.subagents.startContinuable(startSpec(parent))
|
|
|
+ await waitNoActivation(ctx, started.childId)
|
|
|
+ await vi.waitFor(() => { expect(ends).toHaveLength(1) })
|
|
|
+ expect(ends[0]!.stopReason).toBe('completed')
|
|
|
+ expect(warnings.some(warning => warning.includes('settlement notice was not delivered'))).toBe(true)
|
|
|
+ })
|
|
|
+
|
|
|
+ it('stays silent about a child the caller was told does not exist', async () => {
|
|
|
+ const { ctx, parent } = await setup([])
|
|
|
+ const drains: Promise<void>[] = []
|
|
|
+ ctx.on('subagent/start', () => { drains.push(drainManager(ctx)) })
|
|
|
+
|
|
|
+ await expect(ctx.subagents.startContinuable(startSpec(parent)))
|
|
|
+ .rejects.toMatchObject({ code: 'DRAINING' })
|
|
|
+ await Promise.all(drains)
|
|
|
+ expect(settlementNotices(parent)).toEqual([])
|
|
|
+ })
|
|
|
+})
|
|
|
+
|
|
|
describe('continuable lifecycle observation', () => {
|
|
|
it('emits one paired start/end per residency epoch', async () => {
|
|
|
const { ctx, parent } = await setup([textResponse('first'), textResponse('second')])
|
|
|
+ parkParent(ctx, parent)
|
|
|
const starts: SubagentRunInfo[] = []
|
|
|
const ends: SubagentRunEndInfo[] = []
|
|
|
ctx.on('subagent/start', (info) => { starts.push(info) })
|
|
|
@@ -1510,6 +2054,8 @@ describe('continuable lifecycle observation', () => {
|
|
|
expect(starts.map(info => info.provider)).toEqual(['spawn', 'spawn'])
|
|
|
// Each end pairs its own start's runId.
|
|
|
expect(ends.map(info => info.runId)).toEqual(starts.map(info => info.runId))
|
|
|
+ // Both epochs ran their own scripted response; neither exhausted the corpus.
|
|
|
+ expect(ends.map(info => info.stopReason)).toEqual(['completed', 'completed'])
|
|
|
})
|
|
|
})
|
|
|
|