diff --git a/packages/acp-adapter/src/events-map.ts b/packages/acp-adapter/src/events-map.ts index 0448f2eb9cb..9ce23573ad6 100644 --- a/packages/acp-adapter/src/events-map.ts +++ b/packages/acp-adapter/src/events-map.ts @@ -63,6 +63,8 @@ export function assistantDeltaToSessionUpdate( * `blocked` → `refusal`: a prompt hook blocked the turn before the model * ran. ACP has no separate hook-blocked terminal state, so reuse the * refusal channel instead of reporting a clean `end_turn`. + * `interrupted` → `end_turn`: the turn was closed at restore after the + * process died mid-turn; ACP has no interrupted terminal state either. */ export function turnEndReasonToStopReason( reason: TurnEndReason, @@ -78,6 +80,8 @@ export function turnEndReasonToStopReason( return 'end_turn'; case 'blocked': return 'refusal'; + case 'interrupted': + return 'end_turn'; } } diff --git a/packages/acp-server/src/events-map.ts b/packages/acp-server/src/events-map.ts index cb549b73619..5194d59dd8a 100644 --- a/packages/acp-server/src/events-map.ts +++ b/packages/acp-server/src/events-map.ts @@ -55,9 +55,11 @@ export function assistantDeltaToSessionUpdate( * `blocked` → `refusal`: a prompt hook blocked the turn before the model * ran. ACP has no separate hook-blocked terminal state, so reuse the * refusal channel. + * `interrupted` → `end_turn`: the turn was closed at restore after the + * process died mid-turn; ACP has no interrupted terminal state either. */ export function turnEndReasonToStopReason( - reason: TurnEndReason, + reason: TurnEndReason | 'interrupted', error?: { readonly code: string }, ): AcpStopReason { switch (reason) { @@ -70,6 +72,8 @@ export function turnEndReasonToStopReason( return 'end_turn'; case 'blocked': return 'refusal'; + case 'interrupted': + return 'end_turn'; } } diff --git a/packages/agent-core-v2/docs/state-manifest.d.ts b/packages/agent-core-v2/docs/state-manifest.d.ts index bc51251a9a6..9dc9cf802a6 100644 --- a/packages/agent-core-v2/docs/state-manifest.d.ts +++ b/packages/agent-core-v2/docs/state-manifest.d.ts @@ -802,7 +802,7 @@ export interface AgentStateSnapshot { }; readonly lastTurn?: /* ActivityLastTurnState — packages/agent-core-v2/src/agent/activityView/activityView.ts */ { readonly turnId: number; - readonly reason: /* TurnEndReason — packages/agent-core-v2/src/agent/loop/turnEvents.ts */ 'completed' | 'cancelled' | 'failed' | 'blocked'; + readonly reason: /* TurnEndReason — packages/agent-core-v2/src/agent/loop/turnEvents.ts */ 'completed' | 'cancelled' | 'failed' | 'blocked' | 'interrupted'; readonly durationMs?: number; readonly at: number; }; @@ -814,7 +814,7 @@ export interface AgentStateSnapshot { }; 'activityView.lastTurn': /* ActivityLastTurnState — packages/agent-core-v2/src/agent/activityView/activityView.ts */ { readonly turnId: number; - readonly reason: /* TurnEndReason — packages/agent-core-v2/src/agent/loop/turnEvents.ts */ 'completed' | 'cancelled' | 'failed' | 'blocked'; + readonly reason: /* TurnEndReason — packages/agent-core-v2/src/agent/loop/turnEvents.ts */ 'completed' | 'cancelled' | 'failed' | 'blocked' | 'interrupted'; readonly durationMs?: number; readonly at: number; } | undefined; @@ -1192,7 +1192,7 @@ export interface AgentStateSnapshot { readonly cancelledTurnIds: readonly number[]; readonly lastEnded?: { readonly turnId: number; - readonly reason: 'completed' | 'cancelled' | 'failed' | 'blocked'; + readonly reason: 'completed' | 'cancelled' | 'failed' | 'blocked' | 'interrupted'; readonly durationMs?: number; }; }; diff --git a/packages/agent-core-v2/docs/wire-manifest.d.ts b/packages/agent-core-v2/docs/wire-manifest.d.ts index 0ef3e5124d2..5d6030ba707 100644 --- a/packages/agent-core-v2/docs/wire-manifest.d.ts +++ b/packages/agent-core-v2/docs/wire-manifest.d.ts @@ -735,7 +735,7 @@ interface TurnEndedPayload { _name: 'turn.ended'; agentId: string; turnId: number; - reason: 'completed' | 'cancelled' | 'failed' | 'blocked'; + reason: 'completed' | 'cancelled' | 'failed' | 'blocked' | 'interrupted'; /** KimiErrorPayload */ error?: { code: (typeof ErrorCodes)[keyof typeof ErrorCodes]; @@ -781,6 +781,7 @@ interface TurnEndedPayload { }; }; durationMs?: number; + interruptReason?: 'user_cancelled' | 'aborted' | 'max_steps' | 'error' | 'filtered' | 'blocked'; } /** diff --git a/packages/agent-core-v2/src/agent/interruptionReminder/interruptionReminderService.ts b/packages/agent-core-v2/src/agent/interruptionReminder/interruptionReminderService.ts index d5348bb6995..8882c8676ab 100644 --- a/packages/agent-core-v2/src/agent/interruptionReminder/interruptionReminderService.ts +++ b/packages/agent-core-v2/src/agent/interruptionReminder/interruptionReminderService.ts @@ -10,6 +10,7 @@ import { IAgentScopeContext } from '#/agent/scopeContext/scopeContext'; import { IAgentLifecycleService } from '#/session/agentLifecycle/agentLifecycle'; import { IAgentStateService } from '#/agent/state/agentState'; import { IEventBus } from '#/app/event/eventBus'; +import { IEventDispatcher } from '#/state/eventDispatcher'; import { IAgentInterruptionReminderService } from './interruptionReminder'; import { INTERRUPTION_REMINDER_VARIANT, interruptionReminderKey } from './interruptionReminderOps'; @@ -29,6 +30,7 @@ export class AgentInterruptionReminderService constructor( @IEventBus eventBus: IEventBus, @IAgentContextMemoryService private readonly context: IAgentContextMemoryService, + @IEventDispatcher dispatcher: IEventDispatcher, @IAgentLifecycleService agentLifecycle: IAgentLifecycleService, @IAgentScopeContext scopeContext: IAgentScopeContext, @IAgentStateService agentState: IAgentStateService, @@ -37,6 +39,7 @@ export class AgentInterruptionReminderService agentState.contributeState(interruptionReminderKey); this._register( eventBus.subscribe(TurnEnded, (event) => { + if (dispatcher.restorePhase === 'restoring') return; if (event.reason !== 'cancelled' || event.interruptReason !== 'user_cancelled') return; const origin = lastComparableMessage(this.context.get())?.origin; if (origin?.kind === 'injection' && origin.variant === INTERRUPTION_REMINDER_VARIANT) return; diff --git a/packages/agent-core-v2/src/agent/loop/turnEvents.ts b/packages/agent-core-v2/src/agent/loop/turnEvents.ts index 4994ea4a095..0ad81f1bc67 100644 --- a/packages/agent-core-v2/src/agent/loop/turnEvents.ts +++ b/packages/agent-core-v2/src/agent/loop/turnEvents.ts @@ -6,7 +6,7 @@ import type { FinishReason } from '#/kosong/contract/provider'; import type { ContentPart, TextPart } from '#/kosong/contract/message'; import type { TokenUsage } from '#/kosong/contract/usage'; -export type TurnEndReason = 'completed' | 'cancelled' | 'failed' | 'blocked'; +export type TurnEndReason = 'completed' | 'cancelled' | 'failed' | 'blocked' | 'interrupted'; export type TurnInterruptReason = | 'user_cancelled' diff --git a/packages/agent-core-v2/src/agent/loop/turnOps.ts b/packages/agent-core-v2/src/agent/loop/turnOps.ts index c209e9356bd..b5bb9317bdb 100644 --- a/packages/agent-core-v2/src/agent/loop/turnOps.ts +++ b/packages/agent-core-v2/src/agent/loop/turnOps.ts @@ -15,7 +15,7 @@ export interface TurnModelState { readonly cancelledTurnIds: readonly number[]; readonly lastEnded?: { readonly turnId: number; - readonly reason: 'completed' | 'cancelled' | 'failed' | 'blocked'; + readonly reason: 'completed' | 'cancelled' | 'failed' | 'blocked' | 'interrupted'; readonly durationMs?: number; }; } @@ -74,15 +74,16 @@ export interface TurnCancel { const turnEndedSchema = z.object({ agentId: z.string(), turnId: z.number(), - reason: z.enum(['completed', 'cancelled', 'failed', 'blocked']), + reason: z.enum(['completed', 'cancelled', 'failed', 'blocked', 'interrupted']), error: z.custom().optional(), durationMs: z.number().optional(), + interruptReason: z.enum(['user_cancelled', 'aborted', 'max_steps', 'error', 'filtered', 'blocked']).optional(), }); export interface TurnEndedPayload { readonly agentId: string; readonly turnId: number; - readonly reason: 'completed' | 'cancelled' | 'failed' | 'blocked'; + readonly reason: 'completed' | 'cancelled' | 'failed' | 'blocked' | 'interrupted'; readonly error?: KimiErrorPayload; readonly durationMs?: number; readonly interruptReason?: TurnInterruptReason; @@ -103,6 +104,7 @@ export class TurnEnded extends AgentEvent2 { }; if (this.error !== undefined) record['error'] = this.error; if (this.durationMs !== undefined) record['durationMs'] = this.durationMs; + if (this.interruptReason !== undefined) record['interruptReason'] = this.interruptReason; record['time'] = this.time; return record as SerializedEvent2; } @@ -137,7 +139,7 @@ export const turnKey = defineState( lastEnded: { turnId: e.turnId, reason: e.reason, durationMs: e.durationMs }, })); -function advanceTurnClock( +export function advanceTurnClock( state: TurnModelState, nextTurnId: number, cancelledTurnIds: readonly number[] = state.cancelledTurnIds, diff --git a/packages/agent-core-v2/src/app/event/eventBusService.ts b/packages/agent-core-v2/src/app/event/eventBusService.ts index 340bdb16f30..0f1a3adfc05 100644 --- a/packages/agent-core-v2/src/app/event/eventBusService.ts +++ b/packages/agent-core-v2/src/app/event/eventBusService.ts @@ -15,24 +15,45 @@ export class EventBusService extends Service implements ISessionEventBus { private readonly allEmitter = this._register(new Emitter>('*')); private readonly perType = new Map>>(); private readonly agents = new Map(); + private readonly latestGeneration = new Map(); + private readonly retiredGeneration = new Map(); private readonly sources = new WeakMap, AgentContext>(); activateAgent(agent: AgentContext): void { + this.latestGeneration.set( + agent.agentId, + Math.max(this.latestGeneration.get(agent.agentId) ?? 0, agent.generation), + ); this.agents.set(agent.agentId, agent); } deactivateAgent(agent: AgentContext): void { + this.retiredGeneration.set( + agent.agentId, + Math.max(this.retiredGeneration.get(agent.agentId) ?? 0, agent.generation), + ); + this.latestGeneration.set( + agent.agentId, + Math.max(this.latestGeneration.get(agent.agentId) ?? 0, agent.generation), + ); if (this.agents.get(agent.agentId) === agent) this.agents.delete(agent.agentId); } publish(event: Event2, agent?: AgentContext): void { const cls = event.constructor as Event2Class; if (cls.agentDomain) { - if ( - agent === undefined || - this.agents.get(agent.agentId) !== agent || - (event as Event2 & AgentDomainTrait).agentId !== agent.agentId - ) { + const eventAgentId = (event as Event2 & AgentDomainTrait).agentId; + if (agent === undefined || eventAgentId !== agent.agentId) { + throw new Error(`Agent event '${event.type}' has no active lifecycle context`); + } + const active = this.agents.get(agent.agentId); + if (active !== agent) { + const latest = this.latestGeneration.get(agent.agentId); + const retired = this.retiredGeneration.get(agent.agentId); + const known = latest !== undefined || retired !== undefined; + if (known && (agent.generation < (latest ?? 0) || agent.generation <= (retired ?? 0))) { + return; + } throw new Error(`Agent event '${event.type}' has no active lifecycle context`); } } diff --git a/packages/agent-core-v2/src/session/sessionActivity/sessionOutcomeMirrorService.ts b/packages/agent-core-v2/src/session/sessionActivity/sessionOutcomeMirrorService.ts index 1ebd6a49c38..ff5e7e94609 100644 --- a/packages/agent-core-v2/src/session/sessionActivity/sessionOutcomeMirrorService.ts +++ b/packages/agent-core-v2/src/session/sessionActivity/sessionOutcomeMirrorService.ts @@ -64,7 +64,11 @@ export class SessionOutcomeMirror extends Disposable implements ISessionOutcomeM this.write('completed'); return; } - if (event.reason === 'failed' || event.reason === 'blocked') { + if ( + event.reason === 'failed' || + event.reason === 'blocked' || + event.reason === 'interrupted' + ) { this.write('failed'); return; } @@ -86,7 +90,7 @@ export class SessionOutcomeMirror extends Disposable implements ISessionOutcomeM const reason = event.lastTurn?.reason; if (reason === 'completed' || reason === 'cancelled') { this.write(reason, { touchUpdatedAt: false }); - } else if (reason === 'failed' || reason === 'blocked') { + } else if (reason === 'failed' || reason === 'blocked' || reason === 'interrupted') { this.write('failed', { touchUpdatedAt: false }); } }), diff --git a/packages/agent-core-v2/src/state/eventDispatcher.ts b/packages/agent-core-v2/src/state/eventDispatcher.ts index caecbe00c41..2848238c359 100644 --- a/packages/agent-core-v2/src/state/eventDispatcher.ts +++ b/packages/agent-core-v2/src/state/eventDispatcher.ts @@ -10,6 +10,8 @@ export type EventDispatcherHooks = { readonly onDidRestore: Record; }; +export type RestorePhase = 'new' | 'restoring' | 'ready' | 'failed'; + export interface ModelCheckpointDepth { readonly id: string; readonly depth: number; @@ -24,6 +26,8 @@ export interface IEventDispatcher extends DurableRuntimeParticipantHost { readonly hooks: Hooks; + readonly restorePhase: RestorePhase; + dispatch(event: Event2): Promise; history(key: ReplayableStateKey): readonly PatchEntry[]; checkpointDepth(key: ReplayableStateKey): number; diff --git a/packages/agent-core-v2/src/state/eventDispatcherService.ts b/packages/agent-core-v2/src/state/eventDispatcherService.ts index 345e028372d..0f8fe9fe4d4 100644 --- a/packages/agent-core-v2/src/state/eventDispatcherService.ts +++ b/packages/agent-core-v2/src/state/eventDispatcherService.ts @@ -19,6 +19,8 @@ import { type Event2Class, } from '#/app/event/event2'; import { IEventBus } from '#/app/event/eventBus'; +import { ContextAppendLoopEvent } from '#/agent/contextMemory/contextEvents'; +import { advanceTurnClock, TurnCancel, TurnEnded, TurnPrompt, type TurnModelState } from '#/agent/loop/turnOps'; import type { ContentPart } from '#/kosong/contract/message'; import { OrderedHookSlot } from '#/hooks'; import { IWireService } from '#/wire/wire'; @@ -31,7 +33,7 @@ import { type AgentModel, type AgentModelDefinition, } from './agentModel'; -import { IEventDispatcher, type ModelCheckpointDepth } from './eventDispatcher'; +import { IEventDispatcher, type ModelCheckpointDepth, type RestorePhase } from './eventDispatcher'; import { StateError, StateErrors } from './errors'; import { expandedModelAppliers, @@ -109,8 +111,6 @@ interface PreparedParticipant { readonly inversePatches: PatchEntry['inversePatches']; } -type RestorePhase = 'new' | 'restoring' | 'ready' | 'failed'; - class FoldContextImpl implements FoldContext { pendingCheckpoint = false; pendingClear = false; @@ -178,7 +178,7 @@ export class EventDispatcherService extends Service implements IEventDispatcher readLegacyState: (key) => this.agentState.get(key), }; - private restorePhase: RestorePhase = 'new'; + restorePhase: RestorePhase = 'new'; private dispatching = false; private disposed = false; private queue: QueuedEvent[] = []; @@ -718,6 +718,9 @@ export class EventDispatcherService extends Service implements IEventDispatcher } this.restorePhase = 'restoring'; try { + let openTurn: number | undefined; + let openTurnCancel: 'user_cancelled' | 'aborted' | undefined; + let turnClock: TurnModelState = { nextTurnId: 0, cancelledTurnIds: [] }; let recordIndex = 0; for await (const record of this.wire.readJournal()) { if (record.type === 'metadata') continue; @@ -750,9 +753,51 @@ export class EventDispatcherService extends Service implements IEventDispatcher continue; } this.executeEvent(event, true); + if (event instanceof TurnPrompt) { + openTurn = turnClock.nextTurnId; + openTurnCancel = undefined; + turnClock = advanceTurnClock(turnClock, turnClock.nextTurnId + 1); + } else if (event instanceof TurnCancel) { + if ( + event.target !== undefined && + event.turnId !== undefined && + event.turnId >= turnClock.nextTurnId + ) { + turnClock = advanceTurnClock(turnClock, turnClock.nextTurnId, [ + ...turnClock.cancelledTurnIds, + event.turnId, + ]); + } + if (event.turnId !== undefined && event.turnId === openTurn) { + openTurnCancel = event.reason ?? 'aborted'; + } + } else if (event instanceof TurnEnded) { + if (event.turnId === openTurn) { + openTurn = undefined; + openTurnCancel = undefined; + } + } else if (event instanceof ContextAppendLoopEvent) { + const loopEvent = event.event; + if (loopEvent.type !== 'tool.result' && loopEvent.turnId !== undefined) { + const turnId = Number.parseInt(loopEvent.turnId, 10); + if (Number.isInteger(turnId) && turnId >= turnClock.nextTurnId) { + turnClock = advanceTurnClock(turnClock, turnId + 1); + } + } + } recordIndex++; } await this.rehydrateStates(); + if (openTurn !== undefined) { + const event = new TurnEnded({ + agentId: this.agentScope!.agentId, + turnId: openTurn, + reason: openTurnCancel === undefined ? 'interrupted' : 'cancelled', + interruptReason: openTurnCancel ?? 'aborted', + }); + this.executeEvent(event, false); + await this.wire.flush(); + } this.restorePhase = 'ready'; await this.hooks.onDidRestore.run({}); } catch (error) { 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 5fd1356b53c..f3da8d3ec97 100644 --- a/packages/agent-core-v2/test/agent/loop/loop.test.ts +++ b/packages/agent-core-v2/test/agent/loop/loop.test.ts @@ -148,7 +148,7 @@ describe('Agent loop', () => { [emit] agent.activity.updated { "time": "