Merge PR #8920: fix(compaction): bound summarization input and stop retrying overflow (@PaleRoses)

# Conflicts:
#	packages/agent/src/compaction/compaction.ts
This commit is contained in:
can1357
2026-08-19 01:39:07 +02:00
7 changed files with 433 additions and 22 deletions
+168 -20
View File
@@ -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";
@@ -888,6 +889,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,
@@ -900,6 +968,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) {
@@ -909,11 +1050,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) {
@@ -1248,6 +1384,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,
@@ -1257,21 +1419,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;
}
let prevCompactionIndex = findReadableCompactionIndex(pathEntries, settings, activeModel);
// Honor the latest `/clear` reset boundary. `/clear` records a
// `reset_boundary` marker and reports the model context empty, so compaction
@@ -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);
});
});
+1 -1
View File
@@ -107,7 +107,7 @@ const TRANSIENT_ENVELOPE_PATTERN = /anthropic stream envelope error:/i;
const TRANSIENT_ENVELOPE_BEFORE_START_PATTERN = /before message_start/i;
export const STREAM_READ_ERROR_PATTERN = /stream[_ -]?read[_ -]?error/i;
export const TRANSIENT_TRANSPORT_PATTERN =
/\b(?:no[_ -]?capacity|(?:high|peak)[ _-]?demand|(?:at|over|insufficient)[ _-]?capacity|capacity[ _-]?(?:exceeded|exhausted)|peak[ _-]?load)\b|overloaded|provider.?returned.?error|rate.?limit|too many requests|429|500|502|503|504|service.?unavailable|server.?error|internal.?error|retry your request|network.?error|connection.?error|connection.?refused|unable.?to.?connect\.\s*is the computer able to access the url\?|other side closed|fetch failed|upstream.?connect|upstream.?request.?failed|reset before headers|socket hang up|timed? out|timeout|terminated|retry delay|stream stall|no error details in response|HTTP2(?:StreamReset|RefusedStream|EnhanceYourCalm)|nghttp2_(?:internal_error|refused_stream)|stream closed with error code nghttp2_(?:internal_error|refused_stream)|malformed.?function.?call/i;
/\b(?:no[_ -]?capacity|(?:high|peak)[ _-]?demand|(?:at|over|insufficient)[ _-]?capacity|capacity[ _-]?(?:exceeded|exhausted)|peak[ _-]?load)\b|overloaded|provider.?returned.?error|rate.?limit|too many requests|\b(?:429|500|502|503|504)\b|service.?unavailable|server.?error|internal.?error|retry your request|network.?error|connection.?error|connection.?refused|unable.?to.?connect\.\s*is the computer able to access the url\?|other side closed|fetch failed|upstream.?connect|upstream.?request.?failed|reset before headers|socket hang up|timed? out|timeout|terminated|retry delay|stream stall|no error details in response|HTTP2(?:StreamReset|RefusedStream|EnhanceYourCalm)|nghttp2_(?:internal_error|refused_stream)|stream closed with error code nghttp2_(?:internal_error|refused_stream)|malformed.?function.?call/i;
const AUTH_FAILURE_PATTERN =
/\b(?:401|403|unauthorized|forbidden|authentication|auth[_ ]?unavailable|no auth available|(?:invalid|no)[_ ]?api[_ ]?key)\b/i;
const MALFORMED_FUNCTION_CALL_PATTERN = /\bmalformed.?function.?call\b/i;
+4
View File
@@ -102,6 +102,10 @@ function isRetryableOneshotFailure(errorId: number, errorStatus: number | undefi
// Replaying the same prompt produces the same malformed output.
if (AIError.LLAMA_CPP_TOOL_CALL_PARSE_PATTERN.test(errorMessage)) return false;
if (AIError.is(errorId, AIError.Flag.ContentBlocked)) return false;
// A oneshot replays a FIXED prompt, so an input the model cannot fit fails
// identically on every attempt. Retrying burns the caller's deadline instead
// of reaching the fallback that can actually shrink the input.
if (AIError.is(errorId, AIError.Flag.ContextOverflow)) return false;
return (
AIError.isTransientStatus(errorStatus) ||
AIError.is(errorId, AIError.Flag.Transient) ||
@@ -0,0 +1,36 @@
import { describe, expect, it } from "bun:test";
import * as AIError from "@oh-my-pi/pi-ai/error";
/**
* The transient classifier matches bare HTTP status codes in error text. Those
* digits must be a token of their own: omp appends its own
* `raw-http-request=<...>/<random-id>.json` pointer to provider errors, and a
* random id containing `503` used to make a hard 400 look retryable — which
* turned a deterministic oversized-prompt rejection into ten identical retries.
*/
describe("transient status classification", () => {
const overflowWithArtifactPointer =
'Summarization failed: 400 {"type":"error","error":{"type":"invalid_request_error",' +
'"message":"prompt is too long: 3030000 tokens > 1000000 maximum"}}\n' +
"raw-http-request=/home/u/.omp/logs/http-400-requests/1787022540720-3o503gxo48bvb.json";
it("does not call a 400 transient because an artifact id embeds a status code", () => {
const id = AIError.classify(new Error(overflowWithArtifactPointer), "anthropic-messages");
expect(AIError.is(id, AIError.Flag.ContextOverflow)).toBe(true);
expect(AIError.is(id, AIError.Flag.Transient)).toBe(false);
});
it("still classifies a real gateway status as transient", () => {
for (const text of ["503 Service Unavailable", "upstream returned 502", "HTTP 429 from provider"]) {
const id = AIError.classify(new Error(text), "anthropic-messages");
expect(AIError.is(id, AIError.Flag.Transient)).toBe(true);
}
});
it("does not treat status digits inside an identifier as transient", () => {
for (const text of ["model gpt-500x rejected the request", "request req500502 failed validation"]) {
const id = AIError.classify(new Error(text), "anthropic-messages");
expect(AIError.is(id, AIError.Flag.Transient)).toBe(false);
}
});
});
+24
View File
@@ -124,6 +124,30 @@ describe("retryTransientCompletion", () => {
expect(final.stopReason).toBe("error");
});
it("does not retry an input the model cannot fit", async () => {
// A oneshot replays a fixed prompt: the same overflow comes back every
// attempt, so the retries only delay the caller's fallback. Observed live
// as 10 identical 3M-token compaction summarization calls.
let calls = 0;
const final = await retryTransientCompletion(
() => {
calls += 1;
return Promise.resolve(
message({
stopReason: "error",
errorStatus: 400,
errorMessage:
"invalid_request_error: prompt is too long: 3059586 tokens > 1000000 maximum (raw-http-request=/logs/1787022540720-3o503gxo48bvb.json)",
}),
);
},
{ ...fast, maxAttempts: 5 },
);
expect(calls).toBe(1);
expect(final.stopReason).toBe("error");
});
it("does not retry a deterministic llama.cpp tool-call parse failure reported as 500", async () => {
let calls = 0;
const final = await retryTransientCompletion(
@@ -477,7 +477,9 @@ export class SessionMaintenance {
const config = this.#withPlanProtection({
...(opts.config ?? AGGRESSIVE_SHAKE_CONFIG),
// Skip entries summarized away by the latest compaction — shaking them
// only churns persisted history with no prompt/cache effect.
// only churns persisted history with no prompt/cache effect. The cut is
// unconditional on the wire (see `buildSessionContext`), so a compaction
// the active model cannot replay still hides its prefix from the prompt.
keepBoundaryId: latestCompaction?.firstKeptEntryId,
});
const regions = collectShakeRegions(branchEntries, config);
@@ -2729,9 +2731,16 @@ export class SessionMaintenance {
}
const retryAfterMs = this.#host.parseRetryAfterMsFromError(message);
// An input the summarizer cannot fit is deterministic: the same
// prompt fails identically every attempt, so the retry budget is
// pure latency and the next candidate (a larger window) is the
// only move that can succeed. Overflow therefore vetoes the
// transient/usage-limit arms, which a provider blob can trip on
// coincidence alone.
const shouldRetry =
retrySettings.enabled &&
attempt < retrySettings.maxRetries &&
!AIError.is(id, AIError.Flag.ContextOverflow) &&
(retryAfterMs !== undefined ||
AIError.is(id, AIError.Flag.Transient) ||
AIError.is(id, AIError.Flag.UsageLimit));