diff --git a/packages/agent-core-v2/src/agent/externalHooks/externalHooksService.ts b/packages/agent-core-v2/src/agent/externalHooks/externalHooksService.ts index a58c386eb..4107613af 100644 --- a/packages/agent-core-v2/src/agent/externalHooks/externalHooksService.ts +++ b/packages/agent-core-v2/src/agent/externalHooks/externalHooksService.ts @@ -398,17 +398,18 @@ export class AgentExternalHooksService extends Disposable implements IAgentExter } private async runPreCompact(ctx: FullCompactionWillCompactContext): Promise { - ctx.signal.throwIfAborted(); + const signal = ctx.abortController.signal; + signal.throwIfAborted(); const engine = await this.readyEngine(); await engine?.trigger('PreCompact', { matcherValue: ctx.trigger, - signal: ctx.signal, + signal, inputData: { trigger: ctx.trigger, tokenCount: ctx.tokenCount, }, }); - ctx.signal.throwIfAborted(); + signal.throwIfAborted(); } private notifyPostCompact(event: { trigger: CompactionSource; result: CompactionResult }): void { diff --git a/packages/agent-core-v2/src/agent/fullCompaction/fullCompaction.ts b/packages/agent-core-v2/src/agent/fullCompaction/fullCompaction.ts index 358bdfe4b..87622e731 100644 --- a/packages/agent-core-v2/src/agent/fullCompaction/fullCompaction.ts +++ b/packages/agent-core-v2/src/agent/fullCompaction/fullCompaction.ts @@ -13,9 +13,10 @@ export interface CompactInput { } export interface FullCompactionWillCompactContext { + readonly abortController: AbortController; + readonly promise: Promise; readonly trigger: CompactionSource; readonly tokenCount: number; - readonly signal: AbortSignal; } export interface FullCompactionDidCompactContext { @@ -26,9 +27,8 @@ export interface FullCompactionDidCompactContext { export interface IAgentFullCompactionService { readonly _serviceBrand: undefined; - readonly isCompacting: boolean; + readonly compacting: FullCompactionWillCompactContext | null; begin(input: CompactInput): boolean; - cancel(): void; readonly hooks: Hooks<{ onWillCompact: FullCompactionWillCompactContext; diff --git a/packages/agent-core-v2/src/agent/fullCompaction/fullCompactionService.ts b/packages/agent-core-v2/src/agent/fullCompaction/fullCompactionService.ts index 22b0838f9..87d656b92 100644 --- a/packages/agent-core-v2/src/agent/fullCompaction/fullCompactionService.ts +++ b/packages/agent-core-v2/src/agent/fullCompaction/fullCompactionService.ts @@ -34,7 +34,6 @@ import { IAgentFullCompactionService, type CompactInput, type FullCompactionCompleteData, - type FullCompactionDidCompactContext, type FullCompactionWillCompactContext, } from './fullCompaction'; import { @@ -72,9 +71,7 @@ const DEFAULT_COMPACTION_MAX_COMPLETION_TOKENS = 128 * 1024; type CompactionTelemetryProperties = Record; -interface ActiveCompaction { - readonly abortController: AbortController; - promise: Promise; +interface ActiveCompaction extends FullCompactionWillCompactContext { blockedByTurn: boolean; } @@ -98,7 +95,7 @@ export class AgentFullCompactionService extends Disposable implements IAgentFull private readonly strategy: CompactionStrategy; private compactionCountInTurn = 0; - private compacting: ActiveCompaction | null = null; + private _compacting: ActiveCompaction | null = null; // Token count right after the last successful compaction. While nothing new // has been appended, the history is already in its minimal compacted form; // re-compacting would only summarize the summary again, so @@ -147,12 +144,12 @@ export class AgentFullCompactionService extends Disposable implements IAgentFull ); } - get isCompacting(): boolean { - return this.compacting !== null; + get compacting(): FullCompactionWillCompactContext | null { + return this._compacting; } begin(input: CompactInput): boolean { - if (this.compacting) return false; + if (this._compacting) return false; const data: CompactionBeginData = { source: input.source, instruction: input.instruction }; if (data.source === 'manual') { this.compactionCountInTurn = 0; @@ -162,6 +159,7 @@ export class AgentFullCompactionService extends Disposable implements IAgentFull if (this.compactionCountInTurn > this.strategy.maxCompactionPerTurn) return false; const history = this.context.get(); + const tokenCount = estimateTokensForMessages(history); const compactedCount = this.strategy.computeCompactCount(history, data.source); if (compactedCount === 0) { throw new KimiError(ErrorCodes.COMPACTION_UNABLE, 'No prefix that can be compacted in current history.'); @@ -169,29 +167,46 @@ export class AgentFullCompactionService extends Disposable implements IAgentFull this.wire.dispatch(fullCompactionBegin(data)); + const abortController = new AbortController(); + let resolveCompaction!: (result: CompactionResult) => void; + let rejectCompaction!: (reason: unknown) => void; + const promise = new Promise((resolve, reject) => { + resolveCompaction = resolve; + rejectCompaction = reject; + }); const active: ActiveCompaction = { - abortController: new AbortController(), - promise: Promise.resolve(), + abortController, + promise, + trigger: data.source, + tokenCount, blockedByTurn: false, }; - this.compacting = active; - active.promise = this.compactionWorker(active, active.abortController.signal, data, compactedCount); + this._compacting = active; + abortController.signal.addEventListener('abort', () => { + this.cancelActive(active); + }, { once: true }); + void this.compactionWorker(active, data, compactedCount) + .then(resolveCompaction, rejectCompaction); + void active.promise.catch(() => undefined); return true; } - cancel(): void { - const active = this.compacting; - if (active === null) return; + private cancelActive(active: ActiveCompaction): boolean { + if (this._compacting !== active) return false; this.wire.dispatch(fullCompactionCancel({})); - active.abortController.abort(); - this.compacting = null; + this._compacting = null; + if (!active.abortController.signal.aborted) { + active.abortController.abort(); + } this.eventBus.publish({ type: 'compaction.cancelled' }); + return true; } - private markCompleted(result: FullCompactionCompleteData): void { - if (this.compacting === null) return; + private markCompleted(active: ActiveCompaction, result: FullCompactionCompleteData): boolean { + if (this._compacting !== active) return false; this.wire.dispatch(fullCompactionComplete(result)); - this.compacting = null; + this._compacting = null; + return true; } private normalizeAfterReplay(): void { @@ -227,7 +242,7 @@ export class AgentFullCompactionService extends Disposable implements IAgentFull ); } const didStartCompaction = this.beginAutoCompaction(); - if (!didStartCompaction && !this.compacting) { + if (!didStartCompaction && !this._compacting) { await next(); return; } @@ -253,7 +268,7 @@ export class AgentFullCompactionService extends Disposable implements IAgentFull } private checkAutoCompaction(throwOnLimit = true): boolean { - if (this.compacting) return true; + if (this._compacting) return true; if ( this.lastCompactedTokenCount !== null && this.tokenCountWithPending() <= this.lastCompactedTokenCount @@ -265,7 +280,7 @@ export class AgentFullCompactionService extends Disposable implements IAgentFull } private beginAutoCompaction(throwOnLimit = true): boolean { - if (this.compacting) return true; + if (this._compacting) return true; const maxCompactions = this.strategy.maxCompactionPerTurn; if (this.compactionCountInTurn >= maxCompactions) { if (throwOnLimit) { @@ -279,26 +294,30 @@ export class AgentFullCompactionService extends Disposable implements IAgentFull } private async block(signal?: AbortSignal, turnId?: number): Promise { - const active = this.compacting; + const active = this._compacting; if (active === null) return; active.blockedByTurn = true; if (signal !== undefined) { signal.addEventListener('abort', () => { - if (this.compacting === active) { - this.cancel(); + if (this._compacting === active) { + active.abortController.abort(); } }, { once: true }); } this.eventBus.publish({ type: 'compaction.blocked', turnId }); - await active.promise; + try { + await active.promise; + } catch (error) { + if (active.abortController.signal.aborted || isAbortError(error)) return; + throw error; + } } private async compactionWorker( active: ActiveCompaction, - signal: AbortSignal, data: Readonly, initialCompactedCount: number, - ): Promise { + ): Promise { try { const finalResult: CompactionResult = { summary: '', @@ -311,9 +330,8 @@ export class AgentFullCompactionService extends Disposable implements IAgentFull let compactedCount = initialCompactedCount; for (let round = 1; ; round++) { - const result = await this.compactionRound(round, signal, data, compactedCount); - if (result === undefined) return; - if (this.compacting !== active) return; + const result = await this.compactionRound(active, round, data, compactedCount); + if (this._compacting !== active) throw compactionCancelledReason(active); finalResult.summary = result.summary; finalResult.contextSummary = result.contextSummary; @@ -330,17 +348,23 @@ export class AgentFullCompactionService extends Disposable implements IAgentFull if (compactedCount === 0) break; } - if (this.compacting !== active) return; + if (this._compacting !== active) throw compactionCancelledReason(active); this.lastCompactedTokenCount = finalResult.tokensAfter; - this.markCompleted(completeData(finalResult)); + if (!this.markCompleted(active, completeData(finalResult))) { + throw compactionCancelledReason(active); + } const { contextSummary: _contextSummary, ...eventResult } = finalResult; void _contextSummary; this.eventBus.publish({ type: 'compaction.completed', result: eventResult, trigger: data.source }); + return finalResult; } catch (error) { - if (isAbortError(error)) return; - const blockedByTurn = this.compacting === active && active.blockedByTurn; - if (this.compacting === active) { - this.cancel(); + if (active.abortController.signal.aborted || isAbortError(error)) { + this.cancelActive(active); + throw error; + } + const blockedByTurn = this._compacting === active && active.blockedByTurn; + if (this._compacting === active) { + this.cancelActive(active); } if (blockedByTurn) { throw error; @@ -349,15 +373,16 @@ export class AgentFullCompactionService extends Disposable implements IAgentFull type: 'error', ...toKimiErrorPayload(error), }); + throw error; } } private async compactionRound( + active: ActiveCompaction, round: number, - signal: AbortSignal, data: Readonly, initialCompactedCount: number, - ): Promise { + ): Promise { const startedAt = Date.now(); const originalHistory = [...this.context.get()]; const tokensBefore = estimateTokensForMessages(originalHistory); @@ -365,16 +390,13 @@ export class AgentFullCompactionService extends Disposable implements IAgentFull try { let compactedCount = initialCompactedCount; + const signal = active.abortController.signal; signal.throwIfAborted(); // One logical compaction fires the hook once, even when it takes // multiple window-sized rounds to bring the context under the ratio. if (round === 1) { - await this.hooks.onWillCompact.run({ - trigger: data.source, - tokenCount: tokensBefore, - signal, - }); + await this.hooks.onWillCompact.run(active); } const resolvedModel = this.profile.resolveModelContext(); @@ -441,8 +463,11 @@ export class AgentFullCompactionService extends Disposable implements IAgentFull } if (!historySafeToCompact(this.context.get(), originalHistory)) { - this.cancel(); - return undefined; + const active = this._compacting; + if (active !== null) { + this.cancelActive(active); + } + throw compactionCancelledReason(active); } const summary = this.postProcessSummary(attempt.summary); @@ -467,7 +492,7 @@ export class AgentFullCompactionService extends Disposable implements IAgentFull }); return result; } catch (error) { - if (isAbortError(error)) return undefined; + if (isAbortError(error)) throw error; this.telemetry.track('compaction_failed', { source: data.source, tokens_before: tokensBefore, @@ -548,6 +573,14 @@ function usageTelemetry(usage: TokenUsage | null): CompactionTelemetryProperties }; } +function compactionCancelledReason(active: ActiveCompaction | null): Error { + const reason = active?.abortController.signal.reason; + if (reason instanceof Error) return reason; + const error = new Error('Compaction cancelled.'); + error.name = 'AbortError'; + return error; +} + function isTodoItem(value: unknown): value is TodoItem { if (value === null || typeof value !== 'object') return false; const item = value as { title?: unknown; status?: unknown }; diff --git a/packages/agent-core-v2/src/agent/rpc/rpcService.ts b/packages/agent-core-v2/src/agent/rpc/rpcService.ts index 184edffd2..c8a75c59a 100644 --- a/packages/agent-core-v2/src/agent/rpc/rpcService.ts +++ b/packages/agent-core-v2/src/agent/rpc/rpcService.ts @@ -268,10 +268,11 @@ export class AgentRPCService implements IAgentRPCService { } cancelCompaction(_payload: EmptyPayload): void { - if (this.fullCompaction.isCompacting) { + const active = this.fullCompaction.compacting; + if (active !== null) { this.telemetry.track('cancel', { from: 'compacting' }); } - this.fullCompaction.cancel(); + active?.abortController.abort(); } registerTool(payload: RegisterToolPayload): void {