返回源码地图

packages/coding-agent/src/core/compaction/compaction.ts

v1.0.0 · a13d35a742c6 · 07:usage估算、合法切点、split turn与摘要失败

完整原文供逐行核对;页面收录不代表每行都经过人工语义审核。MIT 许可见 许可证。

1/**
2 * Context compaction for long sessions.
3 *
4 * Pure functions for compaction logic. The session manager handles I/O,
5 * and after compaction the session is reloaded.
6 */
7
8import type { AgentMessage, StreamFn, ThinkingLevel } from "@earendil-works/pi-agent-core";
9import {
10 contentText,
11 getCurrentSystemMessage,
12 normalizeContext,
13 type RetryCallbacks,
14 type RetryPolicy,
15 retryAssistantCall,
16 uuidv7,
17} from "@earendil-works/pi-ai";
18import type {
19 AssistantMessage,
20 Model,
21 SimpleStreamOptions,
22 SystemMessage,
23 TranscriptContext,
24 Usage,
25} from "@earendil-works/pi-ai/compat";
26import { completeSimple } from "@earendil-works/pi-ai/compat";
27import { convertToLlm } from "../messages.ts";
28import {
29 buildSessionProjection,
30 type CompactionEntry,
31 type ProjectedSessionEntry,
32 type SessionEntry,
33 type SessionProjection,
34 sessionEntryToContextMessages,
35} from "../session-manager.ts";
36import { combineUsage } from "../usage-totals.ts";
37import {
38 computeFileLists,
39 createFileOps,
40 extractFileOpsFromMessage,
41 type FileOperations,
42 formatFileOperations,
43 SUMMARIZATION_SYSTEM_PROMPT,
44 serializeConversation,
45} from "./utils.ts";
46
47// ============================================================================
48// File Operation Tracking
49// ============================================================================
50
51/** Details stored in CompactionEntry.details for file tracking */
52export interface CompactionDetails {
53 readFiles: string[];
54 modifiedFiles: string[];
55}
56
57/**
58 * Extract file operations from messages and previous compaction entries.
59 */
60function extractFileOperations(
61 messages: AgentMessage[],
62 entries: SessionEntry[],
63 prevCompactionIndex: number,
64): FileOperations {
65 const fileOps = createFileOps();
66
67 // Collect from previous compaction's details (if pi-generated)
68 if (prevCompactionIndex >= 0) {
69 const prevCompaction = entries[prevCompactionIndex] as CompactionEntry;
70 if (!prevCompaction.fromHook && prevCompaction.details) {
71 // fromHook field kept for session file compatibility
72 const details = prevCompaction.details as CompactionDetails;
73 if (Array.isArray(details.readFiles)) {
74 for (const f of details.readFiles) fileOps.read.add(f);
75 }
76 if (Array.isArray(details.modifiedFiles)) {
77 for (const f of details.modifiedFiles) fileOps.edited.add(f);
78 }
79 }
80 }
81
82 // Extract from tool calls in messages
83 for (const msg of messages) {
84 extractFileOpsFromMessage(msg, fileOps);
85 }
86
87 return fileOps;
88}
89
90// ============================================================================
91// Message Extraction
92// ============================================================================
93
94/**
95 * Extract AgentMessage from an entry if it produces one.
96 * Returns undefined for entries that don't contribute to LLM context.
97 */
98function getMessagesFromProjectedEntryForCompaction(entry: ProjectedSessionEntry): AgentMessage[] {
99 if (entry.sourceEntry.type === "compaction") return [];
100 // System messages are prompt state, not conversation; the compaction entry carries their replay.
101 return entry.messages.filter((message) => message.role !== "system");
102}
103
104/** Result from compact() - SessionManager adds uuid/parentUuid when saving */
105export interface CompactionResult<T = unknown> {
106 summary: string;
107 firstKeptEntryId: string;
108 tokensBefore: number;
109 estimatedTokensAfter?: number;
110 /** Usage from the LLM call(s) that generated this summary, if available */
111 usage?: Usage;
112 /** Extension-specific data (e.g., ArtifactIndex, version markers for structured compaction) */
113 details?: T;
114}
115
116// ============================================================================
117// Types
118// ============================================================================
119
120export interface CompactionSettings {
121 enabled: boolean;
122 reserveTokens: number;
123 keepRecentTokens: number;
124}
125
126export const DEFAULT_COMPACTION_SETTINGS: CompactionSettings = {
127 enabled: true,
128 reserveTokens: 16384,
129 keepRecentTokens: 20000,
130};
131
132// ============================================================================
133// Token calculation
134// ============================================================================
135
136/**
137 * Calculate total context tokens from usage.
138 * Uses the native totalTokens field when available, falls back to computing from components.
139 */
140export function calculateContextTokens(usage: Usage): number {
141 return usage.totalTokens || usage.input + usage.output + usage.cacheRead + usage.cacheWrite;
142}
143
144/**
145 * Get usage from an assistant message if available.
146 * Skips aborted, error, and all-zero usage messages as they don't have valid usage data.
147 */
148function getAssistantUsage(msg: AgentMessage): Usage | undefined {
149 if (msg.role === "assistant" && "usage" in msg) {
150 const assistantMsg = msg as AssistantMessage;
151 if (
152 assistantMsg.stopReason !== "aborted" &&
153 assistantMsg.stopReason !== "error" &&
154 assistantMsg.usage &&
155 calculateContextTokens(assistantMsg.usage) > 0
156 ) {
157 return assistantMsg.usage;
158 }
159 }
160 return undefined;
161}
162
163/**
164 * Find the last valid assistant message usage from session entries.
165 */
166export function getLastAssistantUsage(entries: SessionEntry[]): Usage | undefined {
167 for (let i = entries.length - 1; i >= 0; i--) {
168 const entry = entries[i];
169 if (entry.type === "message") {
170 const usage = getAssistantUsage(entry.message);
171 if (usage) return usage;
172 }
173 }
174 return undefined;
175}
176
177export interface ContextUsageEstimate {
178 tokens: number;
179 usageTokens: number;
180 trailingTokens: number;
181 lastUsageIndex: number | null;
182}
183
184function getLastAssistantUsageInfo(messages: AgentMessage[]): { usage: Usage; index: number } | undefined {
185 for (let i = messages.length - 1; i >= 0; i--) {
186 const usage = getAssistantUsage(messages[i]);
187 if (usage) return { usage, index: i };
188 }
189 return undefined;
190}
191
192/**
193 * Estimate context tokens from messages, using the last assistant usage when available.
194 * If there are messages after the last usage, estimate their tokens with estimateTokens.
195 */
196export function estimateContextTokens(messages: AgentMessage[]): ContextUsageEstimate {
197 const usageInfo = getLastAssistantUsageInfo(messages);
198
199 if (!usageInfo) {
200 let estimated = 0;
201 for (const message of messages) {
202 estimated += estimateTokens(message);
203 }
204 return {
205 tokens: estimated,
206 usageTokens: 0,
207 trailingTokens: estimated,
208 lastUsageIndex: null,
209 };
210 }
211
212 const usageTokens = calculateContextTokens(usageInfo.usage);
213 let trailingTokens = 0;
214 for (let i = usageInfo.index + 1; i < messages.length; i++) {
215 trailingTokens += estimateTokens(messages[i]);
216 }
217
218 return {
219 tokens: usageTokens + trailingTokens,
220 usageTokens,
221 trailingTokens,
222 lastUsageIndex: usageInfo.index,
223 };
224}
225
226/** Estimate projected context without trusting usage captured before a later edit or compaction. */
227export function estimateProjectedContextTokens(
228 projection: SessionProjection,
229 branchEntries: SessionEntry[],
230): ContextUsageEstimate {
231 const estimate = estimateContextTokens(projection.messages);
232 if (estimate.lastUsageIndex !== null) {
233 let projectedMessageIndex = 0;
234 let usageEntryId: string | undefined;
235 for (const entry of projection.entries) {
236 const nextMessageIndex = projectedMessageIndex + entry.messages.length;
237 if (estimate.lastUsageIndex < nextMessageIndex) {
238 usageEntryId = entry.sourceEntry.id;
239 break;
240 }
241 projectedMessageIndex = nextMessageIndex;
242 }
243
244 const usageEntryIndex = usageEntryId ? branchEntries.findIndex((entry) => entry.id === usageEntryId) : -1;
245 let latestInvalidatingEntryIndex = -1;
246 for (let i = branchEntries.length - 1; i >= 0; i--) {
247 const entry = branchEntries[i];
248 if (entry.type === "context_edit" || entry.type === "compaction") {
249 latestInvalidatingEntryIndex = i;
250 break;
251 }
252 }
253 if (usageEntryIndex > latestInvalidatingEntryIndex) return estimate;
254 }
255
256 const currentSystem = getCurrentSystemMessage(projection.messages);
257 let tokens = currentSystem ? estimateTokens(currentSystem) : 0;
258 for (const message of projection.messages) {
259 if (message.role !== "system") tokens += estimateTokens(message);
260 }
261 return { tokens, usageTokens: 0, trailingTokens: tokens, lastUsageIndex: null };
262}
263
264/**
265 * Check if compaction should trigger based on context usage.
266 */
267export function shouldCompact(contextTokens: number, contextWindow: number, settings: CompactionSettings): boolean {
268 if (!settings.enabled) return false;
269 return contextTokens > contextWindow - settings.reserveTokens;
270}
271
272// ============================================================================
273// Cut point detection
274// ============================================================================
275
276const ESTIMATED_IMAGE_CHARS = 4800;
277
278function estimateTextAndImageContentChars(content: string | Array<{ type: string; text?: string }>): number {
279 if (typeof content === "string") {
280 return content.length;
281 }
282
283 let chars = 0;
284 for (const block of content) {
285 if (block.type === "text" && block.text) {
286 chars += block.text.length;
287 } else if (block.type === "image") {
288 chars += ESTIMATED_IMAGE_CHARS;
289 }
290 }
291 return chars;
292}
293
294/**
295 * Estimate token count for a message using chars/4 heuristic.
296 * This is conservative (overestimates tokens).
297 */
298export function estimateTokens(message: AgentMessage): number {
299 let chars = 0;
300
301 switch (message.role) {
302 case "system": {
303 const system = message as SystemMessage;
304 chars = estimateTextAndImageContentChars(system.content);
305 if (system.sections) {
306 for (const section of Object.values(system.sections)) {
307 if (section) chars += section.length;
308 }
309 }
310 if (system.toolsAdded) chars += JSON.stringify(system.toolsAdded).length;
311 return Math.ceil(chars / 4);
312 }
313 case "user": {
314 chars = estimateTextAndImageContentChars(
315 (message as { content: string | Array<{ type: string; text?: string }> }).content,
316 );
317 return Math.ceil(chars / 4);
318 }
319 case "assistant": {
320 const assistant = message as AssistantMessage;
321 for (const block of assistant.content) {
322 if (block.type === "text") {
323 chars += block.text.length;
324 } else if (block.type === "thinking") {
325 chars += block.thinking.length;
326 } else if (block.type === "toolCall") {
327 chars += block.name.length + JSON.stringify(block.arguments).length;
328 }
329 }
330 return Math.ceil(chars / 4);
331 }
332 case "custom":
333 case "toolResult": {
334 chars = estimateTextAndImageContentChars(message.content);
335 return Math.ceil(chars / 4);
336 }
337 case "bashExecution": {
338 chars = message.command.length + message.output.length;
339 return Math.ceil(chars / 4);
340 }
341 case "branchSummary":
342 case "compactionSummary": {
343 chars = message.summary.length;
344 return Math.ceil(chars / 4);
345 }
346 }
347
348 return 0;
349}
350
351function isCutPointMessage(message: AgentMessage): boolean {
352 switch (message.role) {
353 case "user":
354 case "assistant":
355 case "bashExecution":
356 case "custom":
357 case "branchSummary":
358 case "compactionSummary":
359 return true;
360 case "toolResult":
361 return false;
362 }
363 return false;
364}
365
366function isTurnStartMessage(message: AgentMessage): boolean {
367 switch (message.role) {
368 case "user":
369 case "bashExecution":
370 case "custom":
371 case "branchSummary":
372 case "compactionSummary":
373 return true;
374 case "assistant":
375 case "toolResult":
376 return false;
377 }
378 return false;
379}
380
381function isTurnStartEntry(entry: SessionEntry): boolean {
382 if (entry.type === "compaction") {
383 return false;
384 }
385 return sessionEntryToContextMessages(entry).some(isTurnStartMessage);
386}
387
388/**
389 * Find valid cut points: indices of context-visible user-like or assistant messages.
390 * Never cut at tool results (they must follow their tool call).
391 * When we cut at an assistant message with tool calls, its tool results follow it
392 * and will be kept.
393 */
394function findValidCutPoints(entries: SessionEntry[], startIndex: number, endIndex: number): number[] {
395 const cutPoints: number[] = [];
396 for (let i = startIndex; i < endIndex; i++) {
397 const entry = entries[i];
398 if (entry.type === "compaction") {
399 continue;
400 }
401 if (sessionEntryToContextMessages(entry).some(isCutPointMessage)) {
402 cutPoints.push(i);
403 }
404 }
405 return cutPoints;
406}
407
408/**
409 * Find the context-visible user-role message that starts the turn containing the given entry index.
410 * Returns -1 if no turn start found before the index.
411 */
412export function findTurnStartIndex(entries: SessionEntry[], entryIndex: number, startIndex: number): number {
413 for (let i = entryIndex; i >= startIndex; i--) {
414 if (isTurnStartEntry(entries[i])) {
415 return i;
416 }
417 }
418 return -1;
419}
420
421export interface CutPointResult {
422 /** Index of first entry to keep */
423 firstKeptEntryIndex: number;
424 /** Index of user message that starts the turn being split, or -1 if not splitting */
425 turnStartIndex: number;
426 /** Whether this cut splits a turn (cut point is not a user message) */
427 isSplitTurn: boolean;
428}
429
430/**
431 * Find the cut point in session entries that keeps approximately `keepRecentTokens`.
432 *
433 * Algorithm: Walk backwards from newest, accumulating estimated message sizes.
434 * Stop when we've accumulated >= keepRecentTokens. Cut at that point.
435 *
436 * Can cut at user OR assistant messages (never tool results). When cutting at an
437 * assistant message with tool calls, its tool results come after and will be kept.
438 *
439 * Returns CutPointResult with:
440 * - firstKeptEntryIndex: the entry index to start keeping from
441 * - turnStartIndex: if cutting mid-turn, the user message that started that turn
442 * - isSplitTurn: whether we're cutting in the middle of a turn
443 *
444 * Only considers entries between `startIndex` and `endIndex` (exclusive).
445 */
446export function findCutPoint(
447 entries: SessionEntry[],
448 startIndex: number,
449 endIndex: number,
450 keepRecentTokens: number,
451): CutPointResult {
452 const cutPoints = findValidCutPoints(entries, startIndex, endIndex);
453
454 if (cutPoints.length === 0) {
455 return { firstKeptEntryIndex: startIndex, turnStartIndex: -1, isSplitTurn: false };
456 }
457
458 // Walk backwards from newest, accumulating estimated message sizes
459 let accumulatedTokens = 0;
460 let cutIndex = cutPoints[0]; // Default: keep from first message (not header)
461
462 for (let i = endIndex - 1; i >= startIndex; i--) {
463 const entry = entries[i];
464 const messageTokens = sessionEntryToContextMessages(entry).reduce(
465 (sum, message) => sum + estimateTokens(message),
466 0,
467 );
468 if (messageTokens === 0) continue;
469 accumulatedTokens += messageTokens;
470
471 // Check if we've exceeded the budget
472 if (accumulatedTokens >= keepRecentTokens) {
473 // Prefer the closest valid cut point at or after this entry. If trailing
474 // tool results exceed the budget by themselves, keep their preceding
475 // assistant tool call instead of falling back to the first message.
476 cutIndex = cutPoints.find((candidate) => candidate >= i) ?? cutPoints[cutPoints.length - 1];
477 break;
478 }
479 }
480
481 // Scan backwards from cutIndex to include adjacent metadata entries that do not affect context.
482 while (cutIndex > startIndex) {
483 const prevEntry = entries[cutIndex - 1];
484 // Stop at compaction boundaries or context-visible entries.
485 if (prevEntry.type === "compaction" || sessionEntryToContextMessages(prevEntry).length > 0) {
486 break;
487 }
488 cutIndex--;
489 }
490
491 // Determine if this is a split turn
492 const cutEntry = entries[cutIndex];
493 const startsTurn = isTurnStartEntry(cutEntry);
494 const turnStartIndex = startsTurn ? -1 : findTurnStartIndex(entries, cutIndex, startIndex);
495
496 return {
497 firstKeptEntryIndex: cutIndex,
498 turnStartIndex,
499 isSplitTurn: !startsTurn && turnStartIndex !== -1,
500 };
501}
502
503// ============================================================================
504// Summarization
505// ============================================================================
506
507const SUMMARIZATION_PROMPT = `The messages above are a conversation to summarize. Create a structured context checkpoint summary that another LLM will use to continue the work.
508
509Use this EXACT format:
510
511## Goal
512[What is the user trying to accomplish? Can be multiple items if the session covers different tasks.]
513
514## Constraints & Preferences
515- [Any constraints, preferences, or requirements mentioned by user]
516- [Or "(none)" if none were mentioned]
517
518## Progress
519### Done
520- [x] [Completed tasks/changes]
521
522### In Progress
523- [ ] [Current work]
524
525### Blocked
526- [Issues preventing progress, if any]
527
528## Key Decisions
529- **[Decision]**: [Brief rationale]
530
531## Next Steps
5321. [Ordered list of what should happen next]
533
534## Critical Context
535- [Any data, examples, or references needed to continue]
536- [Or "(none)" if not applicable]
537
538Keep each section concise. Preserve exact file paths, function names, and error messages.`;
539
540const UPDATE_SUMMARIZATION_INSTRUCTIONS = `Update the existing structured summary with new information. RULES:
541- PRESERVE all existing information from the previous summary
542- ADD new progress, decisions, and context from the new messages
543- UPDATE the Progress section: move items from "In Progress" to "Done" when completed
544- UPDATE "Next Steps" based on what was accomplished
545- PRESERVE exact file paths, function names, and error messages
546- If something is no longer relevant, you may remove it
547
548Use this EXACT format:
549
550## Goal
551[Preserve existing goals, add new ones if the task expanded]
552
553## Constraints & Preferences
554- [Preserve existing, add new ones discovered]
555
556## Progress
557### Done
558- [x] [Include previously done items AND newly completed items]
559
560### In Progress
561- [ ] [Current work - update based on progress]
562
563### Blocked
564- [Current blockers - remove if resolved]
565
566## Key Decisions
567- **[Decision]**: [Brief rationale] (preserve all previous, add new)
568
569## Next Steps
5701. [Update based on current state]
571
572## Critical Context
573- [Preserve important context, add new if needed]
574
575Keep each section concise. Preserve exact file paths, function names, and error messages.`;
576
577const UPDATE_SUMMARIZATION_PROMPT = `The messages above are NEW conversation messages to incorporate into the existing summary provided in <previous-summary> tags.
578
579${UPDATE_SUMMARIZATION_INSTRUCTIONS}`;
580
581/**
582 * Returns an error message when a summarization response cannot safely be persisted.
583 * A length stop contains partial text and must not become a session checkpoint.
584 */
585export function getSummarizationFailure(response: AssistantMessage, label: string): string | undefined {
586 if (response.stopReason === "error") {
587 return `${label} failed: ${response.errorMessage || "Unknown error"}`;
588 }
589 if (response.stopReason === "length") {
590 return `${label} failed: generation hit the token cap and the summary is incomplete`;
591 }
592 return undefined;
593}
594
595function createSummarizationOptions(
596 model: Model<any>,
597 maxTokens: number,
598 apiKey: string | undefined,
599 headers: Record<string, string> | undefined,
600 env: Record<string, string> | undefined,
601 signal: AbortSignal | undefined,
602 thinkingLevel: ThinkingLevel | undefined,
603 sessionId: string | undefined,
604): SimpleStreamOptions {
605 const options: SimpleStreamOptions = { maxTokens, signal, apiKey, headers, env, sessionId };
606 if (model.reasoning && thinkingLevel && thinkingLevel !== "off") {
607 options.reasoning = thinkingLevel;
608 }
609 return options;
610}
611
612/**
613 * Shared choke point for every compaction/branch-summary summarization call. Wraps the
614 * single LLM call in {@link retryAssistantCall} so transient stream drops (e.g.
615 * `terminated`, socket close) honor the configured retry policy instead of failing
616 * the whole compaction on the first attempt. Deterministic errors and aborts return
617 * immediately (see {@link retryAssistantCall}).
618 */
619export async function completeSummarization(
620 model: Model<any>,
621 context: TranscriptContext,
622 options: SimpleStreamOptions,
623 streamFn?: StreamFn,
624 retry?: RetryPolicy,
625 callbacks?: RetryCallbacks,
626): Promise<AssistantMessage> {
627 // Avoid cache writes for one-off summaries. Reuse caller-supplied routing when available;
628 // callers without a session ID, including branch summaries, receive a fresh routing ID.
629 const requestOptions: SimpleStreamOptions = {
630 ...options,
631 cacheRetention: "none",
632 sessionId: options.sessionId ?? uuidv7(),
633 };
634 const produce = async (): Promise<AssistantMessage> =>
635 streamFn
636 ? (await streamFn(model, context, requestOptions)).result()
637 : completeSimple(model, context, requestOptions);
638 return retryAssistantCall(produce, retry, requestOptions.signal, callbacks);
639}
640
641/**
642 * Generate a summary of the conversation using the LLM.
643 * If previousSummary is provided, uses the update prompt to merge.
644 */
645export async function generateSummary(
646 currentMessages: AgentMessage[],
647 model: Model<any>,
648 reserveTokens: number,
649 apiKey: string | undefined,
650 headers?: Record<string, string>,
651 signal?: AbortSignal,
652 customInstructions?: string,
653 previousSummary?: string,
654 thinkingLevel?: ThinkingLevel,
655 streamFn?: StreamFn,
656 env?: Record<string, string>,
657 retry?: RetryPolicy,
658 callbacks?: RetryCallbacks,
659 sessionId?: string,
660): Promise<string> {
661 return (
662 await generateSummaryWithUsage(
663 currentMessages,
664 model,
665 reserveTokens,
666 apiKey,
667 headers,
668 signal,
669 customInstructions,
670 previousSummary,
671 thinkingLevel,
672 streamFn,
673 env,
674 retry,
675 callbacks,
676 sessionId,
677 )
678 ).text;
679}
680
681/** Build the provider context for a standalone summary request. */
682function buildSummarizationContext(promptText: string): TranscriptContext {
683 return normalizeContext({
684 systemPrompt: SUMMARIZATION_SYSTEM_PROMPT,
685 messages: [
686 {
687 role: "user",
688 content: [{ type: "text", text: promptText }],
689 timestamp: Date.now(),
690 },
691 ],
692 });
693}
694
695/** Generate or update a conversation summary and return its provider usage. */
696export async function generateSummaryWithUsage(
697 currentMessages: AgentMessage[],
698 model: Model<any>,
699 reserveTokens: number,
700 apiKey: string | undefined,
701 headers?: Record<string, string>,
702 signal?: AbortSignal,
703 customInstructions?: string,
704 previousSummary?: string,
705 thinkingLevel?: ThinkingLevel,
706 streamFn?: StreamFn,
707 env?: Record<string, string>,
708 retry?: RetryPolicy,
709 callbacks?: RetryCallbacks,
710 sessionId?: string,
711): Promise<{ text: string; usage: Usage }> {
712 const maxTokens = Math.min(
713 Math.floor(0.8 * reserveTokens),
714 model.maxTokens > 0 ? model.maxTokens : Number.POSITIVE_INFINITY,
715 );
716
717 // Use update prompt if we have a previous summary, otherwise initial prompt
718 let basePrompt = previousSummary ? UPDATE_SUMMARIZATION_PROMPT : SUMMARIZATION_PROMPT;
719 if (customInstructions) {
720 basePrompt = `${basePrompt}\n\nAdditional focus: ${customInstructions}`;
721 }
722
723 // Serialize conversation to text so model doesn't try to continue it
724 // Convert to LLM messages first (handles custom types like bashExecution, custom, etc.)
725 const llmMessages = convertToLlm(currentMessages);
726 const conversationText = serializeConversation(llmMessages);
727
728 // Build the prompt with conversation wrapped in tags
729 let promptText = `<conversation>\n${conversationText}\n</conversation>\n\n`;
730 if (previousSummary) {
731 promptText += `<previous-summary>\n${previousSummary}\n</previous-summary>\n\n`;
732 }
733 promptText += basePrompt;
734
735 const completionOptions = createSummarizationOptions(
736 model,
737 maxTokens,
738 apiKey,
739 headers,
740 env,
741 signal,
742 thinkingLevel,
743 sessionId,
744 );
745
746 const response = await completeSummarization(
747 model,
748 buildSummarizationContext(promptText),
749 completionOptions,
750 streamFn,
751 retry,
752 callbacks,
753 );
754
755 const failure = getSummarizationFailure(response, "Summarization");
756 if (failure) {
757 throw new Error(failure);
758 }
759 if (response.content.some((block) => block.type === "toolCall")) {
760 throw new Error("Summarization attempted to call a tool");
761 }
762
763 const textContent = contentText(response.content);
764
765 return { text: textContent, usage: response.usage };
766}
767
768// ============================================================================
769// Compaction Preparation (for extensions)
770// ============================================================================
771
772export interface CompactionPreparation {
773 /** UUID of first entry to keep */
774 firstKeptEntryId: string;
775 /** Messages that will be summarized and discarded */
776 messagesToSummarize: AgentMessage[];
777 /** Messages that will be turned into turn prefix summary (if splitting) */
778 turnPrefixMessages: AgentMessage[];
779 /** Whether this is a split turn (cut point in middle of turn) */
780 isSplitTurn: boolean;
781 tokensBefore: number;
782 /** Summary from previous compaction, for iterative update */
783 previousSummary?: string;
784 /** File operations extracted from messagesToSummarize */
785 fileOps: FileOperations;
786 /** Compaction settions from settings.jsonl */
787 settings: CompactionSettings;
788}
789
790function isProjectedTurnStart(entry: ProjectedSessionEntry): boolean {
791 if (entry.sourceEntry.type === "compaction") return false;
792 return entry.messages.some(isTurnStartMessage);
793}
794
795function findProjectedTurnStartIndex(entries: ProjectedSessionEntry[], entryIndex: number, startIndex: number): number {
796 for (let i = entryIndex; i >= startIndex; i--) {
797 if (isProjectedTurnStart(entries[i])) return i;
798 }
799 return -1;
800}
801
802function findProjectedCutPoint(
803 entries: ProjectedSessionEntry[],
804 startIndex: number,
805 endIndex: number,
806 keepRecentTokens: number,
807): CutPointResult {
808 const cutPoints: number[] = [];
809 for (let i = startIndex; i < endIndex; i++) {
810 const entry = entries[i];
811 if (entry.sourceEntry.type !== "compaction" && entry.messages.some(isCutPointMessage)) cutPoints.push(i);
812 }
813 if (cutPoints.length === 0) {
814 return { firstKeptEntryIndex: startIndex, turnStartIndex: -1, isSplitTurn: false };
815 }
816
817 let accumulatedTokens = 0;
818 let exceededBudget = false;
819 let cutIndex = cutPoints[0];
820 for (let i = endIndex - 1; i >= startIndex; i--) {
821 const messageTokens = entries[i].messages.reduce((sum, message) => sum + estimateTokens(message), 0);
822 if (messageTokens === 0) continue;
823 accumulatedTokens += messageTokens;
824 if (accumulatedTokens >= keepRecentTokens) {
825 exceededBudget = true;
826 cutIndex = cutPoints.find((candidate) => candidate >= i) ?? cutPoints[cutPoints.length - 1];
827 break;
828 }
829 }
830
831 // A recovery attempt and its omission edits are context-invisible after the last
832 // visible input. Advance only for a closed suffix containing an omitted assistant
833 // attempt; arbitrary metadata must not move the cut past unsent input.
834 const suffix = entries.slice(cutIndex + 1, endIndex);
835 const isIntrinsicallyVisible = (entry: ProjectedSessionEntry): boolean =>
836 entry.sourceEntry.type !== "context_edit" && sessionEntryToContextMessages(entry.sourceEntry).length > 0;
837 const isOmitted = (entry: ProjectedSessionEntry): boolean =>
838 isIntrinsicallyVisible(entry) && entry.messages.length === 0;
839 const omittedSuffixIds = new Set(suffix.filter(isOmitted).map((entry) => entry.sourceEntry.id));
840 const hasExternalReplacement = suffix.some(
841 (entry) =>
842 entry.sourceEntry.type === "context_edit" &&
843 entry.sourceEntry.replacement !== null &&
844 !omittedSuffixIds.has(entry.sourceEntry.targetId),
845 );
846 const isRecoveryOmissionSuffix =
847 exceededBudget &&
848 !hasExternalReplacement &&
849 suffix.some(
850 (entry) =>
851 entry.sourceEntry.type === "message" && entry.sourceEntry.message.role === "assistant" && isOmitted(entry),
852 ) &&
853 suffix.every(
854 (entry) => entry.sourceEntry.type !== "compaction" && (!isIntrinsicallyVisible(entry) || isOmitted(entry)),
855 );
856 if (isRecoveryOmissionSuffix) cutIndex++;
857
858 while (cutIndex > startIndex) {
859 const previous = entries[cutIndex - 1];
860 if (previous.sourceEntry.type === "compaction" || previous.messages.length > 0) break;
861 cutIndex--;
862 }
863 const startsTurn = isProjectedTurnStart(entries[cutIndex]);
864 const turnStartIndex = startsTurn ? -1 : findProjectedTurnStartIndex(entries, cutIndex, startIndex);
865 return {
866 firstKeptEntryIndex: cutIndex,
867 turnStartIndex,
868 isSplitTurn: !startsTurn && turnStartIndex !== -1,
869 };
870}
871
872export function prepareCompaction(
873 pathEntries: SessionEntry[],
874 settings: CompactionSettings,
875): CompactionPreparation | undefined {
876 if (pathEntries.length > 0 && pathEntries[pathEntries.length - 1].type === "compaction") {
877 return undefined;
878 }
879
880 const projection = buildSessionProjection(pathEntries);
881 const projectedEntries = projection.entries;
882 const sourceEntries = projectedEntries.map((entry) => entry.sourceEntry);
883 // The newest compaction is projected first. Older compaction entries can still
884 // occur in its retained raw range, but their projected contribution is empty.
885 const prevCompactionIndex = projectedEntries.findIndex(
886 (entry) => entry.sourceEntry.type === "compaction" && entry.messages.length > 0,
887 );
888
889 let previousSummary: string | undefined;
890 let boundaryStart = 0;
891 if (prevCompactionIndex >= 0) {
892 previousSummary = (projectedEntries[prevCompactionIndex].sourceEntry as CompactionEntry).summary;
893 // The canonical projection has already selected the previous compaction's retained tail.
894 boundaryStart = prevCompactionIndex + 1;
895 }
896 const boundaryEnd = projectedEntries.length;
897 const tokensBefore = estimateProjectedContextTokens(projection, pathEntries).tokens;
898 const cutPoint = findProjectedCutPoint(projectedEntries, boundaryStart, boundaryEnd, settings.keepRecentTokens);
899
900 const firstKeptEntry = projectedEntries[cutPoint.firstKeptEntryIndex]?.sourceEntry;
901 if (!firstKeptEntry?.id) return undefined;
902 const firstKeptEntryId = firstKeptEntry.id;
903 const historyEnd = cutPoint.isSplitTurn ? cutPoint.turnStartIndex : cutPoint.firstKeptEntryIndex;
904
905 const messagesToSummarize = projectedEntries
906 .slice(boundaryStart, historyEnd)
907 .flatMap(getMessagesFromProjectedEntryForCompaction);
908 const turnPrefixMessages = cutPoint.isSplitTurn
909 ? projectedEntries
910 .slice(cutPoint.turnStartIndex, cutPoint.firstKeptEntryIndex)
911 .flatMap(getMessagesFromProjectedEntryForCompaction)
912 : [];
913
914 if (messagesToSummarize.length === 0 && turnPrefixMessages.length === 0) return undefined;
915
916 // Extract file operations from edited model-visible messages and the previous compaction.
917 const fileOps = extractFileOperations(messagesToSummarize, sourceEntries, prevCompactionIndex);
918
919 // Also extract file ops from turn prefix if splitting
920 if (cutPoint.isSplitTurn) {
921 for (const msg of turnPrefixMessages) {
922 extractFileOpsFromMessage(msg, fileOps);
923 }
924 }
925
926 return {
927 firstKeptEntryId,
928 messagesToSummarize,
929 turnPrefixMessages,
930 isSplitTurn: cutPoint.isSplitTurn,
931 tokensBefore,
932 previousSummary,
933 fileOps,
934 settings,
935 };
936}
937
938// ============================================================================
939// Main compaction function
940// ============================================================================
941
942const TURN_PREFIX_SUMMARIZATION_PROMPT = `The messages above are earlier context from an ongoing conversation. Later messages are stored separately and do not need to be reconstructed.
943
944Create a concise checkpoint of the user's request and the progress shown above. This checkpoint will be placed before the later messages so the conversation can continue with the necessary context.
945
946## Original Request
947[What did the user ask for?]
948
949## Progress So Far
950- [Key decisions and work completed in these messages]
951
952## Context Needed to Continue
953- [Information from these messages needed to understand the later work]
954
955Only summarize information explicitly present above. Do not infer or recreate later messages.`;
956
957/**
958 * Generate summaries for compaction using prepared data.
959 * Returns CompactionResult - SessionManager adds uuid/parentUuid when saving.
960 *
961 * @param preparation - Pre-calculated preparation from prepareCompaction()
962 * @param customInstructions - Optional custom focus for the summary
963 * @param sessionId - Optional routing session ID forwarded without enabling prompt caching
964 */
965export async function compact(
966 preparation: CompactionPreparation,
967 model: Model<any>,
968 apiKey: string | undefined,
969 headers?: Record<string, string>,
970 customInstructions?: string,
971 signal?: AbortSignal,
972 thinkingLevel?: ThinkingLevel,
973 streamFn?: StreamFn,
974 env?: Record<string, string>,
975 retry?: RetryPolicy,
976 callbacks?: RetryCallbacks,
977 sessionId?: string,
978): Promise<CompactionResult> {
979 const {
980 firstKeptEntryId,
981 messagesToSummarize,
982 turnPrefixMessages,
983 isSplitTurn,
984 tokensBefore,
985 previousSummary,
986 fileOps,
987 settings,
988 } = preparation;
989
990 // Generate summaries and merge into one
991 let summary: string;
992 let summaryUsage: Usage;
993
994 if (isSplitTurn && turnPrefixMessages.length > 0) {
995 let historyText = previousSummary ?? "No prior history.";
996 let historyUsage: Usage | undefined;
997 if (messagesToSummarize.length > 0) {
998 const historyResult = await generateSummaryWithUsage(
999 messagesToSummarize,
1000 model,
1001 settings.reserveTokens,
1002 apiKey,
1003 headers,
1004 signal,
1005 customInstructions,
1006 previousSummary,
1007 thinkingLevel,
1008 streamFn,
1009 env,
1010 retry,
1011 callbacks,
1012 sessionId,
1013 );
1014 historyText = historyResult.text;
1015 historyUsage = historyResult.usage;
1016 }
1017 const turnPrefixResult = await generateTurnPrefixSummary(
1018 turnPrefixMessages,
1019 model,
1020 settings.reserveTokens,
1021 apiKey,
1022 headers,
1023 env,
1024 signal,
1025 thinkingLevel,
1026 streamFn,
1027 retry,
1028 callbacks,
1029 sessionId,
1030 );
1031 // Merge into single summary
1032 summary = `${historyText}\n\n---\n\n**Turn Context (split turn):**\n\n${turnPrefixResult.text}`;
1033 summaryUsage = historyUsage ? combineUsage(historyUsage, turnPrefixResult.usage) : turnPrefixResult.usage;
1034 } else {
1035 // Just generate history summary
1036 const result = await generateSummaryWithUsage(
1037 messagesToSummarize,
1038 model,
1039 settings.reserveTokens,
1040 apiKey,
1041 headers,
1042 signal,
1043 customInstructions,
1044 previousSummary,
1045 thinkingLevel,
1046 streamFn,
1047 env,
1048 retry,
1049 callbacks,
1050 sessionId,
1051 );
1052 summary = result.text;
1053 summaryUsage = result.usage;
1054 }
1055
1056 // Compute file lists and append to summary
1057 const { readFiles, modifiedFiles } = computeFileLists(fileOps);
1058 summary += formatFileOperations(readFiles, modifiedFiles);
1059
1060 if (!firstKeptEntryId) {
1061 throw new Error("First kept entry has no UUID - session may need migration");
1062 }
1063
1064 return {
1065 summary,
1066 firstKeptEntryId,
1067 tokensBefore,
1068 usage: summaryUsage,
1069 details: { readFiles, modifiedFiles } as CompactionDetails,
1070 };
1071}
1072
1073/**
1074 * Generate a summary for a turn prefix (when splitting a turn).
1075 */
1076async function generateTurnPrefixSummary(
1077 messages: AgentMessage[],
1078 model: Model<any>,
1079 reserveTokens: number,
1080 apiKey: string | undefined,
1081 headers?: Record<string, string>,
1082 env?: Record<string, string>,
1083 signal?: AbortSignal,
1084 thinkingLevel?: ThinkingLevel,
1085 streamFn?: StreamFn,
1086 retry?: RetryPolicy,
1087 callbacks?: RetryCallbacks,
1088 sessionId?: string,
1089): Promise<{ text: string; usage: Usage }> {
1090 const maxTokens = Math.min(
1091 Math.floor(0.5 * reserveTokens),
1092 model.maxTokens > 0 ? model.maxTokens : Number.POSITIVE_INFINITY,
1093 ); // Smaller budget for turn prefix
1094 const llmMessages = convertToLlm(messages);
1095 const conversationText = serializeConversation(llmMessages);
1096 const promptText = `# Conversation\n${conversationText}\n\n# Instructions\n${TURN_PREFIX_SUMMARIZATION_PROMPT}`;
1097
1098 const response = await completeSummarization(
1099 model,
1100 buildSummarizationContext(promptText),
1101 createSummarizationOptions(model, maxTokens, apiKey, headers, env, signal, thinkingLevel, sessionId),
1102 streamFn,
1103 retry,
1104 callbacks,
1105 );
1106
1107 const failure = getSummarizationFailure(response, "Turn prefix summarization");
1108 if (failure) {
1109 throw new Error(failure);
1110 }
1111 if (response.content.some((block) => block.type === "toolCall")) {
1112 throw new Error("Turn prefix summarization attempted to call a tool");
1113 }
1114
1115 return {
1116 text: contentText(response.content),
1117 usage: response.usage,
1118 };
1119}