diff --git a/.changeset/interrupt-reminder.md b/.changeset/interrupt-reminder.md new file mode 100644 index 00000000000..2a28b1cf875 --- /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/.changeset/task-output-nonblocking.md b/.changeset/task-output-nonblocking.md new file mode 100644 index 00000000000..4a95c93e1c1 --- /dev/null +++ b/.changeset/task-output-nonblocking.md @@ -0,0 +1,5 @@ +--- +"@moonshot-ai/kimi-code": patch +--- + +Remove the blocking `block`/`timeout` wait from the TaskOutput tool so checking a background task can no longer stall the conversation; it now always returns an immediate snapshot, and completion still arrives via automatic notification. diff --git a/docs/en/reference/tools.md b/docs/en/reference/tools.md index 053ad956020..15f05cfa61c 100644 --- a/docs/en/reference/tools.md +++ b/docs/en/reference/tools.md @@ -111,7 +111,7 @@ Background task tools manage tasks started via `Bash`, `Agent`, or `AskUserQuest **`TaskList`** returns the list of background tasks. Optional parameters: `active_only` (defaults to true; lists only running tasks) and `limit` (defaults to 20; range 1–100). -**`TaskOutput`** returns the status and output of a task given its `task_id`. The inline preview includes at most the most recent 32 KB of content; the full log is saved to disk, and the tool also returns an `output_path` with a suggestion to use `Read` for paginated access. Optional `block` (defaults to false) and `timeout` (seconds to wait; defaults to 30; range 0–3600) parameters allow waiting for the task to complete before returning. +**`TaskOutput`** returns the status and output of a task given its `task_id`. The inline preview includes at most the most recent 32 KB of content; the full log is saved to disk, and the tool also returns an `output_path` with a suggestion to use `Read` for paginated access. The call is always non-blocking — it returns the current snapshot immediately, and task completion is delivered via automatic notification. **`TaskStop`** accepts a `task_id` and optional `reason` (defaults to `Stopped by TaskStop`). Safe to call on tasks that are already in a terminal state. diff --git a/docs/zh/reference/tools.md b/docs/zh/reference/tools.md index 971b826d29b..9c194002616 100644 --- a/docs/zh/reference/tools.md +++ b/docs/zh/reference/tools.md @@ -109,7 +109,7 @@ Plan 模式是一种受约束的工作状态:进入后 `Write` 与 `Edit` 只 **`TaskList`** 返回后台任务列表。可选参数 `active_only`(默认 true,仅列出运行中的任务)和 `limit`(默认 20,取值范围 1–100)。 -**`TaskOutput`** 根据 `task_id` 返回任务状态与输出。内联预览最多包含最近 32 KB 的内容;完整日志保存在磁盘上,工具会一并返回 `output_path` 并提示通过 `Read` 分页读取。可选 `block`(默认 false)和 `timeout`(等待秒数,默认 30,取值范围 0–3600)参数可用于等待任务完成后再返回。 +**`TaskOutput`** 根据 `task_id` 返回任务状态与输出。内联预览最多包含最近 32 KB 的内容;完整日志保存在磁盘上,工具会一并返回 `output_path` 并提示通过 `Read` 分页读取。该调用始终是非阻塞的——立即返回当前快照,任务完成会通过自动通知送达。 **`TaskStop`** 接受 `task_id` 和可选的 `reason`(默认 `Stopped by TaskStop`)。对已处于终止状态的任务也能安全调用。 diff --git a/packages/agent-core-v2/docs/state-manifest.d.ts b/packages/agent-core-v2/docs/state-manifest.d.ts index 88bce478663..98a7cc7407b 100644 --- a/packages/agent-core-v2/docs/state-manifest.d.ts +++ b/packages/agent-core-v2/docs/state-manifest.d.ts @@ -1022,7 +1022,7 @@ export interface AgentStateSnapshot { 'llmRequester.lastConfigLogSignature': string | undefined; 'llmRequester.mediaDegradedTurns': Set; 'llmRequester.mediaStrippedTurns': Map; 'llmRequester.turnConfigs': Map>; + retryable: boolean; + cause?: { + code: (typeof ErrorCodes)[keyof typeof ErrorCodes]; + message: string; + name?: string; + details?: Readonly>; + retryable: boolean; + cause?: { + code: (typeof ErrorCodes)[keyof typeof ErrorCodes]; + message: string; + name?: string; + details?: Readonly>; + retryable: boolean; + cause?: { + code: (typeof ErrorCodes)[keyof typeof ErrorCodes]; + message: string; + name?: string; + details?: Readonly>; + retryable: boolean; + cause?: { + code: (typeof ErrorCodes)[keyof typeof ErrorCodes]; + message: string; + name?: string; + details?: Readonly>; + retryable: boolean; + cause?: { + code: ErrorCode; + message: string; + name?: string; + details?: Readonly>; + retryable: boolean; + cause?: ErrorPayload; + }; + }; + }; + }; + }; + }; + durationMs?: number; } /** @@ -634,6 +701,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; @@ -656,6 +724,7 @@ interface WirePayloadMap { "tools.unregister_user_tool": ToolsUnregisterUserToolPayload; "tools.update_store": ToolsUpdateStorePayload; "turn.cancel": TurnCancelPayload; + "turn.ended": TurnEndedPayload; "turn.prompt": TurnPromptPayload; "turn.steer": TurnSteerPayload; "usage.record": UsageRecordPayload; diff --git a/packages/agent-core-v2/scripts/gen-state-manifest.mts b/packages/agent-core-v2/scripts/gen-state-manifest.mts index 01d0d2b7645..171f854cb48 100644 --- a/packages/agent-core-v2/scripts/gen-state-manifest.mts +++ b/packages/agent-core-v2/scripts/gen-state-manifest.mts @@ -108,6 +108,16 @@ function tsFieldKey(key: string): string { return /^[$A-Z_a-z][$\w]*$/.test(key) ? key : JSON.stringify(key); } +/** + * The checker names a `unique symbol` key `__@@` — + * the numeric id is a compilation-global counter that shifts with unrelated + * edits, so the manifest renders the stable `__@` form instead. + */ +function stableSymbolKey(key: string): string { + const match = /^__@(.+)@\d+$/.exec(key); + return match === null ? key : `__@${match[1]}`; +} + // --------------------------------------------------------------------------- // Static pass — key constants and their register call sites // --------------------------------------------------------------------------- @@ -621,7 +631,7 @@ class TypeRenderer { const propLines = rendered.split('\n'); propLines[propLines.length - 1] += ';'; lines.push( - ` ${readonly}${tsFieldKey(prop.getName())}${optional ? '?' : ''}: ${propLines[0]}`, + ` ${readonly}${tsFieldKey(stableSymbolKey(prop.getName()))}${optional ? '?' : ''}: ${propLines[0]}`, ...propLines.slice(1).map((line) => ` ${line}`), ); } diff --git a/packages/agent-core-v2/src/agent/goal/goalService.ts b/packages/agent-core-v2/src/agent/goal/goalService.ts index 5c38e1b97f1..99428ef6004 100644 --- a/packages/agent-core-v2/src/agent/goal/goalService.ts +++ b/packages/agent-core-v2/src/agent/goal/goalService.ts @@ -625,7 +625,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') { @@ -988,18 +988,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 00000000000..619accfcfef --- /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 00000000000..0c7a39e13b6 --- /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 00000000000..1974cd741bb --- /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 20751f8ffb5..fe1680ae10d 100644 --- a/packages/agent-core-v2/src/agent/loop/loopService.ts +++ b/packages/agent-core-v2/src/agent/loop/loopService.ts @@ -53,12 +53,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 { @@ -92,8 +93,8 @@ import { type TurnSeed, } from './stepRequest'; import { StepRequestQueue, type StepRequestBatch } from './stepRequestQueue'; -import { isDisplayablePromptOrigin, turnPromptText } from './turnEvents'; -import { cancelTurn, promptTurn, TurnModel } from './turnOps'; +import { isDisplayablePromptOrigin, turnPromptText, type TurnInterruptReason } from './turnEvents'; +import { cancelTurn, endTurn, promptTurn, TurnModel } from './turnOps'; export type LoopInterruptReason = 'aborted' | 'max_steps' | 'error'; @@ -283,7 +284,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; } @@ -293,7 +297,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'; @@ -506,20 +510,25 @@ 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({ type: 'turn.ended', turnId: turn.id, reason: result.type, error, - durationMs: Date.now() - startedAt, + 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, @@ -618,6 +627,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, @@ -801,6 +811,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, @@ -808,13 +819,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( @@ -877,6 +895,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, @@ -1031,55 +1069,73 @@ 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': - case 'openai_compaction': - 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; + } + case 'openai_compaction': { + // Opaque compaction replay items are not streamed to the UI. + return; + } + default: { + const _exhaustive: never = part; + return _exhaustive; + } } - } + }, + drainInterruptedContent: () => + partialContent.splice(0).filter((part) => !isVacuousContentPart(part)), }; } } @@ -1138,9 +1194,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 0d3f60ab49f..aff0ecd7014 100644 --- a/packages/agent-core-v2/src/agent/loop/turnEvents.ts +++ b/packages/agent-core-v2/src/agent/loop/turnEvents.ts @@ -21,6 +21,14 @@ import type { TokenUsage } from '#/kosong/contract/usage'; /** Why a turn ended. `blocked` folds into `failed` at the wire edge. */ 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; @@ -50,6 +58,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 0a6714e92c0..02c19beebe4 100644 --- a/packages/agent-core-v2/src/agent/loop/turnOps.ts +++ b/packages/agent-core-v2/src/agent/loop/turnOps.ts @@ -4,11 +4,16 @@ * * Owns the next available turn id, including cancelled queued reservations and * legacy loop-event observations. Consumed by the Agent-scope `loopService`. + * 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. The + * `interruptionReminder` domain projects `turn.cancel` into its own model. */ import { z } from 'zod'; import { defineModel } from '#/wire/model'; +import type { KimiErrorPayload } from '#/_base/errors/serialize'; import type { ContentPart } from '#/kosong/contract/message'; import type { PromptOrigin } from '#/agent/contextMemory/types'; @@ -46,6 +51,7 @@ declare module '#/wire/types' { 'turn.prompt': typeof promptTurn; 'turn.steer': typeof steerTurn; 'turn.cancel': typeof cancelTurn; + 'turn.ended': typeof endTurn; } } @@ -63,13 +69,25 @@ 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]); }, }); +export const endTurn = TurnModel.defineOp('turn.ended', { + schema: z.object({ + turnId: z.number(), + reason: z.enum(['completed', 'cancelled', 'failed', 'blocked']), + error: z.custom().optional(), + durationMs: z.number().optional(), + }), + apply: (s) => s, +}); + function advanceTurnClock( state: TurnModelState, nextTurnId: number, diff --git a/packages/agent-core-v2/src/agent/task/taskService.ts b/packages/agent-core-v2/src/agent/task/taskService.ts index cc454d3dec0..b81ebb938be 100644 --- a/packages/agent-core-v2/src/agent/task/taskService.ts +++ b/packages/agent-core-v2/src/agent/task/taskService.ts @@ -189,7 +189,7 @@ const NOTIFICATION_FALLBACK_PREVIEW_BYTES = 3_000; const ACTIVE_BACKGROUND_TASK_INJECTION_VARIANT = 'background_task_status'; const ACTIVE_BACKGROUND_TASK_GUIDANCE = [ 'The conversation was compacted, so the earlier messages that started these background tasks are gone — but the tasks are still running from before.', - 'Do not start duplicates. Use TaskOutput to fetch a task’s result, TaskList to list them, and TaskStop to cancel one.', + 'Do not start duplicates. Use TaskList to list them, TaskOutput for a non-blocking status/output snapshot, and TaskStop to cancel one — completion arrives via automatic notification.', ].join(' '); export function isAgentTaskTerminal(status: AgentTaskStatus): boolean { diff --git a/packages/agent-core-v2/src/agent/tools/agent/agent-background-enabled.md b/packages/agent-core-v2/src/agent/tools/agent/agent-background-enabled.md index 67cafb5b770..7c2c65f3575 100644 --- a/packages/agent-core-v2/src/agent/tools/agent/agent-background-enabled.md +++ b/packages/agent-core-v2/src/agent/tools/agent/agent-background-enabled.md @@ -1,3 +1,3 @@ When `run_in_background=true`, the subagent runs detached from this turn. The completion arrives in a later turn as a synthetic user-role message containing its result — you do not need to poll, sleep, or check on its progress. Continue with other work or respond to the user. Never fabricate or predict what the result will say. -Default to a foreground subagent (omit `run_in_background`) when your next step needs its result — foreground hands the result straight back. Reach for `run_in_background=true` only when you have other work to do while it runs and do not need its result to proceed. Never launch in the background and then immediately wait on it (with `TaskOutput block=true`, sleeping, or otherwise): that just blocks the turn for no benefit — run it in the foreground instead. +Default to a foreground subagent (omit `run_in_background`) when your next step needs its result — foreground hands the result straight back. Reach for `run_in_background=true` only when you have other work to do while it runs and do not need its result to proceed. Never launch in the background and then immediately wait on it (by polling `TaskOutput`, sleeping, or otherwise): that just blocks the turn for no benefit — run it in the foreground instead. diff --git a/packages/agent-core-v2/src/agent/tools/os/bash/bash.md b/packages/agent-core-v2/src/agent/tools/os/bash/bash.md index 13fc46b9187..6b3ec9c4b4a 100644 --- a/packages/agent-core-v2/src/agent/tools/os/bash/bash.md +++ b/packages/agent-core-v2/src/agent/tools/os/bash/bash.md @@ -13,7 +13,7 @@ The dedicated tools render in the per-tool permission UI and keep raw stdout out **Output:** The stdout and stderr will be combined and returned as a string. The output may be truncated if it is too long. If the command exits non-zero, the output ends with a `Command failed with exit code: N` line; a command killed by its timeout or interrupted by the user ends with its own message instead. -If `run_in_background=true`, the command will be started as a background task and this tool will return a task ID instead of waiting for command completion. When doing that, you must provide a short `description`. Background commands default to a ${DEFAULT_BACKGROUND_TIMEOUT_S}s timeout and `timeout` is capped at ${MAX_BACKGROUND_TIMEOUT_S}s; set `disable_timeout=true` only when the task should run without a timeout. You will be automatically notified when the task completes. After starting one, default to returning control to the user instead of immediately waiting on it. Use `TaskOutput` only for a non-blocking status/output snapshot — do not set `block=true` to wait for a task you just launched, since its completion arrives automatically; reserve `block=true` for when the user explicitly asked you to wait. Use `TaskStop` only if the task must be cancelled. If a human user wants to inspect background tasks themselves, point them to the `/tasks` command, which opens an interactive panel; it has no subcommands. +If `run_in_background=true`, the command will be started as a background task and this tool will return a task ID instead of waiting for command completion. When doing that, you must provide a short `description`. Background commands default to a ${DEFAULT_BACKGROUND_TIMEOUT_S}s timeout and `timeout` is capped at ${MAX_BACKGROUND_TIMEOUT_S}s; set `disable_timeout=true` only when the task should run without a timeout. You will be automatically notified when the task completes. After starting one, default to returning control to the user instead of immediately waiting on it. Use `TaskOutput` only for a non-blocking status/output snapshot — do not wait on a task you just launched, since its completion arrives automatically. Use `TaskStop` only if the task must be cancelled. If a human user wants to inspect background tasks themselves, point them to the `/tasks` command, which opens an interactive panel; it has no subcommands. **Guidelines for safety and security:** - Each shell tool call will be executed in a fresh shell environment. The shell variables, current working directory changes, and the shell history is not preserved between calls. To run a command in a particular directory, pass the `cwd` argument (or use absolute paths) rather than relying on a `cd` from an earlier call. diff --git a/packages/agent-core-v2/src/agent/tools/os/bash/bashTool.ts b/packages/agent-core-v2/src/agent/tools/os/bash/bashTool.ts index 88c3bb66673..0cda986a259 100644 --- a/packages/agent-core-v2/src/agent/tools/os/bash/bashTool.ts +++ b/packages/agent-core-v2/src/agent/tools/os/bash/bashTool.ts @@ -385,7 +385,7 @@ export class BashTool implements IBashTool { if (!output.fullOutputAvailable || output.outputPath === undefined) return result; const taskOutputHint = this.allowBackground() - ? `, or TaskOutput(task_id="${taskId}", block=false)` + ? `, or TaskOutput(task_id="${taskId}")` : ''; const reference = `\n\n[Full output saved]\n` + diff --git a/packages/agent-core-v2/src/agent/tools/task/task-output/task-output.md b/packages/agent-core-v2/src/agent/tools/task/task-output/task-output.md index 7b902943d1f..7e53c035d10 100644 --- a/packages/agent-core-v2/src/agent/tools/task/task-output/task-output.md +++ b/packages/agent-core-v2/src/agent/tools/task/task-output/task-output.md @@ -1,13 +1,11 @@ Retrieve a snapshot of a running or completed background task. -Use this after `Bash(run_in_background=true)` or `Agent(run_in_background=true)` to check progress, or to read the output of a task that has already completed. +Use this after `Bash(run_in_background=true)`, `Agent(run_in_background=true)`, or `AskUserQuestion(background=true)` to check progress, or to read the output of a task that has already completed. Guidelines: - Prefer relying on automatic completion notifications. Use this tool only when you need task output before the automatic notification arrives. -- By default this tool is non-blocking and returns a current status/output snapshot — that is the normal way to use it. +- This tool is always non-blocking: it returns the current status/output snapshot immediately and never waits for the task to finish. - Do not use TaskOutput to wait for a result you need before continuing — if your next step depends on the task's result, run that task in the foreground instead. TaskOutput is for a deliberate progress check you will act on without blocking, not a way to sit and wait for a background task you just launched. -- Use block=true only when the user explicitly asked you to wait for the task. Never block on a task you launched in the current turn — if you need its result right away, it should have been a foreground call. -- If a block=true call returns `retrieval_status: timeout` (the task is still running), do not block on the same task again. Continue with other work or hand back to the user — the completion notification arrives on its own. - This tool returns structured task metadata, a fixed-size output preview, and an output_path for the full log. - For a terminal task, the metadata also explains why it ended. A shell command that runs to completion reports `status: completed` on a zero exit, or `status: failed` with its non-zero `exit_code` — judge that failure from the `exit_code`, because a plain command failure carries no `stop_reason` and no `terminal_reason`. `terminal_reason` is a categorical label emitted only when the end is not an ordinary exit: `timed_out` when the deadline aborted it, `stopped` when it was explicitly stopped, or `failed` when it errored without producing an exit code; the `stopped` and `failed` cases also carry a human-readable `stop_reason`. A task that finished on its own with a clean exit carries neither `stop_reason` nor `terminal_reason`. - The full, never-truncated log is always available at output_path; use the `Read` tool with that path to page through it, whether or not the preview was truncated. diff --git a/packages/agent-core-v2/src/agent/tools/task/task-output/task-output.ts b/packages/agent-core-v2/src/agent/tools/task/task-output/task-output.ts index aa3c00eae72..977db35429c 100644 --- a/packages/agent-core-v2/src/agent/tools/task/task-output/task-output.ts +++ b/packages/agent-core-v2/src/agent/tools/task/task-output/task-output.ts @@ -15,21 +15,6 @@ import { type AgentTool } from '#/tool/toolContract'; export const TaskOutputInputSchema = z.object({ task_id: z.string().describe('The background task ID to inspect.'), - block: z - .boolean() - .default(false) - .describe( - 'Whether to wait for the task to finish before returning. Discouraged — background tasks notify automatically on completion; use only when the user explicitly asked you to wait.', - ) - .optional(), - timeout: z - .number() - .int() - .min(0) - .max(3600) - .default(30) - .describe('Maximum number of seconds to wait when block=true.') - .optional(), }); export type TaskOutputInput = z.infer; diff --git a/packages/agent-core-v2/src/agent/tools/task/task-output/taskOutputTool.ts b/packages/agent-core-v2/src/agent/tools/task/task-output/taskOutputTool.ts index 063a45f914c..d5983a17180 100644 --- a/packages/agent-core-v2/src/agent/tools/task/task-output/taskOutputTool.ts +++ b/packages/agent-core-v2/src/agent/tools/task/task-output/taskOutputTool.ts @@ -39,12 +39,8 @@ const OUTPUT_PREVIEW_BYTES = 32 * 1024; const PAGING_HINT_LINES = 300; -function retrievalStatus( - status: AgentTaskStatus, - block: boolean | undefined, -): 'success' | 'timeout' | 'not_ready' { - if (TERMINAL_STATUSES.has(status)) return 'success'; - return block ? 'timeout' : 'not_ready'; +function retrievalStatus(status: AgentTaskStatus): 'success' | 'not_ready' { + return TERMINAL_STATUSES.has(status) ? 'success' : 'not_ready'; } function terminalReason(info: AgentTaskInfo): 'timed_out' | 'stopped' | 'failed' | undefined { @@ -85,23 +81,11 @@ export class TaskOutputTool implements ITaskOutputTool { description: `Reading output of task ${args.task_id}`, approvalRule: this.name, matchesRule: (ruleArgs) => matchesGlobRuleSubject(ruleArgs, args.task_id), - execute: ({ signal }) => this.execute(args, signal), + execute: () => this.execute(args), }; } - private async execute( - args: TaskOutputInput, - signal: AbortSignal, - ): Promise { - const info = this.tasks.getTask(args.task_id); - if (!info) { - return { isError: true, output: `Task not found: ${args.task_id}` }; - } - - if (args.block && !TERMINAL_STATUSES.has(info.status)) { - await this.tasks.wait(args.task_id, (args.timeout ?? 30) * 1000, signal); - } - + private async execute(args: TaskOutputInput): Promise { const current = this.tasks.getTask(args.task_id); if (!current) { return { isError: true, output: `Task not found: ${args.task_id}` }; @@ -111,7 +95,7 @@ export class TaskOutputTool implements ITaskOutputTool { const lines = [ formatPlainObject({ - retrievalStatus: retrievalStatus(current.status, args.block), + retrievalStatus: retrievalStatus(current.status), ...current, outputPath: output.outputPath, terminalReason: terminalReason(current), @@ -122,10 +106,6 @@ export class TaskOutputTool implements ITaskOutputTool { fullOutputTool: output.fullOutputAvailable && output.outputPath !== undefined ? 'Read' : undefined, fullOutputHint: fullOutputHint(output), - nextStep: - args.block === true && !TERMINAL_STATUSES.has(current.status) - ? 'The task is still running after waiting. Do not block on it again — continue with other work or hand back to the user; you will be notified automatically when it completes.' - : undefined, }), '', ]; diff --git a/packages/agent-core-v2/src/index.ts b/packages/agent-core-v2/src/index.ts index 0365ccf86a3..87116b36711 100644 --- a/packages/agent-core-v2/src/index.ts +++ b/packages/agent-core-v2/src/index.ts @@ -494,6 +494,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 25d89d66803..e302fc13663 100644 --- a/packages/agent-core-v2/test/agent/fullCompaction/fullCompaction.test.ts +++ b/packages/agent-core-v2/test/agent/fullCompaction/fullCompaction.test.ts @@ -1204,6 +1204,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 8c59fd79312..1f24896d5f1 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 a09086dcd1e..90237db691b 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'; @@ -78,6 +81,7 @@ describe('Agent loop', () => { [wire] context.append_loop_event { "event": { "type": "content.part", "uuid": "", "turnId": "0", "step": 1, "stepUuid": "", "part": { "type": "think", "think": "" } }, "time": "