refactor: simplify full compaction contract

This commit is contained in:
_Kerman 2026-07-07 15:03:11 +08:00
parent 08e754401a
commit e2cecf6551
4 changed files with 91 additions and 56 deletions

View file

@ -398,17 +398,18 @@ export class AgentExternalHooksService extends Disposable implements IAgentExter
}
private async runPreCompact(ctx: FullCompactionWillCompactContext): Promise<void> {
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 {

View file

@ -13,9 +13,10 @@ export interface CompactInput {
}
export interface FullCompactionWillCompactContext {
readonly abortController: AbortController;
readonly promise: Promise<CompactionResult>;
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;

View file

@ -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<string, string | number | boolean | undefined>;
interface ActiveCompaction {
readonly abortController: AbortController;
promise: Promise<void>;
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<CompactionResult>((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<void> {
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<CompactionBeginData>,
initialCompactedCount: number,
): Promise<void> {
): Promise<CompactionResult> {
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<CompactionBeginData>,
initialCompactedCount: number,
): Promise<CompactionResult | undefined> {
): Promise<CompactionResult> {
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 };

View file

@ -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 {