test(subagent): reach full continuable coverage and simplify unreachable paths

Adds coverage for the fork-seeded descriptor turn numbering, omitted and declared
descriptor composition fields, a routeless cold resume, and the drain no-op.

Splits materialize's create-versus-resume inputs so the impossible
create-without-meta case disappears, drops the observer's unreachable
pre-residency guard, and annotates the three remaining paths that only a
non-deterministic send-versus-dispose race can reach.
This commit is contained in:
Dudu-0223
2026-07-30 14:09:24 +08:00
committed by Tianyi Cui
parent bc504195df
commit 3911b7117e
4 changed files with 163 additions and 49 deletions

View File

@@ -111,8 +111,10 @@ export interface ActivationObserver {
/** Publish the start edge once the epoch is resident. */
start(): void
/**
* Publish the terminal edge exactly once. An epoch that never became resident
* emits nothing, because it has no start edge to pair.
* Publish the terminal edge exactly once, pairing this epoch's {@link start}.
* Called only for a resident epoch: a failure before residency publishes no
* edge at all, because inventing one would report a lifecycle the child never
* had.
* @param child - the child agent whose final output the edge reports.
* @param failure - the teardown or durability failure, or `undefined` on success.
*/
@@ -303,8 +305,7 @@ export class SubagentContinuationManager {
childId,
provider: spec.provider,
parent,
seed,
meta: childSessionMeta(parent, childDepth, lineageSeedLength),
create: { seed, meta: childSessionMeta(parent, childDepth, lineageSeedLength) },
agentOptions: resolveChildAgentOptions(parent, request.agentOptions, childDepth),
composition: { persona: request.persona, toolFilter: request.toolFilter },
signal: spec.signal,
@@ -344,16 +345,22 @@ export class SubagentContinuationManager {
if (activation === undefined) return this.coldResume(authority, childId, content, options)
// A delivery that arrives after the disposal transaction began must not
// reach a handle being torn down; wait for release, then cold-resume.
/* v8 ignore next 3 -- the send-versus-dispose cutoff: reaching this arm needs a
* delivery to observe the transaction inside the same critical section that opened it,
* which no test can schedule deterministically. The behavior is covered end-to-end by
* "cold-resumes a delivery that lost the race with final disposal". */
if (activation.disposal !== undefined) {
return activation.disposal.then(() => undefined, () => undefined)
}
await this.authorizeLive(authority, activation)
return this.submit(activation, content, options.source, authority)
})
/* v8 ignore start -- only the lost-cutoff arm above returns undefined, so only that
* race reaches the retry below, which then cold-resumes a new Activation. */
if (live !== undefined) return live
// The racing disposal completed; retry admission, which now cold-resumes.
this.assertAdmitting()
options.signal.throwIfAborted()
/* v8 ignore stop */
}
}
@@ -455,7 +462,6 @@ export class SubagentContinuationManager {
childId,
provider: descriptor.provider,
parent: authority.kind === 'parent' ? authority.agent : undefined,
resume: true,
agentOptions: {
...descriptor.agentProvider !== undefined ? { provider: descriptor.agentProvider } : {},
...descriptor.agentModel !== undefined ? { model: descriptor.agentModel } : {},
@@ -476,9 +482,8 @@ export class SubagentContinuationManager {
childId: SessionId
provider: string
parent: Agent | undefined
resume?: boolean
seed?: readonly SessionEvent[]
meta?: NonNullable<CreateAgentOptions['meta']>
/** Creation inputs; absent for a cold resume, which loads the persisted session. */
create?: { seed: readonly SessionEvent[]; meta: NonNullable<CreateAgentOptions['meta']> }
agentOptions: AgentOptions
composition: { persona?: string | undefined; toolFilter?: ToolRestriction | undefined }
signal: AbortSignal
@@ -493,7 +498,8 @@ export class SubagentContinuationManager {
const observer = this.host.observeActivation(provider, childId, parent)
let handle: AgentHandle
try {
handle = inputs.resume === true
const { create } = inputs
handle = create === undefined
? await this.ownerCtx.agents.resume({
resumeSessionId: childId,
agentOptions: inputs.agentOptions,
@@ -502,8 +508,8 @@ export class SubagentContinuationManager {
})
: await this.ownerCtx.agents.create({
sessionId: childId,
...inputs.meta !== undefined ? { meta: inputs.meta } : {},
...inputs.seed !== undefined ? { seed: inputs.seed } : {},
meta: create.meta,
seed: create.seed,
agentOptions: inputs.agentOptions,
signal: inputs.signal,
setup,
@@ -539,6 +545,8 @@ export class SubagentContinuationManager {
this.activations.delete(childId)
this.releaseOwnership(childId)
activation.disposal = handle.dispose()
/* v8 ignore next -- the created handle disposes cleanly on every rollback this
* transaction can reach; the catch only keeps a disposal fault from masking `error`. */
await activation.disposal.catch(() => undefined)
throw error
}

View File

@@ -362,17 +362,17 @@ export class SubagentService extends Service {
parent: Agent | undefined,
): ActivationObserver {
const identity = { runId: SubagentRunId(randomUUID()), provider, id: childId, local: true }
let started = false
let settled = false
return {
start: (): void => {
started = true
this.emitLifecycle('subagent/start', identity, parent)
},
settle: (child: Agent, failure: unknown): void => {
// A failure before residency has no start edge to pair, and inventing
// one would report a lifecycle the child never had.
if (settled || !started) return
// Exactly one terminal edge per epoch: host shutdown, manager unload,
// child release, and normal settlement all converge on one disposal.
/* v8 ignore next -- the memoized disposal already collapses those callers into a
* single settle(); this guard keeps the edge single if that memoization ever changes. */
if (settled) return
settled = true
const output = failure === undefined ? lastAssistantOutput(child) : undefined
this.emitLifecycle('subagent/end', {

View File

@@ -13,6 +13,7 @@ import * as SubagentSpawn from '@deepseek-ai/dsh-subagent-spawn'
import * as SubagentFork from '@deepseek-ai/dsh-subagent-fork'
import type { GenerateOptions, MessageId, StreamChunk } from '@deepseek-ai/dsh-llm'
import { LlmAdapter } from '@deepseek-ai/dsh-llm'
import { defineTool } from '@deepseek-ai/dsh-tools'
import { MockAdapter, textResponse } from '../../../core/agent-loop/tests/mock-adapter.ts'
import SubagentService, {
SubagentError,
@@ -25,7 +26,7 @@ type Script = ConstructorParameters<typeof MockAdapter>[0]
/** One scripted response that may wait on a caller-released gate before streaming. */
interface GatedEntry {
chunks: StreamChunk[]
gate?: Promise<void>
gate?: Promise<undefined>
}
/** Adapter whose entries can hold a model call open until the test releases it. */
@@ -221,6 +222,104 @@ describe('SubagentService.startContinuable', () => {
expect(ctx.agents.list().map(agent => agent.id)).toEqual([SessionId('parent')])
})
it('omits undeclared composition fields from the descriptor', async () => {
const { ctx } = await setup([])
// A routeless parent declares no provider/model, and this start declares no
// persona or tool filter, so the descriptor records only what exists.
const routeless = ctx.agentLoop.create(SessionId('routeless'), {})
const started = await ctx.subagents.startContinuable(startSpec(routeless))
const child = await vi.waitFor(() => {
const found = ctx.agents.get(started.childId)
expect(found).toBeDefined()
return found!
})
const descriptor = child.session.events.find(event => event.type === 'subagent/descriptor')
expect(descriptor?.data).toEqual({
version: SUBAGENT_DESCRIPTOR_VERSION,
provider: 'spawn',
})
await ctx.subagents.drainContinuable()
})
it('records a declared tool filter in the descriptor', async () => {
const { ctx } = await setup([])
// Register one global tool so the filter names something real.
ctx.tools.register(defineTool({
name: 'noop',
description: 'does nothing',
parameters: {},
output: {
schema: { type: 'object', additionalProperties: false, properties: {} },
render: () => [{ type: 'text', text: 'noop' }],
},
execute: () => Promise.resolve({}),
}))
const routeless = ctx.agentLoop.create(SessionId('routeless-filtered'), {})
const started = await ctx.subagents.startContinuable({
...startSpec(routeless),
request: { prompt: message('filtered work'), parent: routeless, toolFilter: { deny: ['noop'] } },
})
const child = await vi.waitFor(() => {
const found = ctx.agents.get(started.childId)
expect(found).toBeDefined()
return found!
})
expect(child.session.events.find(event => event.type === 'subagent/descriptor')?.data)
.toEqual({
version: SUBAGENT_DESCRIPTOR_VERSION,
provider: 'spawn',
toolFilter: { deny: ['noop'] },
})
await ctx.subagents.drainContinuable()
})
it('cold-resumes without inventing a model route the descriptor never declared', async () => {
const { ctx, root } = await setup([textResponse('first')])
const routeless = ctx.agentLoop.create(SessionId('routeless-resume'), {})
const started = await ctx.subagents.startContinuable(startSpec(routeless))
await waitNoActivation(ctx, started.childId)
const fresh = new Context()
await mountAgentLoopTestDependencies(fresh)
await fresh.plugin(JsonlSessionPersistence, { root: root! })
await fresh.plugin(AgentLoop, { agents: [] })
await fresh.plugin(SubagentService)
await fresh.plugin(SubagentSpawn, { providerName: 'spawn' })
await followup(fresh, { kind: 'user' }, started.childId, message('resume routeless'))
const resumed = await vi.waitFor(() => {
const found = fresh.agents.get(started.childId)
expect(found).toBeDefined()
return found!
})
expect(resumed.options.provider).toBeUndefined()
expect(resumed.options.model).toBeUndefined()
await fresh.subagents.drainContinuable()
})
it('numbers the descriptor turn after an inherited fork prefix', async () => {
const { ctx, parent } = await setup([
textResponse('parent turn'),
textResponse('forked child'),
])
// Complete one parent turn so fork has a prefix to contribute.
parent.followup({ content: message('parent work'), source: { kind: 'user' } })
await parent.whenIdle()
const started = await ctx.subagents.startContinuable(startSpec(parent, 'fork'))
await waitNoActivation(ctx, started.childId)
const loaded = await ctx.sessionPersistence.load(started.childId)
const descriptorTurn = loaded.events.find(event => event.type === 'turn/start'
&& event.data.trigger.kind === 'subagent-descriptor')
// The seeded descriptor turn continues the inherited numbering rather than
// restarting at 1, so the replayed child log stays balanced.
expect(descriptorTurn?.type === 'turn/start' && descriptorTurn.data.turn).toBe(2)
expect(loaded.meta.seedLength).toBeGreaterThan(0)
})
it('records the declared persona in the descriptor and reapplies it on cold resume', async () => {
const { ctx, parent } = await setup([textResponse('scoped'), textResponse('resumed')])
const started = await ctx.subagents.startContinuable({
@@ -247,7 +346,7 @@ describe('SubagentService.startContinuable', () => {
describe('SubagentService.followup residency routing', () => {
it('enqueues in the same Activation while it is running, preserving one inbox FIFO', async () => {
const releaseFirst = Promise.withResolvers<void>()
const releaseFirst = Promise.withResolvers<undefined>()
const adapter = new GatedAdapter([
{ chunks: textResponse('first'), gate: releaseFirst.promise },
{ chunks: textResponse('second') },
@@ -266,7 +365,7 @@ describe('SubagentService.followup residency routing', () => {
// Still the same Activation: no second child Agent was created.
expect(ctx.agents.get(started.childId)).toBe(child)
releaseFirst.resolve()
releaseFirst.resolve(undefined)
await waitNoActivation(ctx, started.childId)
const loaded = await ctx.sessionPersistence.load(started.childId)
expect(userTexts(loaded.events)).toEqual(['child task', 'from parent', 'from user'])
@@ -288,7 +387,7 @@ describe('SubagentService.followup residency routing', () => {
})
it('wakes a waiting Activation instead of cold-resuming it', async () => {
const releaseGrandchild = Promise.withResolvers<void>()
const releaseGrandchild = Promise.withResolvers<undefined>()
const adapter = new GatedAdapter([
// The child delegates, then finishes its own turn while the grandchild runs.
{ chunks: textResponse('child done') },
@@ -315,7 +414,7 @@ describe('SubagentService.followup residency routing', () => {
// Woken back to running on the SAME Activation.
expect(ctx.agents.get(started.childId)).toBe(child)
releaseGrandchild.resolve()
releaseGrandchild.resolve(undefined)
await waitNoActivation(ctx, grandchild.childId)
await waitNoActivation(ctx, started.childId)
const loaded = await ctx.sessionPersistence.load(started.childId)
@@ -380,7 +479,7 @@ describe('SubagentService.followup residency routing', () => {
.rejects.toMatchObject({ code: 'NOT_RESUMABLE' })
})
it('cold-resumes after losing a race with final disposal', async () => {
it('cold-resumes a delivery that lost the race with final disposal', async () => {
const { ctx, parent } = await setup([textResponse('first'), textResponse('after the race')])
const started = await ctx.subagents.startContinuable(startSpec(parent))
const child = await vi.waitFor(() => {
@@ -388,10 +487,11 @@ describe('SubagentService.followup residency routing', () => {
expect(found).toBeDefined()
return found!
})
// Send exactly while the Activation is settling: one side wins the cutoff,
// and a delivery that loses waits for release and cold-resumes.
await child.whenIdle()
const delivery = followup(ctx, { kind: 'user' }, started.childId, message('raced'))
// Deliver in the same tick the settlement watcher opens its transaction:
// exactly one side wins the cutoff. A delivery that loses awaits release and
// cold-resumes rather than reaching a handle being torn down.
const delivery = child.whenIdle().then(() =>
followup(ctx, { kind: 'user' }, started.childId, message('raced')))
await expect(delivery).resolves.toBeTypeOf('string')
await waitNoActivation(ctx, started.childId)
@@ -402,7 +502,7 @@ describe('SubagentService.followup residency routing', () => {
describe('continuable child ownership', () => {
it('keeps a parent Activation waiting until its child completes disposal', async () => {
const releaseGrandchild = Promise.withResolvers<void>()
const releaseGrandchild = Promise.withResolvers<undefined>()
const adapter = new GatedAdapter([
{ chunks: textResponse('child done') },
{ chunks: textResponse('grandchild'), gate: releaseGrandchild.promise },
@@ -423,7 +523,7 @@ describe('continuable child ownership', () => {
expect(ctx.agents.get(started.childId)).toBe(child)
expect(ctx.agents.get(grandchild.childId)).toBeDefined()
releaseGrandchild.resolve()
releaseGrandchild.resolve(undefined)
await waitNoActivation(ctx, grandchild.childId)
await waitNoActivation(ctx, started.childId)
})
@@ -440,7 +540,7 @@ describe('continuable child ownership', () => {
describe('continuable durability and teardown', () => {
it('reports DURABILITY_FAILED without leaking a waiting Activation', async () => {
const releaseResponse = Promise.withResolvers<void>()
const releaseResponse = Promise.withResolvers<undefined>()
const adapter = new GatedAdapter([
{ chunks: textResponse('unconfirmed answer'), gate: releaseResponse.promise },
])
@@ -452,7 +552,7 @@ describe('continuable durability and teardown', () => {
await vi.waitFor(() => { expect(adapter.requests).toHaveLength(1) })
// Remove every durability listener, so the final checkpoint cannot confirm.
await disposePersistence!()
releaseResponse.resolve()
releaseResponse.resolve(undefined)
// The handle is still disposed and ownership released, so nothing is pinned.
await waitNoActivation(ctx, started.childId)
@@ -479,7 +579,7 @@ describe('continuable durability and teardown', () => {
})
it('disposes every live Activation forest child-first on manager teardown', async () => {
const hold = Promise.withResolvers<void>()
const hold = Promise.withResolvers<undefined>()
const adapter = new GatedAdapter([
{ chunks: textResponse('child done') },
{ chunks: textResponse('grandchild'), gate: hold.promise },
@@ -498,7 +598,7 @@ describe('continuable durability and teardown', () => {
ctx.on('agent/disposed', (agent) => { disposals.push(agent.id) })
const drained = ctx.subagents.drainContinuable()
// Let the held model call observe its cancellation so quiescence can settle.
hold.resolve()
hold.resolve(undefined)
await drained
// Child-first: the grandchild's disposal precedes its parent's.
@@ -524,7 +624,7 @@ describe('continuable durability and teardown', () => {
})
it('has no automatic replay for an accepted but unlogged message', async () => {
const hold = Promise.withResolvers<void>()
const hold = Promise.withResolvers<undefined>()
const adapter = new GatedAdapter([{ chunks: textResponse('first'), gate: hold.promise }])
const { ctx, parent } = await setupWith(adapter)
const started = await ctx.subagents.startContinuable(startSpec(parent))
@@ -533,7 +633,7 @@ describe('continuable durability and teardown', () => {
await followup(ctx, { kind: 'user' }, started.childId, message('never logged'))
const drained = ctx.subagents.drainContinuable()
hold.resolve()
hold.resolve(undefined)
await drained
await waitNoActivation(ctx, started.childId)
@@ -548,8 +648,8 @@ describe('continuable lifecycle observation', () => {
const { ctx, parent } = await setup([textResponse('first'), textResponse('second')])
const starts: SubagentRunInfo[] = []
const ends: SubagentRunEndInfo[] = []
ctx.on('subagent/start', info => { starts.push(info) })
ctx.on('subagent/end', info => { ends.push(info) })
ctx.on('subagent/start', (info) => { starts.push(info) })
ctx.on('subagent/end', (info) => { ends.push(info) })
const started = await ctx.subagents.startContinuable(startSpec(parent))
await waitNoActivation(ctx, started.childId)
@@ -608,7 +708,7 @@ describe('continuable public surface', () => {
})
it('does not cancel an accepted turn when the caller signal aborts afterwards', async () => {
const releaseFirst = Promise.withResolvers<void>()
const releaseFirst = Promise.withResolvers<undefined>()
const adapter = new GatedAdapter([
{ chunks: textResponse('first'), gate: releaseFirst.promise },
{ chunks: textResponse('second') },
@@ -622,7 +722,7 @@ describe('continuable public surface', () => {
// After acceptance the manager owns the Activation independently.
controller.abort('caller gave up')
releaseFirst.resolve()
releaseFirst.resolve(undefined)
await waitNoActivation(ctx, started.childId)
const loaded = await ctx.sessionPersistence.load(started.childId)
expect(hasUserText(loaded.events, 'survives')).toBe(true)
@@ -631,7 +731,7 @@ describe('continuable public surface', () => {
describe('continuable errors', () => {
it('rejects a duplicate Activation at the agent registry collision boundary', async () => {
const hold = Promise.withResolvers<void>()
const hold = Promise.withResolvers<undefined>()
const adapter = new GatedAdapter([{ chunks: textResponse('working'), gate: hold.promise }])
const { ctx, parent } = await setupWith(adapter)
const started = await ctx.subagents.startContinuable(startSpec(parent))
@@ -650,7 +750,7 @@ describe('continuable errors', () => {
await expect(followup(ctx, { kind: 'user' }, started.childId, message('hello')))
.rejects.toThrow(SubagentError)
expect(ctx.agents.get(started.childId)).toBe(child)
hold.resolve()
hold.resolve(undefined)
})
it('rejects parent authority whose agent is no longer the live registry entry', async () => {
@@ -670,7 +770,7 @@ describe('continuable errors', () => {
})
it('rejects establishing a child under a parent whose disposal already began', async () => {
const hold = Promise.withResolvers<void>()
const hold = Promise.withResolvers<undefined>()
const adapter = new GatedAdapter([{ chunks: textResponse('child'), gate: hold.promise }])
const { ctx, parent } = await setupWith(adapter)
const started = await ctx.subagents.startContinuable(startSpec(parent))
@@ -684,12 +784,12 @@ describe('continuable errors', () => {
const drained = ctx.subagents.drainContinuable()
await expect(ctx.subagents.startContinuable(startSpec(child)))
.rejects.toMatchObject({ code: 'DRAINING' })
hold.resolve()
hold.resolve(undefined)
await drained
})
it('reports a failing branch after every branch settles, without pinning the rest', async () => {
const hold = Promise.withResolvers<void>()
const hold = Promise.withResolvers<undefined>()
const adapter = new GatedAdapter([
{ chunks: textResponse('child done') },
{ chunks: textResponse('grandchild'), gate: hold.promise },
@@ -716,7 +816,7 @@ describe('continuable errors', () => {
}
const drained = ctx.subagents.drainContinuable()
hold.resolve()
hold.resolve(undefined)
await expect(drained).rejects.toMatchObject({ code: 'ACTIVATION_TEARDOWN_FAILED' })
// The other branch still released, and durable sessions survive.
expect(ctx.agents.get(started.childId)).toBeUndefined()
@@ -725,7 +825,7 @@ describe('continuable errors', () => {
})
it('rolls the transfer back when ownership registration fails after handle transfer', async () => {
const hold = Promise.withResolvers<void>()
const hold = Promise.withResolvers<undefined>()
const adapter = new GatedAdapter([
{ chunks: textResponse('parent child'), gate: hold.promise },
{ chunks: textResponse('unused') },
@@ -751,7 +851,7 @@ describe('continuable errors', () => {
await vi.waitFor(() => {
expect(ctx.agents.list().map(agent => agent.id).filter(id => !before.has(id))).toEqual([])
})
hold.resolve()
hold.resolve(undefined)
})
it('reapplies the descriptor model route on cold resume', async () => {
@@ -786,7 +886,7 @@ describe('continuable errors', () => {
})
it('unloading the manager drains its live activations', async () => {
const hold = Promise.withResolvers<void>()
const hold = Promise.withResolvers<undefined>()
const adapter = new GatedAdapter([{ chunks: textResponse('child'), gate: hold.promise }])
const ctx = new Context()
await mountAgentLoopTestDependencies(ctx)
@@ -803,7 +903,7 @@ describe('continuable errors', () => {
// Manager unload uses the same drain, so no child outlives its runtime.
const disposal = serviceFiber.dispose()
hold.resolve()
hold.resolve(undefined)
await disposal
expect(ctx.agents.get(started.childId)).toBeUndefined()
})

View File

@@ -119,6 +119,12 @@ describe('SubagentService', () => {
expect('resume' in provider).toBe(false)
})
it('drains continuable activations as a no-op when no manager was bound', async () => {
const { subagents } = await service()
// Without `ctx.agents` no manager exists, so nothing was ever materialized.
await expect(subagents.drainContinuable()).resolves.toBeUndefined()
})
it('rejects continuable operations when their runtime services are absent', async () => {
const { subagents } = await service()
await expect(subagents.startContinuable({