From 6fdd04812387ea8e53ebaaf453bb7cd302efb05a Mon Sep 17 00:00:00 2001 From: Tianyi Cui <53024+tianyicui@users.noreply.github.com> Date: Wed, 17 Jun 2026 21:25:47 +0800 Subject: [PATCH] fix(agent-loop): harden lifecycle edge cases --- packages/agent-loop/src/agent.ts | 20 +++++--- packages/agent-loop/src/loop.ts | 26 +++++++--- packages/agent-loop/tests/agent.spec.ts | 65 ++++++++++++++++++++++++ packages/agent-loop/tests/loop.spec.ts | 66 ++++++++++++++++++++++++- packages/agent/src/types.ts | 11 +++-- 5 files changed, 167 insertions(+), 21 deletions(-) diff --git a/packages/agent-loop/src/agent.ts b/packages/agent-loop/src/agent.ts index 8040e1b107..96a1158c66 100644 --- a/packages/agent-loop/src/agent.ts +++ b/packages/agent-loop/src/agent.ts @@ -63,7 +63,11 @@ export class LoopAgent implements Agent { // waiter (AGENTS.md "contain callback exceptions" — a lifecycle await must // not hang on one bad listener). if (status !== 'running') this.settleIdleWaiters() - this.ctx.emit('agent/status', this, status) + try { + this.ctx.emit('agent/status', this, status) + } catch (error: unknown) { + this.ctx.logger.warn(`agent "${this.id}": agent/status listener threw on ${status}: ${String(error)}`) + } } /** @@ -177,16 +181,16 @@ export class LoopAgent implements Agent { * `running`. If it is already disposed, awaits {@link done} (the loop-exit * promise) — `agent/status('disposed')` fires in the disposer BEFORE the * driver loop has unwound, so it is NOT itself a quiescence signal. If it is - * idle, resolves immediately. Otherwise queues an internal waiter (see - * {@link idleWaiters}) released on the next running→idle/disposed transition, - * resolving on `idle` directly (the turn fully ended) or chaining {@link done} - * on `disposed` (wait for the loop to actually exit). Implements the - * {@link Agent.whenIdle} contract used by teardown (`abort()` then - * `await whenIdle()`). + * idle AND has no queued work, resolves immediately. Otherwise queues an + * internal waiter (see {@link idleWaiters}) released on the next + * running→idle/disposed transition, resolving on `idle` directly (the turn + * fully ended) or chaining {@link done} on `disposed` (wait for the loop to + * actually exit). Implements the {@link Agent.whenIdle} contract used by + * teardown (`abort()` then `await whenIdle()`). */ whenIdle(): Promise { if (this._status === 'disposed') return this.done - if (this._status !== 'running') return Promise.resolve() + if (this._status !== 'running' && !this.inbox.hasQueued) return Promise.resolve() // Register an internal waiter (resolved by settleIdleWaiters on the next // running→idle/disposed transition), NOT an effect-scoped `ctx.on` listener: // a concurrent fiber disposal runs this agent's listener disposers, which diff --git a/packages/agent-loop/src/loop.ts b/packages/agent-loop/src/loop.ts index 74ccb0f71e..a9b3fbd2c5 100644 --- a/packages/agent-loop/src/loop.ts +++ b/packages/agent-loop/src/loop.ts @@ -203,8 +203,8 @@ async function runTurn(ctx: Context, agent: LoopAgent, handle: LoopHandle, turn: // agent/step-end emit is contained: a throwing step-end listener must not // abort finalization and strand the turn open (turn/end balance > notifying // one bad listener). Appended before the emit (ADR 0003 append-before-emit). - const closeStep = (): void => { - if (!stepOpen) return + const closeStep = (): boolean => { + if (!stepOpen) return false stepOpen = false // Session.append pushes step/end BEFORE notifying session/event listeners, // so a throwing listener leaves step/end in the log (balance holds) but @@ -227,6 +227,7 @@ async function runTurn(ctx: Context, agent: LoopAgent, handle: LoopHandle, turn: // itself succeeded, AND keeps finalization going when closeStep runs from // the outer catch. if (failure !== undefined) failTurn(toError(failure)) + return failure !== undefined } // Record a step/turn failure exactly once: append the single `error` event @@ -358,7 +359,7 @@ async function runTurn(ctx: Context, agent: LoopAgent, handle: LoopHandle, turn: // Steering that arrived during streaming/tool execution. const steered = drainSteering(ctx, agent, turn) - closeStep() + if (closeStep()) break const defaultDecision = stepOutcome.hadToolCalls || steered let shouldContinue: boolean @@ -502,16 +503,27 @@ async function runStep( // tool dispatch actually uses. let message: Message = assembler.message() message = await ctx.waterfall('agent/step-result', agent, turn, step, message, () => Promise.resolve(message)) + const finish = assembler.finish + const messageForLog: Message = finish.kind === 'max-tokens' + ? { ...message, content: message.content.filter(block => block.type !== 'tool-call') } + : message - session.append('assistant/message', { turn, step, content: message.content }) + if (finish.kind !== 'max-tokens' || messageForLog.content.length > 0) { + session.append('assistant/message', { turn, step, content: messageForLog.content }) + } if (assembler.usage) { session.append('usage', { turn, step, usage: assembler.usage }) } // --- Tool execution (sequential; parallel execution is a TODO) --- // ToolRegistry.execute converts tool failures (including aborts) into - // isError results, so abort is re-checked around every call here. - const toolCalls = message.content.filter(block => block.type === 'tool-call') + // isError results, so abort is re-checked around every call here. A + // max-tokens step is cut off: any tool-call block in it may be partial, so it + // is neither dispatched nor recorded in the derived-history assistant message + // above. Raw assistant/chunk events still preserve the exact stream. + const toolCalls = finish.kind === 'max-tokens' + ? [] + : message.content.filter(block => block.type === 'tool-call') for (const call of toolCalls) { /* v8 ignore next -- signal.reason always set by agent.abort() which provides a default */ if (signal.aborted) throw new Error(String(signal.reason ?? 'aborted')) @@ -552,7 +564,7 @@ async function runStep( /* v8 ignore stop */ } - return { hadToolCalls: toolCalls.length > 0, finish: assembler.finish } + return { hadToolCalls: toolCalls.length > 0, finish } } /** The last turn number in a (possibly seeded) session log, or 0. */ diff --git a/packages/agent-loop/tests/agent.spec.ts b/packages/agent-loop/tests/agent.spec.ts index 4e2cd6c046..4fd7237061 100644 --- a/packages/agent-loop/tests/agent.spec.ts +++ b/packages/agent-loop/tests/agent.spec.ts @@ -32,6 +32,17 @@ function waitForIdle(ctx: Context, agent: LoopAgent): Promise { }) } +function waitForStatus(ctx: Context, agent: LoopAgent, expected: LoopAgent['status']): Promise { + return new Promise((resolve) => { + const dispose = ctx.on('agent/status', (subject, status) => { + if (subject === agent && status === expected) { + dispose() + resolve() + } + }) + }) +} + function send(agent: LoopAgent, text: string) { agent.send([{ type: 'text', text }]) } @@ -265,6 +276,24 @@ describe('LoopAgent', () => { expect(agent.status).not.toBe('running') }) + it('whenIdle() waits for queued work that has not flipped status yet', async () => { + const adapter = new MockAdapter(['hang']) + const ctx = await harness(adapter) + const agent = ctx.agentLoop.create('a1', { model: 'mock' }) + + send(agent, 'queued') + let settled = false + const idle = agent.whenIdle().then(() => { settled = true }) + await Promise.resolve() + expect(settled).toBe(false) + + await waitForStatus(ctx, agent, 'running') + agent.abort('done') + await idle + expect(settled).toBe(true) + expect(agent.status).toBe('idle') + }) + it('whenIdle() awaits the running→idle transition, ignoring other subjects/running events', async () => { const adapter = new MockAdapter([textResponse('ok'), textResponse('ok')]) const ctx = await harness(adapter) @@ -366,6 +395,42 @@ describe('LoopAgent', () => { expect(doneResolved).toBe(true) }) + it('contains a throwing agent/status listener on the running transition', async () => { + const adapter = new MockAdapter([textResponse('ok')]) + const ctx = await harness(adapter) + const warn = vi.spyOn(ctx.logger, 'warn').mockImplementation(() => undefined) + const agent = ctx.agentLoop.create('a1', { model: 'mock' }) + ctx.on('agent/status', (_subject, status) => { + if (status === 'running') throw new Error('bad running listener') + }) + + send(agent, 'go') + await agent.whenIdle() + + expect(adapter.requests).toHaveLength(1) + expect(agent.status).toBe('idle') + expect(warn).toHaveBeenCalledWith(expect.stringContaining('agent/status listener threw on running')) + warn.mockRestore() + }) + + it('contains a throwing agent/status listener on the idle transition', async () => { + const adapter = new MockAdapter([textResponse('ok')]) + const ctx = await harness(adapter) + const warn = vi.spyOn(ctx.logger, 'warn').mockImplementation(() => undefined) + const agent = ctx.agentLoop.create('a1', { model: 'mock' }) + ctx.on('agent/status', (_subject, status) => { + if (status === 'idle') throw new Error('bad idle listener') + }) + + send(agent, 'go') + await agent.whenIdle() + + expect(adapter.requests).toHaveLength(1) + expect(agent.status).toBe('idle') + expect(warn).toHaveBeenCalledWith(expect.stringContaining('agent/status listener threw on idle')) + warn.mockRestore() + }) + it('abort() resolves reason to "aborted" when no reason provided', async () => { const adapter = new MockAdapter(['hang']) const ctx = await harness(adapter) diff --git a/packages/agent-loop/tests/loop.spec.ts b/packages/agent-loop/tests/loop.spec.ts index 4d8cdefe6c..3f3f52c48e 100644 --- a/packages/agent-loop/tests/loop.spec.ts +++ b/packages/agent-loop/tests/loop.spec.ts @@ -1,6 +1,6 @@ import { describe, expect, it } from 'vitest' import { Context } from 'cordis' -import LlmService, { StreamChunk } from '@deepseek-ai/dsh-llm' +import LlmService, { CallId, StreamChunk } from '@deepseek-ai/dsh-llm' import SessionStore, { TurnEndReason } from '@deepseek-ai/dsh-session' import SystemPrompt from '@deepseek-ai/dsh-system-prompt' import ToolRegistry, { defineTool } from '@deepseek-ai/dsh-tools' @@ -406,6 +406,70 @@ describe('agent loop', () => { expect(reasons).toEqual([{ kind: 'max-tokens' }, { kind: 'completed' }]) }) + it('does not dispatch tool calls from a max-tokens-truncated step', async () => { + const callId = CallId('c1') + const adapter = new MockAdapter([[ + { type: 'block-start', index: 0, blockType: 'tool-call' }, + { type: 'tool-call-delta', index: 0, id: callId, name: 'echo', argumentsDelta: '{"text":"x"}' }, + { type: 'block-end', index: 0, block: { type: 'tool-call', id: callId, name: 'echo', arguments: '{"text":"x"}' } }, + { type: 'usage', usage: { inputTokens: 10, outputTokens: 5 } }, + { type: 'finish', reason: { kind: 'max-tokens' } }, + ]]) + const ctx = await harness(adapter) + let executions = 0 + ctx.tools.register(defineTool({ + name: 'echo', + description: '', + parameters: { text: { type: 'string' } }, + async execute() { + executions += 1 + return [{ type: 'text', text: 'should not run' }] + }, + })) + const agent = ctx.agentLoop.create('a1', { model: 'mock' }) + + const reasons: TurnEndReason[] = [] + ctx.on('agent/turn-end', (_agent, _turn, reason) => void reasons.push(reason)) + + send(agent, 'go') + await waitForIdle(ctx, agent) + + expect(executions).toBe(0) + expect(agent.session.events.some(e => e.type === 'tool/call')).toBe(false) + expect(agent.session.deriveMessages()).toEqual([ + { role: 'user', content: [{ type: 'text', text: 'go' }] }, + ]) + expect(reasons).toEqual([{ kind: 'max-tokens' }]) + }) + + it('stops the turn when agent/step-end listener failure has recorded an error', async () => { + const adapter = new MockAdapter([ + toolCallResponse('c1', 'echo', { text: 'x' }), + textResponse('should not run'), + ]) + const ctx = await harness(adapter) + ctx.tools.register(defineTool({ + name: 'echo', + description: '', + parameters: { text: { type: 'string' } }, + async execute(args) { + return [{ type: 'text', text: String(args.text) }] + }, + })) + const agent = ctx.agentLoop.create('a1', { model: 'mock' }) + let threw = false + ctx.on('agent/step-end', () => { + if (!threw) { threw = true; throw new Error('bad step-end listener') } + }) + + send(agent, 'go') + await waitForIdle(ctx, agent) + + expect(adapter.requests).toHaveLength(1) + const turnEnd = agent.session.events.findLast(e => e.type === 'turn/end') + expect(turnEnd?.type === 'turn/end' && turnEnd.data.reason.kind).toBe('error') + }) + it('chains queued messages into consecutive turns', async () => { const adapter = new MockAdapter([textResponse('first'), textResponse('second')]) const ctx = await harness(adapter) diff --git a/packages/agent/src/types.ts b/packages/agent/src/types.ts index 0b786df593..9040c45ac5 100644 --- a/packages/agent/src/types.ts +++ b/packages/agent/src/types.ts @@ -82,11 +82,12 @@ export interface Agent { /** * Resolve once the agent has reached quiescence after settling out of - * `running`, or immediately if it is already idle. The quiescence signal a - * teardown awaits: `agent.abort()` then `await agent.whenIdle()` guarantees - * the in-flight turn has fully stopped before the caller proceeds (a closing - * ACP connection, a disposing UI plugin), rather than returning while the - * driver is still streaming. + * `running`, or immediately if it is already idle with no queued work. The + * quiescence signal a teardown awaits: `agent.abort()` then + * `await agent.whenIdle()` guarantees queued/running work has fully stopped + * before the caller proceeds (a closing ACP connection, a disposing UI + * plugin), rather than returning while the driver is still streaming or about + * to start a queued turn. * * "Quiescence", not merely "status changed": a disposed agent emits * `agent/status('disposed')` from inside its disposer, BEFORE the driver loop