fix(compaction): bound summarization input and stop retrying overflow
A session that crossed a provider boundary compacted 90 times in three days without ever succeeding: every attempt asked the summarizer to read the whole re-expanded span in one call (2.33M tokens on 08-15, 3.03M by 08-17, against a 1M cap), and every rejection was retried ten times. Three independent defects: 1. `generateSummary` serialized the entire span into one prompt with no budget check. It now plans windows that fit the summarizer's context and folds them with the update prompt that iterative compaction already uses, so a stranded boundary is recovered instead of rejected. A provider that rejects a window the catalog said would fit (claude-sonnet-4-5 advertises 1M but is beta-gated to 200k on OAuth credentials) halves what was actually sent and re-plans, because only the rejection knows the real cap. 2. `TRANSIENT_TRANSPORT_PATTERN` matched bare status codes, so the random id in the `raw-http-request=.../1787022540720-3o503gxo48bvb.json` pointer omp appends to its own errors classified a deterministic 400 as a transient 503. Statuses are now word-boundaried, matching AUTH_FAILURE_PATTERN. 3. Neither retry layer vetoed ContextOverflow, so one failure became up to 30 identical calls (10 outer x 3 oneshot). A oneshot replays a fixed prompt, so an input that does not fit never fits; both layers now fail fast to the next candidate. The boundary scan that decides which compaction entry a model can actually read is extracted as `findReadableCompactionIndex`, since the fold and `prepareCompaction` both need it. Verified by replaying the session that failed: 7,096 messages summarize in 3 calls with a largest prompt of 773,705 tokens under the real 1M cap, and in 15 calls with a largest prompt of 196,148 tokens under a simulated 200k cap.
This commit is contained in:
@@ -23,6 +23,7 @@ import {
|
||||
type Usage,
|
||||
withAuth,
|
||||
} from "@oh-my-pi/pi-ai";
|
||||
import type { Dialect } from "@oh-my-pi/pi-ai/dialect";
|
||||
import * as AIError from "@oh-my-pi/pi-ai/error";
|
||||
import { createOpenAICodexCompactionRequestContext } from "@oh-my-pi/pi-ai/providers/openai-codex-responses";
|
||||
import { convertTools } from "@oh-my-pi/pi-ai/providers/openai-responses";
|
||||
@@ -887,6 +888,73 @@ function createSnapcompactArchiveMigrationMessage(archiveText: string): Message
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* Fallback window for a model whose catalog entry carries no usable context
|
||||
* window; matches the smallest window any compaction-capable model ships with.
|
||||
*/
|
||||
const DEFAULT_SUMMARY_INPUT_WINDOW = 200_000;
|
||||
|
||||
/** Floor for one summarization window, so a tiny model still makes progress. */
|
||||
const MIN_SUMMARY_INPUT_TOKENS = 16_384;
|
||||
|
||||
/**
|
||||
* Usable conversation input for ONE summarization call: the summarizer's window
|
||||
* minus the summary it must emit, the previous summary it carries forward, and
|
||||
* prompt scaffolding. Providers tokenize differently from the local cl100k
|
||||
* estimate, so the window is discounted before the fixed reserves come off.
|
||||
*/
|
||||
function summaryInputBudgetTokens(model: Model, maxTokens: number): number {
|
||||
const window = model.contextWindow && model.contextWindow > 0 ? model.contextWindow : DEFAULT_SUMMARY_INPUT_WINDOW;
|
||||
// 0.8, not "window minus reserves": provider tokenizers disagree with the
|
||||
// local cl100k estimate by a few percent, and being wrong here is a hard
|
||||
// 400 on the one call that is supposed to rescue an oversized session.
|
||||
return Math.max(MIN_SUMMARY_INPUT_TOKENS, Math.floor(window * 0.8) - maxTokens - MAX_SUMMARY_TOKENS);
|
||||
}
|
||||
|
||||
/**
|
||||
* Clamp one serialized window to the budget. Only reachable when a SINGLE
|
||||
* message serializes above the budget (an oversized paste): the alternative is
|
||||
* a provider rejection that no retry can clear, which strands the session with
|
||||
* a full window forever.
|
||||
*/
|
||||
function clampConversationToBudget(text: string, budgetTokens: number, tokens: number): string {
|
||||
if (tokens <= budgetTokens) return text;
|
||||
const keep = Math.max(1024, Math.floor((text.length * budgetTokens * 0.95) / tokens));
|
||||
if (keep >= text.length) return text;
|
||||
return `${text.slice(0, keep)}\n\n[... ${text.length - keep} more characters truncated]`;
|
||||
}
|
||||
|
||||
/** One planned summarization call: its messages and the budget they were packed for. */
|
||||
interface SummaryWindow {
|
||||
messages: Message[];
|
||||
budgetTokens: number;
|
||||
/** Serialization reused from the fit check, so the common path serializes once. */
|
||||
text?: string;
|
||||
}
|
||||
|
||||
/**
|
||||
* Partition a conversation into windows that each fit `budgetTokens`, splitting
|
||||
* on message boundaries. Only called when the whole conversation does not fit —
|
||||
* the common single-window path never pays this per-message sizing pass.
|
||||
*/
|
||||
function planSummaryWindows(messages: Message[], dialect: Dialect | undefined, budgetTokens: number): Message[][] {
|
||||
const windows: Message[][] = [];
|
||||
let current: Message[] = [];
|
||||
let currentTokens = 0;
|
||||
for (const message of messages) {
|
||||
const tokens = countTokens(serializeConversationForSummary([message], dialect));
|
||||
if (currentTokens > 0 && currentTokens + tokens > budgetTokens) {
|
||||
windows.push(current);
|
||||
current = [];
|
||||
currentTokens = 0;
|
||||
}
|
||||
current.push(message);
|
||||
currentTokens += tokens;
|
||||
}
|
||||
if (current.length > 0) windows.push(current);
|
||||
return windows;
|
||||
}
|
||||
|
||||
export async function generateSummary(
|
||||
currentMessages: AgentMessage[],
|
||||
model: Model,
|
||||
@@ -899,6 +967,79 @@ export async function generateSummary(
|
||||
): Promise<string> {
|
||||
const maxTokens = Math.min(Math.floor(0.8 * reserveTokens), MAX_SUMMARY_TOKENS);
|
||||
|
||||
// Serialize conversation to text so model doesn't try to continue it
|
||||
// Convert to LLM messages first (handles custom app messages when caller provides a transformer).
|
||||
const llmMessages = (options?.convertToLlm ?? defaultConvertToLlm)(currentMessages);
|
||||
const dialect = preferredDialect(model.id);
|
||||
const wholeConversation = serializeConversationForSummary(llmMessages, dialect);
|
||||
const budgetTokens = summaryInputBudgetTokens(model, maxTokens);
|
||||
// A span that outgrew the summarizer's window is summarized as a fold: each
|
||||
// window updates the summary carried out of the previous one, which is the
|
||||
// same contract the update prompt already implements for iterative
|
||||
// compaction. The alternative is a hard provider rejection on a prompt no
|
||||
// retry can shrink — the state a cross-provider compaction boundary
|
||||
// (see `prepareCompaction`) puts a long session into. One window is the
|
||||
// common case and costs exactly the one call it always did.
|
||||
const pending: SummaryWindow[] =
|
||||
countTokens(wholeConversation) <= budgetTokens
|
||||
? [{ messages: llmMessages, budgetTokens, text: wholeConversation }]
|
||||
: planSummaryWindows(llmMessages, dialect, budgetTokens).map(messages => ({ messages, budgetTokens }));
|
||||
|
||||
let carriedSummary = previousSummary;
|
||||
while (pending.length > 0) {
|
||||
const window = pending[0];
|
||||
const text = window.text ?? serializeConversationForSummary(window.messages, dialect);
|
||||
const windowTokens = countTokens(text);
|
||||
try {
|
||||
carriedSummary = await summarizeConversationWindow(
|
||||
clampConversationToBudget(text, window.budgetTokens, windowTokens),
|
||||
carriedSummary,
|
||||
model,
|
||||
maxTokens,
|
||||
apiKey,
|
||||
signal,
|
||||
customInstructions,
|
||||
options,
|
||||
);
|
||||
} catch (error) {
|
||||
// The catalog window can overstate what the provider actually accepts:
|
||||
// `claude-sonnet-4-5` advertises 1M but is beta-gated to 200k on OAuth
|
||||
// credentials (see `anthropic.ts` — the 1M beta is never advertised).
|
||||
// Halve and re-plan rather than failing the whole compaction on a
|
||||
// window size only the provider can tell us is wrong.
|
||||
// Halve what was actually SENT, not the budget it was planned against:
|
||||
// the rejection proves the plan was fiction, so converging on the real
|
||||
// cap must not spend a call per level of an imaginary ladder.
|
||||
const halved = Math.floor(Math.min(window.budgetTokens, windowTokens) / 2);
|
||||
if (!AIError.is(AIError.classify(error), AIError.Flag.ContextOverflow) || halved < MIN_SUMMARY_INPUT_TOKENS) {
|
||||
throw error;
|
||||
}
|
||||
pending.splice(
|
||||
0,
|
||||
1,
|
||||
...planSummaryWindows(window.messages, dialect, halved).map(messages => ({
|
||||
messages,
|
||||
budgetTokens: halved,
|
||||
})),
|
||||
);
|
||||
continue;
|
||||
}
|
||||
pending.shift();
|
||||
}
|
||||
return carriedSummary ?? "";
|
||||
}
|
||||
|
||||
/** One summarization call over a single conversation window. */
|
||||
async function summarizeConversationWindow(
|
||||
conversationText: string,
|
||||
previousSummary: string | undefined,
|
||||
model: Model,
|
||||
maxTokens: number,
|
||||
apiKey: ApiKey,
|
||||
signal: AbortSignal | undefined,
|
||||
customInstructions: string | undefined,
|
||||
options: SummaryOptions | undefined,
|
||||
): Promise<string> {
|
||||
// Use update prompt if we have a previous summary, otherwise initial prompt
|
||||
let basePrompt = previousSummary ? UPDATE_SUMMARIZATION_PROMPT : SUMMARIZATION_PROMPT;
|
||||
if (options?.promptOverride) {
|
||||
@@ -908,11 +1049,6 @@ export async function generateSummary(
|
||||
basePrompt = `${basePrompt}\n\nAdditional focus: ${customInstructions}`;
|
||||
}
|
||||
|
||||
// Serialize conversation to text so model doesn't try to continue it
|
||||
// Convert to LLM messages first (handles custom app messages when caller provides a transformer).
|
||||
const llmMessages = (options?.convertToLlm ?? defaultConvertToLlm)(currentMessages);
|
||||
const conversationText = serializeConversationForSummary(llmMessages, preferredDialect(model.id));
|
||||
|
||||
// Build the prompt with conversation wrapped in tags
|
||||
let promptText = `<conversation>\n${conversationText}\n</conversation>\n\n`;
|
||||
if (previousSummary) {
|
||||
@@ -1247,6 +1383,32 @@ function remotePreserveReusable(
|
||||
return v2Ok || shouldUseOpenAiRemoteCompaction(activeModel);
|
||||
}
|
||||
|
||||
/**
|
||||
* Index of the newest compaction entry the active model can actually read, or
|
||||
* `-1` when none can.
|
||||
*
|
||||
* A provider-native remote compaction (V2 or V1) stores an opaque replay payload
|
||||
* and only a placeholder summary, so for any OTHER provider that entry
|
||||
* summarizes nothing and the history behind it is still live context. Callers
|
||||
* must therefore treat it as absent: `prepareCompaction` re-expands past it and
|
||||
* summarizes those messages locally, and the maintenance ops that use the
|
||||
* compaction boundary to skip "already summarized away" entries must not skip
|
||||
* entries that no summary covers.
|
||||
*/
|
||||
export function findReadableCompactionIndex(
|
||||
pathEntries: SessionEntry[],
|
||||
settings: CompactionSettings,
|
||||
activeModel?: Model,
|
||||
): number {
|
||||
for (let i = pathEntries.length - 1; i >= 0; i--) {
|
||||
if (pathEntries[i].type !== "compaction") continue;
|
||||
const entry = pathEntries[i] as CompactionEntry;
|
||||
if (activeModel && !remotePreserveReusable(entry.preserveData, activeModel, settings)) continue;
|
||||
return i;
|
||||
}
|
||||
return -1;
|
||||
}
|
||||
|
||||
export function prepareCompaction(
|
||||
pathEntries: SessionEntry[],
|
||||
settings: CompactionSettings,
|
||||
@@ -1256,21 +1418,7 @@ export function prepareCompaction(
|
||||
return undefined;
|
||||
}
|
||||
|
||||
let prevCompactionIndex = -1;
|
||||
for (let i = pathEntries.length - 1; i >= 0; i--) {
|
||||
if (pathEntries[i].type !== "compaction") continue;
|
||||
// Skip a prior remote compaction (V2 or V1) whose provider-native replay the
|
||||
// active model cannot read: its summary is only an opaque placeholder, so
|
||||
// re-expand its original messages and summarize them locally rather than
|
||||
// stranding that history. compact() still reuses the payload when the active
|
||||
// model can replay it (same provider, remote enabled).
|
||||
const entry = pathEntries[i] as CompactionEntry;
|
||||
if (activeModel && !remotePreserveReusable(entry.preserveData, activeModel, settings)) {
|
||||
continue;
|
||||
}
|
||||
prevCompactionIndex = i;
|
||||
break;
|
||||
}
|
||||
const prevCompactionIndex = findReadableCompactionIndex(pathEntries, settings, activeModel);
|
||||
const boundaryStart = prevCompactionIndex + 1;
|
||||
const boundaryEnd = pathEntries.length;
|
||||
|
||||
|
||||
@@ -0,0 +1,190 @@
|
||||
import { describe, expect, test, vi } from "bun:test";
|
||||
import type { AgentMessage } from "@oh-my-pi/pi-agent-core";
|
||||
import {
|
||||
DEFAULT_COMPACTION_SETTINGS,
|
||||
findReadableCompactionIndex,
|
||||
generateSummary,
|
||||
type SessionEntry,
|
||||
} from "@oh-my-pi/pi-agent-core/compaction";
|
||||
import type { AssistantMessage, Model } from "@oh-my-pi/pi-ai";
|
||||
import * as ai from "@oh-my-pi/pi-ai";
|
||||
import { getBundledModel } from "@oh-my-pi/pi-catalog/models";
|
||||
|
||||
function createAssistantMessage(text: string): AssistantMessage {
|
||||
return {
|
||||
role: "assistant",
|
||||
content: [{ type: "text", text }],
|
||||
timestamp: Date.now(),
|
||||
provider: "mock",
|
||||
model: "mock",
|
||||
api: "mock",
|
||||
usage: {
|
||||
input: 0,
|
||||
output: 0,
|
||||
cacheRead: 0,
|
||||
cacheWrite: 0,
|
||||
totalTokens: 0,
|
||||
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
|
||||
},
|
||||
stopReason: "stop",
|
||||
};
|
||||
}
|
||||
|
||||
function getModel(contextWindow: number): Model {
|
||||
const model = getBundledModel("anthropic", "claude-sonnet-4-5");
|
||||
if (!model) throw new Error("Expected built-in anthropic/claude-sonnet-4-5 to exist");
|
||||
return { ...model, contextWindow };
|
||||
}
|
||||
|
||||
/** ~4 chars per cl100k token, so each turn is roughly `tokens` tokens of input. */
|
||||
function turn(index: number, tokens: number): AgentMessage[] {
|
||||
return [
|
||||
{ role: "user", content: `turn ${index} ${"work ".repeat(tokens / 2)}`, timestamp: index },
|
||||
createAssistantMessage(`did ${index}`),
|
||||
];
|
||||
}
|
||||
|
||||
function promptTextOf(call: unknown[]): string {
|
||||
const context = call[1] as { messages: { content: { type: string; text: string }[] }[] };
|
||||
return context.messages[0].content[0].text;
|
||||
}
|
||||
|
||||
describe("summarization input budget", () => {
|
||||
test("summarizes a fitting conversation in one call", async () => {
|
||||
const spy = vi.spyOn(ai, "completeSimple").mockResolvedValue(createAssistantMessage("summary"));
|
||||
try {
|
||||
const summary = await generateSummary(turn(1, 200), getModel(200_000), 16_384, "test-key");
|
||||
expect(spy.mock.calls.length).toBe(1);
|
||||
expect(summary).toBe("summary");
|
||||
} finally {
|
||||
spy.mockRestore();
|
||||
}
|
||||
});
|
||||
|
||||
test("folds a conversation larger than the summarizer window across calls", async () => {
|
||||
let call = 0;
|
||||
const spy = vi
|
||||
.spyOn(ai, "completeSimple")
|
||||
.mockImplementation(async () => createAssistantMessage(`summary ${++call}`));
|
||||
try {
|
||||
// 40k-token window leaves ~4k of conversation budget after the summary
|
||||
// reserve, so ~48k tokens of conversation cannot be one prompt.
|
||||
const messages = Array.from({ length: 12 }, (_, i) => turn(i, 4_000)).flat();
|
||||
const summary = await generateSummary(messages, getModel(40_000), 16_384, "test-key");
|
||||
|
||||
expect(spy.mock.calls.length).toBeGreaterThan(1);
|
||||
expect(summary).toBe(`summary ${spy.mock.calls.length}`);
|
||||
|
||||
// Every window is inside the budget, and every window after the first
|
||||
// carries the summary of the ones before it.
|
||||
const prompts = spy.mock.calls.map(promptTextOf);
|
||||
for (const prompt of prompts) {
|
||||
expect(prompt.length).toBeLessThan(40_000 * 4);
|
||||
}
|
||||
expect(prompts[0]).not.toContain("<previous-summary>");
|
||||
expect(prompts[1]).toContain("<previous-summary>\nsummary 1\n</previous-summary>");
|
||||
|
||||
// The fold covers the whole span: first and last turns both reach a call.
|
||||
expect(prompts[0]).toContain("turn 0");
|
||||
expect(prompts[prompts.length - 1]).toContain("turn 11");
|
||||
} finally {
|
||||
spy.mockRestore();
|
||||
}
|
||||
});
|
||||
|
||||
test("shrinks windows when the provider rejects a prompt the catalog said would fit", async () => {
|
||||
// claude-sonnet-4-5 advertises a 1M window but is beta-gated to 200k on
|
||||
// OAuth credentials (`anthropic.ts` never advertises the 1M beta), so the
|
||||
// only authority on the real cap is the rejection itself.
|
||||
const providerCapChars = 160_000;
|
||||
const rejected: number[] = [];
|
||||
let call = 0;
|
||||
const spy = vi.spyOn(ai, "completeSimple").mockImplementation(async (_model, context) => {
|
||||
const prompt = promptTextOf([_model, context]);
|
||||
if (prompt.length > providerCapChars) {
|
||||
rejected.push(prompt.length);
|
||||
throw new Error(`400 prompt is too long: ${prompt.length} tokens > ${providerCapChars} maximum`);
|
||||
}
|
||||
return createAssistantMessage(`summary ${++call}`);
|
||||
});
|
||||
try {
|
||||
const messages = Array.from({ length: 60 }, (_, i) => turn(i, 4_000)).flat();
|
||||
const summary = await generateSummary(messages, getModel(400_000), 16_384, "test-key");
|
||||
|
||||
// The first plan trusted the catalog and was rejected; the fold halved
|
||||
// the window instead of failing the compaction.
|
||||
expect(rejected.length).toBeGreaterThan(0);
|
||||
expect(rejected.length).toBeLessThan(4);
|
||||
expect(summary).toBe(`summary ${call}`);
|
||||
const accepted = spy.mock.calls.map(promptTextOf).filter(p => p.length <= providerCapChars);
|
||||
expect(accepted[0]).toContain("turn 0");
|
||||
expect(accepted[accepted.length - 1]).toContain("turn 59");
|
||||
} finally {
|
||||
spy.mockRestore();
|
||||
}
|
||||
});
|
||||
|
||||
test("propagates a non-overflow failure instead of shrinking", async () => {
|
||||
let calls = 0;
|
||||
const spy = vi.spyOn(ai, "completeSimple").mockImplementation(async () => {
|
||||
calls++;
|
||||
throw new Error("provider exploded");
|
||||
});
|
||||
try {
|
||||
const messages = Array.from({ length: 12 }, (_, i) => turn(i, 4_000)).flat();
|
||||
await expect(generateSummary(messages, getModel(40_000), 16_384, "test-key")).rejects.toThrow(
|
||||
"provider exploded",
|
||||
);
|
||||
expect(calls).toBe(1);
|
||||
} finally {
|
||||
spy.mockRestore();
|
||||
}
|
||||
});
|
||||
|
||||
test("carries a caller-supplied previous summary into the first window", async () => {
|
||||
const spy = vi.spyOn(ai, "completeSimple").mockResolvedValue(createAssistantMessage("merged"));
|
||||
try {
|
||||
await generateSummary(turn(1, 200), getModel(200_000), 16_384, "test-key", undefined, undefined, "earlier");
|
||||
expect(promptTextOf(spy.mock.calls[0])).toContain("<previous-summary>\nearlier\n</previous-summary>");
|
||||
} finally {
|
||||
spy.mockRestore();
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
function compactionEntry(id: string, preserveData?: Record<string, unknown>): SessionEntry {
|
||||
return {
|
||||
type: "compaction",
|
||||
id,
|
||||
parentId: null,
|
||||
timestamp: new Date().toISOString(),
|
||||
summary: "summary",
|
||||
firstKeptEntryId: `${id}-kept`,
|
||||
tokensBefore: 1,
|
||||
preserveData,
|
||||
} satisfies SessionEntry;
|
||||
}
|
||||
|
||||
describe("readable compaction boundary", () => {
|
||||
const local = compactionEntry("local");
|
||||
const remote = compactionEntry("remote", {
|
||||
openaiRemoteCompaction: { provider: "openai-codex", replacementHistory: [] },
|
||||
});
|
||||
const entries = [local, remote];
|
||||
|
||||
test("skips a provider-native compaction another provider cannot replay", () => {
|
||||
const model = getBundledModel("anthropic", "claude-sonnet-4-5");
|
||||
if (!model) throw new Error("Expected built-in anthropic/claude-sonnet-4-5 to exist");
|
||||
expect(findReadableCompactionIndex(entries, DEFAULT_COMPACTION_SETTINGS, model)).toBe(0);
|
||||
});
|
||||
|
||||
test("keeps a provider-native compaction the same provider can replay", () => {
|
||||
const model = getBundledModel("openai-codex", "gpt-5.6-sol");
|
||||
if (!model) throw new Error("Expected built-in openai-codex/gpt-5.6-sol to exist");
|
||||
expect(findReadableCompactionIndex(entries, DEFAULT_COMPACTION_SETTINGS, model)).toBe(1);
|
||||
});
|
||||
|
||||
test("without an active model the newest entry is the boundary", () => {
|
||||
expect(findReadableCompactionIndex(entries, DEFAULT_COMPACTION_SETTINGS)).toBe(1);
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user