diff --git a/packages/subagent/subagent/src/continuation.ts b/packages/subagent/subagent/src/continuation.ts index e1a03489d7..3ab4d342ef 100644 --- a/packages/subagent/subagent/src/continuation.ts +++ b/packages/subagent/subagent/src/continuation.ts @@ -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 + /** Creation inputs; absent for a cold resume, which loads the persisted session. */ + create?: { seed: readonly SessionEvent[]; meta: NonNullable } 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 } diff --git a/packages/subagent/subagent/src/index.ts b/packages/subagent/subagent/src/index.ts index 5f5f0d6fab..0789113e7e 100644 --- a/packages/subagent/subagent/src/index.ts +++ b/packages/subagent/subagent/src/index.ts @@ -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', { diff --git a/packages/subagent/subagent/tests/continuation.spec.ts b/packages/subagent/subagent/tests/continuation.spec.ts index 2728415fde..612c6763e8 100644 --- a/packages/subagent/subagent/tests/continuation.spec.ts +++ b/packages/subagent/subagent/tests/continuation.spec.ts @@ -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[0] /** One scripted response that may wait on a caller-released gate before streaming. */ interface GatedEntry { chunks: StreamChunk[] - gate?: Promise + gate?: Promise } /** 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() + const releaseFirst = Promise.withResolvers() 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() + const releaseGrandchild = Promise.withResolvers() 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() + const releaseGrandchild = Promise.withResolvers() 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() + const releaseResponse = Promise.withResolvers() 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() + const hold = Promise.withResolvers() 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() + const hold = Promise.withResolvers() 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() + const releaseFirst = Promise.withResolvers() 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() + const hold = Promise.withResolvers() 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() + const hold = Promise.withResolvers() 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() + const hold = Promise.withResolvers() 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() + const hold = Promise.withResolvers() 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() + const hold = Promise.withResolvers() 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() }) diff --git a/packages/subagent/subagent/tests/service.spec.ts b/packages/subagent/subagent/tests/service.spec.ts index 8b68e9554f..2260aaf119 100644 --- a/packages/subagent/subagent/tests/service.spec.ts +++ b/packages/subagent/subagent/tests/service.spec.ts @@ -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({