Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions packages/acp-adapter/src/events-map.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -78,6 +80,8 @@ export function turnEndReasonToStopReason(
return 'end_turn';
case 'blocked':
return 'refusal';
case 'interrupted':
return 'end_turn';
}
}

Expand Down
6 changes: 5 additions & 1 deletion packages/acp-server/src/events-map.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand All @@ -70,6 +72,8 @@ export function turnEndReasonToStopReason(
return 'end_turn';
case 'blocked':
return 'refusal';
case 'interrupted':
return 'end_turn';
}
}

Expand Down
6 changes: 3 additions & 3 deletions packages/agent-core-v2/docs/state-manifest.d.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
};
Expand All @@ -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;
Expand Down Expand Up @@ -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;
};
};
Expand Down
3 changes: 2 additions & 1 deletion packages/agent-core-v2/docs/wire-manifest.d.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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];
Expand Down Expand Up @@ -781,6 +781,7 @@ interface TurnEndedPayload {
};
};
durationMs?: number;
interruptReason?: 'user_cancelled' | 'aborted' | 'max_steps' | 'error' | 'filtered' | 'blocked';
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand All @@ -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,
Expand All @@ -37,6 +39,7 @@ export class AgentInterruptionReminderService
agentState.contributeState(interruptionReminderKey);
this._register(
eventBus.subscribe(TurnEnded, (event) => {
if (dispatcher.restorePhase === 'restoring') return;

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Preserve reminders for restored user cancellations

If the process crashes after persisting an active turn.cancel with reason: 'user_cancelled' but before the normal turn.ended handler appends its interruption reminder, restore() now synthesizes the cancelled ending while restorePhase is still restoring, so this guard discards the only notification that would add the reminder. The next user prompt can therefore send the model partial output from the cancelled turn without warning that it is incomplete; suppress replayed historical endings without suppressing the newly synthesized recovery ending, or append its reminder after restoration.

Useful? React with 👍 / 👎.

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;
Expand Down
2 changes: 1 addition & 1 deletion packages/agent-core-v2/src/agent/loop/turnEvents.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand Down
10 changes: 6 additions & 4 deletions packages/agent-core-v2/src/agent/loop/turnOps.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
};
}
Expand Down Expand Up @@ -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<KimiErrorPayload>().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;
Expand All @@ -103,6 +104,7 @@ export class TurnEnded extends AgentEvent2<TurnEndedPayload> {
};
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;
}
Expand Down Expand Up @@ -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,
Expand Down
31 changes: 26 additions & 5 deletions packages/agent-core-v2/src/app/event/eventBusService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,24 +15,45 @@ export class EventBusService extends Service implements ISessionEventBus {
private readonly allEmitter = this._register(new Emitter<Event2<any>>('*'));
private readonly perType = new Map<string, Emitter<Event2<any>>>();
private readonly agents = new Map<string, AgentContext>();
private readonly latestGeneration = new Map<string, number>();
private readonly retiredGeneration = new Map<string, number>();
private readonly sources = new WeakMap<Event2<any>, 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<any>, agent?: AgentContext): void {
const cls = event.constructor as Event2Class;
if (cls.agentDomain) {
if (
agent === undefined ||
this.agents.get(agent.agentId) !== agent ||
(event as Event2<any> & AgentDomainTrait).agentId !== agent.agentId
) {
const eventAgentId = (event as Event2<any> & 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`);
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}
Expand All @@ -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 });
}
}),
Expand Down
4 changes: 4 additions & 0 deletions packages/agent-core-v2/src/state/eventDispatcher.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,8 @@ export type EventDispatcherHooks = {
readonly onDidRestore: Record<string, never>;
};

export type RestorePhase = 'new' | 'restoring' | 'ready' | 'failed';

export interface ModelCheckpointDepth {
readonly id: string;
readonly depth: number;
Expand All @@ -24,6 +26,8 @@ export interface IEventDispatcher extends DurableRuntimeParticipantHost {

readonly hooks: Hooks<EventDispatcherHooks>;

readonly restorePhase: RestorePhase;

dispatch(event: Event2<any>): Promise<void>;
history<S>(key: ReplayableStateKey<S>): readonly PatchEntry[];
checkpointDepth(key: ReplayableStateKey<any>): number;
Expand Down
53 changes: 49 additions & 4 deletions packages/agent-core-v2/src/state/eventDispatcherService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand All @@ -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,
Expand Down Expand Up @@ -109,8 +111,6 @@ interface PreparedParticipant {
readonly inversePatches: PatchEntry['inversePatches'];
}

type RestorePhase = 'new' | 'restoring' | 'ready' | 'failed';

class FoldContextImpl implements FoldContext {
pendingCheckpoint = false;
pendingClear = false;
Expand Down Expand Up @@ -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[] = [];
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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) {
Expand Down
2 changes: 1 addition & 1 deletion packages/agent-core-v2/test/agent/loop/loop.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -148,7 +148,7 @@ describe('Agent loop', () => {
[emit] agent.activity.updated { "time": "<time>", "lifecycle": "ready", "turn": { "turnId": 0, "origin": { "kind": "user" }, "phase": "running", "step": 1, "ending": false, "pendingApprovals": [], "activeToolCalls": [], "since": "<time>" }, "background": [], "agentId": "main" }
[wire] context.append_loop_event { "agentId": "main", "event": { "type": "content.part", "uuid": "<uuid-2>", "turnId": "0", "step": 1, "stepUuid": "<uuid-1>", "part": { "type": "text", "text": "blocked" } }, "time": "<time>" }
[wire] context.append_loop_event { "agentId": "main", "event": { "type": "step.end", "uuid": "<uuid-1>", "turnId": "0", "step": 1, "finishReason": "filtered", "usage": { "inputOther": 3, "output": 5, "inputCacheRead": 0, "inputCacheCreation": 0 }, "messageId": "mock-1", "providerFinishReason": "filtered", "rawFinishReason": "filtered" }, "time": "<time>" }
[wire] turn.ended { "agentId": "main", "turnId": 0, "reason": "failed", "error": { "code": "provider.filtered", "message": "Provider safety policy blocked the response.", "name": "ProviderFilteredError", "details": { "finishReason": "filtered" }, "retryable": false }, "time": "<time>" }
[wire] turn.ended { "agentId": "main", "turnId": 0, "reason": "failed", "error": { "code": "provider.filtered", "message": "Provider safety policy blocked the response.", "name": "ProviderFilteredError", "details": { "finishReason": "filtered" }, "retryable": false }, "interruptReason": "filtered", "time": "<time>" }
[emit] turn.ended { "time": "<time>", "agentId": "main", "turnId": 0, "reason": "failed", "error": { "code": "provider.filtered", "message": "Provider safety policy blocked the response.", "name": "ProviderFilteredError", "details": { "finishReason": "filtered" }, "retryable": false }, "interruptReason": "filtered" }
`);

Expand Down
3 changes: 2 additions & 1 deletion packages/agent-core-v2/test/agent/loop/turnOps.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -68,7 +68,7 @@ describe('turnKey lastEnded', () => {
});

describe('TurnEnded serialization', () => {
it('emits the op record shape without the bus-only interruptReason', () => {
it('persists the interruptReason in the op record shape', () => {
const event = new TurnEnded(
{
agentId: 'main',
Expand All @@ -85,6 +85,7 @@ describe('TurnEnded serialization', () => {
turnId: 3,
reason: 'cancelled',
durationMs: 12,
interruptReason: 'user_cancelled',
time: 42,
});
});
Expand Down
Loading
Loading