diff --git a/packages/agent-core-v2/src/agent/contextMemory/loopEventFold.ts b/packages/agent-core-v2/src/agent/contextMemory/loopEventFold.ts index f10abbd93..cf68b8c9f 100644 --- a/packages/agent-core-v2/src/agent/contextMemory/loopEventFold.ts +++ b/packages/agent-core-v2/src/agent/contextMemory/loopEventFold.ts @@ -152,6 +152,7 @@ function createLoopEventFoldWithState( return; } case 'step.end': { + if (event.finishReason === 'interrupted' || event.finishReason === 'error') return; settleOpen(time); flushDeferred(); return; diff --git a/packages/agent-core-v2/src/agent/loop/loopService.ts b/packages/agent-core-v2/src/agent/loop/loopService.ts index c8fe3bfe9..74cece5d5 100644 --- a/packages/agent-core-v2/src/agent/loop/loopService.ts +++ b/packages/agent-core-v2/src/agent/loop/loopService.ts @@ -805,40 +805,56 @@ export class AgentLoopService extends Disposable implements IAgentLoopService { this.activeRequestTrace = undefined; await this.hooks.onWillBeginStep.run({ turnId, step: currentStep, firstStepOfTurn, signal }); const markStepStarted = this.beginStep(turnId, signal, currentStep, stepUuid, onStarted); - const streamParts = this.createStreamPartHandler(turnId, markStepStarted); - const request = this.llmRequester.start( - { source: { type: 'turn', turnId, step: currentStep } }, - streamParts.handle, - signal, - ); - this.activeRequestTrace = request.trace; - let response: AgentLLMRequestFinish; + let stepEndAppended = false; try { - response = await request.result; + const streamParts = this.createStreamPartHandler(turnId, markStepStarted); + const request = this.llmRequester.start( + { source: { type: 'turn', turnId, step: currentStep } }, + streamParts.handle, + signal, + ); + this.activeRequestTrace = request.trace; + let response: AgentLLMRequestFinish; + try { + response = await request.result; + } catch (error) { + this.appendInterruptedStreamContent(turnId, currentStep, stepUuid, streamParts, turnSignal); + throw error; + } + this.lastRequestTraceId = request.trace.traceId; + this.appendResponseContent(turnId, currentStep, stepUuid, response); + const finishReason = await this.executeStepTools( + turnId, + signal, + currentStep, + stepUuid, + response, + request.trace, + ); + this.finishStep(turnId, signal, currentStep, stepUuid, response, finishReason, markStepStarted); + stepEndAppended = true; + const hookStopTurn = await this.runAfterStep( + turnId, + signal, + currentStep, + firstStepOfTurn, + response.usage, + finishReason, + ); + return { stopReason: finishReason, hookStopTurn }; } catch (error) { - this.appendInterruptedStreamContent(turnId, currentStep, stepUuid, streamParts, turnSignal); + if (!stepEndAppended) { + this.context.appendLoopEvent({ + type: 'step.end', + uuid: stepUuid, + turnId: String(turnId), + step: currentStep, + finishReason: + isAbortError(error) || signal.aborted || turnSignal.aborted ? 'interrupted' : 'error', + }); + } throw error; } - this.lastRequestTraceId = request.trace.traceId; - this.appendResponseContent(turnId, currentStep, stepUuid, response); - const finishReason = await this.executeStepTools( - turnId, - signal, - currentStep, - stepUuid, - response, - request.trace, - ); - this.finishStep(turnId, signal, currentStep, stepUuid, response, finishReason, markStepStarted); - const hookStopTurn = await this.runAfterStep( - turnId, - signal, - currentStep, - firstStepOfTurn, - response.usage, - finishReason, - ); - return { stopReason: finishReason, hookStopTurn }; } private beginStep( diff --git a/packages/agent-core-v2/test/agent/contextMemory/loopEventFold.test.ts b/packages/agent-core-v2/test/agent/contextMemory/loopEventFold.test.ts index 968969e51..74bca6df3 100644 --- a/packages/agent-core-v2/test/agent/contextMemory/loopEventFold.test.ts +++ b/packages/agent-core-v2/test/agent/contextMemory/loopEventFold.test.ts @@ -216,6 +216,54 @@ describe('loop-event fold parity', () => { expect(folded).toEqual([]); }); + it('keeps the open assistant untouched when step.end reports an interruption', () => { + const folded = foldAll([], [ + { type: 'step.begin', uuid: 's1' }, + { + type: 'content.part', + stepUuid: 's1', + part: { type: 'text', text: 'partial' }, + }, + { type: 'step.end', uuid: 's1', finishReason: 'interrupted' }, + ]); + + expect(shapes(folded)).toEqual([ + { + role: 'assistant', + content: [{ type: 'text', text: 'partial' }], + toolCalls: [], + toolCallId: undefined, + isError: undefined, + partial: true, + }, + ]); + }); + + it('settles a failed step at the next step.begin as before', () => { + const folded = foldAll([], [ + { type: 'step.begin', uuid: 's1' }, + { type: 'step.end', uuid: 's1', finishReason: 'error' }, + { type: 'step.begin', uuid: 's2' }, + { + type: 'content.part', + stepUuid: 's2', + part: { type: 'text', text: 'recovered' }, + }, + { type: 'step.end', uuid: 's2' }, + ]); + + expect(shapes(folded)).toEqual([ + { + role: 'assistant', + content: [{ type: 'text', text: 'recovered' }], + toolCalls: [], + toolCallId: undefined, + isError: undefined, + partial: undefined, + }, + ]); + }); + it('drops an assistant whose only recorded part is an empty thinking block at step.end', () => { const folded = foldAll([], [ { type: 'step.begin', uuid: 's1' }, diff --git a/packages/agent-core-v2/test/agent/loop/loop.test.ts b/packages/agent-core-v2/test/agent/loop/loop.test.ts index ab86bfd24..15c1ae7e1 100644 --- a/packages/agent-core-v2/test/agent/loop/loop.test.ts +++ b/packages/agent-core-v2/test/agent/loop/loop.test.ts @@ -274,6 +274,37 @@ describe('Agent loop', () => { ); }); + it('records step.end with finishReason error when a step fails', async () => { + profile.update({ activeToolNames: [] }); + + await ctx.rpc.prompt({ input: [{ type: 'text', text: 'Hello' }] }); + await ctx.untilTurnEnd(); + + const begins = wireLoopEvents(ctx, 'step.begin'); + const ends = wireLoopEvents(ctx, 'step.end'); + expect(begins).toHaveLength(1); + expect(ends).toEqual([ + expect.objectContaining({ uuid: begins[0]!['uuid'], finishReason: 'error' }), + ]); + }); + + it('records step.end with finishReason interrupted when the turn is cancelled mid-step', async () => { + ctx.mockNextResponse({ type: 'text', text: 'partial answer' }, { type: 'text', text: ' more' }); + const subscription = ctx.get(IEventBus).subscribe(AssistantDelta, () => { + loop.cancel(); + }); + const turn = (await loop.enqueue(nextTurnMessage('Hello')).assigned).turn; + await expect(turn.result).resolves.toMatchObject({ type: 'cancelled' }); + subscription.dispose(); + + const begins = wireLoopEvents(ctx, 'step.begin'); + const ends = wireLoopEvents(ctx, 'step.end'); + expect(begins).toHaveLength(1); + expect(ends).toEqual([ + expect.objectContaining({ uuid: begins[0]!['uuid'], finishReason: 'interrupted' }), + ]); + }); + it('does not run loop error handlers for aborted turns', async () => { let called = false; loop.registerLoopErrorHandler({ @@ -1536,6 +1567,20 @@ function nextTurnMessage(text: string): MessageStepRequest { ); } +function wireLoopEvents( + target: TestAgentContext, + eventType: string, +): Array> { + return target.allEvents + .filter( + (entry) => + entry.type === '[wire]' && + entry.event === 'context.append_loop_event' && + (entry.args as { event?: { type?: string } }).event?.type === eventType, + ) + .map((entry) => (entry.args as { event: Record }).event); +} + function createTimingRequester(): IAgentLLMRequesterService { const timing: ModelRequestTiming = { firstTokenLatencyMs: 100, diff --git a/packages/agent-core-v2/test/agent/stepRetry/stepRetry.test.ts b/packages/agent-core-v2/test/agent/stepRetry/stepRetry.test.ts index 7ed2b8744..d01d2c8e0 100644 --- a/packages/agent-core-v2/test/agent/stepRetry/stepRetry.test.ts +++ b/packages/agent-core-v2/test/agent/stepRetry/stepRetry.test.ts @@ -33,6 +33,17 @@ describe('stepRetry plugin', () => { return ctx.allEvents.filter((event) => event.type === '[rpc]' && event.event === name); } + function wireLoopEvents(eventType: string): Array> { + return ctx.allEvents + .filter( + (entry) => + entry.type === '[wire]' && + entry.event === 'context.append_loop_event' && + (entry.args as { event?: { type?: string } }).event?.type === eventType, + ) + .map((entry) => (entry.args as { event: Record }).event); + } + async function runTurn(turnId: number, signal?: AbortSignal) { void ctx.dispatcher.dispatch(new TurnStarted({ turnId, origin: { kind: 'user' } })); const loop = ctx.get(IAgentLoopService); @@ -108,6 +119,37 @@ describe('stepRetry plugin', () => { ]); }); + it('pairs every retried step.begin with a step.end in the wire', async () => { + vi.useFakeTimers(); + let calls = 0; + ctx = createTestAgent( + llmGenerateServices(async () => { + calls += 1; + if (calls === 1) throw new APIConnectionError('terminated'); + return { + id: 'retry-response', + message: { + role: 'assistant', + content: [{ type: 'text', text: 'recovered' }], + toolCalls: [], + }, + usage: emptyUsage(), + finishReason: 'completed', + rawFinishReason: 'stop', + }; + }), + ); + + const result = await runTurn(1); + + expect(result).toEqual({ type: 'completed', steps: 2, truncated: false }); + const begins = wireLoopEvents('step.begin'); + const ends = wireLoopEvents('step.end'); + expect(begins).toHaveLength(2); + expect(ends.map((event) => event['finishReason'])).toEqual(['error', 'end_turn']); + expect(ends.map((event) => event['uuid'])).toEqual(begins.map((event) => event['uuid'])); + }); + it('fails the turn after maxAttempts and reports the interruption only then', async () => { vi.useFakeTimers(); let calls = 0;