diff --git a/.changeset/interrupt-reminder.md b/.changeset/interrupt-reminder.md new file mode 100644 index 0000000000..2a28b1cf87 --- /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 9f79e4be23..049d4bb86b 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 43cb118bef..0b3633e294 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 0000000000..619accfcfe --- /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 0000000000..0c7a39e13b --- /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 0000000000..1974cd741b --- /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 110e1c76f5..df36717a2e 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; - } - 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; + 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; + } } - } + }, + 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 06614f3f28..fed477dd99 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 4082a3421a..ad96e5b638 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 169cf99901..650e5e0beb 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 6160bf672e..d2f0490542 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 8c59fd7931..1f24896d5f 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 89bddbe559..493e579781 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": "