Download packages/agent-core-v2/src/agent/fullCompaction/fullCompactionService.ts from SaylorTwift/kimi-code: direct link, hf CLI and curl.
- Browser
- Download file 34.6 kB
-
https://huggingface.co/SaylorTwift/kimi-code/resolve/main/packages/agent-core-v2/src/agent/fullCompaction/fullCompactionService.ts
- Command line
-
hf download hf://SaylorTwift/kimi-code/packages/agent-core-v2/src/agent/fullCompaction/fullCompactionService.ts
-
curl -L -o fullCompactionService.ts https://huggingface.co/SaylorTwift/kimi-code/resolve/main/packages/agent-core-v2/src/agent/fullCompaction/fullCompactionService.ts
34.6 kB
| import type { IDisposable } from '#/_base/di/lifecycle'; | |
| import { Service } from "#/_base/di/service"; | |
| import { LifecycleScope } from '#/app/scopes'; | |
| import { ScopeActivation, registerScopedService } from '#/_base/di/scope'; | |
| import { defineState } from '#/state/state'; | |
| import { estimateTokensForMessage } from "#/llm-adapter/contract/tokens"; | |
| import { buildCompactionSummaryText, isRealUserInput } from '#/agent/contextMemory/compactionHandoff'; | |
| import { IAgentContextMemoryService } from '#/agent/contextMemory/contextMemory'; | |
| import type { ContextMessage } from '#/agent/contextMemory/types'; | |
| import { ISessionTokenCountingService } from '#/session/tokenCounting/sessionTokenCounting'; | |
| import { IAgentLLMRequesterService, type AgentLLMRequestFinish } from '#/agent/llmRequester/llmRequester'; | |
| import type { LLMRequestTrace } from '#/llm-adapter/contract/request-trace'; | |
| import { retryBackoffDelay, sleepForRetry } from '#/_base/utils/retry'; | |
| import { runWithCredentialRecovery } from '#/llm-adapter/model/credential-recovery'; | |
| import { IAgentLoopService, type LoopErrorContext } from '#/agent/loop/loop'; | |
| import { TurnStarted } from '#/agent/loop/turnEvents'; | |
| import { TurnEnded } from '#/agent/loop/turnOps'; | |
| import { isAbortError } from '#/_base/utils/abort'; | |
| import { IAgentProfileService, type ProfileModelContext } from '#/agent/profile/profile'; | |
| import { | |
| agentContextOfScope, | |
| IAgentScopeContext, | |
| } from '#/agent/scopeContext/scopeContext'; | |
| import { IAgentStateService } from '#/agent/state/agentState'; | |
| import { IAgentToolRegistryService } from '#/agent/toolRegistry/toolRegistry'; | |
| import { stripDynamicToolContext } from '#/agent/toolSelect/dynamicTools'; | |
| import { IAgentToolSelectService } from '#/agent/toolSelect/toolSelect'; | |
| import { IAgentTodoService } from '#/features/todo/todoService'; | |
| import { renderTodoList } from '#/features/todo/todoItem'; | |
| import { onUnexpectedError } from '#/_base/errors/unexpectedError'; | |
| import type { WireLineRange } from '#/wire/record'; | |
| import { IWireService } from '#/wire/wire'; | |
| import { | |
| APIContextOverflowError, | |
| APIEmptyResponseError, | |
| APIStatusError, | |
| isRetryableGenerateError, | |
| } from '#/llm-adapter/contract/errors'; | |
| import { createUserMessage, type Message } from '#/llm-adapter/contract/message'; | |
| import type { ToolDescription as Tool } from '#human/llm/message'; | |
| import { inputTotal, type TokenUsage } from '#human/llm/usage'; | |
| import { IEventBus } from '#/app/event/eventBus'; | |
| import type { CompactionFailedEvent, CompactionFinishedEvent } from '#/app/telemetry/events'; | |
| import { ITelemetryService } from '#/app/telemetry/telemetry'; | |
| import { ErrorCodes, Error2, isCodedError, isError2, toKimiErrorPayload, unwrapErrorCause } from "#/errors"; | |
| import { AgentErrorEvent } from '#/agent/mcp/mcpEvents'; | |
| import { IEventDispatcher } from '#/state/eventDispatcher'; | |
| import { renderCompactionInstruction } from './compactionInstruction'; | |
| import { renderContextRecoveryPointer } from './contextRecovery'; | |
| import { | |
| IAgentFullCompactionService, | |
| type FullCompactionInput, | |
| type FullCompactionTask, | |
| } from './fullCompaction'; | |
| import { | |
| RuntimeCompactionStrategy, | |
| type CompactionStrategy, | |
| } from './strategy'; | |
| import { | |
| CompactionBlocked, | |
| CompactionCancelled, | |
| CompactionCompleted, | |
| fullCompactionKey, | |
| fullCompactionWireRangesKey, | |
| FullCompactionBegin, | |
| FullCompactionCancel, | |
| FullCompactionComplete, | |
| } from './compactionOps'; | |
| import { | |
| type CompactionBeginData, | |
| type CompactionResult, | |
| } from './types'; | |
| import { Emitter, type Event } from '#/_base/event'; | |
| import { OrderedHookSlot } from '#/hooks'; | |
| export const MAX_COMPACTION_RETRY_ATTEMPTS = 5; | |
| const DEFAULT_COMPACTION_MAX_COMPLETION_TOKENS = 128 * 1024; | |
| const OVERFLOW_CONTEXT_SAFETY_RATIO = 0.85; | |
| const OVERFLOW_STATUS_RECOVERY_RATIO = 0.5; | |
| const MAX_COMPACTION_OVERFLOW_SHRINK_ATTEMPTS = 3; | |
| const COMPACTION_OVERFLOW_SHRINK_RATIOS = [0.7, 0.5, 0.35] as const; | |
| const EMPTY_TOOL_PARAMETERS: Record<string, unknown> = { | |
| type: 'object', | |
| properties: {}, | |
| }; | |
| type CompactionTelemetryProperties = Pick< | |
| CompactionFinishedEvent, | |
| 'input_tokens' | 'output_tokens' | 'input_cache_read' | 'input_cache_creation' | |
| >; | |
| interface ActiveCompaction extends FullCompactionTask { | |
| readonly originTurnId?: number; | |
| readonly quiescence?: IDisposable; | |
| trace?: LLMRequestTrace; | |
| blockedByTurn: boolean; | |
| } | |
| interface CompactionAttemptResult { | |
| readonly summary: string; | |
| readonly usage: TokenUsage | null; | |
| readonly traceId?: string; | |
| } | |
| class CompactionTruncatedError extends Error { | |
| constructor() { | |
| super('Compaction response was truncated before producing a complete summary.'); | |
| this.name = 'CompactionTruncatedError'; | |
| } | |
| } | |
| export const fullCompactionCompactionCountInTurnKey = defineState<number>( | |
| 'fullCompaction.compactionCountInTurn', | |
| () => 0, | |
| ); | |
| export const fullCompactionObservedMaxContextTokensByModelKey = defineState<Map<string, number>>( | |
| 'fullCompaction.observedMaxContextTokensByModel', | |
| () => new Map(), | |
| ); | |
| export const fullCompactionLastCompactedTokenCountKey = defineState<number | null>( | |
| 'fullCompaction.lastCompactedTokenCount', | |
| () => null, | |
| ); | |
| export const fullCompactionConsecutiveOverflowCompactionsKey = defineState<number>( | |
| 'fullCompaction.consecutiveOverflowCompactions', | |
| () => 0, | |
| ); | |
| export const fullCompactionActiveTurnIdKey = defineState<number | undefined>( | |
| 'fullCompaction.activeTurnId', | |
| () => undefined as number | undefined, | |
| ); | |
| export class AgentFullCompactionService extends Service implements IAgentFullCompactionService { | |
| declare readonly _serviceBrand: undefined; | |
| readonly hooks: IAgentFullCompactionService['hooks'] = { | |
| onWillCompact: new OrderedHookSlot<FullCompactionTask>(), | |
| }; | |
| private readonly _onDidFinishCompaction = this._register(new Emitter<FullCompactionTask>()); | |
| readonly onDidFinishCompaction: Event<FullCompactionTask> = this._onDidFinishCompaction.event; | |
| private readonly strategy: CompactionStrategy; | |
| private _compacting: ActiveCompaction | null = null; | |
| constructor( | |
| private readonly context: IAgentContextMemoryService, | |
| private readonly tokenCounting: ISessionTokenCountingService, | |
| private readonly llmRequester: IAgentLLMRequesterService, | |
| private readonly profile: IAgentProfileService, | |
| private readonly toolRegistry: IAgentToolRegistryService, | |
| private readonly toolSelect: IAgentToolSelectService, | |
| private readonly agent: IAgentScopeContext, | |
| private readonly todo: IAgentTodoService, | |
| private readonly telemetry: ITelemetryService, | |
| private readonly dispatcher: IEventDispatcher, | |
| private readonly eventBus: IEventBus, | |
| private readonly loopService: IAgentLoopService, | |
| private readonly states: IAgentStateService, | |
| private readonly wire: IWireService, | |
| ) { | |
| super(); | |
| this.states.contributeState(fullCompactionKey); | |
| this.states.contributeState(fullCompactionWireRangesKey); | |
| this.states.contributeState(fullCompactionCompactionCountInTurnKey); | |
| this.states.contributeState(fullCompactionObservedMaxContextTokensByModelKey); | |
| this.states.contributeState(fullCompactionLastCompactedTokenCountKey); | |
| this.states.contributeState(fullCompactionConsecutiveOverflowCompactionsKey); | |
| this.states.contributeState(fullCompactionActiveTurnIdKey); | |
| this.strategy = new RuntimeCompactionStrategy( | |
| () => this.resolveModelContextWithEffectiveMax(), | |
| (message) => this.tokenCounting.estimateMessage(message), | |
| ); | |
| this._register( | |
| this.dispatcher.hooks.onDidRestore.register('full-compaction', async (_ctx, next) => { | |
| this.normalizeAfterReplay(); | |
| await next(); | |
| }), | |
| ); | |
| this._register( | |
| this.eventBus.subscribe(TurnStarted, () => this.resetForTurn()), | |
| ); | |
| this._register( | |
| this.eventBus.subscribe(TurnEnded, () => { | |
| this.activeTurnId = undefined; | |
| }), | |
| ); | |
| this._register( | |
| this.loopService.hooks.onWillBeginStep.register('full-compaction', async (ctx, next) => { | |
| await this.beforeStep(ctx.signal, ctx.turnId); | |
| await next(); | |
| }), | |
| ); | |
| this._register( | |
| this.loopService.hooks.onDidFinishStep.register('full-compaction', async (_ctx, next) => { | |
| await this.afterStep(); | |
| await next(); | |
| }), | |
| ); | |
| this._register( | |
| this.loopService.registerLoopErrorHandler({ | |
| id: 'full-compaction', | |
| match: (context) => this.shouldRecoverFromContextOverflow(context.error), | |
| handle: (context) => this.recoverFromContextOverflow(context), | |
| }), | |
| ); | |
| } | |
| private get compactionCountInTurn(): number { | |
| return this.states.get(fullCompactionCompactionCountInTurnKey); | |
| } | |
| private set compactionCountInTurn(value: number) { | |
| this.states.set(fullCompactionCompactionCountInTurnKey, value); | |
| } | |
| private get observedMaxContextTokensByModel(): Map<string, number> { | |
| return this.states.get(fullCompactionObservedMaxContextTokensByModelKey); | |
| } | |
| private get lastCompactedTokenCount(): number | null { | |
| return this.states.get(fullCompactionLastCompactedTokenCountKey); | |
| } | |
| private set lastCompactedTokenCount(value: number | null) { | |
| this.states.set(fullCompactionLastCompactedTokenCountKey, value); | |
| } | |
| private get consecutiveOverflowCompactions(): number { | |
| return this.states.get(fullCompactionConsecutiveOverflowCompactionsKey); | |
| } | |
| private set consecutiveOverflowCompactions(value: number) { | |
| this.states.set(fullCompactionConsecutiveOverflowCompactionsKey, value); | |
| } | |
| private get activeTurnId(): number | undefined { | |
| return this.states.get(fullCompactionActiveTurnIdKey); | |
| } | |
| private set activeTurnId(value: number | undefined) { | |
| this.states.set(fullCompactionActiveTurnIdKey, value); | |
| } | |
| get compacting(): FullCompactionTask | null { | |
| return this._compacting; | |
| } | |
| cancel(): void { | |
| const active = this._compacting; | |
| if (active !== null) { | |
| this.telemetry.track2('cancel', { | |
| from: 'compacting', | |
| trace_id: active.traceId, | |
| }); | |
| } | |
| active?.abortController.abort(); | |
| } | |
| private getEffectiveMaxContextTokens(): number { | |
| const capability = this.profile.data().modelCapabilities; | |
| const configured = capability.max_input_tokens ?? capability.max_context_tokens; | |
| const modelAlias = this.profile.data().modelAlias; | |
| const observed = | |
| modelAlias === undefined ? undefined : this.observedMaxContextTokensByModel.get(modelAlias); | |
| if (observed === undefined) return configured; | |
| if (configured <= 0) return observed; | |
| return Math.min(configured, observed); | |
| } | |
| private resolveModelContextWithEffectiveMax(): ProfileModelContext { | |
| const resolved = this.profile.resolveModelContext(); | |
| const effectiveMax = this.getEffectiveMaxContextTokens(); | |
| return { | |
| ...resolved, | |
| modelCapabilities: { | |
| ...resolved.modelCapabilities, | |
| max_context_tokens: effectiveMax, | |
| max_input_tokens: effectiveMax, | |
| }, | |
| }; | |
| } | |
| private currentRequestTokens(): number { | |
| return this.requestTokens(this.context.get()); | |
| } | |
| private requestTokens(messages: readonly Message[]): number { | |
| return this.tokenCounting.requestSize({ | |
| systemPrompt: this.profile.getSystemPrompt(), | |
| tools: this.defaultTools().filter((tool) => tool.deferred !== true), | |
| messages, | |
| }); | |
| } | |
| private defaultTools(): readonly Tool[] { | |
| return this.toolSelect | |
| .shapeTools(this.toolRegistry.list()) | |
| .map((tool) => ({ | |
| name: tool.name, | |
| description: tool.description, | |
| parameters: tool.parameters ?? EMPTY_TOOL_PARAMETERS, | |
| deferred: tool.deferred, | |
| })); | |
| } | |
| private shouldRecoverFromContextOverflow( | |
| error: unknown, | |
| estimatedRequestTokens = this.currentRequestTokens(), | |
| ): boolean { | |
| if (isCodedError(error) && error.code === ErrorCodes.CONTEXT_OVERFLOW) return true; | |
| const statusError = findAPIStatusError(error); | |
| if (statusError instanceof APIContextOverflowError) return true; | |
| if (statusError === undefined || statusError.statusCode !== 413) return false; | |
| const effectiveMax = this.getEffectiveMaxContextTokens(); | |
| return ( | |
| effectiveMax > 0 && | |
| estimatedRequestTokens >= effectiveMax * OVERFLOW_STATUS_RECOVERY_RATIO | |
| ); | |
| } | |
| private observeContextOverflow(estimatedRequestTokens: number): void { | |
| if (!Number.isFinite(estimatedRequestTokens) || estimatedRequestTokens <= 0) return; | |
| const modelAlias = this.profile.data().modelAlias; | |
| if (modelAlias === undefined) return; | |
| const observed = Math.max( | |
| 1, | |
| Math.floor(estimatedRequestTokens * OVERFLOW_CONTEXT_SAFETY_RATIO), | |
| ); | |
| const current = this.getEffectiveMaxContextTokens(); | |
| if (current > 0 && observed >= current) return; | |
| this.observedMaxContextTokensByModel.set(modelAlias, observed); | |
| } | |
| begin(input: FullCompactionInput): boolean { | |
| if (this._compacting) return false; | |
| const data: CompactionBeginData = { source: input.source, instruction: input.instruction }; | |
| if (!this.reserveCompactionSlot(data.source)) return false; | |
| const tokenCount = this.validateCompactionStart(data.source); | |
| const quiescence = data.source === 'manual' | |
| ? this.loopService.tryAcquireQuiescence() | |
| : undefined; | |
| if (data.source === 'manual' && quiescence === undefined) { | |
| throw new Error2( | |
| ErrorCodes.COMPACTION_UNABLE, | |
| 'Cannot compact while a turn is active or another context change is running. Wait for it to finish, then retry.', | |
| ); | |
| } | |
| try { | |
| void this.dispatcher.dispatch( | |
| new FullCompactionBegin({ ...data, agentId: this.agent.agentId }), | |
| ); | |
| const active = this.createActiveCompaction( | |
| data.source, | |
| tokenCount, | |
| data.source === 'auto' ? this.activeTurnId : undefined, | |
| quiescence, | |
| ); | |
| this._compacting = active.task; | |
| active.task.abortController.signal.addEventListener( | |
| 'abort', | |
| () => this.cancelActive(active.task), | |
| { once: true }, | |
| ); | |
| void this.compactionWorker(active.task, data).then(active.resolve, active.reject); | |
| void active.task.promise.catch(() => undefined); | |
| return true; | |
| } catch (error) { | |
| quiescence?.dispose(); | |
| throw error; | |
| } | |
| } | |
| private reserveCompactionSlot(source: CompactionBeginData['source']): boolean { | |
| if (source === 'manual') { | |
| this.compactionCountInTurn = 0; | |
| } else { | |
| this.compactionCountInTurn += 1; | |
| } | |
| return this.compactionCountInTurn <= this.strategy.maxCompactionPerTurn; | |
| } | |
| private validateCompactionStart(source: CompactionBeginData['source']): number { | |
| const history = this.context.get(); | |
| if (history.length === 0) { | |
| throw new Error2(ErrorCodes.COMPACTION_UNABLE, 'No messages to compact in current history.'); | |
| } | |
| if (source === 'manual' && this.loopService.snapshot().state !== 'idle') { | |
| throw new Error2( | |
| ErrorCodes.COMPACTION_UNABLE, | |
| 'Cannot compact while a turn is active. Wait for it to finish, then retry.', | |
| ); | |
| } | |
| return this.requestTokens(history); | |
| } | |
| private createActiveCompaction( | |
| trigger: CompactionBeginData['source'], | |
| tokenCount: number, | |
| originTurnId: number | undefined, | |
| quiescence: IDisposable | undefined, | |
| ): { | |
| readonly task: ActiveCompaction; | |
| readonly resolve: (result: CompactionResult) => void; | |
| readonly reject: (reason: unknown) => void; | |
| } { | |
| const abortController = new AbortController(); | |
| let resolve!: (result: CompactionResult) => void; | |
| let reject!: (reason: unknown) => void; | |
| const promise = new Promise<CompactionResult>((onResolve, onReject) => { | |
| resolve = onResolve; | |
| reject = onReject; | |
| }); | |
| return { | |
| task: { | |
| abortController, | |
| promise, | |
| trigger, | |
| tokenCount, | |
| originTurnId, | |
| quiescence, | |
| get traceId() { | |
| return this.trace?.traceId; | |
| }, | |
| blockedByTurn: false, | |
| }, | |
| resolve, | |
| reject, | |
| }; | |
| } | |
| override dispose(): void { | |
| if (this._compacting !== null && !this._compacting.abortController.signal.aborted) { | |
| this._compacting.abortController.abort(); | |
| } | |
| super.dispose(); | |
| } | |
| private cancelActive(active: ActiveCompaction): boolean { | |
| if (this._compacting !== active) return false; | |
| void this.dispatcher.dispatch(new FullCompactionCancel({ agentId: this.agent.agentId })); | |
| this._compacting = null; | |
| if (!active.abortController.signal.aborted) { | |
| active.abortController.abort(); | |
| } | |
| void this.dispatcher.dispatch(new CompactionCancelled({ agentId: this.agent.agentId })); | |
| return true; | |
| } | |
| private markCompleted(active: ActiveCompaction): boolean { | |
| if (this._compacting !== active) return false; | |
| void this.dispatcher.dispatch(new FullCompactionComplete({ agentId: this.agent.agentId })); | |
| this._compacting = null; | |
| return true; | |
| } | |
| private normalizeAfterReplay(): void { | |
| if (this.states.get(fullCompactionKey).phase !== 'running') return; | |
| void this.dispatcher.dispatch(new FullCompactionCancel({ agentId: this.agent.agentId })); | |
| } | |
| private resetForTurn(): void { | |
| this.compactionCountInTurn = 0; | |
| this.lastCompactedTokenCount = null; | |
| this.consecutiveOverflowCompactions = 0; | |
| } | |
| private async recoverFromContextOverflow( | |
| context: LoopErrorContext, | |
| ): Promise<boolean> { | |
| this.recordOverflowRecovery(context.error); | |
| const didStartCompaction = this.beginAutoCompaction(); | |
| if (!didStartCompaction && !this._compacting) return false; | |
| await this.block(context.signal, context.turnId); | |
| return this.retryFailedDriver(context); | |
| } | |
| private recordOverflowRecovery(error: unknown): void { | |
| this.observeContextOverflow(this.currentRequestTokens()); | |
| this.consecutiveOverflowCompactions += 1; | |
| const maxAttempts = this.strategy.maxOverflowCompactionAttempts; | |
| if (this.consecutiveOverflowCompactions <= maxAttempts) return; | |
| throw new Error2( | |
| ErrorCodes.CONTEXT_OVERFLOW, | |
| `Compaction failed to bring the context under the model window after ${String(maxAttempts)} attempts.`, | |
| { cause: error instanceof Error ? error : undefined }, | |
| ); | |
| } | |
| private retryFailedDriver(context: LoopErrorContext): boolean { | |
| if (context.signal.aborted) return false; | |
| context.retry(); | |
| return true; | |
| } | |
| private async beforeStep(signal: AbortSignal, turnId?: number): Promise<void> { | |
| this.activeTurnId = turnId; | |
| this.checkAutoCompaction(); | |
| if (this.strategy.shouldBlock(this.tokenCountWithPending())) { | |
| await this.block(signal, turnId); | |
| } | |
| } | |
| private async afterStep(): Promise<void> { | |
| this.consecutiveOverflowCompactions = 0; | |
| if (this.strategy.checkAfterStep) { | |
| this.checkAutoCompaction(false); | |
| } | |
| } | |
| private checkAutoCompaction(throwOnLimit = true): boolean { | |
| if (this._compacting) return true; | |
| if ( | |
| this.lastCompactedTokenCount !== null && | |
| this.tokenCountWithPending() <= this.lastCompactedTokenCount | |
| ) { | |
| return false; | |
| } | |
| if (!this.strategy.shouldCompact(this.tokenCountWithPending())) return false; | |
| return this.beginAutoCompaction(throwOnLimit); | |
| } | |
| private beginAutoCompaction(throwOnLimit = true): boolean { | |
| if (this._compacting) return true; | |
| const maxCompactions = this.strategy.maxCompactionPerTurn; | |
| if (this.compactionCountInTurn >= maxCompactions) { | |
| if (throwOnLimit) { | |
| throw new Error2(ErrorCodes.CONTEXT_OVERFLOW, `Compaction limit exceeded (${String(maxCompactions)})`, { | |
| details: { maxCompactions }, | |
| }); | |
| } | |
| return false; | |
| } | |
| return this.begin({ source: 'auto' }); | |
| } | |
| private async block(signal?: AbortSignal, turnId?: number): Promise<void> { | |
| const active = this._compacting; | |
| if (active === null) return; | |
| active.blockedByTurn = true; | |
| this.propagateBlockingAbort(active, signal); | |
| void this.dispatcher.dispatch( | |
| new CompactionBlocked({ agentId: this.agent.agentId, turnId }), | |
| ); | |
| try { | |
| await active.promise; | |
| } catch (error) { | |
| if (this.wasBlockingWaitAborted(active, signal, error)) return; | |
| throw error; | |
| } | |
| } | |
| private propagateBlockingAbort(active: ActiveCompaction, signal: AbortSignal | undefined): void { | |
| signal?.addEventListener( | |
| 'abort', | |
| () => { | |
| if (this._compacting === active) active.abortController.abort(); | |
| }, | |
| { once: true }, | |
| ); | |
| } | |
| private wasBlockingWaitAborted( | |
| active: ActiveCompaction, | |
| signal: AbortSignal | undefined, | |
| error: unknown, | |
| ): boolean { | |
| return ( | |
| signal?.aborted === true && | |
| (active.abortController.signal.aborted || isAbortError(error)) | |
| ); | |
| } | |
| private async compactionWorker( | |
| active: ActiveCompaction, | |
| data: Readonly<CompactionBeginData>, | |
| ): Promise<CompactionResult> { | |
| try { | |
| const result = await this.compactionRound(active, data); | |
| if (this._compacting !== active) throw compactionCancelledReason(active); | |
| this.lastCompactedTokenCount = result.tokensAfter; | |
| if (!this.markCompleted(active)) { | |
| throw compactionCancelledReason(active); | |
| } | |
| const { contextSummary: _contextSummary, ...eventResult } = result; | |
| void _contextSummary; | |
| void this.dispatcher.dispatch( | |
| new CompactionCompleted({ agentId: this.agent.agentId, result: eventResult }), | |
| ); | |
| return result; | |
| } catch (error) { | |
| 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; | |
| } | |
| void this.dispatcher.dispatch( | |
| new AgentErrorEvent({ ...toKimiErrorPayload(error), agentId: this.agent.agentId }), | |
| ); | |
| throw error; | |
| } finally { | |
| try { | |
| this._onDidFinishCompaction.fire(active); | |
| } finally { | |
| active.quiescence?.dispose(); | |
| } | |
| } | |
| } | |
| private async compactionRound( | |
| active: ActiveCompaction, | |
| data: Readonly<CompactionBeginData>, | |
| ): Promise<CompactionResult> { | |
| const startedAt = Date.now(); | |
| const originalHistory = [...this.context.get()]; | |
| const tokensBefore = this.requestTokens(originalHistory); | |
| let retryCount = 0; | |
| let thinkingEffort = this.profile.data().thinkingLevel; | |
| try { | |
| const signal = active.abortController.signal; | |
| signal.throwIfAborted(); | |
| await this.hooks.onWillCompact.run(active); | |
| const resolvedModel = this.profile.resolveModelContext(); | |
| thinkingEffort = resolvedModel.thinkingLevel; | |
| const maxContextTokens = resolvedModel.modelCapabilities.max_context_tokens; | |
| const defaultCompactionCap = | |
| maxContextTokens > 0 | |
| ? Math.min(maxContextTokens, DEFAULT_COMPACTION_MAX_COMPLETION_TOKENS) | |
| : undefined; | |
| const compactionMaxOutputSize = resolvedModel.maxOutputSize ?? defaultCompactionCap; | |
| const instruction = renderCompactionInstruction({ customInstruction: data.instruction }); | |
| const maxAttempts = resolvedModel.compactionMaxAttempts ?? MAX_COMPACTION_RETRY_ATTEMPTS; | |
| let attempt: CompactionAttemptResult | undefined; | |
| let historyForModel: readonly ContextMessage[] = stripDynamicToolContext(originalHistory); | |
| let droppedCount = 0; | |
| let overflowShrinkCount = 0; | |
| let requestAttempts = 0; | |
| while (true) { | |
| const messagesToCompact = historyForModel; | |
| const messages: Message[] = [...messagesToCompact, createUserMessage(instruction)]; | |
| const estimatedCompactionRequestTokens = this.requestTokens(messages); | |
| requestAttempts += 1; | |
| try { | |
| const runRequest = async () => { | |
| const request = this.llmRequester.start( | |
| { | |
| messages, | |
| maxOutputSize: compactionMaxOutputSize, | |
| source: { | |
| type: 'operation', | |
| turnId: active.originTurnId, | |
| requestKind: 'full_compaction', | |
| logFields: { droppedCount }, | |
| }, | |
| }, | |
| undefined, | |
| signal, | |
| ); | |
| active.trace = request.trace; | |
| return request.result; | |
| }; | |
| const result = await runWithCredentialRecovery( | |
| this.llmRequester.currentCredentialProvider(), | |
| runRequest, | |
| signal, | |
| ); | |
| attempt = collectSummary(result); | |
| break; | |
| } catch (error) { | |
| const isContextOverflow = this.shouldRecoverFromContextOverflow( | |
| error, | |
| estimatedCompactionRequestTokens, | |
| ); | |
| if (isContextOverflow) { | |
| this.observeContextOverflow(estimatedCompactionRequestTokens); | |
| overflowShrinkCount += 1; | |
| if ( | |
| overflowShrinkCount > MAX_COMPACTION_OVERFLOW_SHRINK_ATTEMPTS || | |
| requestAttempts >= maxAttempts || | |
| messagesToCompact.length <= 1 | |
| ) { | |
| throw error; | |
| } | |
| const before = messagesToCompact.length; | |
| historyForModel = shrinkCompactionHistoryAfterOverflow( | |
| messagesToCompact, | |
| overflowShrinkCount, | |
| (message) => this.tokenCounting.estimateMessage(message), | |
| ); | |
| if (historyForModel.length === 0) throw error; | |
| droppedCount += before - historyForModel.length; | |
| retryCount = 0; | |
| continue; | |
| } | |
| const unwrappedError = unwrapErrorCause(error); | |
| if ( | |
| (error instanceof CompactionTruncatedError || | |
| (unwrappedError instanceof APIEmptyResponseError && | |
| unwrappedError.finishReason !== 'filtered')) && | |
| messagesToCompact.length > 1 | |
| ) { | |
| if (requestAttempts >= maxAttempts) { | |
| throw error; | |
| } | |
| const reduced = dropOldestMessageAndLeadingToolResults(messagesToCompact); | |
| droppedCount += messagesToCompact.length - reduced.length; | |
| historyForModel = reduced; | |
| retryCount = 0; | |
| continue; | |
| } | |
| if (!isRetryableGenerateError(unwrappedError)) { | |
| throw error; | |
| } | |
| if (requestAttempts >= maxAttempts) { | |
| throw error; | |
| } | |
| await sleepForRetry(retryBackoffDelay(retryCount), signal); | |
| retryCount += 1; | |
| } | |
| } | |
| if (attempt === undefined) { | |
| throw new APIEmptyResponseError( | |
| 'The compaction response did not contain a usable summary.', | |
| ); | |
| } | |
| if (!historySafeToCompact(this.context.get(), originalHistory)) { | |
| const active = this._compacting; | |
| if (active !== null) { | |
| this.cancelActive(active); | |
| } | |
| throw compactionCancelledReason(active); | |
| } | |
| const summary = await this.postProcessSummary(attempt.summary); | |
| const wireLines = await this.captureWireLines(); | |
| const recoveryFooter = this.renderRecoveryFooter(wireLines); | |
| const summaryText = buildCompactionSummaryText(summary); | |
| const result = this.context.applyCompaction({ | |
| summary, | |
| contextSummary: | |
| recoveryFooter === undefined ? summaryText : `${summaryText}\n\n${recoveryFooter}`, | |
| compactedCount: originalHistory.length, | |
| tokensBefore, | |
| summaryOutputTokens: | |
| attempt.usage === null | |
| ? undefined | |
| : attempt.usage.output + | |
| (recoveryFooter === undefined ? 0 : this.tokenCounting.estimateText(recoveryFooter)), | |
| requestOverheadTokens: this.requestTokens([]), | |
| droppedCount: droppedCount === 0 ? undefined : droppedCount, | |
| wireLines, | |
| }); | |
| const properties: CompactionFinishedEvent = { | |
| turn_id: active.originTurnId, | |
| source: data.source, | |
| tokens_before: result.tokensBefore, | |
| tokens_after: result.tokensAfter, | |
| duration_ms: Date.now() - startedAt, | |
| compacted_count: result.compactedCount, | |
| dropped_count: result.droppedCount, | |
| retry_count: retryCount, | |
| round: 1, | |
| thinking_effort: thinkingEffort, | |
| trace_id: attempt.traceId, | |
| ...usageTelemetry(attempt.usage), | |
| }; | |
| this.telemetry.track2('compaction_finished', properties); | |
| return result; | |
| } catch (error) { | |
| if (isAbortError(error)) throw error; | |
| const properties: CompactionFailedEvent = { | |
| turn_id: active.originTurnId, | |
| source: data.source, | |
| tokens_before: tokensBefore, | |
| duration_ms: Date.now() - startedAt, | |
| round: 1, | |
| retry_count: retryCount, | |
| thinking_effort: thinkingEffort, | |
| error_type: error instanceof Error ? error.name : 'Unknown', | |
| trace_id: findAPIStatusError(error)?.traceId ?? active.traceId, | |
| }; | |
| this.telemetry.track2('compaction_failed', properties); | |
| if ( | |
| isError2(error) && | |
| (error.code === ErrorCodes.AUTH_LOGIN_REQUIRED || | |
| error.code === ErrorCodes.PROVIDER_AUTH_ERROR) | |
| ) { | |
| throw error; | |
| } | |
| throw new Error2(ErrorCodes.COMPACTION_FAILED, String(error), { cause: error }); | |
| } | |
| } | |
| private async postProcessSummary(summary: string): Promise<string> { | |
| const todos = this.todo.get(); | |
| if (todos.length === 0) { | |
| return summary; | |
| } | |
| return `${summary.trim()}\n\n${renderTodoList(todos, '## TODO List')}`; | |
| } | |
| private async captureWireLines(): Promise<WireLineRange | undefined> { | |
| try { | |
| await this.wire.flush(); | |
| } catch (error) { | |
| onUnexpectedError(error); | |
| return undefined; | |
| } | |
| const end = this.wire.lineCount(); | |
| const previous = this.states.get(fullCompactionWireRangesKey).at(-1); | |
| const start = Math.max(previous?.end ?? 0, this.wire.lastContextClearLine() ?? 0) + 1; | |
| if (end < start) return undefined; | |
| return { start, end }; | |
| } | |
| private renderRecoveryFooter(wireLines: WireLineRange | undefined): string | undefined { | |
| if (wireLines === undefined) return undefined; | |
| const journalPath = this.wire.journalPath(); | |
| if (journalPath === undefined) return undefined; | |
| const windows = [...this.states.get(fullCompactionWireRangesKey), wireLines]; | |
| return renderContextRecoveryPointer({ journalPath, windows }); | |
| } | |
| private tokenCountWithPending(): number { | |
| return this.tokenCounting.get(agentContextOfScope(this.agent)).size; | |
| } | |
| } | |
| function findAPIStatusError(error: unknown): APIStatusError | undefined { | |
| let current: unknown = error; | |
| const seen = new Set<unknown>(); | |
| while (current !== undefined && current !== null && !seen.has(current)) { | |
| if (current instanceof APIStatusError) return current; | |
| seen.add(current); | |
| current = current instanceof Error ? current.cause : undefined; | |
| } | |
| return undefined; | |
| } | |
| function collectSummary(finish: AgentLLMRequestFinish): CompactionAttemptResult { | |
| if (finish.providerFinishReason === 'truncated') { | |
| throw new CompactionTruncatedError(); | |
| } | |
| const summary = finish.message.content | |
| .filter((part) => part.type === 'text') | |
| .map((part) => part.text) | |
| .join('') | |
| .trim(); | |
| if (summary.length === 0) { | |
| throw new APIEmptyResponseError( | |
| 'The compaction response did not contain a non-empty summary.', | |
| ); | |
| } | |
| return { summary, usage: finish.usage, traceId: finish.traceId }; | |
| } | |
| function historySafeToCompact( | |
| current: readonly ContextMessage[], | |
| original: readonly ContextMessage[], | |
| ): boolean { | |
| if (current.length < original.length) return false; | |
| if (!original.every((message, index) => message === current[index])) return false; | |
| return current.slice(original.length).every(isRealUserInput); | |
| } | |
| function shrinkCompactionHistoryAfterOverflow<T extends Message>( | |
| messages: readonly T[], | |
| attempt: number, | |
| estimateMessage: (message: T) => number = estimateTokensForMessage, | |
| ): T[] { | |
| if (messages.length <= 1) return messages.slice(); | |
| const ratio = COMPACTION_OVERFLOW_SHRINK_RATIOS[ | |
| Math.min(attempt - 1, COMPACTION_OVERFLOW_SHRINK_RATIOS.length - 1) | |
| ]!; | |
| let totalTokens = 0; | |
| for (const message of messages) totalTokens += estimateMessage(message); | |
| const tokenBudget = Math.floor(totalTokens * ratio); | |
| return takeRecentMessagesWithinTokenBudget(messages, tokenBudget, estimateMessage); | |
| } | |
| function takeRecentMessagesWithinTokenBudget<T extends Message>( | |
| messages: readonly T[], | |
| tokenBudget: number, | |
| estimateMessage: (message: T) => number = estimateTokensForMessage, | |
| ): T[] { | |
| let start = messages.length; | |
| let tokens = 0; | |
| for (let i = messages.length - 1; i >= 0; i--) { | |
| const messageTokens = estimateMessage(messages[i]!); | |
| if (tokens + messageTokens > tokenBudget) break; | |
| tokens += messageTokens; | |
| start = i; | |
| } | |
| if (start === 0) start = 1; | |
| return dropLeadingToolResults(messages.slice(start)); | |
| } | |
| function dropOldestMessageAndLeadingToolResults<T extends { readonly role: string }>( | |
| messages: readonly T[], | |
| ): T[] { | |
| if (messages.length <= 1) return messages.slice(); | |
| return dropLeadingToolResults(messages.slice(1)); | |
| } | |
| function dropLeadingToolResults<T extends { readonly role: string }>(messages: readonly T[]): T[] { | |
| let start = 0; | |
| while (start < messages.length && messages[start]!.role === 'tool') { | |
| start += 1; | |
| } | |
| return messages.slice(start); | |
| } | |
| function usageTelemetry(usage: TokenUsage | null): CompactionTelemetryProperties { | |
| if (usage === null) return {}; | |
| return { | |
| input_tokens: inputTotal(usage), | |
| output_tokens: usage.output, | |
| input_cache_read: usage.inputCacheRead, | |
| input_cache_creation: usage.inputCacheCreation, | |
| }; | |
| } | |
| 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; | |
| } | |
| registerScopedService( | |
| LifecycleScope.Agent, | |
| IAgentFullCompactionService, | |
| AgentFullCompactionService, | |
| ScopeActivation.OnScopeCreated, | |
| 'fullCompaction', | |
| ); | |