import type { AgentEndEvent, ExtensionAPI, ExtensionContext, MessageStartEvent, SessionBeforeCompactEvent, SessionCompactEvent, TurnEndEvent, TurnStartEvent, } from "@earendil-works/pi-coding-agent"; import type { AutoCompactConfig, ConfigLoadResult } from "./config.ts"; import { CoordinationClient, type CoordinationTransaction } from "./coordination-client.ts"; import type { CompactionOutcome } from "./coordination-protocol.ts"; export const CONTINUATION_CUSTOM_TYPE = "wj-pi-auto-compact-continuation"; export const CONTINUATION_CONTENT = "[wj-pi-auto-compact/v1]\nAuto-compaction completed. Continue the interrupted task from the compacted context."; export const CONTINUATION_START_TIMEOUT_MS = 5_000; export interface AutoCompactControllerOptions { readonly continuationStartTimeoutMs?: number; } interface ContextUsageSnapshot { readonly tokens: number; readonly contextWindow: number; readonly percent: number; } type PendingContinuationState = "awaiting_start" | "start_unconfirmed" | "cancelling"; interface PendingContinuation { readonly transaction: CoordinationTransaction; readonly requestId: string; readonly epoch: number; timer: NodeJS.Timeout | undefined; state: PendingContinuationState; cancellation: Promise | undefined; } interface ActiveCompaction { readonly transaction: CoordinationTransaction; readonly baselineCompactionId: string | undefined; readonly needsContinuation: boolean; nativeAttempted: boolean; nativeAttemptSignal: AbortSignal | undefined; nativeWillRetry: boolean; compactionObserved: boolean; externalManualObserved: boolean; manualInFlight: boolean; cancelManual: (() => void) | undefined; } export class AutoCompactController { private readonly pi: ExtensionAPI; private readonly configResult: ConfigLoadResult; private readonly continuationStartTimeoutMs: number; private config: AutoCompactConfig; private coordination: CoordinationClient; private sessionEpoch = 0; private thresholdPending = false; private thresholdUsage: ContextUsageSnapshot | undefined; private toolContinuationCandidate = false; private prepareFailedForRun = false; private preparing: Promise | undefined; private active: ActiveCompaction | undefined; private readonly pendingContinuations = new Map(); private warningShown = false; constructor( pi: ExtensionAPI, configResult: ConfigLoadResult, options: AutoCompactControllerOptions = {}, ) { this.pi = pi; this.configResult = configResult; this.continuationStartTimeoutMs = options.continuationStartTimeoutMs ?? CONTINUATION_START_TIMEOUT_MS; this.config = configResult.config; this.coordination = new CoordinationClient(pi.events); } async startSession(ctx: ExtensionContext): Promise { await this.endCurrentSession(); this.coordination = new CoordinationClient(this.pi.events); void this.coordination.discover(); if (this.configResult.kind === "invalid" && !this.warningShown && ctx.hasUI) { this.warningShown = true; ctx.ui.notify(`Auto-compaction is disabled: ${this.configResult.reason}`, "warning"); } } onTurnEnd(event: TurnEndEvent, ctx: ExtensionContext): void { if (!this.config.enabled || this.active || this.prepareFailedForRun) return; const usage = ctx.getContextUsage(); if (usage === undefined || usage.tokens === null || usage.percent === null) return; if (usage.percent < this.config.maxContextPercent) return; this.thresholdPending = true; this.thresholdUsage = { tokens: usage.tokens, contextWindow: usage.contextWindow, percent: usage.percent, }; this.toolContinuationCandidate = isSuccessfulToolTurn(event); } async onTurnStart(_event: TurnStartEvent, ctx: ExtensionContext): Promise { if ( !this.config.enabled || !this.thresholdPending || !this.toolContinuationCandidate || this.active || this.prepareFailedForRun ) return; const prepared = await this.establishBarrier(ctx, true); if (!prepared) return; ctx.abort(); } async onAgentEnd(_event: AgentEndEvent, ctx: ExtensionContext): Promise { if (!this.config.enabled || this.active || !this.thresholdPending || this.prepareFailedForRun) return; await this.establishBarrier(ctx, false); } onSessionBeforeCompact(event: SessionBeforeCompactEvent, ctx: ExtensionContext): void { const active = this.active; if (!active || active.compactionObserved || !sameBranch(event.branchEntries, ctx.sessionManager.getBranch())) return; if (event.reason === "manual") { if (!active.manualInFlight) active.externalManualObserved = true; return; } active.nativeAttempted = true; active.nativeAttemptSignal = event.signal; active.nativeWillRetry ||= event.reason === "overflow" && event.willRetry; } onSessionCompact(event: SessionCompactEvent, ctx: ExtensionContext): void { const active = this.active; if ( !active || event.compactionEntry.id === active.baselineCompactionId || !ctx.sessionManager.getBranch().some((entry) => entry.id === event.compactionEntry.id) ) return; if (event.reason === "manual") { if (!active.manualInFlight) active.externalManualObserved = true; return; } active.compactionObserved = true; active.nativeWillRetry ||= event.reason === "overflow" && event.willRetry; } onMessageStart(event: MessageStartEvent): void { if ( event.message.role !== "custom" || event.message.customType !== CONTINUATION_CUSTOM_TYPE ) return; const requestId = continuationRequestId(event.message.details); if (requestId === undefined) return; const pending = this.pendingContinuations.get(requestId); if (!pending || pending.epoch !== this.sessionEpoch || pending.state === "cancelling") return; if (pending.timer !== undefined) clearTimeout(pending.timer); this.pendingContinuations.delete(requestId); } async onAgentSettled(ctx: ExtensionContext): Promise { const active = this.active; if (!active) { this.prepareFailedForRun = false; return; } if (active.compactionObserved) { await this.finish(active, "succeeded"); return; } if (active.externalManualObserved || this.hasNewCompaction(ctx, active)) { await this.finish(active, "not_started"); return; } if (active.nativeAttempted) { const outcome: CompactionOutcome = active.nativeAttemptSignal?.aborted ? "cancelled" : "failed"; await this.finish(active, outcome); return; } const outcome = await this.runManualCompaction(ctx, active); await this.finish(active, outcome); } async shutdown(): Promise { await this.endCurrentSession(); } private async endCurrentSession(): Promise { this.sessionEpoch += 1; this.coordination.close(); const pendingContinuations = [...this.pendingContinuations.values()]; for (const pending of pendingContinuations) { if (pending.timer !== undefined) clearTimeout(pending.timer); } this.pendingContinuations.clear(); await Promise.all(pendingContinuations.flatMap((pending) => pending.cancellation === undefined ? [] : [pending.cancellation] )); const preparing = this.preparing; if (preparing !== undefined) { try { await preparing; } catch { // prepare 自身会把不确定参与者回滚为 not_started。 } } const active = this.active; this.active = undefined; if (active !== undefined) { active.cancelManual?.(); try { await active.transaction.cancel(); } catch { // 业务 ACK 已有期限;关闭流程仍需清理本地监听器和状态。 } } this.resetTriggerState(); } private async establishBarrier(ctx: ExtensionContext, needsContinuation: boolean): Promise { if (this.preparing) return await this.preparing; const usage = this.thresholdUsage; if (usage === undefined) { this.thresholdPending = false; this.toolContinuationCandidate = false; return false; } const epoch = this.sessionEpoch; const baselineCompactionId = latestCompactionId(ctx); this.preparing = (async () => { const result = await this.coordination.prepare(); if (epoch !== this.sessionEpoch || this.active) { if (result.prepared) await result.transaction.cancel(); return false; } if (!result.prepared) { this.prepareFailedForRun = true; this.thresholdPending = false; this.thresholdUsage = undefined; this.toolContinuationCandidate = false; return false; } this.active = { transaction: result.transaction, baselineCompactionId, needsContinuation, nativeAttempted: false, nativeAttemptSignal: undefined, nativeWillRetry: false, compactionObserved: false, externalManualObserved: false, manualInFlight: false, cancelManual: undefined, }; this.thresholdPending = false; this.thresholdUsage = undefined; this.toolContinuationCandidate = false; if (ctx.hasUI) { ctx.ui.notify(formatCompactionNotice(usage, this.config.maxContextPercent), "info"); } return true; })().finally(() => { this.preparing = undefined; }); return await this.preparing; } private async runManualCompaction(ctx: ExtensionContext, active: ActiveCompaction): Promise { active.manualInFlight = true; return await new Promise((resolve) => { let settled = false; const finish = (outcome: CompactionOutcome): void => { if (settled) return; settled = true; active.manualInFlight = false; active.cancelManual = undefined; resolve(outcome); }; active.cancelManual = () => finish("cancelled"); const customInstructions = this.config.customInstructions; try { ctx.compact({ ...(customInstructions === undefined ? {} : { customInstructions }), onComplete: () => finish("succeeded"), onError: (error) => { if (active.compactionObserved) { finish("succeeded"); return; } finish(isCancellation(error) ? "cancelled" : "failed"); }, }); } catch (error) { if (active.compactionObserved) { finish("succeeded"); return; } finish(isCancellation(error) ? "cancelled" : "failed"); } }); } private async finish(active: ActiveCompaction, outcome: CompactionOutcome): Promise { if (this.active !== active) return; const continuationExpected = outcome === "succeeded" && active.needsContinuation && !active.nativeWillRetry; const completed = await active.transaction.complete(outcome, continuationExpected); if (this.active !== active) return; this.active = undefined; this.prepareFailedForRun = false; if (continuationExpected && completed) { await this.startContinuation(active.transaction, continuationExpected); } } private async startContinuation( transaction: CoordinationTransaction, continuationExpected: boolean, ): Promise { if (!continuationExpected) return; const pending: PendingContinuation = { transaction, requestId: transaction.requestId, epoch: this.sessionEpoch, timer: undefined, state: "awaiting_start", cancellation: undefined, }; this.pendingContinuations.set(pending.requestId, pending); try { this.pi.sendMessage( { customType: CONTINUATION_CUSTOM_TYPE, content: CONTINUATION_CONTENT, display: true, details: Object.freeze({ coordinationRequestId: transaction.requestId }), }, { triggerTurn: true }, ); } catch { await this.cancelPendingContinuation(pending); return; } if ( this.pendingContinuations.get(pending.requestId) === pending && pending.state === "awaiting_start" ) { pending.timer = setTimeout(() => { this.markContinuationStartUnconfirmed(pending); }, this.continuationStartTimeoutMs); } } private markContinuationStartUnconfirmed(pending: PendingContinuation): void { if ( this.pendingContinuations.get(pending.requestId) !== pending || pending.state !== "awaiting_start" ) return; pending.state = "start_unconfirmed"; pending.timer = undefined; } private cancelPendingContinuation(pending: PendingContinuation): Promise { if (this.pendingContinuations.get(pending.requestId) !== pending) { return pending.cancellation ?? Promise.resolve(); } if (pending.cancellation !== undefined) return pending.cancellation; pending.state = "cancelling"; if (pending.timer !== undefined) { clearTimeout(pending.timer); pending.timer = undefined; } pending.cancellation = (async () => { try { await pending.transaction.cancel(); } catch { // 同步发送失败是确定的本地事实;ACK 失败只影响远端确认。 } finally { if (this.pendingContinuations.get(pending.requestId) === pending) { this.pendingContinuations.delete(pending.requestId); } } })(); return pending.cancellation; } private hasNewCompaction(ctx: ExtensionContext, active: ActiveCompaction): boolean { return active.compactionObserved || latestCompactionId(ctx) !== active.baselineCompactionId; } private resetTriggerState(): void { this.thresholdPending = false; this.thresholdUsage = undefined; this.toolContinuationCandidate = false; this.prepareFailedForRun = false; this.preparing = undefined; this.active = undefined; } } const PERCENT_FORMATTER = new Intl.NumberFormat("en-US", { maximumFractionDigits: 2 }); const TOKEN_FORMATTER = new Intl.NumberFormat("en-US", { maximumFractionDigits: 0 }); function formatCompactionNotice(usage: ContextUsageSnapshot, thresholdPercent: number): string { const thresholdTokens = Math.ceil((usage.contextWindow * thresholdPercent) / 100); return `Context usage is ${PERCENT_FORMATTER.format(usage.percent)}% (${TOKEN_FORMATTER.format(usage.tokens)} / ${TOKEN_FORMATTER.format(usage.contextWindow)} tokens). Auto-compaction threshold is ${PERCENT_FORMATTER.format(thresholdPercent)}% (${TOKEN_FORMATTER.format(thresholdTokens)} tokens). Preparing to compact.`; } function latestCompactionId(ctx: ExtensionContext): string | undefined { const branch = ctx.sessionManager.getBranch(); for (let index = branch.length - 1; index >= 0; index -= 1) { const entry = branch[index]; if (entry?.type === "compaction") return entry.id; } return undefined; } function isSuccessfulToolTurn(event: TurnEndEvent): boolean { if (event.message.role !== "assistant") return false; if (event.message.stopReason === "error" || event.message.stopReason === "aborted") return false; return event.toolResults.length > 0 && event.message.content.some((content) => content.type === "toolCall"); } function sameBranch( eventBranch: SessionBeforeCompactEvent["branchEntries"], currentBranch: ReturnType, ): boolean { if (eventBranch.length !== currentBranch.length) return false; return eventBranch.every((entry, index) => entry.id === currentBranch[index]?.id); } function continuationRequestId(details: unknown): string | undefined { if (typeof details !== "object" || details === null || Array.isArray(details)) return undefined; const requestId = (details as Record).coordinationRequestId; return typeof requestId === "string" ? requestId : undefined; } function isCancellation(error: unknown): boolean { if (!(error instanceof Error)) return false; return error.name === "AbortError" || error.message === "Compaction cancelled"; }