diff --git a/packages/core/agent-loop/src/agent.ts b/packages/core/agent-loop/src/agent.ts index 5eb29db8f6..aab8e2610b 100644 --- a/packages/core/agent-loop/src/agent.ts +++ b/packages/core/agent-loop/src/agent.ts @@ -150,7 +150,7 @@ export class ReactLoopAgent implements Agent { try { while (await this.turn()) {} } catch (_error) { - // Admission and turn boundaries report before rethrowing; the driver only contains the rejection. + // Reported failures and cancellation are contained at the driver boundary. } finally { if (this.phase.kind === 'running') { this.setPhase({ kind: 'idle', lastTurn: this.phase.turn }) @@ -190,15 +190,14 @@ export class ReactLoopAgent implements Agent { const lastTurn = this.phase.kind === 'collecting' ? this.phase.lastTurn : this.phase.turn const phase = { kind: 'running' as const, abort, turn: lastTurn, step: 0 } this.setPhase(phase) - if (signal.aborted) return this.inbox.hasPending + signal.throwIfAborted() let admission: Admission try { admission = await this.admit(true) if (admission.kind !== 'admitted') return false signal.throwIfAborted() } catch (error: unknown) { - // oxlint-disable-next-line typescript/no-unnecessary-condition -- cancel may abort while admission awaits - if (signal.aborted) return this.inbox.hasPending + if (signal.aborted) throw error this.throwError(error) } const turn = ++phase.turn @@ -237,16 +236,15 @@ export class ReactLoopAgent implements Agent { if (admission.kind === 'empty' && turnEnds) break } } catch (error: unknown) { - // oxlint-disable-next-line typescript/no-unnecessary-condition -- cancel may abort during any awaited turn operation if (signal.aborted) { turnEnds = { kind: 'aborted', reason: signal.reason as AgentCancelCause } - } else { - turnEnds = { - kind: 'error', - error: error instanceof LlmError ? error.failure : errorChain(error), - } - this.throwError(error) + throw error } + turnEnds = { + kind: 'error', + error: error instanceof LlmError ? error.failure : errorChain(error), + } + this.throwError(error) } finally { try { // oxlint-disable-next-line typescript/no-non-null-assertion -- every exit assigns a turn ending diff --git a/packages/core/agent-loop/tests/cancel.spec.ts b/packages/core/agent-loop/tests/cancel.spec.ts index e86d1f9636..4461ec9cb7 100644 --- a/packages/core/agent-loop/tests/cancel.spec.ts +++ b/packages/core/agent-loop/tests/cancel.spec.ts @@ -73,7 +73,7 @@ describe('Agent.cancel()', () => { }) it('cancel({ keepInbox: true }) preserves queued work and emits no discard', async () => { - const adapter = new MockAdapter([textResponse('reply')]) + const adapter = new MockAdapter([textResponse('preserved reply'), textResponse('wake reply')]) const ctx = await harness(adapter) const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' }) @@ -85,11 +85,43 @@ describe('Agent.cancel()', () => { agent.cancel({ kind: 'user' }, { keepInbox: true }) expect(agent.session.events.some(event => event.type === 'agent/inbox/spliced' && event.data.outcome === 'canceled')).toBe(false) + await agent.whenIdle() + expect(agent.inbox.nextTurn).toHaveLength(1) + expect(userTexts(agent)).toEqual([]) + expect(adapter.requests).toHaveLength(0) // The preserved item still runs once a later follow-up wakes the driver. + const idle = waitForIdle(ctx, agent) send(agent, 'wake it') - await waitForIdle(ctx, agent) + await idle expect(userTexts(agent)).toEqual(['preserved', 'wake it']) + expect(adapter.requests).toHaveLength(2) + }) + + it('cancel({ keepInbox: true }) parks queued work after an active turn aborts', async () => { + const adapter = new MockAdapter([ + 'hang', + textResponse('preserved reply'), + textResponse('wake reply'), + ]) + const ctx = await harness(adapter) + const agent = ctx.agentLoop.create(SessionId('keep-after-abort'), { provider: 'mock', model: 'mock' }) + + send(agent, 'active') + await new Promise(resolve => setTimeout(resolve, 30)) + send(agent, 'preserved') + agent.cancel({ kind: 'user' }, { keepInbox: true }) + await agent.whenIdle() + + expect(userTexts(agent)).toEqual(['active']) + expect(agent.inbox.nextTurn).toHaveLength(1) + expect(adapter.requests).toHaveLength(1) + + const idle = waitForIdle(ctx, agent) + send(agent, 'wake it') + await idle + expect(userTexts(agent)).toEqual(['active', 'preserved', 'wake it']) + expect(adapter.requests).toHaveLength(3) }) it('pre-step cancel drops the about-to-start turn (no turn is opened)', async () => { @@ -195,8 +227,12 @@ describe('Agent.cancel()', () => { expect(userTexts(agent)).toEqual(['first', 'later']) }) - it('replacement work queued after idle-listener cancellation still runs', async () => { - const adapter = new MockAdapter([textResponse('first reply'), textResponse('replacement reply')]) + it('replacement work queued after idle-listener cancellation waits for another wakeup', async () => { + const adapter = new MockAdapter([ + textResponse('first reply'), + textResponse('replacement reply'), + textResponse('wake reply'), + ]) const ctx = await harness(adapter) const agent = ctx.agentLoop.create(SessionId('idle-listener-post-cancel-send'), { provider: 'mock', model: 'mock' }) @@ -216,8 +252,15 @@ describe('Agent.cancel()', () => { if (replacementIdle === undefined) throw new Error('idle listener did not register replacement work') await replacementIdle - expect(adapter.requests).toHaveLength(2) - expect(userTexts(agent)).toEqual(['first', 'surviving replacement']) + expect(adapter.requests).toHaveLength(1) + expect(userTexts(agent)).toEqual(['first']) + expect(agent.inbox.nextTurn).toHaveLength(1) + + const idle = waitForIdle(ctx, agent) + send(agent, 'wake it') + await idle + expect(adapter.requests).toHaveLength(3) + expect(userTexts(agent)).toEqual(['first', 'surviving replacement', 'wake it']) }) it('cancel() mid-step aborts the active turn and drops every queued tail item', async () => { @@ -462,8 +505,7 @@ describe('Agent.cancel()', () => { expect(agent.session.events.some(e => e.type === 'turn/start')).toBe(false) }) - it('window 2: whenIdle() does NOT resolve early when a running listener cancels then queues replacement work', async () => { - // Cancellation must not settle idle while replacement work remains queued. + it('a running-listener cancellation parks replacement work until another wakeup', async () => { const adapter = new MockAdapter([textResponse('A reply'), textResponse('B reply')]) const ctx = await harness(adapter) const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' }) @@ -481,16 +523,17 @@ describe('Agent.cancel()', () => { await idle dispose() - // whenIdle() resolved only AFTER B's turn ran: B's user message + a turn/end - // are in the log, and A was dropped. - expect(userTexts(agent)).toContain('B') - expect(userTexts(agent)).not.toContain('A') - expect(agent.session.events.some(e => e.type === 'turn/end')).toBe(true) + expect(userTexts(agent)).toEqual([]) + expect(agent.inbox.nextTurn).toHaveLength(1) + + const replacementIdle = waitForIdle(ctx, agent) + send(agent, 'C') + await replacementIdle + expect(userTexts(agent)).toEqual(['B', 'C']) + expect(agent.session.events.filter(event => event.type === 'turn/end')).toHaveLength(2) }) - it('whenIdle() does NOT resolve early when a new prompt is queued during a pre-step cancel', async () => { - // The subtle race: a whenIdle() waiter is registered for prompt A; cancel() clears A; - // prompt B is queued before the loop resumes from the idle wait. + it('a prompt queued during pre-step cancellation waits for another wakeup', async () => { const adapter = new MockAdapter([textResponse('A reply'), textResponse('B reply')]) const ctx = await harness(adapter) const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' }) @@ -500,13 +543,15 @@ describe('Agent.cancel()', () => { agent.cancel({ kind: 'user' }) // arms marker, clears A send(agent, 'B') // B races in before the loop resumes - // whenIdle() must resolve only after B's turn fully ran — by which point B's user message - // and a turn/end are in the log. await idle - expect(userTexts(agent)).toContain('B') - expect(agent.session.events.some(e => e.type === 'turn/end')).toBe(true) - // A was dropped (never ran); only B's turn is recorded. - expect(userTexts(agent)).not.toContain('A') + expect(userTexts(agent)).toEqual([]) + expect(agent.inbox.nextTurn).toHaveLength(1) + + const replacementIdle = waitForIdle(ctx, agent) + send(agent, 'C') + await replacementIdle + expect(userTexts(agent)).toEqual(['B', 'C']) + expect(agent.session.events.filter(event => event.type === 'turn/end')).toHaveLength(2) }) it("cancel clears the turn's steering — it is not re-enqueued as a fresh turn", async () => { @@ -537,8 +582,12 @@ describe('Agent.cancel()', () => { expect(flat).not.toContain('steer text') }) - it('keeps replacement work queued synchronously by an abort observer', async () => { - const adapter = new MockAdapter(['hang', textResponse('replacement reply')]) + it('parks replacement work queued synchronously by an abort observer', async () => { + const adapter = new MockAdapter([ + 'hang', + textResponse('replacement reply'), + textResponse('wake reply'), + ]) const ctx = await harness(adapter) const agent = ctx.agentLoop.create(SessionId('abort-observer-replacement'), { provider: 'mock', model: 'mock' }) @@ -547,7 +596,7 @@ describe('Agent.cancel()', () => { const signal = adapter.requests[0]?.signal if (signal === undefined) throw new Error('model request omitted its turn signal') signal.addEventListener('abort', () => { send(agent, 'replacement') }, { once: true }) - const idle = waitForIdle(ctx, agent) + const idle = agent.whenIdle() agent.cancel({ kind: 'user' }) await Promise.race([ idle, @@ -563,12 +612,19 @@ describe('Agent.cancel()', () => { }), ]) - expect(adapter.requests).toHaveLength(2) - expect(userTexts(agent)).toEqual(['original', 'replacement']) + expect(adapter.requests).toHaveLength(1) + expect(userTexts(agent)).toEqual(['original']) + expect(agent.inbox.nextTurn).toHaveLength(1) const reasons = agent.session.events .filter(event => event.type === 'turn/end') .map(event => event.type === 'turn/end' ? event.data.reason : undefined) - expect(reasons).toEqual([{ kind: 'aborted', reason: { kind: 'user' } }, { kind: 'completed' }]) + expect(reasons).toEqual([{ kind: 'aborted', reason: { kind: 'user' } }]) + + const replacementIdle = waitForIdle(ctx, agent) + send(agent, 'wake it') + await replacementIdle + expect(adapter.requests).toHaveLength(3) + expect(userTexts(agent)).toEqual(['original', 'replacement', 'wake it']) }) it('keeps the first typed cause for an active turn', async () => { diff --git a/packages/core/agent-loop/tests/contract-regressions.spec.ts b/packages/core/agent-loop/tests/contract-regressions.spec.ts index acfd5b7184..76bee1dd5b 100644 --- a/packages/core/agent-loop/tests/contract-regressions.spec.ts +++ b/packages/core/agent-loop/tests/contract-regressions.spec.ts @@ -128,8 +128,11 @@ describe('assistant replay provenance', () => { }) describe('abort during tool execution ends the turn', () => { - it('records context finalized after a tool-step abort in the next turn', async () => { - const adapter = new MockAdapter([toolCallResponse('c1', 'aborter', {})]) + it('parks context finalized after a tool-step abort until another wakeup', async () => { + const adapter = new MockAdapter([ + toolCallResponse('c1', 'aborter', {}), + textResponse('after wake'), + ]) const ctx = await harness(adapter) const agent = ctx.agentLoop.create(SessionId('a-abort-injection'), { provider: 'mock', model: 'mock' }) ctx.tools.register(defineContentToolFixture({ @@ -153,14 +156,20 @@ describe('abort during tool execution ends the turn', () => { send(agent, 'go') await waitForIdle(ctx, agent) - const events = [...agent.session.events] - expect(events + expect(agent.session.events .filter(event => event.type === 'tool/result' || (event.type === 'user/message' && event.data.source.kind === 'plugin') || event.type === 'step/end' || event.type === 'turn/end') .map(event => event.type)) - .toEqual(['tool/result', 'step/end', 'turn/end', 'user/message', 'step/end', 'turn/end']) - expect(events + .toEqual(['tool/result', 'step/end', 'turn/end']) + expect(agent.inbox.nextStep.map(inboxText)) + .toEqual(['accepted result context after abort']) + + const idle = waitForIdle(ctx, agent) + send(agent, 'wake') + await idle + + expect(agent.session.events .flatMap(event => event.type === 'user/message' && event.data.source.kind === 'plugin' ? [event.data.content] : [])) @@ -224,7 +233,7 @@ describe('abort during tool execution ends the turn', () => { .toBeUndefined() }) - it('records result context finalized after disposal cancellation', async () => { + it('parks result context finalized after disposal cancellation without opening another turn', async () => { const adapter = new MockAdapter([toolCallResponse('c1', 'waiter', {})]) const ctx = await harness(adapter) const started = Promise.withResolvers() @@ -264,9 +273,11 @@ describe('abort during tool execution ends the turn', () => { .flatMap(event => event.type === 'user/message' && event.data.source.kind === 'plugin' ? [event.data.content] : [])) - .toEqual([ - [{ type: 'text', text: 'accepted result context during disposal' }], - ]) + .toEqual([]) + expect(agent.inbox.nextStep.map(inboxText)) + .toEqual(['accepted result context during disposal']) + expect(agent.session.events.filter(event => event.type === 'turn/start')) + .toHaveLength(1) expect(agent.session.events.find(event => event.type === 'turn/end')?.data.reason) .toEqual({ kind: 'aborted', reason: { kind: 'disposed' } }) }) diff --git a/packages/core/agent-loop/tests/tool-calls.spec.ts b/packages/core/agent-loop/tests/tool-calls.spec.ts index 98ac46270f..7393122109 100644 --- a/packages/core/agent-loop/tests/tool-calls.spec.ts +++ b/packages/core/agent-loop/tests/tool-calls.spec.ts @@ -518,10 +518,10 @@ describe('tool-call scheduler: abort handling', () => { ]) }) - it('stops replenishing after abort, commits started results, and drains accepted additional contexts', async () => { + it('stops replenishing after abort, commits started results, and parks accepted additional contexts', async () => { const adapter = new MockAdapter([ multiCall([1, 2, 3, 4].map(n => ({ id: `c${n}`, name: 'p', args: { id: String(n) } }))), - textResponse('should never be requested'), + textResponse('after wake'), ]) const ctx = await harness(adapter, 2) const gated = gatedParallelTool('p') @@ -566,9 +566,23 @@ describe('tool-call scheduler: abort handling', () => { const settled = events(agent).filter(e => e.type === 'tool/result' || (e.type === 'user/message' && e.data.source.kind === 'plugin')) expect(settled.map(e => e.type)) - .toEqual(['tool/result', 'tool/result', 'tool/result', 'tool/result', 'user/message', 'user/message']) - expect(settled.filter(e => e.type === 'user/message') - .map(e => (e.data.content[0] as { text: string }).text)) + .toEqual(['tool/result', 'tool/result', 'tool/result', 'tool/result']) + expect(agent.inbox.nextStep.map(message => message.content[0])) + .toEqual([ + { type: 'text', text: 'ctx-c1' }, + { type: 'text', text: 'ctx-c2' }, + ]) + + const idle = waitForIdle(ctx, agent) + agent.followup(createUserMessage({ content: [{ type: 'text', text: 'wake' }], source: { kind: 'user' } })) + await idle + + expect(events(agent).flatMap(e => + e.type === 'user/message' + && e.data.source.kind === 'plugin' + && e.data.content[0]?.type === 'text' + ? [e.data.content[0].text] + : [])) .toEqual(['ctx-c1', 'ctx-c2']) })