From b692de0241db81faa1ef1daf4e7d5a3377eaa4e3 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Wed, 23 Sep 2026 16:38:40 -0700 Subject: [PATCH 1/8] fix(compaction): limit idle recompress to the Anthropic cache Guessed windows for other providers fold a cache that is still warm. Only Anthropic publishes a 5-minute expiry this client actually uses. --- docs/ARCHITECTURE.md | 12 ++-- src/agent/compaction.test.ts | 68 ++++++++++-------- src/provider/cache-ttl.test.ts | 49 +++++-------- src/provider/cache-ttl.ts | 128 ++++++++++----------------------- 4 files changed, 99 insertions(+), 158 deletions(-) diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index a6a1b2c27..9d549b1b3 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -193,14 +193,10 @@ When a cycle's input tokens cross a threshold, the director compacts the inferen - **Idle (end-of-turn)** — An interactive turn can end with a reply and then sit idle with no tool batch to intercept; the governor requests a continuation at that pause and compacts when it arrives. An operator message that races the continuation still compacts first, then re-enters inference to answer it. - **Idle recompress past the provider TTL** — The fold is a re-compress, not a cache play: provider KV caches expire on their own schedule, so a session compresses _after_ that expiry (the next turn is a cheaper write and later reads compound on the shrunk context) instead of only at 60% tokens. It is not limited to under-threshold sessions: once growth hysteresis has cleared `pending`, an over-threshold session in that gap can fire the same fold. Any live re-entry past the window fires it — an empty stall ping re-enters meter-only, a raced operator message re-infers — under reason `cache-ttl-recompress` through the same compactor, so the fresh tail stays raw exactly as in the threshold path. The interval follows provider cache economics (`src/provider/cache-ttl.ts`), not a global N minutes: -| Provider segment | Idle recompress allowed after | Why | -| ----------------------------------------------------------------- | ----------------------------- | ------------------------------------------------------------------------------------ | -| `anthropic`, `zen-messages`, `opencode-go-messages` | 5 min | Default ephemeral cache TTL; hour-long breakpoints are never set | -| `openai-responses`, `codex-responses`, `openai-compatible`, `xai` | 10 min | In-memory prefixes typically evicted after 5–10 min idle; xAI assumed same economics | -| `gemini` | 15 min | No published implicit-cache eviction window; conservative | -| `deepseek` | 60 min | On-disk context cache persists for hours-to-days | -| `ollama` | never | Local inference has no remote cache to expire | -| anything else | 10 min | Assumed OpenAI-style in-memory economics; tune per upstream | +| Provider segment | Idle recompress allowed after | Why | +| -------------------------------------------------------------- | ----------------------------- | ---------------------------------------------------------------------------------------------------------------- | +| `anthropic`, `zen-messages`, `opencode-go-messages` | 5 min | Published default ephemeral cache TTL; the 1-hour TTL is never set | +| anything else, including OpenAI, xAI, Gemini, DeepSeek, ollama | never | No published 5-minute expiry, or no remote cache. Guessing a shorter window folds a cache that may still be warm | The `ollama` table key is the bare provider segment; the harness stamps slash-form on `sourceId` (`ollama/default` / `ollama/`), with `provider: openai-compatible` (`buildOpenAISource`) and a bare model (`llama3` / `qwen3`). Idle-recompress disable matches via `isOllamaProviderId` on that `sourceId` (`src/provider/ollama.ts`), not by looking up slash-form as a table key. diff --git a/src/agent/compaction.test.ts b/src/agent/compaction.test.ts index f517d7f4a..caa0b4bb0 100644 --- a/src/agent/compaction.test.ts +++ b/src/agent/compaction.test.ts @@ -67,6 +67,7 @@ function turnsOfLength(count: number, textLength: number): ConversationTurn[] { function inferenceDone( input: number, text = "", + provider = "p", ): Extract { return { type: "inference.done", @@ -75,7 +76,7 @@ function inferenceDone( content: text.length > 0 ? [{ type: "text", text }] : [], }, usage: usage(input), - source: { sourceId: "s", provider: "p", model: "m" }, + source: { sourceId: "s", provider, model: "m" }, } as unknown as Extract; } @@ -946,7 +947,11 @@ describe("provider-aware idle recompress (CL-8745)", () => { ): Extract { const source = typeof modelOrSource === "string" - ? { sourceId: "s", provider: "p", model: modelOrSource } + ? { + sourceId: "s", + provider: modelOrSource.split("/")[0] ?? "p", + model: modelOrSource, + } : { sourceId: modelOrSource.sourceId ?? "s", provider: modelOrSource.provider ?? "p", @@ -1021,11 +1026,12 @@ describe("provider-aware idle recompress (CL-8745)", () => { expect(governor.resumeAfterCompact(emptyMessage())).toBe("meter"); }); - test("production LastCycleSource: anthropic fires at 5m, codex stays quiet until 10m", () => { + test("production LastCycleSource: anthropic fires at 5m, codex never does", () => { // Harness stamps { sourceId, provider, model } with a bare model, not // slash-form "anthropic/claude-opus-4-6". Anthropic's 5-minute window - // must come from provider, not a dummy model string; Codex must not - // inherit that 5-minute fire from sourceId "codex/work". + // must come from provider, not a dummy model string. Codex has no + // published 5-minute expiry, so it must stay quiet past the old 10-minute + // guess as well. let nowMs = 15_000_000; const clock = () => nowMs; const anthropic = createCompactionGovernor(() => undefined, "", [], clock); @@ -1057,13 +1063,13 @@ describe("provider-aware idle recompress (CL-8745)", () => { codex.interceptIdleContinuation(emptyMessage(), capabilities), ).toBeNull(); - nowMs += 5 * MINUTE_MS; + nowMs += 30 * MINUTE_MS; expect( codex.interceptIdleContinuation(emptyMessage(), capabilities), - ).toEqual(ttlCompact); + ).toBeNull(); }); - test("follows provider economics: deepseek waits out its long window, ollama never fires", () => { + test("does not idle-recompress DeepSeek or ollama", () => { let nowMs = 20_000_000; const clock = () => nowMs; const deepseek = createCompactionGovernor(() => undefined, "", [], clock); @@ -1077,30 +1083,20 @@ describe("provider-aware idle recompress (CL-8745)", () => { tenTurns, ); - // Past Anthropic/OpenAI windows but inside DeepSeek's hour: neither fires. - nowMs += 30 * MINUTE_MS; + nowMs += 90 * MINUTE_MS; expect( deepseek.interceptIdleContinuation(emptyMessage(), capabilities), ).toBeNull(); expect( local.interceptIdleContinuation(emptyMessage(), capabilities), ).toBeNull(); - - // Past DeepSeek's hour: recompress fires; local inference still never - // does — no remote cache means no cache benefit. - nowMs += 31 * MINUTE_MS; - expect( - deepseek.interceptIdleContinuation(emptyMessage(), capabilities), - ).toEqual(ttlCompact); - expect( - local.interceptIdleContinuation(emptyMessage(), capabilities), - ).toBeNull(); }); test("production ollama LastCycleSource never fires cache-ttl-recompress", () => { // Harness stamps { sourceId, provider, model } with a bare model. Ollama is // buildOpenAISource: sourceId "ollama/default", provider openai-compatible, - // model llama3. Keying TTL only off model would take the 10-minute default. + // Keying TTL off the bare model would miss the ollama sourceId. Local + // inference stays disabled. let nowMs = 25_000_000; const governor = createCompactionGovernor( () => undefined, @@ -1120,7 +1116,7 @@ describe("provider-aware idle recompress (CL-8745)", () => { tenTurns, ); - nowMs += 11 * MINUTE_MS; + nowMs += 5 * MINUTE_MS + 1; expect( governor.interceptIdleContinuation(emptyMessage(), capabilities), ).toBeNull(); @@ -1142,8 +1138,8 @@ describe("provider-aware idle recompress (CL-8745)", () => { () => nowMs, ); governor.noteInferenceDone(inferenceDone(overThreshold), tenTurns); - // The "m" fixture model carries the 10-minute default TTL; advance past it. - nowMs += 11 * MINUTE_MS; + // Fixture provider "p" has no TTL. Threshold arming still owns this session. + nowMs += 5 * MINUTE_MS + 1; // Threshold arming owns the over-threshold session: no TTL double-fold. expect( governor.interceptIdleContinuation(emptyMessage(), capabilities), @@ -1165,18 +1161,24 @@ describe("provider-aware idle recompress (CL-8745)", () => { [], () => nowMs, ); - governor.noteInferenceDone(inferenceDone(overThreshold), tenTurns); + governor.noteInferenceDone( + inferenceDone(overThreshold, "", "anthropic"), + tenTurns, + ); expect( governor.interceptActions(toolDone(), inferAction, capabilities), ).not.toBeNull(); expect(governor.resumeAfterCompact(emptyMessage())).toBe("infer"); - governor.noteInferenceDone(inferenceDone(overThreshold), tenTurns); + governor.noteInferenceDone( + inferenceDone(overThreshold, "", "anthropic"), + tenTurns, + ); expect( governor.interceptActions(toolDone(), inferAction, capabilities), ).toBeNull(); - nowMs += 11 * MINUTE_MS; + nowMs += 5 * MINUTE_MS + 1; expect( governor.interceptIdleContinuation(emptyMessage(), capabilities), ).toEqual(ttlCompact); @@ -1193,14 +1195,20 @@ describe("provider-aware idle recompress (CL-8745)", () => { [], () => nowMs, ); - governor.noteInferenceDone(inferenceDone(overThreshold), tenTurns); + governor.noteInferenceDone( + inferenceDone(overThreshold, "", "anthropic"), + tenTurns, + ); expect( governor.interceptActions(toolDone(), inferAction, capabilities), ).not.toBeNull(); expect(governor.resumeAfterCompact(emptyMessage())).toBe("infer"); - governor.noteInferenceDone(inferenceDone(overThreshold), tenTurns); - nowMs += 11 * MINUTE_MS; + governor.noteInferenceDone( + inferenceDone(overThreshold, "", "anthropic"), + tenTurns, + ); + nowMs += 5 * MINUTE_MS + 1; expect( governor.interceptIdleContinuation(emptyMessage(), capabilities), ).toEqual(ttlCompact); diff --git a/src/provider/cache-ttl.test.ts b/src/provider/cache-ttl.test.ts index 0bf58fbf4..ceb383731 100644 --- a/src/provider/cache-ttl.test.ts +++ b/src/provider/cache-ttl.test.ts @@ -4,24 +4,21 @@ import { cacheTtlMsFor } from "./cache-ttl.js"; const MINUTE_MS = 60_000; describe("cacheTtlMsFor", () => { - test("maps providers to their documented cache TTL windows", () => { + test("allows idle recompress only for the published Anthropic 5-minute window", () => { expect(cacheTtlMsFor("anthropic/claude-opus-4-6")).toBe(5 * MINUTE_MS); - expect(cacheTtlMsFor("openai-responses/gpt-5.6")).toBe(10 * MINUTE_MS); - expect(cacheTtlMsFor("codex-responses/gpt-5.6")).toBe(10 * MINUTE_MS); - expect(cacheTtlMsFor("openai-compatible/custom")).toBe(10 * MINUTE_MS); - expect(cacheTtlMsFor("xai/thegreataxios")).toBe(10 * MINUTE_MS); - expect(cacheTtlMsFor("gemini/gemini-3-pro")).toBe(15 * MINUTE_MS); - expect(cacheTtlMsFor("deepseek/deepseek-chat")).toBe(60 * MINUTE_MS); - }); - - test("covers Anthropic-protocol adapters with the 5-minute window", () => { expect(cacheTtlMsFor("zen-messages/claude-opus-4-6")).toBe(5 * MINUTE_MS); expect(cacheTtlMsFor("opencode-go-messages/claude-opus-4-6")).toBe( 5 * MINUTE_MS, ); }); - test("disables idle recompress for local inference and missing model ids", () => { + test("disables providers whose expiry is unpublished or longer than 5 minutes", () => { + expect(cacheTtlMsFor("openai-responses/gpt-5.6")).toBeUndefined(); + expect(cacheTtlMsFor("codex-responses/gpt-5.6")).toBeUndefined(); + expect(cacheTtlMsFor("openai-compatible/custom")).toBeUndefined(); + expect(cacheTtlMsFor("xai/thegreataxios")).toBeUndefined(); + expect(cacheTtlMsFor("gemini/gemini-3-pro")).toBeUndefined(); + expect(cacheTtlMsFor("deepseek/deepseek-chat")).toBeUndefined(); expect(cacheTtlMsFor("ollama/llama3.1")).toBeUndefined(); expect(cacheTtlMsFor(undefined)).toBeUndefined(); expect(cacheTtlMsFor("")).toBeUndefined(); @@ -37,7 +34,7 @@ describe("cacheTtlMsFor", () => { ).toBeUndefined(); }); - test("maps bare LastCycleSource ids through provider and family", () => { + test("maps a bare Anthropic LastCycleSource and leaves Codex quiet", () => { expect( cacheTtlMsFor({ provider: "anthropic", @@ -49,25 +46,17 @@ describe("cacheTtlMsFor", () => { provider: "codex-responses", model: "gpt-5.6-luna", }), - ).toBe(10 * MINUTE_MS); - }); - - test("falls back to model family for unrecognized provider prefixes", () => { - expect(cacheTtlMsFor("proxy-acme/grok-4")).toBe(10 * MINUTE_MS); - expect(cacheTtlMsFor("proxy-acme/gemini-3-pro")).toBe(15 * MINUTE_MS); - expect(cacheTtlMsFor("proxy-acme/ollama-qwen")).toBeUndefined(); - // An exact provider-segment match wins over the model family: an - // openai-compatible account fronting Claude keeps the generic window. - expect(cacheTtlMsFor("openai-compatible/claude-opus-4-6")).toBe( - 10 * MINUTE_MS, - ); + ).toBeUndefined(); }); - test("assumes OpenAI-style economics for unrecognized providers", () => { - expect(cacheTtlMsFor("bifrost/some-model")).toBe(10 * MINUTE_MS); - expect(cacheTtlMsFor("totally-new-provider/model-x")).toBe(10 * MINUTE_MS); - // A truly unknown id (no provider segment, no family substring) still - // gets the 10-minute default — not the local-inference disable. - expect(cacheTtlMsFor("unknown-id")).toBe(10 * MINUTE_MS); + test("does not inherit a window from the model family or an unknown provider", () => { + expect(cacheTtlMsFor("proxy-acme/grok-4")).toBeUndefined(); + expect(cacheTtlMsFor("proxy-acme/gemini-3-pro")).toBeUndefined(); + expect(cacheTtlMsFor("proxy-acme/claude-opus-4-6")).toBeUndefined(); + // The provider segment wins: an openai-compatible account fronting + // Claude is not the Anthropic messages protocol. + expect(cacheTtlMsFor("openai-compatible/claude-opus-4-6")).toBeUndefined(); + expect(cacheTtlMsFor("bifrost/some-model")).toBeUndefined(); + expect(cacheTtlMsFor("unknown-id")).toBeUndefined(); }); }); diff --git a/src/provider/cache-ttl.ts b/src/provider/cache-ttl.ts index a3803b52f..62104c3f0 100644 --- a/src/provider/cache-ttl.ts +++ b/src/provider/cache-ttl.ts @@ -1,36 +1,21 @@ // Provider prompt-cache TTLs for idle recompression (CL-8745). // -// The fold is a re-compress, not a cache play: provider KV caches expire on -// their own schedule, and idle compression fires *after* that expiry so the -// next turn is a cheaper write (a smaller prefix to re-process) and later -// reads compound on the shrunk context. The interval therefore follows -// provider cache economics, not a single global N minutes. +// The fold is a re-compress, not a cache play. It is allowed only after a +// published cache expiry, so the next turn rewrites a smaller prefix instead +// of compacting while a warm cache would still have been a cheap read. // -// Mapping — idle recompress is allowed once `now - lastCacheWrite >= ttl`: -// - Anthropic (`anthropic`, plus `zen-messages` / `opencode-go-messages`, -// which speak the Anthropic messages protocol): default 5-minute ephemeral -// cache TTL. Hour-long breakpoints exist but we never set them. 5 min. -// - OpenAI family (`openai-responses`, `codex-responses`, -// `openai-compatible`): in-memory prefixes are typically evicted after -// 5-10 minutes idle (retained up to an hour off-peak); GPT-5.6+ families -// default to a 30-minute minimum TTL instead. 10 min splits the range. -// - xAI (`xai`, `grok-*`): per-server cache, no published TTL; assumed -// OpenAI-style in-memory economics. 10 min. -// - Gemini (`gemini`, `gemini-*`): implicit caching on 2.5+ with no published -// eviction window; explicit caches default to 1 hour. Conservative. 15 min. -// - DeepSeek (`deepseek`): on-disk context cache cleared only after -// hours-to-days idle, so an early recompress would fold a still-warm cache. -// 60 min. -// - ollama: local inference has no remote cache to expire, so an idle -// recompress would burn local compute for no cache benefit. Disabled. -// Production LastCycleSource is `{ sourceId, provider, model }` with a -// bare model; Ollama is `buildOpenAISource` (`provider: openai-compatible`, -// `sourceId` like `ollama/default`, model `llama3` / `qwen3`). Identity is -// `sourceId` via `isOllamaProviderId`, not a slash-form string the harness -// never stamps. -// - Unknown/custom providers (bifrost proxy, `openai-compatible` fronting an -// unlisted upstream, unrecognized strings): assumed OpenAI-style in-memory -// economics. 10 min. +// Anthropic is the only provider with a published idle expiry we actually +// use: the default ephemeral cache lasts 5 minutes and refreshes on each +// hit. The 1-hour TTL exists, costs more to write, and this client never +// sets it. `zen-messages` and `opencode-go-messages` speak that same +// messages protocol, so they share the 5-minute window. +// +// Everyone else is disabled. OpenAI GPT-5.6+ stays eligible for at least 30 +// minutes. Gemini's implicit cache has no published eviction window and its +// explicit cache defaults to 1 hour. DeepSeek's disk cache is cleared over +// hours to days. xAI publishes no TTL. Guessing a shorter window folds a +// cache that is still warm. Ollama has no remote cache. Unknown providers +// stay off rather than inheriting a default. // // Deliberate non-goal: the summary call itself carries no `cache_control` — // prompt-caching the fold is not attempted. @@ -39,38 +24,14 @@ import { isOllamaProviderId } from "./ollama.js"; const MINUTE_MS = 60_000; -const PROVIDER_TTLS_MS: Record = { - anthropic: 5 * MINUTE_MS, - "zen-messages": 5 * MINUTE_MS, - "opencode-go-messages": 5 * MINUTE_MS, - "openai-responses": 10 * MINUTE_MS, - "codex-responses": 10 * MINUTE_MS, - "openai-compatible": 10 * MINUTE_MS, - gemini: 15 * MINUTE_MS, - deepseek: 60 * MINUTE_MS, - xai: 10 * MINUTE_MS, - // No remote cache to expire; idle recompress would only burn local compute. - ollama: undefined, -}; - -const DEFAULT_TTL_MS = 10 * MINUTE_MS; +/** Published Anthropic ephemeral TTL. The only idle-recompress window. */ +const ANTHROPIC_TTL_MS = 5 * MINUTE_MS; -// Model-family fallback for unrecognized provider prefixes (custom proxies, -// account names). An exact provider-segment match always wins — e.g. an -// `openai-compatible` account fronting Claude keeps the generic 10-minute -// window rather than Claude's 5. Checked in order; first substring wins. -const FAMILY_TTLS_MS: readonly (readonly [string, number | undefined])[] = [ - ["claude", 5 * MINUTE_MS], - ["anthropic", 5 * MINUTE_MS], - ["gpt", 10 * MINUTE_MS], - ["codex", 10 * MINUTE_MS], - ["openai", 10 * MINUTE_MS], - ["grok", 10 * MINUTE_MS], - ["xai", 10 * MINUTE_MS], - ["gemini", 15 * MINUTE_MS], - ["deepseek", 60 * MINUTE_MS], - ["ollama", undefined], -]; +const ANTHROPIC_PROTOCOL = new Set([ + "anthropic", + "zen-messages", + "opencode-go-messages", +]); type CacheTtlIdentity = { sourceId?: string; @@ -78,11 +39,9 @@ type CacheTtlIdentity = { model?: string; }; -// Canonical provider segment of a `provider:model` string: the account or +// Canonical provider segment of a `provider/model` string: the account or // adapter name before the first "/" (custom names like `xai/thegreataxios` -// carry the provider there), else the head before ":". Mirrors the -// segmentation in provider/context-window.ts without its window-table -// fallback semantics. +// carry the provider there), else the head before ":". function canonicalSegment(model: string): string { const lower = model.toLowerCase(); const slash = lower.indexOf("/"); @@ -91,46 +50,35 @@ function canonicalSegment(model: string): string { return colon >= 0 ? head.slice(0, colon) : head; } -function ttlForModelString(model: string | undefined): number | undefined { - if (model === undefined) return undefined; - const segment = canonicalSegment(model); - if (segment.length === 0) return undefined; - if (Object.hasOwn(PROVIDER_TTLS_MS, segment)) - return PROVIDER_TTLS_MS[segment]; - const lower = model.toLowerCase(); - for (const [family, ttl] of FAMILY_TTLS_MS) { - if (lower.includes(family)) return ttl; - } - return DEFAULT_TTL_MS; +function ttlForSegment(segment: string): number | undefined { + if (!ANTHROPIC_PROTOCOL.has(segment)) return undefined; + return ANTHROPIC_TTL_MS; } /** * Milliseconds of provider-cache idle after which a recompress is allowed, - * or `undefined` when the identity is missing/empty or the provider has no - * remote cache (local inference). Never throws; unknown strings get the - * default 10-minute window. + * or `undefined` when idle recompress must not run. Never throws. * * Accepts a slash-form `provider/model` string or a LastCycleSource-shaped - * identity. Production events stamp a bare `model`; Ollama is recognized - * from `sourceId` (`isOllamaProviderId`), not from a combined string the - * harness never stamps. + * identity. An explicit provider wins: only the Anthropic messages protocol + * returns the 5-minute window. A non-Anthropic provider stays disabled even + * when the model id contains "claude". Ollama is recognized from `sourceId` + * (`isOllamaProviderId`) before that, because production stamps + * `provider: openai-compatible` for local inference. */ export function cacheTtlMsFor( identity: string | CacheTtlIdentity | undefined, ): number | undefined { if (identity === undefined) return undefined; - if (typeof identity === "string") return ttlForModelString(identity); + if (typeof identity === "string") + return ttlForSegment(canonicalSegment(identity)); if (identity.sourceId !== undefined && isOllamaProviderId(identity.sourceId)) return undefined; if (identity.provider !== undefined && isOllamaProviderId(identity.provider)) return undefined; - if (identity.provider !== undefined) { - const provider = identity.provider.toLowerCase(); - if (Object.hasOwn(PROVIDER_TTLS_MS, provider)) - return PROVIDER_TTLS_MS[provider]; - } - - return ttlForModelString(identity.model); + if (identity.provider !== undefined && identity.provider.length > 0) + return ttlForSegment(canonicalSegment(identity.provider)); + return ttlForSegment(canonicalSegment(identity.model ?? "")); } From 5840a36563626daac9bc7947e465eec98fc0b7da Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Wed, 23 Sep 2026 18:40:51 -0700 Subject: [PATCH 2/8] feat(session): compact expired anthropic history before the next infer A resumed process only knew the cache write time in memory, so the first request rewrote the full prefix. Persist the Anthropic stamp and fold through the existing compactor before that infer. --- docs/ARCHITECTURE.md | 2 +- src/agent/compaction.ts | 38 ++++ src/agent/director.ts | 14 ++ src/exec/runner.ts | 35 +++- src/process-handlers.ts | 3 + src/provider/cache-ttl.test.ts | 46 ++++- src/provider/cache-ttl.ts | 32 ++++ src/session/active-run.ts | 6 + src/session/assemble-runtime.ts | 28 ++- src/session/cache-ttl-resume.test.ts | 267 +++++++++++++++++++++++++++ src/session/run-sink.ts | 8 +- src/session/state.ts | 4 + src/tui/runner/exit.ts | 5 + src/tui/runner/session.ts | 19 +- src/tui/session-start.ts | 22 +++ 15 files changed, 521 insertions(+), 8 deletions(-) create mode 100644 src/session/cache-ttl-resume.test.ts diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 9d549b1b3..2b99be445 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -203,7 +203,7 @@ The `ollama` table key is the bare provider segment; the harness stamps slash-fo - **Operator `/compact`** — A first-class slash command folds now without waiting for the occupancy governor. Optional trailing instructions go to the summarizer (empty uses the default structured fold) and are stored on the compact record so later auto-folds still see them. Director rebuild (`/model`, interrupt, resume) restores those instructions from the latest compact record's parameters. An in-flight turn uses the same pair-safe compact-then-continue hop as threshold compact; an idle compact shows the fold and does not start a new turn. - **Overflow recovery** — A `context_overflow` inference error would otherwise become a terminal error reply; the governor compacts and retries instead, bounded so a history the compactor cannot shrink does not loop forever. Overflow ignores hysteresis for the compact itself. -The TTL window is measured from the last `inference.done` (every inference rewrites the provider's prefix cache) with one fire per window per compact of any kind, and it never fires with a tool batch outstanding — stall pings can arrive mid-work. The consecutive-compact cap shared with the threshold path still bounds compact→infer→compact. Prompt-caching the summary call itself is an explicit non-goal: the fold carries no cache options. +The TTL window is measured from the last `inference.done` (every inference rewrites the provider's prefix cache) with one fire per window per compact of any kind, and it never fires with a tool batch outstanding — stall pings can arrive mid-work. The consecutive-compact cap shared with the threshold path still bounds compact→infer→compact. Prompt-caching the summary call itself is an explicit non-goal: the fold carries no cache options. That write time is stored on the session run record for Anthropic-protocol identities only. A new process restores it and, once the stamp is at least 5 minutes old and the stored turn count is above the compact floor, issues the same `cache-ttl-recompress` fold before the first infer so the request reads the trimmed turns. The compaction control flow is shaped by a reactor invariant: a `compact` action runs in its own cycle (it cannot be paired with `infer`), and **the reactor delivers no event after a compact cycle**. A director that simply emitted `compact` in place of the follow-up `infer` would leave the loop idle forever — the cause of an earlier stall. Instead the governor pairs `compact` with an emit action (`custom.compaction.continue`), and the host answers that emission by delivering a content-less inbound message (`buildCompactionContinuationMessage()`) through the serial operation queue, generation-guarded like every other deliver. A hop superseded because interrupt already bumped generation and enqueued a rebuild is re-queued onto the replacement agent rather than consumed on the outgoing liveAgent. That message adds no turn (`createInboundTurn` returns `null` for empty content) but re-enters the loop, where the director issues the follow-up `infer` against the freshly truncated history. Each emission is answered at most once (a replayed duplicate of an already-answered emission is ignored), and an unsolicited continuation with neither resume flag set answers `wait` rather than burning a billable inference. The legacy `requestContinuation` closure survives only on the sub-agent path, which delivers directly with no host emit hop; idle arming keeps the two channels exclusive (closure fires and the arming flag stays false, or no closure and the flag reports the arming) so a caller honoring both cannot double-deliver. diff --git a/src/agent/compaction.ts b/src/agent/compaction.ts index 4cff59a80..b981cbc14 100644 --- a/src/agent/compaction.ts +++ b/src/agent/compaction.ts @@ -103,6 +103,28 @@ export function compactFloorNoopNotice(instructions: string): string { : "Nothing to compact yet."; } +/** Run-record `provider:model` (or slash form) as a LastCycleSource. */ +export function lastCycleSourceFromRunModel( + model: string | undefined, +): LastCycleSource | undefined { + if (model === undefined || model.length === 0) return undefined; + const colon = model.indexOf(":"); + if (colon > 0) { + const provider = model.slice(0, colon); + const rest = model.slice(colon + 1); + return { + sourceId: model, + provider, + model: rest.length > 0 ? rest : model, + }; + } + const slash = model.indexOf("/"); + if (slash > 0) { + return { sourceId: model, provider: model.slice(0, slash), model }; + } + return { sourceId: model, provider: model, model }; +} + /** Build the continuation re-entry action for a compacted governor cycle. */ export function compactionContinuationAction( capabilities: ReactorCapabilities, @@ -561,6 +583,21 @@ export function createCompactionGovernor( extraInstructions = trimmed; } + // A new process has no in-memory cache write. Seed the stored stamp, the + // identity that wrote it, and the stored turns so the next message.received + // can fold before its infer. Turn count is the stored set: the inbound + // message is appended after this, and the floor is measured without it. + function restoreCacheWrite(args: { + at: number; + source: LastCycleSource; + turns: readonly ConversationTurn[]; + }): void { + syncFromTurns(args.turns); + lastCacheWriteAt = args.at; + lastCycleSource = args.source; + lastModel = args.source.model; + } + // `/handoff` folds through the same operator pipeline as above, then starts // the next turn immediately: unlike an idle auto-compact (empty synthetic // continuation → meter, no infer), the caller delivers the pivot message @@ -622,6 +659,7 @@ export function createCompactionGovernor( }, requestManual, restoreExtraInstructions, + restoreCacheWrite, requestHandoff, cancelManual, syncFromTurns, diff --git a/src/agent/director.ts b/src/agent/director.ts index 75cef9c41..a6086f480 100644 --- a/src/agent/director.ts +++ b/src/agent/director.ts @@ -12,6 +12,7 @@ import type { ReactorAction, ToolDefinition, ConversationTurn, + LastCycleSource, RetryPolicy, } from "@intx/types/runtime"; import { isCompactSpacerEchoTurn } from "../session/compactor.js"; @@ -821,6 +822,14 @@ class ChatDirectorImpl extends DefaultDirector { this.compaction.restoreExtraInstructions(value); } + restoreCacheWrite(args: { + at: number; + source: LastCycleSource; + turns: readonly ConversationTurn[]; + }): void { + this.compaction.restoreCacheWrite(args); + } + getCompactTurnCount(): number { return this.compaction.compactTurnCount; } @@ -1437,6 +1446,11 @@ export interface ChatDirector extends ReactorDirector { ): ManualCompactArming; getCompactInstructions(): string | undefined; restoreCompactInstructions(value: string | undefined): void; + restoreCacheWrite(args: { + at: number; + source: LastCycleSource; + turns: readonly ConversationTurn[]; + }): void; getCompactTurnCount(): number; /** `/handoff`: arm the shared operator fold for the pivot the caller sends. */ requestHandoff(instructions: string): HandoffArming; diff --git a/src/exec/runner.ts b/src/exec/runner.ts index 267036e18..a6f671acc 100644 --- a/src/exec/runner.ts +++ b/src/exec/runner.ts @@ -107,6 +107,7 @@ import { } from "../session/active-host.js"; import { finalizeRunState, + loadState, saveState, type ConnectedMcpServer, } from "../session/state.js"; @@ -141,6 +142,10 @@ import { tryReadPriorHandoffFile } from "../session/compaction-handoff.js"; import { emitPluginWarningSummary } from "../plugins/diagnostics.js"; import { createModelSummarizer } from "../session/summarizer.js"; import { ID_PREFIX, LOG_NAMESPACE_ROOT } from "../branding.js"; +import { + anthropicCacheWriteAt, + resumeCacheWriteSeed, +} from "../provider/cache-ttl.js"; import type { ReactorEmittedEvent } from "@intx/inference"; import { setAgentSourceUnlessClosed } from "../tui/agent-source-sync.js"; import { ensureFreshInferenceSource } from "../subagent/refresh-inference-source.js"; @@ -433,6 +438,11 @@ export async function runExec(config: Config): Promise { const startedAt = Date.now(); const workdir = sessionContextDir(config.cwd, sessionId); await initSessionDir(config.cwd, sessionId); + const prior = + config.sessionId.length > 0 + ? await loadState(config.cwd, config.sessionId) + : undefined; + const priorState = prior?.kind === "ok" ? prior.state : undefined; let connectedMcp: ConnectedMcpServer[] = []; let agent: Agent | null = null; @@ -455,6 +465,13 @@ export async function runExec(config: Config): Promise { startedAt, turnsUsed: 0, model: `${config.providerName}:${config.model}`, + ...(priorState?.lastCacheWriteAt !== undefined + ? { lastCacheWriteAt: priorState.lastCacheWriteAt } + : {}), + ...(priorState?.lastCacheWriteAt !== undefined && + priorState.model !== undefined + ? { cacheWriteModel: priorState.model } + : {}), }; let stopHeartbeat: (() => void) | undefined; @@ -486,6 +503,9 @@ export async function runExec(config: Config): Promise { model, mcpServers: connectedMcp, ...(activatedTools.length > 0 ? { activatedTools } : {}), + ...(activeRunHandle.lastCacheWriteAt !== undefined + ? { lastCacheWriteAt: activeRunHandle.lastCacheWriteAt } + : {}), ...(status !== "running" ? { finishedAt: Date.now() } : {}), ...(extra?.error !== undefined ? { error: extra.error } : {}), }; @@ -866,6 +886,12 @@ export async function runExec(config: Config): Promise { } }, }), + getCacheWriteSeed: () => + resumeCacheWriteSeed({ + at: activeRunHandle.lastCacheWriteAt, + storedModel: activeRunHandle.cacheWriteModel, + liveProvider: config.providerName, + }), onBuilt: (agent, storage) => { currentAgent = agent; currentStorage = storage; @@ -913,7 +939,14 @@ export async function runExec(config: Config): Promise { getTelemetry: () => liveTelemetry, getSessionId: () => sessionId, getSource: () => liveSource, - onTurnBoundarySnapshot: () => { + onTurnBoundarySnapshot: (event) => { + if (event?.type === "inference.done") { + const at = anthropicCacheWriteAt(event.data.source, Date.now()); + if (at !== undefined) { + activeRunHandle.lastCacheWriteAt = at; + activeRunHandle.cacheWriteModel = `${event.data.source.provider}:${event.data.source.model}`; + } + } void persist("running"); }, resolveContextDir: () => workdir, diff --git a/src/process-handlers.ts b/src/process-handlers.ts index 356ca6df9..4264a5088 100644 --- a/src/process-handlers.ts +++ b/src/process-handlers.ts @@ -132,6 +132,9 @@ async function finalizeActiveRun( ...(run.activatedTools !== undefined ? { activatedTools: run.activatedTools } : {}), + ...(run.lastCacheWriteAt !== undefined + ? { lastCacheWriteAt: run.lastCacheWriteAt } + : {}), }); } catch (saveErr: unknown) { process.stderr.write( diff --git a/src/provider/cache-ttl.test.ts b/src/provider/cache-ttl.test.ts index ceb383731..fbceb7934 100644 --- a/src/provider/cache-ttl.test.ts +++ b/src/provider/cache-ttl.test.ts @@ -1,5 +1,9 @@ import { describe, expect, test } from "bun:test"; -import { cacheTtlMsFor } from "./cache-ttl.js"; +import { + anthropicCacheWriteAt, + cacheTtlMsFor, + resumeCacheWriteSeed, +} from "./cache-ttl.js"; const MINUTE_MS = 60_000; @@ -60,3 +64,43 @@ describe("cacheTtlMsFor", () => { expect(cacheTtlMsFor("unknown-id")).toBeUndefined(); }); }); + +describe("anthropic cache-write stamp", () => { + test("stamps only Anthropic-protocol identities", () => { + expect(anthropicCacheWriteAt("anthropic:claude-opus-4-6", 10)).toBe(10); + expect(anthropicCacheWriteAt({ provider: "zen-messages" }, 10)).toBe(10); + expect(anthropicCacheWriteAt({ provider: "openai" }, 10)).toBeUndefined(); + expect(anthropicCacheWriteAt(undefined, 10)).toBeUndefined(); + }); + + test("resume seed requires both the stored model and the live provider", () => { + expect( + resumeCacheWriteSeed({ + at: 10, + storedModel: "anthropic:claude-opus-4-6", + liveProvider: "anthropic", + }), + ).toEqual({ at: 10, model: "anthropic:claude-opus-4-6" }); + expect( + resumeCacheWriteSeed({ + at: 10, + storedModel: "openai:gpt-5.6", + liveProvider: "openai", + }), + ).toBeUndefined(); + expect( + resumeCacheWriteSeed({ + at: 10, + storedModel: "anthropic:claude-opus-4-6", + liveProvider: "openai", + }), + ).toBeUndefined(); + expect( + resumeCacheWriteSeed({ + at: undefined, + storedModel: "anthropic:claude-opus-4-6", + liveProvider: "anthropic", + }), + ).toBeUndefined(); + }); +}); diff --git a/src/provider/cache-ttl.ts b/src/provider/cache-ttl.ts index 62104c3f0..9dcb1d9ae 100644 --- a/src/provider/cache-ttl.ts +++ b/src/provider/cache-ttl.ts @@ -82,3 +82,35 @@ export function cacheTtlMsFor( return ttlForSegment(canonicalSegment(identity.provider)); return ttlForSegment(canonicalSegment(identity.model ?? "")); } + +/** + * Wall time to record as the last Anthropic-protocol cache write, or + * `undefined` when this inference must not stamp one. Other providers stay + * unstamped so a later resume does not treat their history as an expired + * Anthropic prefix. + */ +export function anthropicCacheWriteAt( + identity: string | CacheTtlIdentity | undefined, + nowMs: number, +): number | undefined { + if (cacheTtlMsFor(identity) === undefined) return undefined; + return nowMs; +} + +/** + * Seed for a resumed session's pre-infer fold. Absent unless both the + * identity that wrote the cache and the provider about to be called are + * Anthropic-protocol. A stamp alone is not enough: a non-Anthropic model + * on the run record must not fold. + */ +export function resumeCacheWriteSeed(args: { + at: number | undefined; + storedModel: string | undefined; + liveProvider: string; +}): { at: number; model: string } | undefined { + if (args.at === undefined) return undefined; + if (args.storedModel === undefined) return undefined; + if (cacheTtlMsFor(args.storedModel) === undefined) return undefined; + if (cacheTtlMsFor(args.liveProvider) === undefined) return undefined; + return { at: args.at, model: args.storedModel }; +} diff --git a/src/session/active-run.ts b/src/session/active-run.ts index dd8495fde..c537c8dd7 100644 --- a/src/session/active-run.ts +++ b/src/session/active-run.ts @@ -26,6 +26,12 @@ export interface RunStateHandle { // crash/signal terminal write can carry them into run.json for the resume // seed. activatedTools?: string[]; + // Last Anthropic-protocol cache write, and the run-record model that wrote + // it. The crash path copies the stamp so a killed process can still fold + // before the next infer. cacheWriteModel is process-local: the on-disk + // model field is the live provider, which resume reads back as this value. + lastCacheWriteAt?: number; + cacheWriteModel?: string; } // Keep the crash/signal handle in step with every persisted snapshot so a diff --git a/src/session/assemble-runtime.ts b/src/session/assemble-runtime.ts index 1e3d88e51..e1bbc6f9b 100644 --- a/src/session/assemble-runtime.ts +++ b/src/session/assemble-runtime.ts @@ -14,6 +14,7 @@ // unifying further would be abstraction for its own sake. import type { EventEmitter } from "node:events"; +import type { ReactorEmittedEvent } from "@intx/inference"; import { createDirectorRegistry, defineAgent, @@ -50,7 +51,10 @@ import { } from "../agent/tool-search.js"; import { normalizeToolDefinitionsForProvider } from "../agent/tool-schema-normalize.js"; import { resolveModelFamilyPolicy } from "../agent/model-family-policy.js"; -import { stickyExtraInstructionsFromRecords } from "../agent/compaction.js"; +import { + stickyExtraInstructionsFromRecords, + lastCycleSourceFromRunModel, +} from "../agent/compaction.js"; import { createChatDirector, type ChatDirector } from "../agent/director.js"; import { createDoomLoopCorrectiveNote } from "../agent/doom-loop-note.js"; import type { AgentToolset } from "../agent/tools.js"; @@ -503,6 +507,13 @@ export interface ChatAgentWiring { getDefaultSource: () => string; /** Read at each build so a compaction-mode toggle is visible on rebuild. */ getCompactor: () => Compactor; + /** + * Present when a resumed run record has an Anthropic-protocol cache write + * and the provider about to be called is the same protocol. Read at each + * build so an interrupt rebuild of the same session still folds before the + * next infer. A new session omits it. + */ + getCacheWriteSeed?: () => { at: number; model: string } | undefined; /** Assigns the runner's live agent/storage holders; keeps call sites unchanged. */ onBuilt: (agent: Agent, storage: ContextStore) => void; /** @@ -718,6 +729,17 @@ export function assembleChatAgent(wiring: ChatAgentWiring): AssembledChatAgent { ), }, }); + const seed = wiring.getCacheWriteSeed?.(); + const source = + seed === undefined ? undefined : lastCycleSourceFromRunModel(seed.model); + if (seed !== undefined && source !== undefined) { + const loaded = await storage.load(); + directorHolder.instance?.restoreCacheWrite({ + at: seed.at, + source, + turns: loaded.turns, + }); + } const admittedAgent = primaryArchive === undefined ? agent @@ -742,7 +764,9 @@ export interface SessionLifecycleWiring { getSessionId: () => string; getSource: () => InferenceSource; initialTurnCount?: number | undefined; - onTurnBoundarySnapshot: () => void; + onTurnBoundarySnapshot: ( + event: Extract, + ) => void; hookEnabled?: Record | undefined; onHookEvent?: ((event: LifecycleHookEvent) => void) | undefined; resolveContextDir: () => string; diff --git a/src/session/cache-ttl-resume.test.ts b/src/session/cache-ttl-resume.test.ts new file mode 100644 index 000000000..94697e61d --- /dev/null +++ b/src/session/cache-ttl-resume.test.ts @@ -0,0 +1,267 @@ +import { mkdtemp, rm } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { describe, expect, test } from "bun:test"; +import { createInboundMessage } from "@intx/mime"; +import { createReactor, type ReactorEmittedEvent } from "@intx/inference"; +import { createDefaultDependencies } from "@intx/inference/providers"; +import type { + ConversationTurn, + InferenceEvent, + ReactorAction, +} from "@intx/types/runtime"; +import { createChatDirector } from "../agent/director.js"; +import { + COMPACTION_CONTINUATION_EVENT, + lastCycleSourceFromRunModel, +} from "../agent/compaction.js"; +import { HANDOFF_LATEST_KEY } from "./compaction-handoff.js"; +import { COMPACTED_PREFIX, createPruningCompactor } from "./compactor.js"; +import { createOptimizedContextStore } from "./optimized-context-store.js"; +import { buildCompactionContinuationMessage } from "./runtime-assembly.js"; + +const MINUTE_MS = 60_000; +const USAGE = { + input: 1, + output: 1, + cacheRead: 0, + cacheWrite: 0, + thinking: 0, +}; + +function userTurn(text: string, timestamp: number): ConversationTurn { + return { role: "user", content: [{ type: "text", text }], timestamp }; +} + +function textLength(turns: readonly ConversationTurn[]): number { + let total = 0; + for (const turn of turns) { + for (const block of turn.content) { + if (block.type === "text") total += block.text.length; + } + } + return total; +} + +function texts(turns: readonly ConversationTurn[]): string[] { + return turns.flatMap((turn) => + turn.content.flatMap((block) => + block.type === "text" ? [block.text] : [], + ), + ); +} + +function seedTurns(): ConversationTurn[] { + const filler = "prior-detail ".repeat(400); + return Array.from({ length: 12 }, (_, index) => + userTurn(`turn-${index} ${filler}`, index + 1), + ); +} + +function actionsOf(result: ReactorAction | ReactorAction[]): ReactorAction[] { + return Array.isArray(result) ? result : [result]; +} + +async function resumeAndInfer(args: { at: number; model: string }): Promise<{ + first: ReactorAction[]; + stored: ConversationTurn[]; + prompt: ConversationTurn[]; + handoff: string | undefined; + seed: ConversationTurn[]; +}> { + const seed = seedTurns(); + const dir = await mkdtemp(join(tmpdir(), "cache-ttl-resume-")); + try { + const store = await createOptimizedContextStore(dir); + await store.writeTurns(seed); + await store.writeMetadata({ + pendingOperations: [], + tokenUsage: USAGE, + }); + await store.commit({ message: "seed" }); + + const source = lastCycleSourceFromRunModel(args.model); + if (source === undefined) throw new Error("missing source"); + const director = createChatDirector("test", [], {}); + director.restoreCacheWrite({ at: args.at, source, turns: seed }); + const original = director.decide.bind(director); + let first: ReactorAction[] | undefined; + director.decide = async (event, state, capabilities) => { + const result = await original(event, state, capabilities); + if (first === undefined) first = actionsOf(result); + return result; + }; + + let prompt: ConversationTurn[] | undefined; + let storedAtInfer: ConversationTurn[] | undefined; + let handoff: string | undefined; + const events: ReactorEmittedEvent[] = []; + let resolveDone: (() => void) | undefined; + const done = new Promise((resolve) => { + resolveDone = resolve; + }); + + const reactor = createReactor({ + sessionId: "cache-ttl-resume", + director, + source: { + id: "anthropic:claude-opus-4-6", + provider: "anthropic", + model: "claude-opus-4-6", + baseURL: "https://example.invalid", + credentialId: "test", + }, + toolRunner: { + async run(call) { + return { callId: call.id, content: "ok" }; + }, + }, + contextStore: store, + deps: createDefaultDependencies(), + doomLoopThreshold: false, + compactors: { + "pruning-compactor": createPruningCompactor({ + summarize: async (turns) => { + const body = texts(turns); + return `${(body[0] ?? "").slice(0, 400)}\n${(body.at(-1) ?? "").slice(0, 200)}`; + }, + }), + }, + inferenceRunner(opts) { + return (async function* () { + prompt = opts.turns; + const loaded = await store.load(); + storedAtInfer = loaded.turns; + try { + const bytes = await store.readBlob(HANDOFF_LATEST_KEY); + handoff = new TextDecoder().decode(bytes); + } catch { + handoff = undefined; + } + const event: InferenceEvent = { + type: "inference.done", + seq: opts.nextSeq(), + data: { + turn: { + role: "assistant", + content: [{ type: "text", text: "ok" }], + model: "claude-opus-4-6", + timestamp: 2, + }, + usage: USAGE, + source: { + sourceId: source.sourceId, + provider: source.provider, + model: source.model, + }, + }, + }; + yield event; + })(); + }, + onEvent(event) { + events.push(event); + if (event.type === COMPACTION_CONTINUATION_EVENT) { + reactor.deliver(buildCompactionContinuationMessage()); + } + if (event.type === "inference.done" || event.type === "reactor.error") { + resolveDone?.(); + } + }, + }); + + reactor.start(); + const started = events.some((event) => event.type === "reactor.start"); + if (!started) { + await new Promise((resolve, reject) => { + const timer = setTimeout( + () => reject(new Error("reactor did not start")), + 5_000, + ); + const poll = (): void => { + if (events.some((event) => event.type === "reactor.start")) { + clearTimeout(timer); + resolve(); + return; + } + setTimeout(poll, 10); + }; + poll(); + }); + } + reactor.deliver( + createInboundMessage({ + from: "user@local", + to: "agent@local", + content: "continue the task", + }), + ); + await Promise.race([ + done, + new Promise((_, reject) => { + setTimeout(() => reject(new Error("resume infer timed out")), 15_000); + }), + ]); + reactor.abort("admin_kill"); + if (first === undefined) throw new Error("director never decided"); + if (prompt === undefined || storedAtInfer === undefined) { + throw new Error( + `infer did not run: ${events.map((event) => event.type).join(",")}`, + ); + } + return { first, stored: storedAtInfer, prompt, handoff, seed }; + } finally { + await rm(dir, { recursive: true, force: true }); + } +} + +describe("resumed Anthropic cache write compacts before infer", () => { + test("an expired Anthropic stamp folds into the context store first", async () => { + const result = await resumeAndInfer({ + at: Date.now() - 6 * MINUTE_MS, + model: "anthropic:claude-opus-4-6", + }); + + expect(result.first[0]).toMatchObject({ + type: "compact", + compactor: "pruning-compactor", + reason: "cache-ttl-recompress", + }); + expect(result.stored.length).toBeLessThan(result.seed.length); + expect(textLength(result.stored)).toBeLessThan(textLength(result.seed)); + expect(textLength(result.prompt)).toBeLessThan(textLength(result.seed)); + expect(texts(result.prompt)).toEqual(texts(result.stored)); + expect( + texts(result.stored).some((text) => text.startsWith(COMPACTED_PREFIX)), + ).toBe(true); + expect(result.handoff?.length ?? 0).toBeGreaterThan(0); + expect(result.handoff).toContain(HANDOFF_LATEST_KEY); + expect(result.handoff).toContain("# Compaction handoff"); + }); + + test("a write inside 5 minutes leaves the stored turns in place", async () => { + const result = await resumeAndInfer({ + at: Date.now() - 2 * MINUTE_MS, + model: "anthropic:claude-opus-4-6", + }); + + expect(result.first.some((action) => action.type === "compact")).toBe( + false, + ); + expect(texts(result.stored)).toEqual(texts(result.seed)); + expect(result.handoff).toBeUndefined(); + }); + + test("a non-Anthropic stored stamp does not compact", async () => { + const result = await resumeAndInfer({ + at: Date.now() - 6 * MINUTE_MS, + model: "openai:gpt-5.6", + }); + + expect(result.first.some((action) => action.type === "compact")).toBe( + false, + ); + expect(texts(result.stored)).toEqual(texts(result.seed)); + expect(result.handoff).toBeUndefined(); + }); +}); diff --git a/src/session/run-sink.ts b/src/session/run-sink.ts index efd3c85fa..0187cb396 100644 --- a/src/session/run-sink.ts +++ b/src/session/run-sink.ts @@ -57,7 +57,11 @@ export interface RunSinkArgs { // alongside the turn count it reports, rather than in a second // subscription to the same event stream in a renderer: the renderer has // already been swapped out from under this constraint three times. - onTurnBoundarySnapshot?: () => void; + // The event is the inference that just finished, so the snapshot can stamp + // an Anthropic cache write before the director's own bookkeeping runs. + onTurnBoundarySnapshot?: ( + event: Extract, + ) => void; } export interface RunSink { @@ -192,7 +196,7 @@ export function createRunSink(args: RunSinkArgs): RunSink { turnInFlight = false; pendingInferenceError = undefined; runError = undefined; - onTurnBoundarySnapshot?.(); + onTurnBoundarySnapshot?.(event); } if (event.type === "reactor.error" && isReactorErrorFatal(event.data)) { const data = event.data as { error: string }; diff --git a/src/session/state.ts b/src/session/state.ts index 13a372245..231ddb03e 100644 --- a/src/session/state.ts +++ b/src/session/state.ts @@ -50,6 +50,10 @@ const RunStateSchema = type({ // so a resume can re-activate them before the first post-resume inference — // the transcript still tells the model they are callable. "activatedTools?": "string[]", + // Wall time of the last Anthropic-protocol inference. A later process + // folds before its first infer once this is at least the published TTL old. + // Absent for other providers and for records written before this field. + "lastCacheWriteAt?": "number", }); export type RunState = typeof RunStateSchema.infer; diff --git a/src/tui/runner/exit.ts b/src/tui/runner/exit.ts index 76da25a6c..f5d07c840 100644 --- a/src/tui/runner/exit.ts +++ b/src/tui/runner/exit.ts @@ -232,6 +232,9 @@ function createRunPersistence(state: RunnerState, services: RunnerServices) { model, mcpServers: state.connectedMcpServers, ...(activatedTools.length > 0 ? { activatedTools } : {}), + ...(services.activeRunHandle.lastCacheWriteAt !== undefined + ? { lastCacheWriteAt: services.activeRunHandle.lastCacheWriteAt } + : {}), ...extra, }; if (clearsActiveRun(kind)) { @@ -715,6 +718,8 @@ export async function createRunLifecycle( model: `${rotatedBundle.selected.id}:${rotatedBundle.selected.model}`, activatedTools: [], }); + delete services.activeRunHandle.lastCacheWriteAt; + delete services.activeRunHandle.cacheWriteModel; services.emitter.emit( "session.title", state.runTaskTitle.trim().length > 0 diff --git a/src/tui/runner/session.ts b/src/tui/runner/session.ts index ff76613b6..2333eb235 100644 --- a/src/tui/runner/session.ts +++ b/src/tui/runner/session.ts @@ -56,6 +56,10 @@ import { resolveLiveSessionSources, type LiveSessionSources, } from "../../session/assemble-runtime.js"; +import { + anthropicCacheWriteAt, + resumeCacheWriteSeed, +} from "../../provider/cache-ttl.js"; import type { CompactionArchive } from "../../session/compaction-archive.js"; import { tryReadPriorHandoffFile } from "../../session/compaction-handoff.js"; import { @@ -156,7 +160,14 @@ export async function assembleTUISession( // persistRunSnapshot onto the state this closure reads it from. getSource: () => state.liveSource, initialTurnCount: start.resumeSeed.turnsUsed, - onTurnBoundarySnapshot: () => { + onTurnBoundarySnapshot: (event) => { + if (event?.type === "inference.done") { + const at = anthropicCacheWriteAt(event.data.source, Date.now()); + if (at !== undefined) { + start.activeRunHandle.lastCacheWriteAt = at; + start.activeRunHandle.cacheWriteModel = `${event.data.source.provider}:${event.data.source.model}`; + } + } void state.persistRunSnapshot?.("running"); }, hookEnabled: initialHookEnabled, @@ -722,6 +733,12 @@ export async function assembleTUISession( }, }), ), + getCacheWriteSeed: () => + resumeCacheWriteSeed({ + at: start.activeRunHandle.lastCacheWriteAt, + storedModel: start.activeRunHandle.cacheWriteModel, + liveProvider: state.config.providerName, + }), onBuilt: (agent, storage) => { state.currentAgent = agent; state.currentStorage = storage; diff --git a/src/tui/session-start.ts b/src/tui/session-start.ts index bd4d5ef2e..e1b429e89 100644 --- a/src/tui/session-start.ts +++ b/src/tui/session-start.ts @@ -47,6 +47,9 @@ export interface ResumeSeed { // tool_search-promoted tool names from the prior run, re-activated before // the first post-resume inference so the wire matches the transcript. activatedTools: string[]; + // Present only when the prior run stamped an Anthropic-protocol cache write. + lastCacheWriteAt?: number; + cacheWriteModel?: string; } const FRESH_RESUME_SEED: ResumeSeed = { @@ -69,6 +72,14 @@ export function resolveResumeSeed(pickedState: RunState | null): ResumeSeed { turnsUsed: pickedState.turnsUsed, mcpServers: pickedState.mcpServers ?? [], activatedTools: pickedState.activatedTools ?? [], + ...(pickedState.lastCacheWriteAt !== undefined + ? { + lastCacheWriteAt: pickedState.lastCacheWriteAt, + ...(pickedState.model !== undefined + ? { cacheWriteModel: pickedState.model } + : {}), + } + : {}), }; } @@ -122,6 +133,7 @@ export function createTUICrashGuard( // too to close that earlier window. const turnsUsed = getActiveRun()?.turnsUsed ?? 0; const activatedTools = getActiveRun()?.activatedTools; + const lastCacheWriteAt = getActiveRun()?.lastCacheWriteAt; clearActiveRun(); clearActiveDisposeHost(); await flushPartialOnCrash().catch((flushErr: unknown) => { @@ -151,6 +163,7 @@ export function createTUICrashGuard( model: `${live.providerName}:${live.model}`, mcpServers: [], ...(activatedTools !== undefined ? { activatedTools } : {}), + ...(lastCacheWriteAt !== undefined ? { lastCacheWriteAt } : {}), }).catch((saveErr: unknown) => { const saveMessage = saveErr instanceof Error ? saveErr.message : String(saveErr); @@ -294,6 +307,9 @@ export async function prepareTUISession( model: `${config.providerName}:${config.model}`, mcpServers: resumeSeed.mcpServers, activatedTools: resumeSeed.activatedTools, + ...(resumeSeed.lastCacheWriteAt !== undefined + ? { lastCacheWriteAt: resumeSeed.lastCacheWriteAt } + : {}), }); // Registered the moment a run starts so the top-level uncaughtException / @@ -311,6 +327,12 @@ export async function prepareTUISession( turnsUsed: resumeSeed.turnsUsed, model: `${config.providerName}:${config.model}`, activatedTools: resumeSeed.activatedTools, + ...(resumeSeed.lastCacheWriteAt !== undefined + ? { lastCacheWriteAt: resumeSeed.lastCacheWriteAt } + : {}), + ...(resumeSeed.cacheWriteModel !== undefined + ? { cacheWriteModel: resumeSeed.cacheWriteModel } + : {}), }; setActiveRun(activeRunHandle); From 4ab98d88afcf4248a27c2aefb6da4ebc3d7dc0c4 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Wed, 23 Sep 2026 19:04:07 -0700 Subject: [PATCH 3/8] fix(provider): resume anthropic cache stamps from the protocol Runners persist the catalog id, which is not the messages adapter name. Requiring that id to match the protocol dropped an expired stamp, so the next request sent the untrimmed history. --- src/exec/runner.ts | 1 + src/provider/cache-ttl.test.ts | 65 +++++++++++++++++++ src/provider/cache-ttl.ts | 45 +++++++++++--- src/session/cache-ttl-resume.test.ts | 93 +++++++++++++++++++++++----- src/tui/runner/session.ts | 1 + 5 files changed, 180 insertions(+), 25 deletions(-) diff --git a/src/exec/runner.ts b/src/exec/runner.ts index a6f671acc..ef09f4f1e 100644 --- a/src/exec/runner.ts +++ b/src/exec/runner.ts @@ -891,6 +891,7 @@ export async function runExec(config: Config): Promise { at: activeRunHandle.lastCacheWriteAt, storedModel: activeRunHandle.cacheWriteModel, liveProvider: config.providerName, + liveProtocol: liveSource.provider, }), onBuilt: (agent, storage) => { currentAgent = agent; diff --git a/src/provider/cache-ttl.test.ts b/src/provider/cache-ttl.test.ts index fbceb7934..9582ce93e 100644 --- a/src/provider/cache-ttl.test.ts +++ b/src/provider/cache-ttl.test.ts @@ -1,4 +1,9 @@ import { describe, expect, test } from "bun:test"; +import { + buildAnthropicSource, + buildGoSource, + buildZenSource, +} from "../config/index.js"; import { anthropicCacheWriteAt, cacheTtlMsFor, @@ -79,6 +84,7 @@ describe("anthropic cache-write stamp", () => { at: 10, storedModel: "anthropic:claude-opus-4-6", liveProvider: "anthropic", + liveProtocol: "anthropic", }), ).toEqual({ at: 10, model: "anthropic:claude-opus-4-6" }); expect( @@ -86,6 +92,7 @@ describe("anthropic cache-write stamp", () => { at: 10, storedModel: "openai:gpt-5.6", liveProvider: "openai", + liveProtocol: "openai", }), ).toBeUndefined(); expect( @@ -93,6 +100,7 @@ describe("anthropic cache-write stamp", () => { at: 10, storedModel: "anthropic:claude-opus-4-6", liveProvider: "openai", + liveProtocol: "openai", }), ).toBeUndefined(); expect( @@ -100,7 +108,64 @@ describe("anthropic cache-write stamp", () => { at: undefined, storedModel: "anthropic:claude-opus-4-6", liveProvider: "anthropic", + liveProtocol: "anthropic", }), ).toBeUndefined(); }); + + test("catalog ids that speak the Anthropic protocol still return the stamp", () => { + const at = 10; + const sources = [ + buildZenSource({ + id: "zen", + model: "claude-opus-4-6", + sessionId: "sess-zen", + }), + buildGoSource({ + id: "opencode-go", + model: "minimax-m3", + sessionId: "sess-go", + }), + buildAnthropicSource({ + id: "acme", + baseURL: "https://example.invalid", + model: "claude-opus-4-6", + }), + ]; + for (const source of sources) { + expect(source.id).not.toBe(source.provider); + const storedModel = `${source.id}:${source.model}`; + expect( + resumeCacheWriteSeed({ + at, + storedModel, + liveProvider: source.id, + liveProtocol: source.provider, + }), + ).toEqual({ at, model: `${source.provider}:${source.model}` }); + } + + const gemini = buildZenSource({ + id: "zen", + model: "gemini-3-flash", + sessionId: "sess-zen", + }); + expect( + resumeCacheWriteSeed({ + at, + storedModel: `${gemini.id}:${gemini.model}`, + liveProvider: gemini.id, + liveProtocol: gemini.provider, + }), + ).toBeUndefined(); + + expect( + resumeCacheWriteSeed({ + at, + storedModel: "zen-messages:claude-opus-4-6", + liveProvider: "zen", + liveProtocol: "zen-messages", + }), + ).toEqual({ at, model: "zen-messages:claude-opus-4-6" }); + }); }); diff --git a/src/provider/cache-ttl.ts b/src/provider/cache-ttl.ts index 9dcb1d9ae..8cf8621e8 100644 --- a/src/provider/cache-ttl.ts +++ b/src/provider/cache-ttl.ts @@ -98,19 +98,48 @@ export function anthropicCacheWriteAt( } /** - * Seed for a resumed session's pre-infer fold. Absent unless both the - * identity that wrote the cache and the provider about to be called are - * Anthropic-protocol. A stamp alone is not enough: a non-Anthropic model - * on the run record must not fold. + * Seed for a resumed session's pre-infer fold. Absent unless the provider + * about to be called speaks the Anthropic messages protocol. A stamp alone + * is not enough: a non-Anthropic live adapter must not fold. + * + * `liveProvider` is the catalog id (`source.id`). Runners persist + * `${id}:${model}`, and that id is `zen`, `opencode-go`, or a custom account + * name — not `zen-messages`, `opencode-go-messages`, or `anthropic`. + * `liveProtocol` is `source.provider`, the adapter that actually wrote the + * cache. A catalog-form stamp is accepted only for this same catalog id, and + * the returned model uses the protocol name so the governor's TTL check + * follows the adapter. A protocol-form stamp (the in-memory `provider:model` + * recorded at inference.done) is kept as-is. */ export function resumeCacheWriteSeed(args: { at: number | undefined; storedModel: string | undefined; liveProvider: string; + liveProtocol: string; }): { at: number; model: string } | undefined { if (args.at === undefined) return undefined; - if (args.storedModel === undefined) return undefined; - if (cacheTtlMsFor(args.storedModel) === undefined) return undefined; - if (cacheTtlMsFor(args.liveProvider) === undefined) return undefined; - return { at: args.at, model: args.storedModel }; + if (args.storedModel === undefined || args.storedModel.length === 0) { + return undefined; + } + if (cacheTtlMsFor(args.liveProtocol) === undefined) return undefined; + if (cacheTtlMsFor(args.storedModel) !== undefined) { + return { at: args.at, model: args.storedModel }; + } + + const stored = splitRunModel(args.storedModel); + if (stored.provider !== args.liveProvider) return undefined; + if (stored.model.length === 0) return undefined; + return { at: args.at, model: `${args.liveProtocol}:${stored.model}` }; +} + +function splitRunModel(value: string): { provider: string; model: string } { + const colon = value.indexOf(":"); + if (colon > 0) { + return { provider: value.slice(0, colon), model: value.slice(colon + 1) }; + } + const slash = value.indexOf("/"); + if (slash > 0) { + return { provider: value.slice(0, slash), model: value.slice(slash + 1) }; + } + return { provider: value, model: "" }; } diff --git a/src/session/cache-ttl-resume.test.ts b/src/session/cache-ttl-resume.test.ts index 94697e61d..45bef180a 100644 --- a/src/session/cache-ttl-resume.test.ts +++ b/src/session/cache-ttl-resume.test.ts @@ -8,6 +8,7 @@ import { createDefaultDependencies } from "@intx/inference/providers"; import type { ConversationTurn, InferenceEvent, + InferenceSource, ReactorAction, } from "@intx/types/runtime"; import { createChatDirector } from "../agent/director.js"; @@ -15,6 +16,13 @@ import { COMPACTION_CONTINUATION_EVENT, lastCycleSourceFromRunModel, } from "../agent/compaction.js"; +import { + buildAnthropicSource, + buildGoSource, + buildOpenAISource, + buildZenSource, +} from "../config/index.js"; +import { resumeCacheWriteSeed } from "../provider/cache-ttl.js"; import { HANDOFF_LATEST_KEY } from "./compaction-handoff.js"; import { COMPACTED_PREFIX, createPruningCompactor } from "./compactor.js"; import { createOptimizedContextStore } from "./optimized-context-store.js"; @@ -62,7 +70,10 @@ function actionsOf(result: ReactorAction | ReactorAction[]): ReactorAction[] { return Array.isArray(result) ? result : [result]; } -async function resumeAndInfer(args: { at: number; model: string }): Promise<{ +async function resumeAndInfer(args: { + at: number; + source: InferenceSource; +}): Promise<{ first: ReactorAction[]; stored: ConversationTurn[]; prompt: ConversationTurn[]; @@ -80,10 +91,21 @@ async function resumeAndInfer(args: { at: number; model: string }): Promise<{ }); await store.commit({ message: "seed" }); - const source = lastCycleSourceFromRunModel(args.model); - if (source === undefined) throw new Error("missing source"); + const storedModel = `${args.source.id}:${args.source.model}`; + const cacheSeed = resumeCacheWriteSeed({ + at: args.at, + storedModel, + liveProvider: args.source.id, + liveProtocol: args.source.provider, + }); + const source = + cacheSeed === undefined + ? undefined + : lastCycleSourceFromRunModel(cacheSeed.model); const director = createChatDirector("test", [], {}); - director.restoreCacheWrite({ at: args.at, source, turns: seed }); + if (cacheSeed !== undefined && source !== undefined) { + director.restoreCacheWrite({ at: cacheSeed.at, source, turns: seed }); + } const original = director.decide.bind(director); let first: ReactorAction[] | undefined; director.decide = async (event, state, capabilities) => { @@ -104,13 +126,7 @@ async function resumeAndInfer(args: { at: number; model: string }): Promise<{ const reactor = createReactor({ sessionId: "cache-ttl-resume", director, - source: { - id: "anthropic:claude-opus-4-6", - provider: "anthropic", - model: "claude-opus-4-6", - baseURL: "https://example.invalid", - credentialId: "test", - }, + source: args.source, toolRunner: { async run(call) { return { callId: call.id, content: "ok" }; @@ -150,9 +166,9 @@ async function resumeAndInfer(args: { at: number; model: string }): Promise<{ }, usage: USAGE, source: { - sourceId: source.sourceId, - provider: source.provider, - model: source.model, + sourceId: args.source.id, + provider: args.source.provider, + model: args.source.model, }, }, }; @@ -216,10 +232,16 @@ async function resumeAndInfer(args: { at: number; model: string }): Promise<{ } describe("resumed Anthropic cache write compacts before infer", () => { + const native = buildAnthropicSource({ + id: "anthropic", + baseURL: "https://example.invalid", + model: "claude-opus-4-6", + }); + test("an expired Anthropic stamp folds into the context store first", async () => { const result = await resumeAndInfer({ at: Date.now() - 6 * MINUTE_MS, - model: "anthropic:claude-opus-4-6", + source: native, }); expect(result.first[0]).toMatchObject({ @@ -242,7 +264,7 @@ describe("resumed Anthropic cache write compacts before infer", () => { test("a write inside 5 minutes leaves the stored turns in place", async () => { const result = await resumeAndInfer({ at: Date.now() - 2 * MINUTE_MS, - model: "anthropic:claude-opus-4-6", + source: native, }); expect(result.first.some((action) => action.type === "compact")).toBe( @@ -255,7 +277,11 @@ describe("resumed Anthropic cache write compacts before infer", () => { test("a non-Anthropic stored stamp does not compact", async () => { const result = await resumeAndInfer({ at: Date.now() - 6 * MINUTE_MS, - model: "openai:gpt-5.6", + source: buildOpenAISource({ + id: "openai", + baseURL: "https://example.invalid/v1", + model: "gpt-5.6", + }), }); expect(result.first.some((action) => action.type === "compact")).toBe( @@ -264,4 +290,37 @@ describe("resumed Anthropic cache write compacts before infer", () => { expect(texts(result.stored)).toEqual(texts(result.seed)); expect(result.handoff).toBeUndefined(); }); + + test("expired Zen, OpenCode Go, and custom Anthropic catalog ids fold before infer", async () => { + const sources = [ + buildZenSource({ + id: "zen", + model: "claude-opus-4-6", + sessionId: "sess-zen", + }), + buildGoSource({ + id: "opencode-go", + model: "minimax-m3", + sessionId: "sess-go", + }), + buildAnthropicSource({ + id: "acme", + baseURL: "https://example.invalid", + model: "claude-opus-4-6", + }), + ]; + for (const source of sources) { + expect(source.id).not.toBe(source.provider); + const result = await resumeAndInfer({ + at: Date.now() - 6 * MINUTE_MS, + source, + }); + expect(result.first[0]).toMatchObject({ + type: "compact", + compactor: "pruning-compactor", + reason: "cache-ttl-recompress", + }); + expect(result.stored.length).toBeLessThan(result.seed.length); + } + }); }); diff --git a/src/tui/runner/session.ts b/src/tui/runner/session.ts index 2333eb235..d44eb90e8 100644 --- a/src/tui/runner/session.ts +++ b/src/tui/runner/session.ts @@ -738,6 +738,7 @@ export async function assembleTUISession( at: start.activeRunHandle.lastCacheWriteAt, storedModel: start.activeRunHandle.cacheWriteModel, liveProvider: state.config.providerName, + liveProtocol: state.liveSource.provider, }), onBuilt: (agent, storage) => { state.currentAgent = agent; From e4744a0b331fb4ae741ea8cabc341d623c3b0a00 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Wed, 23 Sep 2026 19:58:39 -0700 Subject: [PATCH 4/8] feat(session): shrink anthropic prompts once the cache write expires --- src/session/anthropic-cache-prompt.test.ts | 147 ++++++++++++++++++ src/session/anthropic-cache-prompt.ts | 166 +++++++++++++++++++++ src/session/assemble-runtime.ts | 12 ++ 3 files changed, 325 insertions(+) create mode 100644 src/session/anthropic-cache-prompt.test.ts create mode 100644 src/session/anthropic-cache-prompt.ts diff --git a/src/session/anthropic-cache-prompt.test.ts b/src/session/anthropic-cache-prompt.test.ts new file mode 100644 index 000000000..78762c778 --- /dev/null +++ b/src/session/anthropic-cache-prompt.test.ts @@ -0,0 +1,147 @@ +import { mkdtemp, readFile } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { describe, expect, test } from "bun:test"; +import type { ConversationTurn } from "@intx/types/runtime"; + +import { createAnthropicCachePromptTransform } from "./anthropic-cache-prompt.js"; +import { createOptimizedContextStore } from "./optimized-context-store.js"; + +const MINUTE_MS = 60_000; +const NOW = 1_000_000_000_000; +const BODY = "file-body-".repeat(2_000); + +const CTX = { state: {} as never, trigger: "test" }; + +function transformFor(args: { + protocol: string; + cacheWriteAt: number | undefined; + nowMs?: number; +}) { + return createAnthropicCachePromptTransform({ + nowMs: () => args.nowMs ?? NOW, + cacheWriteAt: () => args.cacheWriteAt, + protocol: () => args.protocol, + }); +} + +function history(): ConversationTurn[] { + return [ + { + role: "user", + content: [{ type: "text", text: "read the sources" }], + timestamp: 1, + }, + { + role: "assistant", + content: [ + { + type: "tool_call", + id: "read-1", + name: "read_file", + arguments: { path: "src/big.ts" }, + }, + ], + timestamp: 2, + }, + { + role: "user", + content: [ + { + type: "tool_result", + callId: "read-1", + content: [{ type: "text", text: BODY }], + }, + ], + timestamp: 3, + }, + { + role: "assistant", + content: [ + { + type: "tool_call", + id: "bash-1", + name: "bash", + arguments: { command: "wc -l src/big.ts" }, + }, + ], + timestamp: 4, + }, + { + role: "user", + content: [ + { + type: "tool_result", + callId: "bash-1", + content: [{ type: "text", text: BODY }], + }, + ], + timestamp: 5, + }, + { + role: "assistant", + content: [{ type: "text", text: "newest assistant" }], + timestamp: 6, + }, + { + role: "user", + content: [{ type: "text", text: "newest user" }], + timestamp: 7, + }, + ]; +} + +describe("anthropic cache prompt transform", () => { + test("expired or missing Anthropic stamp stubs old tool bodies and leaves stored turns", async () => { + const dir = await mkdtemp(join(tmpdir(), "anthropic-cache-prompt-")); + const store = await createOptimizedContextStore(dir); + const turns = history(); + await store.writeTurns(turns); + const stored = (await store.load()).turns; + const diskBefore = await readFile(join(dir, "turns.jsonl"), "utf8"); + const storedJson = JSON.stringify(stored); + + for (const cacheWriteAt of [undefined, NOW - 5 * MINUTE_MS]) { + const result = await transformFor({ + protocol: "anthropic", + cacheWriteAt, + }).apply(stored, CTX); + expect(JSON.stringify(result.output).length).toBeLessThan( + storedJson.length, + ); + expect(result.output).toHaveLength(stored.length); + expect(JSON.stringify(result.output)).not.toContain(BODY); + expect(JSON.stringify(result.output)).toContain("newest assistant"); + expect(JSON.stringify(result.output)).toContain("newest user"); + expect(JSON.stringify(result.output)).toContain( + "[read_file result omitted]", + ); + expect(JSON.stringify(result.output)).toContain("[bash result omitted]"); + expect(result.record.reason).toBe("stubbed-tool-results"); + } + + expect(JSON.stringify((await store.load()).turns)).toBe(storedJson); + expect(await readFile(join(dir, "turns.jsonl"), "utf8")).toBe(diskBefore); + expect(storedJson).toContain(BODY); + }); + + test("a stamp two minutes old returns the same turns", async () => { + const turns = history(); + const result = await transformFor({ + protocol: "anthropic", + cacheWriteAt: NOW - 2 * MINUTE_MS, + }).apply(turns, CTX); + expect(result.output).toBe(turns); + expect(result.record.reason).toBe("cache-warm"); + }); + + test("openai leaves the prompt unchanged when the stamp is old", async () => { + const turns = history(); + const result = await transformFor({ + protocol: "openai", + cacheWriteAt: NOW - 60 * MINUTE_MS, + }).apply(turns, CTX); + expect(result.output).toBe(turns); + expect(result.record.reason).toBe("non-anthropic"); + }); +}); diff --git a/src/session/anthropic-cache-prompt.ts b/src/session/anthropic-cache-prompt.ts new file mode 100644 index 000000000..768e28f42 --- /dev/null +++ b/src/session/anthropic-cache-prompt.ts @@ -0,0 +1,166 @@ +// Prompt-only shrink after an Anthropic ephemeral cache write expires. +// Compaction rewrites turns.jsonl and keeps a raw tail, so the next infer +// still cache-writes those tool bodies. This transform runs inside +// executeInfer: its output is what gets written to prompt.jsonl. It never +// writes turns and never calls the compaction governor. + +import type { + ContentBlock, + ContextTransform, + ConversationTurn, +} from "@intx/types/runtime"; + +import { cacheTtlMsFor } from "../provider/cache-ttl.js"; + +export type AnthropicCachePromptDeps = { + nowMs: () => number; + /** Wall time of the last Anthropic-protocol cache write, if one was stamped. */ + cacheWriteAt: () => number | undefined; + /** Live adapter protocol (`InferenceSource.provider`). */ + protocol: () => string | undefined; +}; + +const STRATEGY = "anthropic-cache-prompt"; + +type ToolResultBlock = Extract; + +function passthrough( + turns: ConversationTurn[], + reason: string, +): Awaited> { + return { + output: turns, + record: { + strategy: STRATEGY, + version: "1", + parameters: {}, + reason, + decisions: { stubbed: 0 }, + }, + }; +} + +function newestAssistantTextIndex(turns: readonly ConversationTurn[]): number { + for (let i = turns.length - 1; i >= 0; i--) { + const turn = turns[i]; + if (turn === undefined || turn.role !== "assistant") continue; + if ( + turn.content.some( + (block) => block.type === "text" && block.text.length > 0, + ) + ) { + return i; + } + } + return -1; +} + +// Tool results before the newest assistant text are already consumed. +// When the model has not written that text yet, the newest result turn is +// the unconsumed suffix and stays whole. +function stubBoundary(turns: readonly ConversationTurn[]): number { + const assistantText = newestAssistantTextIndex(turns); + if (assistantText >= 0) return assistantText; + for (let i = turns.length - 1; i >= 0; i--) { + const turn = turns[i]; + if (turn?.content.some((block) => block.type === "tool_result")) return i; + } + return turns.length; +} + +function toolNames(turns: readonly ConversationTurn[]): Map { + const names = new Map(); + for (const turn of turns) { + for (const block of turn.content) { + if (block.type === "tool_call") names.set(block.id, block.name); + } + } + return names; +} + +function stubText(name: string | undefined): string { + return name === undefined + ? "[tool result omitted]" + : `[${name} result omitted]`; +} + +function stubToolResult( + block: ToolResultBlock, + names: ReadonlyMap, +): { block: ToolResultBlock; stubbed: boolean } { + const text = stubText(names.get(block.callId)); + if (JSON.stringify(block.content).length <= text.length) { + return { block, stubbed: false }; + } + const next: ToolResultBlock = { + type: "tool_result", + callId: block.callId, + content: [{ type: "text", text }], + }; + if (block.isError === true) next.isError = true; + return { block: next, stubbed: true }; +} + +function stubTurn( + turn: ConversationTurn, + names: ReadonlyMap, +): { turn: ConversationTurn; stubbed: number } { + let stubbed = 0; + let changed = false; + const content = turn.content.map((block) => { + if (block.type !== "tool_result") return block; + const next = stubToolResult(block, names); + if (!next.stubbed) return block; + changed = true; + stubbed += 1; + return next.block; + }); + return { turn: changed ? { ...turn, content } : turn, stubbed }; +} + +function shrinkPrompt(turns: ConversationTurn[]): { + output: ConversationTurn[]; + stubbed: number; +} { + const boundary = stubBoundary(turns); + const names = toolNames(turns); + let stubbed = 0; + let changed = false; + const output = turns.map((turn, index) => { + if (index >= boundary) return turn; + const next = stubTurn(turn, names); + stubbed += next.stubbed; + if (next.turn !== turn) changed = true; + return next.turn; + }); + if (!changed) return { output: turns, stubbed: 0 }; + return { output, stubbed }; +} + +export function createAnthropicCachePromptTransform( + deps: AnthropicCachePromptDeps, +): ContextTransform { + return { + name: STRATEGY, + version: "1", + async apply(turns, _ctx) { + const ttl = cacheTtlMsFor(deps.protocol()); + if (ttl === undefined) return passthrough(turns, "non-anthropic"); + const at = deps.cacheWriteAt(); + if (at !== undefined && deps.nowMs() - at < ttl) { + return passthrough(turns, "cache-warm"); + } + const shrunk = shrinkPrompt(turns); + return { + output: shrunk.output, + record: { + strategy: STRATEGY, + version: "1", + parameters: {}, + reason: shrunk.stubbed > 0 ? "stubbed-tool-results" : "noop", + decisions: { stubbed: shrunk.stubbed }, + }, + }; + }, + }; +} diff --git a/src/session/assemble-runtime.ts b/src/session/assemble-runtime.ts index e1bbc6f9b..71e8277d4 100644 --- a/src/session/assemble-runtime.ts +++ b/src/session/assemble-runtime.ts @@ -60,7 +60,9 @@ import { createDoomLoopCorrectiveNote } from "../agent/doom-loop-note.js"; import type { AgentToolset } from "../agent/tools.js"; import { createAgentWithLiveToolDispatch } from "../agent/live-tool-dispatch.js"; import { createSessionStores } from "./optimized-context-store.js"; +import { getActiveRun } from "./active-run.js"; import { createAttachmentRehydrateTransform } from "./attachment-store.js"; +import { createAnthropicCachePromptTransform } from "./anthropic-cache-prompt.js"; import { applyRecordingPolicyToText, createCompactionArchive, @@ -701,6 +703,16 @@ export function assembleChatAgent(wiring: ChatAgentWiring): AssembledChatAgent { createAttachmentRehydrateTransform((key) => storageForAgent.readBlob(key), ), + createAnthropicCachePromptTransform({ + nowMs: () => Date.now(), + cacheWriteAt: () => getActiveRun()?.lastCacheWriteAt, + protocol: () => { + const sources = wiring.getSources(); + const preferred = wiring.getDefaultSource(); + const match = sources.find((source) => source.id === preferred); + return (match ?? sources[0])?.provider; + }, + }), ], }, audit, From 7ba40872ed4ec5afcf64acc881268122f3d3c0ea Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Wed, 23 Sep 2026 21:36:24 -0700 Subject: [PATCH 5/8] feat(session): gate anthropic prompt shrink behind settings Cache expiry no longer compacts stored turns. The prompt transform runs only when anthropicCachePrompt is set. --- src/agent/compaction.test.ts | 64 +++------------------- src/agent/compaction.ts | 51 ++--------------- src/config/index.ts | 4 ++ src/config/settings.ts | 10 ++++ src/exec/runner.ts | 1 + src/session/anthropic-cache-prompt.test.ts | 13 +++++ src/session/anthropic-cache-prompt.ts | 5 ++ src/session/assemble-runtime.ts | 3 + src/session/cache-ttl-resume.test.ts | 48 +++++----------- src/tui/runner/session.ts | 1 + 10 files changed, 64 insertions(+), 136 deletions(-) diff --git a/src/agent/compaction.test.ts b/src/agent/compaction.test.ts index caa0b4bb0..06f9c7417 100644 --- a/src/agent/compaction.test.ts +++ b/src/agent/compaction.test.ts @@ -1020,10 +1020,8 @@ describe("provider-aware idle recompress (CL-8745)", () => { nowMs += MINUTE_MS + 1; expect( governor.interceptIdleContinuation(emptyMessage(), capabilities), - ).toEqual(ttlCompact); - expect(continuations).toBe(1); - // Empty continuation adopts the shrunk turns without a new inference. - expect(governor.resumeAfterCompact(emptyMessage())).toBe("meter"); + ).toBeNull(); + expect(continuations).toBe(0); }); test("production LastCycleSource: anthropic fires at 5m, codex never does", () => { @@ -1058,7 +1056,7 @@ describe("provider-aware idle recompress (CL-8745)", () => { nowMs += 5 * MINUTE_MS + 1; expect( anthropic.interceptIdleContinuation(emptyMessage(), capabilities), - ).toEqual(ttlCompact); + ).toBeNull(); expect( codex.interceptIdleContinuation(emptyMessage(), capabilities), ).toBeNull(); @@ -1181,7 +1179,7 @@ describe("provider-aware idle recompress (CL-8745)", () => { nowMs += 5 * MINUTE_MS + 1; expect( governor.interceptIdleContinuation(emptyMessage(), capabilities), - ).toEqual(ttlCompact); + ).toBeNull(); }); test("after threshold compact and gap TTL, growth-armed compact stays blocked until a tool_call", () => { @@ -1211,27 +1209,7 @@ describe("provider-aware idle recompress (CL-8745)", () => { nowMs += 5 * MINUTE_MS + 1; expect( governor.interceptIdleContinuation(emptyMessage(), capabilities), - ).toEqual(ttlCompact); - expect(governor.resumeAfterCompact(emptyMessage())).toBe("meter"); - - // Post-TTL snapshot, then growth past resumeDelta — pending re-arms, - // but the cap is full so interceptActions stays null. - governor.noteInferenceDone(inferenceDone(overThreshold), tenTurns); - governor.noteInferenceDone( - inferenceDone(overThreshold + resumeDelta), - tenTurns, - ); - expect( - governor.interceptActions(toolDone(), inferAction, capabilities), ).toBeNull(); - - governor.noteInferenceDone( - inferenceDoneWithTools(overThreshold + 2 * resumeDelta), - tenTurns, - ); - expect( - governor.interceptActions(toolDone(), inferAction, capabilities), - ).not.toBeNull(); }); test("does not fold under an outstanding tool batch, fires once it settles", () => { @@ -1257,8 +1235,8 @@ describe("provider-aware idle recompress (CL-8745)", () => { ).toBeNull(); expect( governor.interceptIdleContinuation(emptyMessage(), capabilities), - ).toEqual(ttlCompact); - expect(continuations).toBe(1); + ).toBeNull(); + expect(continuations).toBe(0); }); test("one fire per window, then the shared consecutive-compact cap stops the spiral", () => { @@ -1276,31 +1254,10 @@ describe("provider-aware idle recompress (CL-8745)", () => { ); nowMs += 5 * MINUTE_MS + 1; - expect( - governor.interceptIdleContinuation(emptyMessage(), capabilities), - ).toEqual(ttlCompact); - expect(governor.resumeAfterCompact(emptyMessage())).toBe("meter"); - governor.notePostCompact(tenTurns); - - // Inside the next window: no second fire. - nowMs += 4 * MINUTE_MS; expect( governor.interceptIdleContinuation(emptyMessage(), capabilities), ).toBeNull(); - - nowMs += MINUTE_MS + 1; - expect( - governor.interceptIdleContinuation(emptyMessage(), capabilities), - ).toEqual(ttlCompact); - expect(governor.resumeAfterCompact(emptyMessage())).toBe("meter"); - governor.notePostCompact(tenTurns); - - // Cap reached: the third window stays quiet. - nowMs += 5 * MINUTE_MS + 1; - expect( - governor.interceptIdleContinuation(emptyMessage(), capabilities), - ).toBeNull(); - expect(continuations).toBe(2); + expect(continuations).toBe(0); }); test("a raced operator message past the TTL still folds, then re-infers", () => { @@ -1319,11 +1276,8 @@ describe("provider-aware idle recompress (CL-8745)", () => { nowMs += 6 * MINUTE_MS; expect( governor.interceptIdleContinuation(racedMessage(), capabilities), - ).toEqual(ttlCompact); - expect(continuations).toBe(1); - // The follow-up empty continuation carries the infer that answers the - // raced question (mirrors the threshold raced path). - expect(governor.resumeAfterCompact(emptyMessage())).toBe("infer"); + ).toBeNull(); + expect(continuations).toBe(0); }); }); diff --git a/src/agent/compaction.ts b/src/agent/compaction.ts index b981cbc14..9525553fb 100644 --- a/src/agent/compaction.ts +++ b/src/agent/compaction.ts @@ -11,7 +11,6 @@ import { compactionThresholdFor, contextTokensFromUsage, } from "../provider/context-window.js"; -import { cacheTtlMsFor } from "../provider/cache-ttl.js"; import { COMPACTOR_KEEP_RECENT_TURNS, assistantTextIsCompactSpacerEcho, @@ -162,13 +161,7 @@ export function createCompactionGovernor( // can flag the number as approximate instead of implying provider-grade // precision. let usingEstimate = false; - // Last-cycle source of the last inference.done, kept for live re-checks - // between inference cycles (see interceptActions) where the event carries - // no source. Threshold sizing still keys off `model`; TTL identity needs - // `sourceId` / `provider` as well — production LastCycleSource stamps a - // bare model, and Ollama is `openai-compatible` with id `ollama/…`. let lastModel: string | undefined; - let lastCycleSource: LastCycleSource | undefined; let turnCount = 0; // Wall-clock of the last inference.done: the provider (re)wrote its prefix // cache for this session on that turn, so the provider TTL window in @@ -284,7 +277,6 @@ export function createCompactionGovernor( consecutiveThresholdCompacts = 0; } syncFromTurns(turns); - lastCycleSource = event.source; lastModel = event.source?.model; lastCacheWriteAt = now(); // The terminal reply ends the previous tool batch (its results are @@ -390,37 +382,6 @@ export function createCompactionGovernor( return true; } - // Provider-aware idle recompress (CL-8745): the fold is a re-compress, not - // a cache play. Provider KV caches expire on their own schedule - // (provider/cache-ttl.ts); compressing after that expiry makes the next - // turn a cheaper write and later reads compound on the shrunk context. This - // fires on any live re-entry once `now - lastCacheWrite >= ttl`, including - // over-threshold sessions whose threshold path has disarmed (`pending` is - // false via growth hysteresis — the threshold path owns only armed - // over-threshold; the fold is still window- and cap-bounded). Guards, in - // order: threshold arming defers (pending), the fresh-tail floor (turns at - // or under it are all kept, so a fold would shrink nothing), the - // consecutive-compact cap (existing death-spiral bound, shared with the - // threshold path), providers with no TTL (undefined/empty model, local - // inference), no observed cache write yet, the TTL window itself, and one - // fire per window (a fresh fold rewrites the prefix; the summary call - // bypasses this governor so lastCacheWriteAt cannot observe it — - // lastCompactAt covers that). Never fires with a tool batch outstanding: - // the stall ping that triggers this can arrive mid-work. - function isTtlRecompressDue(nowMs: number): boolean { - if (pending) return false; - if (turnCount <= MIN_TURNS_TO_COMPACT) return false; - if (atThresholdCompactCap()) return false; - if (outstandingToolCalls > 0) return false; - const ttl = cacheTtlMsFor(lastCycleSource); - if (ttl === undefined) return false; - if (lastCacheWriteAt === undefined) return false; - if (nowMs - lastCacheWriteAt < ttl) return false; - if (lastCompactAt !== undefined && nowMs - lastCompactAt < ttl) - return false; - return true; - } - function inboundText(event: ReactorInboundEvent): string { if (event.type !== "message.received") return ""; return typeof event.message.content === "string" @@ -479,12 +440,11 @@ export function createCompactionGovernor( // In-flight `/compact` (manualPending without idlePending) waits on // interceptActions; do not steal that hop with a TTL fold. if (manualPending) return null; - if (!isTtlRecompressDue(now())) return null; - return issueIdleFold( - inboundText(event), - capabilities, - "cache-ttl-recompress", - ); + // Cache expiry is a prompt transform, not a fold. Compacting here rewrites + // turns.jsonl and drops the history the transform is supposed to leave + // stored. The Anthropic prompt transform stubs tool bodies on the request + // when this stamp is expired. + return null; } // A context-overflow inference error would otherwise terminate the loop @@ -594,7 +554,6 @@ export function createCompactionGovernor( }): void { syncFromTurns(args.turns); lastCacheWriteAt = args.at; - lastCycleSource = args.source; lastModel = args.source.model; } diff --git a/src/config/index.ts b/src/config/index.ts index eb6da1a91..cbcca2f94 100644 --- a/src/config/index.ts +++ b/src/config/index.ts @@ -576,6 +576,9 @@ export interface Config { cwd: string; task: string; dangerouslySkipPermissions: boolean; + // Experimental prompt shrink after an Anthropic cache expiry. Off unless + // settings set anthropicCachePrompt. + anthropicCachePrompt: boolean; // True when dangerouslySkipPermissions came from the persisted global // default rather than this invocation's CLI flag. Entry points use this to // surface a startup notice since the persisted default is otherwise silent. @@ -1139,6 +1142,7 @@ export async function loadConfig( cwd, task: resumeTask, dangerouslySkipPermissions, + anthropicCachePrompt: settings?.anthropicCachePrompt === true, skipPermissionsFromSettings, auto, command, diff --git a/src/config/settings.ts b/src/config/settings.ts index f8a582cbc..f80ae1a59 100644 --- a/src/config/settings.ts +++ b/src/config/settings.ts @@ -182,6 +182,10 @@ export interface Settings { showPromptCost?: boolean; // User-global YOLO default; `/yolo` writes it. dangerouslySkipPermissions?: boolean; + // Experimental. When true, an expired Anthropic prompt cache stubs old tool + // results on the outgoing prompt only. Default off: sessions do not depend + // on that shrink. + anthropicCachePrompt?: boolean; } function modelRefKey(ref: ModelRef): string { @@ -594,6 +598,7 @@ const SettingsSchema = type({ "favoriteModels?": ModelRefSchema.array(), "showPromptCost?": "boolean", "dangerouslySkipPermissions?": "boolean", + "anthropicCachePrompt?": "boolean", }); // Per-entry MCP shape without the name key. The "exactly one transport" rule is @@ -797,6 +802,7 @@ export const GLOBAL_SETTINGS_OPTIONAL_KEYS = [ "recentModels", "favoriteModels", "dangerouslySkipPermissions", + "anthropicCachePrompt", ] as const satisfies readonly (keyof OptionalSettingsFields)[]; /** Optional local settings keys the load path is required to consider. */ @@ -964,6 +970,10 @@ function normalizeParsedSettings(path: string, parsed: unknown): Settings { s.dangerouslySkipPermissions !== undefined ? Boolean(s.dangerouslySkipPermissions) : undefined, + anthropicCachePrompt: + s.anthropicCachePrompt !== undefined + ? Boolean(s.anthropicCachePrompt) + : undefined, }; return { providers: s.providers as Settings["providers"], diff --git a/src/exec/runner.ts b/src/exec/runner.ts index ef09f4f1e..8757ae19c 100644 --- a/src/exec/runner.ts +++ b/src/exec/runner.ts @@ -862,6 +862,7 @@ export async function runExec(config: Config): Promise { }, getDefaultSource: () => liveDefaultSource.length > 0 ? liveDefaultSource : liveSource.id, + anthropicCachePrompt: () => config.anthropicCachePrompt, getCompactor: () => createSessionPruningCompactor({ summarize: summarizeForCompaction, diff --git a/src/session/anthropic-cache-prompt.test.ts b/src/session/anthropic-cache-prompt.test.ts index 78762c778..55dad43da 100644 --- a/src/session/anthropic-cache-prompt.test.ts +++ b/src/session/anthropic-cache-prompt.test.ts @@ -92,6 +92,19 @@ function history(): ConversationTurn[] { } describe("anthropic cache prompt transform", () => { + test("settings off leaves an expired Anthropic prompt unchanged", async () => { + const turns = history(); + const transform = createAnthropicCachePromptTransform({ + nowMs: () => NOW, + cacheWriteAt: () => NOW - 6 * MINUTE_MS, + protocol: () => "anthropic", + enabled: () => false, + }); + const result = await transform.apply(turns, CTX); + expect(result.output).toBe(turns); + expect(result.record.reason).toBe("disabled"); + }); + test("expired or missing Anthropic stamp stubs old tool bodies and leaves stored turns", async () => { const dir = await mkdtemp(join(tmpdir(), "anthropic-cache-prompt-")); const store = await createOptimizedContextStore(dir); diff --git a/src/session/anthropic-cache-prompt.ts b/src/session/anthropic-cache-prompt.ts index 768e28f42..2d588073a 100644 --- a/src/session/anthropic-cache-prompt.ts +++ b/src/session/anthropic-cache-prompt.ts @@ -18,6 +18,8 @@ export type AnthropicCachePromptDeps = { cacheWriteAt: () => number | undefined; /** Live adapter protocol (`InferenceSource.provider`). */ protocol: () => string | undefined; + /** When omitted, the shrink is on. Settings pass false unless opted in. */ + enabled?: () => boolean; }; const STRATEGY = "anthropic-cache-prompt"; @@ -144,6 +146,9 @@ export function createAnthropicCachePromptTransform( name: STRATEGY, version: "1", async apply(turns, _ctx) { + if (deps.enabled !== undefined && !deps.enabled()) { + return passthrough(turns, "disabled"); + } const ttl = cacheTtlMsFor(deps.protocol()); if (ttl === undefined) return passthrough(turns, "non-anthropic"); const at = deps.cacheWriteAt(); diff --git a/src/session/assemble-runtime.ts b/src/session/assemble-runtime.ts index 71e8277d4..a40296a48 100644 --- a/src/session/assemble-runtime.ts +++ b/src/session/assemble-runtime.ts @@ -509,6 +509,8 @@ export interface ChatAgentWiring { getDefaultSource: () => string; /** Read at each build so a compaction-mode toggle is visible on rebuild. */ getCompactor: () => Compactor; + /** Experimental Anthropic prompt shrink. Default off when omitted. */ + anthropicCachePrompt?: () => boolean; /** * Present when a resumed run record has an Anthropic-protocol cache write * and the provider about to be called is the same protocol. Read at each @@ -706,6 +708,7 @@ export function assembleChatAgent(wiring: ChatAgentWiring): AssembledChatAgent { createAnthropicCachePromptTransform({ nowMs: () => Date.now(), cacheWriteAt: () => getActiveRun()?.lastCacheWriteAt, + enabled: () => wiring.anthropicCachePrompt?.() === true, protocol: () => { const sources = wiring.getSources(); const preferred = wiring.getDefaultSource(); diff --git a/src/session/cache-ttl-resume.test.ts b/src/session/cache-ttl-resume.test.ts index 45bef180a..1e1c36685 100644 --- a/src/session/cache-ttl-resume.test.ts +++ b/src/session/cache-ttl-resume.test.ts @@ -24,7 +24,7 @@ import { } from "../config/index.js"; import { resumeCacheWriteSeed } from "../provider/cache-ttl.js"; import { HANDOFF_LATEST_KEY } from "./compaction-handoff.js"; -import { COMPACTED_PREFIX, createPruningCompactor } from "./compactor.js"; +import { createPruningCompactor } from "./compactor.js"; import { createOptimizedContextStore } from "./optimized-context-store.js"; import { buildCompactionContinuationMessage } from "./runtime-assembly.js"; @@ -41,16 +41,6 @@ function userTurn(text: string, timestamp: number): ConversationTurn { return { role: "user", content: [{ type: "text", text }], timestamp }; } -function textLength(turns: readonly ConversationTurn[]): number { - let total = 0; - for (const turn of turns) { - for (const block of turn.content) { - if (block.type === "text") total += block.text.length; - } - } - return total; -} - function texts(turns: readonly ConversationTurn[]): string[] { return turns.flatMap((turn) => turn.content.flatMap((block) => @@ -231,34 +221,24 @@ async function resumeAndInfer(args: { } } -describe("resumed Anthropic cache write compacts before infer", () => { +describe("resumed Anthropic cache write does not compact", () => { const native = buildAnthropicSource({ id: "anthropic", baseURL: "https://example.invalid", model: "claude-opus-4-6", }); - test("an expired Anthropic stamp folds into the context store first", async () => { + test("an expired Anthropic stamp does not rewrite the stored turns", async () => { const result = await resumeAndInfer({ at: Date.now() - 6 * MINUTE_MS, source: native, }); - expect(result.first[0]).toMatchObject({ - type: "compact", - compactor: "pruning-compactor", - reason: "cache-ttl-recompress", - }); - expect(result.stored.length).toBeLessThan(result.seed.length); - expect(textLength(result.stored)).toBeLessThan(textLength(result.seed)); - expect(textLength(result.prompt)).toBeLessThan(textLength(result.seed)); - expect(texts(result.prompt)).toEqual(texts(result.stored)); - expect( - texts(result.stored).some((text) => text.startsWith(COMPACTED_PREFIX)), - ).toBe(true); - expect(result.handoff?.length ?? 0).toBeGreaterThan(0); - expect(result.handoff).toContain(HANDOFF_LATEST_KEY); - expect(result.handoff).toContain("# Compaction handoff"); + expect(result.first.some((action) => action.type === "compact")).toBe( + false, + ); + expect(texts(result.stored)).toEqual(texts(result.seed)); + expect(result.handoff).toBeUndefined(); }); test("a write inside 5 minutes leaves the stored turns in place", async () => { @@ -291,7 +271,7 @@ describe("resumed Anthropic cache write compacts before infer", () => { expect(result.handoff).toBeUndefined(); }); - test("expired Zen, OpenCode Go, and custom Anthropic catalog ids fold before infer", async () => { + test("expired Zen, OpenCode Go, and custom Anthropic catalog ids do not compact", async () => { const sources = [ buildZenSource({ id: "zen", @@ -315,12 +295,10 @@ describe("resumed Anthropic cache write compacts before infer", () => { at: Date.now() - 6 * MINUTE_MS, source, }); - expect(result.first[0]).toMatchObject({ - type: "compact", - compactor: "pruning-compactor", - reason: "cache-ttl-recompress", - }); - expect(result.stored.length).toBeLessThan(result.seed.length); + expect(result.first.some((action) => action.type === "compact")).toBe( + false, + ); + expect(texts(result.stored)).toEqual(texts(result.seed)); } }); }); diff --git a/src/tui/runner/session.ts b/src/tui/runner/session.ts index d44eb90e8..5cfb73f61 100644 --- a/src/tui/runner/session.ts +++ b/src/tui/runner/session.ts @@ -706,6 +706,7 @@ export async function assembleTUISession( state.liveDefaultSource.length > 0 ? state.liveDefaultSource : state.liveSource.id, + anthropicCachePrompt: () => config.anthropicCachePrompt, getCompactor: () => compactionLifecycle.wrapCompactor( createSessionPruningCompactor({ From d3780b4b5bca1114733e87eb2615dbb9527ca38a Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Wed, 23 Sep 2026 21:48:14 -0700 Subject: [PATCH 6/8] fix(session): keep cache expiry from compacting stored turns The settings fixture covers anthropicCachePrompt. Unused governor clocks are gone. --- docs/ARCHITECTURE.md | 9 +-------- src/agent/compaction.test.ts | 8 -------- src/agent/compaction.ts | 23 +++++------------------ tests/unit/config.test.ts | 1 + 4 files changed, 7 insertions(+), 34 deletions(-) diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 2b99be445..ede6207bb 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -191,20 +191,13 @@ When a cycle's input tokens cross a threshold, the director compacts the inferen - **Threshold at a tool pause** — Once over threshold, the follow-up `infer` after a tool batch is swapped for a `compact` cycle, and inference resumes via a host continuation message. After a compact that remains over the high watermark, the governor uses **growth hysteresis** (wait for usage to grow by ~10% of the window) instead of re-arming on every cycle; dropping under 60% is not required. - **Idle (end-of-turn)** — An interactive turn can end with a reply and then sit idle with no tool batch to intercept; the governor requests a continuation at that pause and compacts when it arrives. An operator message that races the continuation still compacts first, then re-enters inference to answer it. -- **Idle recompress past the provider TTL** — The fold is a re-compress, not a cache play: provider KV caches expire on their own schedule, so a session compresses _after_ that expiry (the next turn is a cheaper write and later reads compound on the shrunk context) instead of only at 60% tokens. It is not limited to under-threshold sessions: once growth hysteresis has cleared `pending`, an over-threshold session in that gap can fire the same fold. Any live re-entry past the window fires it — an empty stall ping re-enters meter-only, a raced operator message re-infers — under reason `cache-ttl-recompress` through the same compactor, so the fresh tail stays raw exactly as in the threshold path. The interval follows provider cache economics (`src/provider/cache-ttl.ts`), not a global N minutes: - -| Provider segment | Idle recompress allowed after | Why | -| -------------------------------------------------------------- | ----------------------------- | ---------------------------------------------------------------------------------------------------------------- | -| `anthropic`, `zen-messages`, `opencode-go-messages` | 5 min | Published default ephemeral cache TTL; the 1-hour TTL is never set | -| anything else, including OpenAI, xAI, Gemini, DeepSeek, ollama | never | No published 5-minute expiry, or no remote cache. Guessing a shorter window folds a cache that may still be warm | +- **Idle recompress past the provider TTL** — Cache expiry does not compact. Stored turns stay. When settings set `anthropicCachePrompt`, and the saved Anthropic-protocol cache write is at least 5 minutes old, the pre-inference context transform stubs old tool-result bodies on the outgoing prompt only. `turns.jsonl` is not rewritten. The setting defaults off. The `ollama` table key is the bare provider segment; the harness stamps slash-form on `sourceId` (`ollama/default` / `ollama/`), with `provider: openai-compatible` (`buildOpenAISource`) and a bare model (`llama3` / `qwen3`). Idle-recompress disable matches via `isOllamaProviderId` on that `sourceId` (`src/provider/ollama.ts`), not by looking up slash-form as a table key. - **Operator `/compact`** — A first-class slash command folds now without waiting for the occupancy governor. Optional trailing instructions go to the summarizer (empty uses the default structured fold) and are stored on the compact record so later auto-folds still see them. Director rebuild (`/model`, interrupt, resume) restores those instructions from the latest compact record's parameters. An in-flight turn uses the same pair-safe compact-then-continue hop as threshold compact; an idle compact shows the fold and does not start a new turn. - **Overflow recovery** — A `context_overflow` inference error would otherwise become a terminal error reply; the governor compacts and retries instead, bounded so a history the compactor cannot shrink does not loop forever. Overflow ignores hysteresis for the compact itself. -The TTL window is measured from the last `inference.done` (every inference rewrites the provider's prefix cache) with one fire per window per compact of any kind, and it never fires with a tool batch outstanding — stall pings can arrive mid-work. The consecutive-compact cap shared with the threshold path still bounds compact→infer→compact. Prompt-caching the summary call itself is an explicit non-goal: the fold carries no cache options. That write time is stored on the session run record for Anthropic-protocol identities only. A new process restores it and, once the stamp is at least 5 minutes old and the stored turn count is above the compact floor, issues the same `cache-ttl-recompress` fold before the first infer so the request reads the trimmed turns. - The compaction control flow is shaped by a reactor invariant: a `compact` action runs in its own cycle (it cannot be paired with `infer`), and **the reactor delivers no event after a compact cycle**. A director that simply emitted `compact` in place of the follow-up `infer` would leave the loop idle forever — the cause of an earlier stall. Instead the governor pairs `compact` with an emit action (`custom.compaction.continue`), and the host answers that emission by delivering a content-less inbound message (`buildCompactionContinuationMessage()`) through the serial operation queue, generation-guarded like every other deliver. A hop superseded because interrupt already bumped generation and enqueued a rebuild is re-queued onto the replacement agent rather than consumed on the outgoing liveAgent. That message adds no turn (`createInboundTurn` returns `null` for empty content) but re-enters the loop, where the director issues the follow-up `infer` against the freshly truncated history. Each emission is answered at most once (a replayed duplicate of an already-answered emission is ignored), and an unsolicited continuation with neither resume flag set answers `wait` rather than burning a billable inference. The legacy `requestContinuation` closure survives only on the sub-agent path, which delivers directly with no host emit hop; idle arming keeps the two channels exclusive (closure fires and the arming flag stays false, or no closure and the flag reports the arming) so a caller honoring both cannot double-deliver. Compaction replaces older turns with a structured, workflow-aware summary rather than a stats blob: sections for **What Happened / What We're Doing / Relevant Links / Action Items / Next Steps**, with the active workflow and step woven in so compacting mid-`/build` or mid-`/plan` preserves the contract. The summary is produced by a one-shot model call; on any failure it falls back to a deterministic summary so a compaction cycle never breaks the session. diff --git a/src/agent/compaction.test.ts b/src/agent/compaction.test.ts index 06f9c7417..4a109be93 100644 --- a/src/agent/compaction.test.ts +++ b/src/agent/compaction.test.ts @@ -984,14 +984,6 @@ describe("provider-aware idle recompress (CL-8745)", () => { } as ReactorInboundEvent; } - const ttlCompact = [ - { - type: "compact", - compactor: "pruning-compactor", - reason: "cache-ttl-recompress", - }, - ] as ReactorAction[]; - test("fires past the provider TTL while under threshold, meter-only on empty", () => { let continuations = 0; let nowMs = 10_000_000; diff --git a/src/agent/compaction.ts b/src/agent/compaction.ts index 9525553fb..aba7164af 100644 --- a/src/agent/compaction.ts +++ b/src/agent/compaction.ts @@ -135,7 +135,7 @@ export function createCompactionGovernor( requestContinuation?: () => void, systemPrompt = "", toolDefinitions: readonly ToolDefinition[] = [], - now: () => number = Date.now, + _now: () => number = Date.now, ) { let pending = false; let idlePending = false; @@ -163,16 +163,6 @@ export function createCompactionGovernor( let usingEstimate = false; let lastModel: string | undefined; let turnCount = 0; - // Wall-clock of the last inference.done: the provider (re)wrote its prefix - // cache for this session on that turn, so the provider TTL window in - // provider/cache-ttl.ts is measured from here. Stamped on every - // inference.done — estimated-usage providers wrote a cache entry too. - let lastCacheWriteAt: number | undefined; - // Wall-clock of the last issued compact of any kind. A fresh fold rewrites - // the session prefix, so the TTL recompress must not fire again inside the - // same window even when its own cache write has not been observed yet (the - // summary call bypasses this governor). - let lastCompactAt: number | undefined; // Tool calls issued by the last inference.done and not yet settled. A TTL // recompress must never fold while a batch is outstanding — the stall ping // that triggers it can arrive mid-work, and folding under it would rewrite @@ -218,7 +208,6 @@ export function createCompactionGovernor( function noteCompactIssued(): void { awaitingPostCompactMeasurement = true; - lastCompactAt = now(); } function atThresholdCompactCap(): boolean { @@ -278,7 +267,6 @@ export function createCompactionGovernor( } syncFromTurns(turns); lastModel = event.source?.model; - lastCacheWriteAt = now(); // The terminal reply ends the previous tool batch (its results are // already in the turns) and opens the batch the reply just issued. A TTL // recompress must never fold while a batch is outstanding — the stall @@ -543,17 +531,16 @@ export function createCompactionGovernor( extraInstructions = trimmed; } - // A new process has no in-memory cache write. Seed the stored stamp, the - // identity that wrote it, and the stored turns so the next message.received - // can fold before its infer. Turn count is the stored set: the inbound - // message is appended after this, and the floor is measured without it. + // A new process has no in-memory turn count. Seed the stored turns so + // threshold and manual compact see the resumed history. The cache-write + // time lives on the run record for the prompt transform, not here. function restoreCacheWrite(args: { at: number; source: LastCycleSource; turns: readonly ConversationTurn[]; }): void { + void args.at; syncFromTurns(args.turns); - lastCacheWriteAt = args.at; lastModel = args.source.model; } diff --git a/tests/unit/config.test.ts b/tests/unit/config.test.ts index dd93b121a..0f28570b4 100644 --- a/tests/unit/config.test.ts +++ b/tests/unit/config.test.ts @@ -324,6 +324,7 @@ test("loadSettings cannot silently drop a known optional key", async () => { recentModels: [{ provider: "p", model: "m" }], favoriteModels: [{ provider: "p", model: "m" }], dangerouslySkipPermissions: true, + anthropicCachePrompt: true, }; await writeFile(globalPath, JSON.stringify(fixture)); const loaded = await loadSettings(globalPath); From dee359c5d0dedcfbb3b987f20ea09400f693f402 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Wed, 23 Sep 2026 21:51:40 -0700 Subject: [PATCH 7/8] docs(architecture): drop the stale idle-recompress table note The provider table is gone, so the sentence that explained its ollama row is gone too. --- docs/ARCHITECTURE.md | 2 -- 1 file changed, 2 deletions(-) diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index ede6207bb..e30636f6a 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -193,8 +193,6 @@ When a cycle's input tokens cross a threshold, the director compacts the inferen - **Idle (end-of-turn)** — An interactive turn can end with a reply and then sit idle with no tool batch to intercept; the governor requests a continuation at that pause and compacts when it arrives. An operator message that races the continuation still compacts first, then re-enters inference to answer it. - **Idle recompress past the provider TTL** — Cache expiry does not compact. Stored turns stay. When settings set `anthropicCachePrompt`, and the saved Anthropic-protocol cache write is at least 5 minutes old, the pre-inference context transform stubs old tool-result bodies on the outgoing prompt only. `turns.jsonl` is not rewritten. The setting defaults off. -The `ollama` table key is the bare provider segment; the harness stamps slash-form on `sourceId` (`ollama/default` / `ollama/`), with `provider: openai-compatible` (`buildOpenAISource`) and a bare model (`llama3` / `qwen3`). Idle-recompress disable matches via `isOllamaProviderId` on that `sourceId` (`src/provider/ollama.ts`), not by looking up slash-form as a table key. - - **Operator `/compact`** — A first-class slash command folds now without waiting for the occupancy governor. Optional trailing instructions go to the summarizer (empty uses the default structured fold) and are stored on the compact record so later auto-folds still see them. Director rebuild (`/model`, interrupt, resume) restores those instructions from the latest compact record's parameters. An in-flight turn uses the same pair-safe compact-then-continue hop as threshold compact; an idle compact shows the fold and does not start a new turn. - **Overflow recovery** — A `context_overflow` inference error would otherwise become a terminal error reply; the governor compacts and retries instead, bounded so a history the compactor cannot shrink does not loop forever. Overflow ignores hysteresis for the compact itself. From 7ab47112134b36df070cd8690fb8ec0e31c2d19d Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Thu, 24 Sep 2026 17:23:27 -0700 Subject: [PATCH 8/8] fix(compaction): drop dead cache-ttl fold state and stale tests MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ## Summary The TTL fold path is gone, so the governor's outstandingToolCalls tracking was write-only and the stall-ping test merged from main still expected a cache-ttl-recompress compact. Rewrite the provider-aware describe block to pin the new contract — idle re-entry never folds, whatever the provider window. ## Verification - bun test src/agent/compaction.test.ts src/subagent/nudge-director.test.ts (117 pass) - bun run check (the one local failure is a /tmp symlink artifact of the worktree path; the same test passes on a real checkout) Refs CL-8914 --- src/agent/compaction.test.ts | 239 ++++++---------------------- src/agent/compaction.ts | 36 +---- src/subagent/nudge-director.test.ts | 33 ++-- 3 files changed, 61 insertions(+), 247 deletions(-) diff --git a/src/agent/compaction.test.ts b/src/agent/compaction.test.ts index 4a109be93..90d585e31 100644 --- a/src/agent/compaction.test.ts +++ b/src/agent/compaction.test.ts @@ -936,7 +936,7 @@ describe("compaction governor", () => { }); }); -describe("provider-aware idle recompress (CL-8745)", () => { +describe("cache expiry never folds (CL-8914)", () => { const MINUTE_MS = 60_000; function ttlInferenceDone( @@ -984,142 +984,48 @@ describe("provider-aware idle recompress (CL-8745)", () => { } as ReactorInboundEvent; } - test("fires past the provider TTL while under threshold, meter-only on empty", () => { + test("idle pings past any provider cache window return null", () => { + // The governor holds no TTL table: an unarmed idle re-entry never + // produces a compact, whatever the provider's cache economics. Staleness + // on the outgoing prompt is the anthropic-cache-prompt transform's job. let continuations = 0; let nowMs = 10_000_000; - const governor = createCompactionGovernor( - () => continuations++, - "", - [], - () => nowMs, - ); - governor.noteInferenceDone( - ttlInferenceDone( - { provider: "anthropic", model: "claude-opus-4-6" }, - false, - ), - tenTurns, - ); - - // Inside the 5-minute Anthropic window: no fire. - nowMs += 4 * MINUTE_MS; - expect( - governor.interceptIdleContinuation(emptyMessage(), capabilities), - ).toBeNull(); - - // Past the window: the same fold as the threshold path (same compactor, - // so the fresh tail stays raw) with an attributable reason. - nowMs += MINUTE_MS + 1; - expect( - governor.interceptIdleContinuation(emptyMessage(), capabilities), - ).toBeNull(); - expect(continuations).toBe(0); - }); - - test("production LastCycleSource: anthropic fires at 5m, codex never does", () => { - // Harness stamps { sourceId, provider, model } with a bare model, not - // slash-form "anthropic/claude-opus-4-6". Anthropic's 5-minute window - // must come from provider, not a dummy model string. Codex has no - // published 5-minute expiry, so it must stay quiet past the old 10-minute - // guess as well. - let nowMs = 15_000_000; - const clock = () => nowMs; - const anthropic = createCompactionGovernor(() => undefined, "", [], clock); - const codex = createCompactionGovernor(() => undefined, "", [], clock); - anthropic.noteInferenceDone( - ttlInferenceDone( - { provider: "anthropic", model: "claude-opus-4-6" }, - false, - ), - tenTurns, - ); - codex.noteInferenceDone( - ttlInferenceDone( - { - sourceId: "codex/work", - provider: "codex-responses", - model: "gpt-5.6-luna", - }, - false, - ), - tenTurns, - ); - - nowMs += 5 * MINUTE_MS + 1; - expect( - anthropic.interceptIdleContinuation(emptyMessage(), capabilities), - ).toBeNull(); - expect( - codex.interceptIdleContinuation(emptyMessage(), capabilities), - ).toBeNull(); - - nowMs += 30 * MINUTE_MS; - expect( - codex.interceptIdleContinuation(emptyMessage(), capabilities), - ).toBeNull(); - }); - - test("does not idle-recompress DeepSeek or ollama", () => { - let nowMs = 20_000_000; const clock = () => nowMs; - const deepseek = createCompactionGovernor(() => undefined, "", [], clock); - const local = createCompactionGovernor(() => undefined, "", [], clock); - deepseek.noteInferenceDone( - ttlInferenceDone("deepseek/deepseek-chat", false), - tenTurns, - ); - local.noteInferenceDone( - ttlInferenceDone("ollama/llama3.1", false), - tenTurns, - ); - - nowMs += 90 * MINUTE_MS; - expect( - deepseek.interceptIdleContinuation(emptyMessage(), capabilities), - ).toBeNull(); - expect( - local.interceptIdleContinuation(emptyMessage(), capabilities), - ).toBeNull(); - }); - - test("production ollama LastCycleSource never fires cache-ttl-recompress", () => { - // Harness stamps { sourceId, provider, model } with a bare model. Ollama is - // buildOpenAISource: sourceId "ollama/default", provider openai-compatible, - // Keying TTL off the bare model would miss the ollama sourceId. Local - // inference stays disabled. - let nowMs = 25_000_000; - const governor = createCompactionGovernor( - () => undefined, - "", - [], - () => nowMs, - ); - governor.noteInferenceDone( - ttlInferenceDone( - { - sourceId: "ollama/default", - provider: "openai-compatible", - model: "llama3", - }, - false, - ), - tenTurns, - ); - - nowMs += 5 * MINUTE_MS + 1; - expect( - governor.interceptIdleContinuation(emptyMessage(), capabilities), - ).toBeNull(); - }); - - test("stays inert with no observed cache write", () => { - const governor = createCompactionGovernor(() => undefined); - expect( - governor.interceptIdleContinuation(emptyMessage(), capabilities), - ).toBeNull(); + const sources = [ + { provider: "anthropic", model: "claude-opus-4-6" }, + { + sourceId: "codex/work", + provider: "codex-responses", + model: "gpt-5.6-luna", + }, + { provider: "deepseek", model: "deepseek-chat" }, + { + sourceId: "ollama/default", + provider: "openai-compatible", + model: "llama3", + }, + { provider: "custom-proxy", model: "unknown-model" }, + ]; + for (const source of sources) { + const governor = createCompactionGovernor( + () => continuations++, + "", + [], + clock, + ); + governor.noteInferenceDone(ttlInferenceDone(source, false), tenTurns); + expect( + governor.interceptIdleContinuation(emptyMessage(), capabilities), + ).toBeNull(); + nowMs += 90 * MINUTE_MS; + expect( + governor.interceptIdleContinuation(emptyMessage(), capabilities), + ).toBeNull(); + } + expect(continuations).toBe(0); }); - test("defers to the armed threshold path while over threshold", () => { + test("an armed threshold fold is not disturbed by idle re-entry", () => { let nowMs = 30_000_000; const governor = createCompactionGovernor( () => undefined, @@ -1128,9 +1034,9 @@ describe("provider-aware idle recompress (CL-8745)", () => { () => nowMs, ); governor.noteInferenceDone(inferenceDone(overThreshold), tenTurns); - // Fixture provider "p" has no TTL. Threshold arming still owns this session. + // Threshold arming owns the over-threshold session and fires at the tool + // pause; a message.received re-entry is not its trigger. nowMs += 5 * MINUTE_MS + 1; - // Threshold arming owns the over-threshold session: no TTL double-fold. expect( governor.interceptIdleContinuation(emptyMessage(), capabilities), ).toBeNull(); @@ -1139,11 +1045,10 @@ describe("provider-aware idle recompress (CL-8745)", () => { ).toBeNull(); }); - test("fires cache-ttl-recompress in the hysteresis gap (over-threshold, no growth)", () => { - // Characterization, not a bug: after a threshold compact, a post-compact - // infer at the same usage clears `pending` via growth hysteresis, so the - // threshold path no longer owns the session. Idle past the provider TTL - // still folds — same window- and cap-bounded path as under-threshold. + test("the hysteresis gap does not fold on cache expiry", () => { + // After a threshold compact, a post-compact infer at the same usage + // clears `pending` via growth hysteresis — the exact re-entry where the + // removed TTL path used to fire. let nowMs = 70_000_000; const governor = createCompactionGovernor( () => undefined, @@ -1174,37 +1079,7 @@ describe("provider-aware idle recompress (CL-8745)", () => { ).toBeNull(); }); - test("after threshold compact and gap TTL, growth-armed compact stays blocked until a tool_call", () => { - // Threshold compact (consecutive=1) plus TTL fire in the hysteresis gap - // (consecutive=2) fills the shared cap. Later growth that would re-arm - // the threshold path stays blocked until a tool_call occupancy resets it. - let nowMs = 80_000_000; - const governor = createCompactionGovernor( - () => undefined, - "", - [], - () => nowMs, - ); - governor.noteInferenceDone( - inferenceDone(overThreshold, "", "anthropic"), - tenTurns, - ); - expect( - governor.interceptActions(toolDone(), inferAction, capabilities), - ).not.toBeNull(); - expect(governor.resumeAfterCompact(emptyMessage())).toBe("infer"); - - governor.noteInferenceDone( - inferenceDone(overThreshold, "", "anthropic"), - tenTurns, - ); - nowMs += 5 * MINUTE_MS + 1; - expect( - governor.interceptIdleContinuation(emptyMessage(), capabilities), - ).toBeNull(); - }); - - test("does not fold under an outstanding tool batch, fires once it settles", () => { + test("an outstanding tool batch does not change idle re-entry", () => { let continuations = 0; let nowMs = 40_000_000; const governor = createCompactionGovernor( @@ -1221,7 +1096,6 @@ describe("provider-aware idle recompress (CL-8745)", () => { expect( governor.interceptIdleContinuation(emptyMessage(), capabilities), ).toBeNull(); - // The batch settles (threshold path uninvolved: under threshold). expect( governor.interceptActions(toolDone(), inferAction, capabilities), ).toBeNull(); @@ -1231,28 +1105,7 @@ describe("provider-aware idle recompress (CL-8745)", () => { expect(continuations).toBe(0); }); - test("one fire per window, then the shared consecutive-compact cap stops the spiral", () => { - let continuations = 0; - let nowMs = 50_000_000; - const governor = createCompactionGovernor( - () => continuations++, - "", - [], - () => nowMs, - ); - governor.noteInferenceDone( - ttlInferenceDone("anthropic/claude-opus-4-6", false), - tenTurns, - ); - - nowMs += 5 * MINUTE_MS + 1; - expect( - governor.interceptIdleContinuation(emptyMessage(), capabilities), - ).toBeNull(); - expect(continuations).toBe(0); - }); - - test("a raced operator message past the TTL still folds, then re-infers", () => { + test("a raced operator message past the window does not fold", () => { let continuations = 0; let nowMs = 60_000_000; const governor = createCompactionGovernor( diff --git a/src/agent/compaction.ts b/src/agent/compaction.ts index aba7164af..90bf14ee2 100644 --- a/src/agent/compaction.ts +++ b/src/agent/compaction.ts @@ -163,12 +163,6 @@ export function createCompactionGovernor( let usingEstimate = false; let lastModel: string | undefined; let turnCount = 0; - // Tool calls issued by the last inference.done and not yet settled. A TTL - // recompress must never fold while a batch is outstanding — the stall ping - // that triggers it can arrive mid-work, and folding under it would rewrite - // turns the pending results still belong to. Assigned (not incremented) on - // every inference.done so the serial loop self-heals a miscount. - let outstandingToolCalls = 0; // Growth hysteresis after a compact that remained over the high watermark: // snapshot the post-compact infer's usage, then do not re-arm until usage // grows by resumeDelta. Cleared once usage drops back to or under high. @@ -267,15 +261,6 @@ export function createCompactionGovernor( } syncFromTurns(turns); lastModel = event.source?.model; - // The terminal reply ends the previous tool batch (its results are - // already in the turns) and opens the batch the reply just issued. A TTL - // recompress must never fold while a batch is outstanding — the stall - // ping that triggers it can arrive mid-work, and folding under it would - // rewrite turns the pending results still belong to. Assigned, not - // incremented, so the serial loop self-heals a miscount. - outstandingToolCalls = event.turn.content.filter( - (block) => block.type === "tool_call", - ).length; const reportedTokens = contextTokensFromUsage(event.usage); usingEstimate = reportedTokens <= 0; const contextTokens = usingEstimate ? estimate.tokens : reportedTokens; @@ -315,7 +300,6 @@ export function createCompactionGovernor( capabilities: ReactorCapabilities, ): ReactorAction[] | null { if (event.type !== "tool.done") return null; - if (outstandingToolCalls > 0) outstandingToolCalls -= 1; const operator = manualPending; if ( !operator && @@ -421,17 +405,11 @@ export function createCompactionGovernor( } return issueIdleFold(content, capabilities, THRESHOLD_COMPACT_REASON); } - // Unarmed idle re-entry past the provider TTL: same fold, same - // keep-recent tail, same cap — but a "cache-ttl-recompress" reason so the - // fold is attributable. No arming: every live re-entry re-checks the - // window, so a sub-agent stall ping or operator message is the trigger. - // In-flight `/compact` (manualPending without idlePending) waits on - // interceptActions; do not steal that hop with a TTL fold. - if (manualPending) return null; - // Cache expiry is a prompt transform, not a fold. Compacting here rewrites - // turns.jsonl and drops the history the transform is supposed to leave - // stored. The Anthropic prompt transform stubs tool bodies on the request - // when this stamp is expired. + // Cache expiry is a prompt transform, not a fold. Compacting on an + // unarmed idle re-entry rewrites turns.jsonl and drops the history the + // transform is supposed to leave stored; the Anthropic prompt transform + // stubs tool bodies on the outgoing request instead. In-flight `/compact` + // (manualPending without idlePending) waits on interceptActions. return null; } @@ -486,10 +464,6 @@ export function createCompactionGovernor( function notePostCompact(turns: readonly ConversationTurn[]): void { syncFromTurns(turns); usingEstimate = true; - // A fold rewrites the turns: results already applied vanish from the - // live set, and post-compact stall pings (empty continuations) carry no - // tool traffic. Reset so a stale count cannot pin the TTL window shut. - outstandingToolCalls = 0; } // True while the governor expects the host to answer a continuation emit. diff --git a/src/subagent/nudge-director.test.ts b/src/subagent/nudge-director.test.ts index 16287c405..8b19c9302 100644 --- a/src/subagent/nudge-director.test.ts +++ b/src/subagent/nudge-director.test.ts @@ -1665,9 +1665,10 @@ describe("SubAgentDirector idle stall ping", () => { expect(folded.some((action) => action.type === "wait")).toBe(false); }); - test("cache-ttl recompress still folds on an in-window empty ping", async () => { - // test-model takes the 10-minute default TTL. The stall window is longer - // so the due fold is still an in-window ping, not a stall nudge. + test("cache expiry does not fold on an in-window empty ping", async () => { + // test-model used to take the 10-minute default TTL. Cache expiry is now + // a prompt transform, not a fold: the in-window ping waits, and the stall + // nudge still fires off the original activity clock. const activityAt = 11_000_000; const cacheTtlMs = 10 * 60_000; const stallTimeoutMs = 15 * 60_000; @@ -1687,29 +1688,15 @@ describe("SubAgentDirector idle stall ping", () => { await director.decide(inferenceDoneText("working"), longState, caps); now += cacheTtlMs + 1; - const folded = actions( - await director.decide(messageReceived(""), longState, caps), - ); - expect(folded).toEqual([ - { - type: "compact", - compactor: "pruning-compactor", - reason: "cache-ttl-recompress", - }, - ]); - expect(continuations).toBe(1); - expect(folded.some((action) => action.type === "infer")).toBe(false); - expect(folded.some((action) => action.type === "wait")).toBe(false); - - // Meter-only resume of the empty fold. Same clock: still inside the - // stall window, and this wait must not count as activity either. - const resumed = actions( + const ping = actions( await director.decide(messageReceived(""), longState, caps), ); - expect(resumed).toEqual([{ type: "wait" }]); - expect(resumed.some((action) => action.type === "infer")).toBe(false); + expect(ping).toEqual([{ type: "wait" }]); + expect(continuations).toBe(0); + expect(ping.some((action) => action.type === "infer")).toBe(false); + expect(ping.some((action) => action.type === "compact")).toBe(false); - // The fold must not stamp lastActivityAt or clear stallNudgeAt. One + // The wait must not stamp lastActivityAt or clear stallNudgeAt. One // stall timeout from the original activity still nudges, once. now = activityAt + stallTimeoutMs; const nudge = actions(