From 1f3f5dadaaa4a1d705cc98aee1dbbef13680502c Mon Sep 17 00:00:00 2001 From: 7Sageer Date: Fri, 31 Jul 2026 18:17:14 +0800 Subject: [PATCH] feat(agent-core-v2): interruption reminder for user-cancelled turns (#2400) * feat(agent-core-v2): interruption reminder for user-cancelled turns When the user interrupts a turn with Esc, append a durable (origin: injection/interruption) to the agent context via a new loop aspect watching turn.ended, so the model learns the previous turn was deliberately cut off. The marker persists to the wire, replays on resume, stays hidden from transcripts, skips non-user aborts and steer, and does not stack on repeated cancels. Two supporting fixes: - An aborted LLM stream now persists its accumulated partial text/thinking as content.part loop events instead of dropping every produced token; gated on the turn signal so retried or step-cancelled attempts keep their partial output out of the record. - The turn.cancel wire op carries an optional reason ('user_cancelled' | 'aborted') so cold readers can tell deliberate interrupts from programmatic aborts. Goal-lifecycle cancels now pass an explicit programmatic reason to keep that field honest. * feat(transcript): mark user-cancelled turns with an interruption marker Project the deliberate user interrupt onto the transcript timeline: the live projector emits an 'interruption' marker when a turn ends with interruptReason 'user_cancelled', and the cold fold consumes the persisted turn.cancel reason into the same marker. Programmatic aborts keep surfacing through their own outlets (errors, goal/task state), and queued cancels that left no visible residue are skipped. * fix(agent-core-v2): make user-turn cancellation idempotent and reconcile interruption reminders on restore * fix(transcript): dedupe user-cancelled interruption markers by turn in the cold fold * chore(agent-core-v2): regenerate state manifest after merging main * refactor(agent-core-v2): split interruptionReminder out of the loop domain The loop domain owns turn execution mechanics; whether an interrupted turn should produce a model-visible reminder is a model-context policy. Move it into its own L4 domain with its own wire model that cross-reduces the loop's turn.cancel fact, and rename the op to interruptionReminder.recorded. --------- Signed-off-by: Haozhe Co-authored-by: Haozhe --- .changeset/interrupt-reminder.md | 5 + .../agent-core-v2/docs/wire-manifest.d.ts | 106 +++--- .../src/agent/goal/goalService.ts | 18 +- .../interruptionReminder.ts | 15 + .../interruptionReminderOps.ts | 43 +++ .../interruptionReminderService.ts | 108 ++++++ .../src/agent/loop/loopService.ts | 168 +++++++--- .../src/agent/loop/turnEvents.ts | 9 + .../agent-core-v2/src/agent/loop/turnOps.ts | 7 +- packages/agent-core-v2/src/index.ts | 3 + .../fullCompaction/fullCompaction.test.ts | 1 + .../test/agent/goal/goal.test.ts | 9 +- .../test/agent/loop/loop.test.ts | 317 +++++++++++++++++- packages/agent-core-v2/test/index.test.ts | 1 + .../agent-core-v2/test/wire/resume.test.ts | 63 ++++ .../kap-server/src/protocol/events-zod.ts | 3 + .../src/services/transcript/coreEventMap.ts | 10 + .../test/services/transcript.test.ts | 24 ++ packages/klient/src/contract/agent/events.ts | 4 + .../test/e2e/invalid-input-matrix.test.ts | 13 +- packages/transcript/src/history/foldFacts.ts | 39 ++- packages/transcript/test/layers.test.ts | 30 ++ 22 files changed, 859 insertions(+), 137 deletions(-) create mode 100644 .changeset/interrupt-reminder.md create mode 100644 packages/agent-core-v2/src/agent/interruptionReminder/interruptionReminder.ts create mode 100644 packages/agent-core-v2/src/agent/interruptionReminder/interruptionReminderOps.ts create mode 100644 packages/agent-core-v2/src/agent/interruptionReminder/interruptionReminderService.ts diff --git a/.changeset/interrupt-reminder.md b/.changeset/interrupt-reminder.md new file mode 100644 index 000000000..2a28b1cf8 --- /dev/null +++ b/.changeset/interrupt-reminder.md @@ -0,0 +1,5 @@ +--- +"@moonshot-ai/kimi-code": patch +--- + +Preserve the assistant's partial output when a turn is interrupted with Esc, and remind the model that the previous turn was deliberately interrupted. diff --git a/packages/agent-core-v2/docs/wire-manifest.d.ts b/packages/agent-core-v2/docs/wire-manifest.d.ts index 9f79e4be2..049d4bb86 100644 --- a/packages/agent-core-v2/docs/wire-manifest.d.ts +++ b/packages/agent-core-v2/docs/wire-manifest.d.ts @@ -21,52 +21,53 @@ // owning model offloads inline media to blob storage), cross-reducers // (foreign models that also reduce this record on dispatch and replay). -// Index (45 record types) -// config.update profile persisted src/agent/profile/profileOps.ts -// context_size.measured contextSize transient src/agent/contextSize/contextSizeOps.ts -// context.append_loop_event contextMemory persisted src/agent/contextMemory/contextOps.ts -// context.append_message contextMemory persisted src/agent/contextMemory/contextOps.ts -// context.apply_compaction contextMemory persisted src/agent/contextMemory/contextOps.ts -// context.clear contextMemory persisted src/agent/contextMemory/contextOps.ts -// context.undo contextMemory persisted src/agent/contextMemory/contextOps.ts -// cron.add cron transient src/session/cron/cronOps.ts -// cron.cursor cron transient src/session/cron/cronOps.ts -// cron.delete cron transient src/session/cron/cronOps.ts -// forked goal persisted src/agent/goal/goalOps.ts -// full_compaction.begin fullCompaction persisted src/agent/fullCompaction/compactionOps.ts -// full_compaction.cancel fullCompaction persisted src/agent/fullCompaction/compactionOps.ts -// full_compaction.complete fullCompaction persisted src/agent/fullCompaction/compactionOps.ts -// goal.clear goal persisted src/agent/goal/goalOps.ts -// goal.create goal persisted src/agent/goal/goalOps.ts -// goal.update goal persisted src/agent/goal/goalOps.ts -// interaction.request interaction persisted src/session/interaction/interactionOps.ts -// interaction.resolved interaction persisted src/session/interaction/interactionOps.ts -// llm.request llm.requestTrace persisted src/agent/llmRequester/llmRequestOps.ts -// llm.tools_snapshot llm.requestTrace persisted src/agent/llmRequester/llmRequestOps.ts -// mcp.tools_discovered mcp.discovery persisted src/agent/mcp/mcpDiscoveryOps.ts -// permission.record_approval_result permissionRules persisted src/agent/permissionRules/permissionRulesOps.ts -// permission.rules.add permissionRules transient src/agent/permissionRules/permissionRulesOps.ts -// permission.set_mode permissionMode persisted src/agent/permissionMode/permissionModeOps.ts -// plan_mode.cancel plan persisted src/agent/plan/planOps.ts -// plan_mode.enter plan persisted src/agent/plan/planOps.ts -// plan_mode.exit plan persisted src/agent/plan/planOps.ts -// plan.revision plan persisted src/agent/plan/planOps.ts -// profile.bind profile persisted src/agent/profile/profileOps.ts -// skill.activate skill transient src/agent/skill/skillOps.ts -// swarm_mode.enter swarm persisted src/agent/swarm/swarmOps.ts -// swarm_mode.exit swarm persisted src/agent/swarm/swarmOps.ts -// task.started task persisted src/agent/task/taskOps.ts -// task.terminated task persisted src/agent/task/taskOps.ts -// tools.register_user_tool userTool persisted src/agent/userTool/userToolOps.ts -// tools.reset_active_tools profile.activeTools persisted src/agent/profile/profileOps.ts -// tools.set_active_tools profile.activeTools persisted src/agent/profile/profileOps.ts -// tools.unregister_user_tool userTool persisted src/agent/userTool/userToolOps.ts -// tools.update_store todo persisted src/session/todo/todoOps.ts -// turn.cancel turn persisted src/agent/loop/turnOps.ts -// turn.ended turn persisted src/agent/loop/turnOps.ts -// turn.prompt turn persisted src/agent/loop/turnOps.ts -// turn.steer turn persisted src/agent/loop/turnOps.ts -// usage.record usage persisted src/agent/usage/usageOps.ts +// Index (46 record types) +// config.update profile persisted src/agent/profile/profileOps.ts +// context_size.measured contextSize transient src/agent/contextSize/contextSizeOps.ts +// context.append_loop_event contextMemory persisted src/agent/contextMemory/contextOps.ts +// context.append_message contextMemory persisted src/agent/contextMemory/contextOps.ts +// context.apply_compaction contextMemory persisted src/agent/contextMemory/contextOps.ts +// context.clear contextMemory persisted src/agent/contextMemory/contextOps.ts +// context.undo contextMemory persisted src/agent/contextMemory/contextOps.ts +// cron.add cron transient src/session/cron/cronOps.ts +// cron.cursor cron transient src/session/cron/cronOps.ts +// cron.delete cron transient src/session/cron/cronOps.ts +// forked goal persisted src/agent/goal/goalOps.ts +// full_compaction.begin fullCompaction persisted src/agent/fullCompaction/compactionOps.ts +// full_compaction.cancel fullCompaction persisted src/agent/fullCompaction/compactionOps.ts +// full_compaction.complete fullCompaction persisted src/agent/fullCompaction/compactionOps.ts +// goal.clear goal persisted src/agent/goal/goalOps.ts +// goal.create goal persisted src/agent/goal/goalOps.ts +// goal.update goal persisted src/agent/goal/goalOps.ts +// interaction.request interaction persisted src/session/interaction/interactionOps.ts +// interaction.resolved interaction persisted src/session/interaction/interactionOps.ts +// interruptionReminder.recorded interruptionReminder persisted src/agent/interruptionReminder/interruptionReminderOps.ts +// llm.request llm.requestTrace persisted src/agent/llmRequester/llmRequestOps.ts +// llm.tools_snapshot llm.requestTrace persisted src/agent/llmRequester/llmRequestOps.ts +// mcp.tools_discovered mcp.discovery persisted src/agent/mcp/mcpDiscoveryOps.ts +// permission.record_approval_result permissionRules persisted src/agent/permissionRules/permissionRulesOps.ts +// permission.rules.add permissionRules transient src/agent/permissionRules/permissionRulesOps.ts +// permission.set_mode permissionMode persisted src/agent/permissionMode/permissionModeOps.ts +// plan_mode.cancel plan persisted src/agent/plan/planOps.ts +// plan_mode.enter plan persisted src/agent/plan/planOps.ts +// plan_mode.exit plan persisted src/agent/plan/planOps.ts +// plan.revision plan persisted src/agent/plan/planOps.ts +// profile.bind profile persisted src/agent/profile/profileOps.ts +// skill.activate skill transient src/agent/skill/skillOps.ts +// swarm_mode.enter swarm persisted src/agent/swarm/swarmOps.ts +// swarm_mode.exit swarm persisted src/agent/swarm/swarmOps.ts +// task.started task persisted src/agent/task/taskOps.ts +// task.terminated task persisted src/agent/task/taskOps.ts +// tools.register_user_tool userTool persisted src/agent/userTool/userToolOps.ts +// tools.reset_active_tools profile.activeTools persisted src/agent/profile/profileOps.ts +// tools.set_active_tools profile.activeTools persisted src/agent/profile/profileOps.ts +// tools.unregister_user_tool userTool persisted src/agent/userTool/userToolOps.ts +// tools.update_store todo persisted src/session/todo/todoOps.ts +// turn.cancel turn persisted src/agent/loop/turnOps.ts +// turn.ended turn persisted src/agent/loop/turnOps.ts +// turn.prompt turn persisted src/agent/loop/turnOps.ts +// turn.steer turn persisted src/agent/loop/turnOps.ts +// usage.record usage persisted src/agent/usage/usageOps.ts /** * model: profile · persisted @@ -306,6 +307,15 @@ interface InteractionResolvedPayload { response: any; } +/** + * model: interruptionReminder · persisted + * owner: src/agent/interruptionReminder/interruptionReminderOps.ts + */ +interface InterruptionReminderRecordedPayload { + _name: 'interruptionReminder.recorded'; + turnId: number; +} + /** * model: llm.requestTrace · persisted * owner: src/agent/llmRequester/llmRequestOps.ts @@ -563,13 +573,14 @@ interface ToolsUpdateStorePayload { } /** - * model: turn · persisted + * model: turn · persisted · cross-reducers: interruptionReminder * owner: src/agent/loop/turnOps.ts */ interface TurnCancelPayload { _name: 'turn.cancel'; turnId?: number; target?: 'active' | 'queued'; + reason?: 'user_cancelled' | 'aborted'; } /** @@ -688,6 +699,7 @@ interface WirePayloadMap { "goal.update": GoalUpdatePayload; "interaction.request": InteractionRequestPayload; "interaction.resolved": InteractionResolvedPayload; + "interruptionReminder.recorded": InterruptionReminderRecordedPayload; "llm.request": LlmRequestPayload; "llm.tools_snapshot": LlmToolsSnapshotPayload; "mcp.tools_discovered": McpToolsDiscoveredPayload; diff --git a/packages/agent-core-v2/src/agent/goal/goalService.ts b/packages/agent-core-v2/src/agent/goal/goalService.ts index 43cb118be..0b3633e29 100644 --- a/packages/agent-core-v2/src/agent/goal/goalService.ts +++ b/packages/agent-core-v2/src/agent/goal/goalService.ts @@ -622,7 +622,7 @@ export class AgentGoalService extends Disposable implements IAgentGoalService { const state = this.requireState(); const snapshot = this.toSnapshot(state); if (state.status === 'active' && this.liveTurnId !== undefined) { - this.loopService.cancel(this.liveTurnId); + this.loopService.cancel(this.liveTurnId, abortError('Goal cancelled')); } this.clearInternal(actor); if (actor === 'user') { @@ -985,18 +985,10 @@ export class AgentGoalService extends Disposable implements IAgentGoalService { const pending = this.pendingContinuation; if (preserveLiveContinuation && pending?.turnId === this.liveTurnId) return; this.pendingContinuation = undefined; - const aborted = - reason === undefined ? pending?.receipt.abort() : pending?.receipt.abort(reason); - if ( - pending !== undefined && - !aborted && - pending.turnId !== undefined - ) { - if (reason === undefined) { - this.loopService.cancel(pending.turnId); - } else { - this.loopService.cancel(pending.turnId, reason); - } + const cancellation = reason ?? abortError('Goal continuation cancelled'); + const aborted = pending?.receipt.abort(cancellation); + if (pending !== undefined && !aborted && pending.turnId !== undefined) { + this.loopService.cancel(pending.turnId, cancellation); } } diff --git a/packages/agent-core-v2/src/agent/interruptionReminder/interruptionReminder.ts b/packages/agent-core-v2/src/agent/interruptionReminder/interruptionReminder.ts new file mode 100644 index 000000000..619accfcf --- /dev/null +++ b/packages/agent-core-v2/src/agent/interruptionReminder/interruptionReminder.ts @@ -0,0 +1,15 @@ +/** + * `interruptionReminder` domain (L4) — user-interruption reminder contract. + * + * Defines the Agent-scoped aspect that records a model-visible reminder after + * a user-cancelled turn. Bound at Agent scope. + */ + +import { createDecorator } from '#/_base/di/instantiation'; + +export interface IAgentInterruptionReminderService { + readonly _serviceBrand: undefined; +} + +export const IAgentInterruptionReminderService = + createDecorator('agentInterruptionReminderService'); diff --git a/packages/agent-core-v2/src/agent/interruptionReminder/interruptionReminderOps.ts b/packages/agent-core-v2/src/agent/interruptionReminder/interruptionReminderOps.ts new file mode 100644 index 000000000..0c7a39e13 --- /dev/null +++ b/packages/agent-core-v2/src/agent/interruptionReminder/interruptionReminderOps.ts @@ -0,0 +1,43 @@ +/** + * `interruptionReminder` domain (L4) — persists and restores pending + * user-interruption reminders. + * + * Projects the `loop` domain's `turn.cancel` fact into the set of turns whose + * interruption reminder still has to reach the conversation, and owns the op + * that records a reminder's delivery. Consumed by the Agent-scope + * `interruptionReminderService`. + */ + +import { z } from 'zod'; + +import { defineModel } from '#/wire/model'; + +export const InterruptionReminderModel = defineModel( + 'interruptionReminder', + () => [], + { + reducers: { + 'turn.cancel': (state, { turnId, target, reason }) => { + if (target !== 'active' || reason !== 'user_cancelled' || turnId === undefined) { + return state; + } + if (state.includes(turnId)) return state; + return [...state, turnId].toSorted((a, b) => a - b); + }, + }, + }, +); + +declare module '#/wire/types' { + interface PersistedOpMap { + 'interruptionReminder.recorded': typeof interruptionReminderRecorded; + } +} + +export const interruptionReminderRecorded = InterruptionReminderModel.defineOp( + 'interruptionReminder.recorded', + { + schema: z.object({ turnId: z.number().int().nonnegative() }), + apply: (state, { turnId }) => state.filter((pendingTurnId) => pendingTurnId !== turnId), + }, +); diff --git a/packages/agent-core-v2/src/agent/interruptionReminder/interruptionReminderService.ts b/packages/agent-core-v2/src/agent/interruptionReminder/interruptionReminderService.ts new file mode 100644 index 000000000..1974cd741 --- /dev/null +++ b/packages/agent-core-v2/src/agent/interruptionReminder/interruptionReminderService.ts @@ -0,0 +1,108 @@ +/** + * `interruptionReminder` domain (L4) — `IAgentInterruptionReminderService` implementation. + * + * Observes turn completion through `event`, persists reminder completion through + * its own wire model, reads conversation history through `contextMemory`, and + * appends model-visible notices through `systemReminder`. Reconciles reminders + * left pending by an interrupted restore. Bound at Agent scope. + */ + +import { Disposable } from '#/_base/di/lifecycle'; +import { LifecycleScope, ScopeActivation, registerScopedService } from '#/_base/di/scope'; +import { IAgentContextMemoryService } from '#/agent/contextMemory/contextMemory'; +import type { ContextMessage } from '#/agent/contextMemory/types'; +import { isVacuousContentPart } from '#/agent/contextMemory/vacuousContent'; +import { IAgentSystemReminderService } from '#/agent/systemReminder/systemReminder'; +import { IEventBus } from '#/app/event/eventBus'; +import { IWireService } from '#/wire/wire'; + +import { IAgentInterruptionReminderService } from './interruptionReminder'; +import { interruptionReminderRecorded, InterruptionReminderModel } from './interruptionReminderOps'; + +export const INTERRUPTION_REMINDER_VARIANT = 'interruption'; + +const INTERRUPTION_REMINDER = [ + 'The previous turn was interrupted by the user before completion;', + 'any partial output shown above is incomplete.', + "The user's next message continues the conversation.", +].join(' '); + +export class AgentInterruptionReminderService + extends Disposable + implements IAgentInterruptionReminderService +{ + declare readonly _serviceBrand: undefined; + + constructor( + @IEventBus eventBus: IEventBus, + @IAgentContextMemoryService private readonly context: IAgentContextMemoryService, + @IAgentSystemReminderService private readonly reminders: IAgentSystemReminderService, + @IWireService private readonly wire: IWireService, + ) { + super(); + this._register( + this.wire.hooks.onDidRestore.register('interruption-reminder', async (_ctx, next) => { + this.reconcilePendingReminders(); + await next(); + }), + ); + this._register( + eventBus.subscribe('turn.ended', (event) => { + if (event.reason !== 'cancelled' || event.interruptReason !== 'user_cancelled') return; + this.recordReminder(event.turnId, true); + }), + ); + } + + private reconcilePendingReminders(): void { + const pending = this.wire.getModel(InterruptionReminderModel); + for (const turnId of pending) this.recordReminder(turnId); + } + + private recordReminder(turnId: number, allowUntracked = false): void { + const pending = this.wire.getModel(InterruptionReminderModel).includes(turnId); + if (!pending && !allowUntracked) return; + if (!this.appendInterruptionReminder()) return; + if (pending) this.wire.dispatch(interruptionReminderRecorded({ turnId })); + } + + private appendInterruptionReminder(): boolean { + const before = this.context.get(); + const origin = lastDurableMessageOrigin(before); + if (origin?.kind === 'injection' && origin.variant === INTERRUPTION_REMINDER_VARIANT) return true; + this.reminders.appendSystemReminder(INTERRUPTION_REMINDER, { + kind: 'injection', + variant: INTERRUPTION_REMINDER_VARIANT, + }); + const after = this.context.get(); + if (after === before) return false; + const appended = lastDurableMessageOrigin(after); + return appended?.kind === 'injection' && appended.variant === INTERRUPTION_REMINDER_VARIANT; + } +} + +function lastDurableMessageOrigin( + messages: readonly ContextMessage[], +): ContextMessage['origin'] | undefined { + for (let i = messages.length - 1; i >= 0; i--) { + const message = messages[i]!; + if ( + message.role === 'assistant' && + message.partial === true && + message.toolCalls.length === 0 && + message.content.every(isVacuousContentPart) + ) { + continue; + } + return message.origin; + } + return undefined; +} + +registerScopedService( + LifecycleScope.Agent, + IAgentInterruptionReminderService, + AgentInterruptionReminderService, + ScopeActivation.OnScopeCreated, + 'interruptionReminder', +); diff --git a/packages/agent-core-v2/src/agent/loop/loopService.ts b/packages/agent-core-v2/src/agent/loop/loopService.ts index 110e1c76f..df36717a2 100644 --- a/packages/agent-core-v2/src/agent/loop/loopService.ts +++ b/packages/agent-core-v2/src/agent/loop/loopService.ts @@ -44,12 +44,13 @@ import { IAgentToolExecutorService } from '#/agent/toolExecutor/toolExecutor'; import { IConfigService } from '#/app/config/config'; import { IEventBus } from '#/app/event/eventBus'; import { type FinishReason } from '#/kosong/contract/provider'; -import { type StreamedMessagePart } from '#/kosong/contract/message'; +import { mergeInPlace, type ContentPart, type StreamedMessagePart } from '#/kosong/contract/message'; import { type TokenUsage } from '#/kosong/contract/usage'; import { BugIndicatingError, ErrorCodes, Error2, isError2, toKimiErrorPayload } from '#/errors'; import { OrderedHookSlot } from '#/hooks'; import { IAgentContextMemoryService } from '#/agent/contextMemory/contextMemory'; +import { isVacuousContentPart } from '#/agent/contextMemory/vacuousContent'; import { IAgentStateService } from '#/agent/state/agentState'; import { IAgentTelemetryContextService } from '#/app/telemetry/agentTelemetryContext'; import type { @@ -83,7 +84,7 @@ import { type TurnSeed, } from './stepRequest'; import { StepRequestQueue, type StepRequestBatch } from './stepRequestQueue'; -import { isDisplayablePromptOrigin, turnPromptText } from './turnEvents'; +import { isDisplayablePromptOrigin, turnPromptText, type TurnInterruptReason } from './turnEvents'; import { cancelTurn, endTurn, promptTurn, TurnModel } from './turnOps'; export type LoopInterruptReason = 'aborted' | 'max_steps' | 'error'; @@ -274,7 +275,10 @@ export class AgentLoopService extends Disposable implements IAgentLoopService { private cancelActiveTurn(turnId: number | undefined, cancellation: unknown): boolean { const job = this.activeTurnJob; if (job === undefined || (turnId !== undefined && job.turn.id !== turnId)) return false; - this.wire.dispatch(cancelTurn({ turnId: job.turn.id, target: 'active' })); + if (job.controller.signal.aborted) return true; + this.wire.dispatch( + cancelTurn({ turnId: job.turn.id, target: 'active', reason: cancelReasonFor(cancellation) }), + ); job.controller.abort(cancellation); return true; } @@ -284,7 +288,7 @@ export class AgentLoopService extends Disposable implements IAgentLoopService { if (index < 0) return false; const [job] = this.pendingTurns.splice(index, 1); if (job === undefined || job.turn.state !== 'queued') return false; - this.wire.dispatch(cancelTurn({ turnId, target: 'queued' })); + this.wire.dispatch(cancelTurn({ turnId, target: 'queued', reason: cancelReasonFor(cancellation) })); for (const step of job.steps.values()) step.cancel(cancellation); job.controller.abort(cancellation); job.turn.state = 'cancelled'; @@ -495,6 +499,8 @@ export class AgentLoopService extends Disposable implements IAgentLoopService { : this.activeRequestTrace?.traceId; if (result !== undefined) { const error = result.type === 'failed' ? toKimiErrorPayload(result.error) : undefined; + const interruptReason = + result.type === 'completed' ? undefined : interruptReasonFor(result); const durationMs = Date.now() - startedAt; this.wire.dispatch(endTurn({ turnId: turn.id, reason: result.type, error, durationMs })); this.eventBus.publish({ @@ -503,14 +509,15 @@ export class AgentLoopService extends Disposable implements IAgentLoopService { reason: result.type, error, durationMs, + interruptReason, }); if (error !== undefined) this.eventBus.publish({ type: 'error', ...error }); - if (result.type !== 'completed') { + if (interruptReason !== undefined) { const interrupted: TurnInterruptedEvent = { turn_id: turn.id, at_step: result.steps, mode, - interrupt_reason: interruptReasonFor(result), + interrupt_reason: interruptReason, provider_type, protocol, thinking_effort: thinkingEffort, @@ -609,6 +616,7 @@ export class AgentLoopService extends Disposable implements IAgentLoopService { const result = await this.executeLoopStep( runtime.turnId, begun.step.signal, + runtime.turnSignal, begun.step.number, begun.step.uuid, options.onStarted, @@ -792,6 +800,7 @@ export class AgentLoopService extends Disposable implements IAgentLoopService { private async executeLoopStep( turnId: number, signal: AbortSignal, + turnSignal: AbortSignal, currentStep: number, stepUuid: string, onStarted: ((step: number) => void) | undefined, @@ -799,13 +808,20 @@ export class AgentLoopService extends Disposable implements IAgentLoopService { this.activeRequestTrace = undefined; await this.hooks.onWillBeginStep.run({ turnId, step: currentStep, 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 } }, - this.createStreamPartHandler(turnId, markStepStarted), + streamParts.handle, signal, ); this.activeRequestTrace = request.trace; - const response = await request.result; + 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( @@ -868,6 +884,26 @@ export class AgentLoopService extends Disposable implements IAgentLoopService { } } + private appendInterruptedStreamContent( + turnId: number, + currentStep: number, + stepUuid: string, + streamParts: StreamPartCollector, + turnSignal: AbortSignal, + ): void { + if (!turnSignal.aborted) return; + for (const part of streamParts.drainInterruptedContent()) { + this.context.appendLoopEvent({ + type: 'content.part', + uuid: randomUUID(), + turnId: String(turnId), + step: currentStep, + stepUuid, + part, + }); + } + } + private async executeStepTools( turnId: number, signal: AbortSignal, @@ -1022,54 +1058,69 @@ export class AgentLoopService extends Disposable implements IAgentLoopService { private createStreamPartHandler( turnId: number, onResponseEvent: () => void, - ): (part: StreamedMessagePart) => void { + ): StreamPartCollector { const callsByIndex = new Map(); + const partialContent: ContentPart[] = []; + let forceContentPartBoundary = false; + const accumulate = (part: ContentPart): void => { + const last = partialContent.at(-1); + if (!forceContentPartBoundary && last !== undefined && mergeInPlace(last, part)) return; + forceContentPartBoundary = false; + partialContent.push({ ...part }); + }; - return (part) => { - switch (part.type) { - case 'text': - onResponseEvent(); - this.eventBus.publish({ type: 'assistant.delta', turnId, delta: part.text }); - return; - case 'think': - onResponseEvent(); - this.eventBus.publish({ type: 'thinking.delta', turnId, delta: part.think }); - return; - case 'image_url': - case 'audio_url': - case 'video_url': - return; - case 'function': { - onResponseEvent(); - callsByIndex.set(part._streamIndex, { id: part.id, name: part.name }); - this.eventBus.publish({ - type: 'tool.call.delta', - turnId, - toolCallId: part.id, - name: part.name, - argumentsPart: part.arguments ?? undefined, - }); - return; + return { + handle: (part) => { + switch (part.type) { + case 'text': + onResponseEvent(); + accumulate(part); + this.eventBus.publish({ type: 'assistant.delta', turnId, delta: part.text }); + return; + case 'think': + onResponseEvent(); + accumulate(part); + this.eventBus.publish({ type: 'thinking.delta', turnId, delta: part.think }); + return; + case 'image_url': + case 'audio_url': + case 'video_url': + return; + case 'function': { + onResponseEvent(); + forceContentPartBoundary = true; + callsByIndex.set(part._streamIndex, { id: part.id, name: part.name }); + this.eventBus.publish({ + type: 'tool.call.delta', + turnId, + toolCallId: part.id, + name: part.name, + argumentsPart: part.arguments ?? undefined, + }); + return; + } + case 'tool_call_part': { + if (part.argumentsPart === null) return; + const toolCall = callsByIndex.get(part.index); + if (toolCall === undefined) return; + onResponseEvent(); + this.eventBus.publish({ + type: 'tool.call.delta', + turnId, + toolCallId: toolCall.id, + name: toolCall.name, + argumentsPart: part.argumentsPart, + }); + return; + } + default: { + const _exhaustive: never = part; + return _exhaustive; + } } - case 'tool_call_part': { - if (part.argumentsPart === null) return; - const toolCall = callsByIndex.get(part.index); - if (toolCall === undefined) return; - onResponseEvent(); - this.eventBus.publish({ - type: 'tool.call.delta', - turnId, - toolCallId: toolCall.id, - name: toolCall.name, - argumentsPart: part.argumentsPart, - }); - return; - } - default: { - const _exhaustive: never = part; - return _exhaustive; - } - } + }, + drainInterruptedContent: () => + partialContent.splice(0).filter((part) => !isVacuousContentPart(part)), }; } } @@ -1128,9 +1179,18 @@ interface StepRuntime { type BeginStepResult = { readonly step: StepRuntime } | { readonly result: LoopRunResult }; +interface StreamPartCollector { + readonly handle: (part: StreamedMessagePart) => void; + drainInterruptedContent(): ContentPart[]; +} + +function cancelReasonFor(cancellation: unknown): 'user_cancelled' | 'aborted' { + return isUserCancellation(cancellation) ? 'user_cancelled' : 'aborted'; +} + function interruptReasonFor( result: Extract, -): TurnInterruptedEvent['interrupt_reason'] { +): TurnInterruptReason { if (result.type === 'cancelled') { return isUserCancellation(result.reason) ? 'user_cancelled' : 'aborted'; } diff --git a/packages/agent-core-v2/src/agent/loop/turnEvents.ts b/packages/agent-core-v2/src/agent/loop/turnEvents.ts index 06614f3f2..fed477dd9 100644 --- a/packages/agent-core-v2/src/agent/loop/turnEvents.ts +++ b/packages/agent-core-v2/src/agent/loop/turnEvents.ts @@ -20,6 +20,14 @@ import type { TokenUsage } from '#/kosong/contract/usage'; export type TurnEndReason = 'completed' | 'cancelled' | 'failed' | 'blocked'; +export type TurnInterruptReason = + | 'user_cancelled' + | 'aborted' + | 'max_steps' + | 'error' + | 'filtered' + | 'blocked'; + export interface TurnStartedEvent { readonly type: 'turn.started'; readonly turnId: number; @@ -49,6 +57,7 @@ export interface TurnEndedEvent { readonly reason: TurnEndReason; readonly error?: KimiErrorPayload; readonly durationMs?: number; + readonly interruptReason?: TurnInterruptReason; } export interface TurnStepStartedEvent { diff --git a/packages/agent-core-v2/src/agent/loop/turnOps.ts b/packages/agent-core-v2/src/agent/loop/turnOps.ts index 4082a3421..ad96e5b63 100644 --- a/packages/agent-core-v2/src/agent/loop/turnOps.ts +++ b/packages/agent-core-v2/src/agent/loop/turnOps.ts @@ -6,7 +6,8 @@ * legacy loop-event observations. Also persists the terminal `turn.ended` * record (reason / error / durationMs) so downstream history rebuilds can * recover how a turn ended; the record carries no engine-restorable state, so - * its `apply` is a no-op. + * its `apply` is a no-op. Consumed by the Agent-scope `loopService`; the + * `interruptionReminder` domain projects `turn.cancel` into its own model. */ import { z } from 'zod'; @@ -68,9 +69,11 @@ export const cancelTurn = TurnModel.defineOp('turn.cancel', { schema: z.object({ turnId: z.number().optional(), target: z.enum(['active', 'queued']).optional(), + reason: z.enum(['user_cancelled', 'aborted']).optional(), }), apply: (s, { turnId, target }) => { - if (target === undefined || turnId === undefined || turnId < s.nextTurnId) return s; + if (target === undefined || turnId === undefined) return s; + if (turnId < s.nextTurnId) return s; return advanceTurnClock(s, s.nextTurnId, [...s.cancelledTurnIds, turnId]); }, }); diff --git a/packages/agent-core-v2/src/index.ts b/packages/agent-core-v2/src/index.ts index 169cf9990..650e5e0be 100644 --- a/packages/agent-core-v2/src/index.ts +++ b/packages/agent-core-v2/src/index.ts @@ -520,6 +520,9 @@ export * from '#/agent/loop/loop'; export * from '#/agent/loop/loopService'; export * from '#/agent/loop/loopContinuation'; export * from '#/agent/loop/loopContinuationService'; +export * from '#/agent/interruptionReminder/interruptionReminder'; +export * from '#/agent/interruptionReminder/interruptionReminderService'; +export * from '#/agent/interruptionReminder/interruptionReminderOps'; export * from '#/agent/mcp/mcp'; export * from '#/agent/mcp/mcpService'; export * from '#/agent/mcp/mcpDiscoveryOps'; diff --git a/packages/agent-core-v2/test/agent/fullCompaction/fullCompaction.test.ts b/packages/agent-core-v2/test/agent/fullCompaction/fullCompaction.test.ts index 6160bf672..d2f049054 100644 --- a/packages/agent-core-v2/test/agent/fullCompaction/fullCompaction.test.ts +++ b/packages/agent-core-v2/test/agent/fullCompaction/fullCompaction.test.ts @@ -1122,6 +1122,7 @@ describe('FullCompaction', () => { code: 'compaction.failed', message: 'APIStatusError: Bad request', }), + interruptReason: 'error', }, }), ); diff --git a/packages/agent-core-v2/test/agent/goal/goal.test.ts b/packages/agent-core-v2/test/agent/goal/goal.test.ts index 8c59fd793..1f24896d5 100644 --- a/packages/agent-core-v2/test/agent/goal/goal.test.ts +++ b/packages/agent-core-v2/test/agent/goal/goal.test.ts @@ -6,6 +6,7 @@ */ import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; +import { isUserCancellation } from '#/_base/utils/abort'; import type { TurnEndedEvent } from '#/agent/loop/turnEvents'; import type { IDisposable } from '#/_base/di/lifecycle'; @@ -1160,7 +1161,8 @@ describe('AgentGoalService core workflow hooks', () => { await goals.cancelGoal(); expect(abort).toHaveBeenCalledOnce(); - expect(cancel).toHaveBeenCalledWith(41); + expect(cancel).toHaveBeenCalledWith(41, expect.any(Error)); + expect(isUserCancellation(cancel.mock.calls[0]?.[1])).toBe(false); }); it.each(['turn', 'token', 'wall-clock'] as const)( @@ -1928,7 +1930,7 @@ describe('AgentGoalService hard wall-clock deadline', () => { } }); - it('keeps user cancellation authoritative when it precedes the wall-clock deadline', async () => { + it('keeps the goal-cancellation abort authoritative when it precedes the wall-clock deadline', async () => { const clock = new ManualGoalDeadlineScheduler(); const llm = blockingGenerate(); const ctx = createTestAgent(appService(IGoalDeadlineScheduler, clock), { @@ -1946,8 +1948,9 @@ describe('AgentGoalService hard wall-clock deadline', () => { await ctx.rpc.cancelGoal({}); expect(llm.signal()).toMatchObject({ aborted: true, - reason: expect.objectContaining({ userCancelled: true }), + reason: expect.objectContaining({ message: 'Goal cancelled' }), }); + expect(isUserCancellation(llm.signal().reason)).toBe(false); clock.advanceBy(1_000); await ctx.untilTurnEnd(); 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 89bddbe55..493e57978 100644 --- a/packages/agent-core-v2/test/agent/loop/loop.test.ts +++ b/packages/agent-core-v2/test/agent/loop/loop.test.ts @@ -2,12 +2,15 @@ import { type ToolCall } from '#/kosong/contract/message'; import { emptyUsage } from '#/kosong/contract/usage'; import { afterEach, beforeEach, describe, expect, it } from 'vitest'; +import type { IDisposable } from '#/_base/di/lifecycle'; import { IAgentProfileService } from '#/index'; import { IAgentLLMRequesterService } from '#/agent/llmRequester/llmRequester'; import type { ModelRequestTiming } from '#/kosong/model/modelRequester'; +import type { ContextMessage } from '#/agent/contextMemory/types'; import { IAgentGoalService } from '#/agent/goal/goal'; import { IAgentLoopService, type Turn } from '#/agent/loop/loop'; import { ContinuationStepRequest, MessageStepRequest } from '#/agent/loop/stepRequest'; +import { RetryStepRequest } from '#/agent/prompt/promptStepRequests'; import type { ExecutableTool } from '#/tool/toolContract'; import { IAgentToolRegistryService } from '#/agent/toolRegistry/toolRegistry'; import { IAgentUsageService } from '#/agent/usage/usage'; @@ -130,7 +133,7 @@ describe('Agent loop', () => { [wire] context.append_loop_event { "event": { "type": "content.part", "uuid": "", "turnId": "0", "step": 1, "stepUuid": "", "part": { "type": "text", "text": "blocked" } }, "time": "