diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index f62e633bb..0d264445c 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -191,24 +191,11 @@ 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 | 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 | - -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. +- **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. - **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 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 f517d7f4a..90d585e31 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; } @@ -935,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( @@ -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", @@ -979,161 +984,48 @@ 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", () => { + 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), - ).toEqual(ttlCompact); - expect(continuations).toBe(1); - // Empty continuation adopts the shrunk turns without a new inference. - expect(governor.resumeAfterCompact(emptyMessage())).toBe("meter"); - }); - - test("production LastCycleSource: anthropic fires at 5m, codex stays quiet until 10m", () => { - // 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". - 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), - ).toEqual(ttlCompact); - expect( - codex.interceptIdleContinuation(emptyMessage(), capabilities), - ).toBeNull(); - - nowMs += 5 * MINUTE_MS; - expect( - codex.interceptIdleContinuation(emptyMessage(), capabilities), - ).toEqual(ttlCompact); - }); - - test("follows provider economics: deepseek waits out its long window, ollama never fires", () => { - 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, - ); - - // Past Anthropic/OpenAI windows but inside DeepSeek's hour: neither fires. - nowMs += 30 * 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. - let nowMs = 25_000_000; - const governor = createCompactionGovernor( - () => undefined, - "", - [], - () => nowMs, - ); - governor.noteInferenceDone( - ttlInferenceDone( - { - sourceId: "ollama/default", - provider: "openai-compatible", - model: "llama3", - }, - false, - ), - tenTurns, - ); - - nowMs += 11 * MINUTE_MS; - 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, @@ -1142,9 +1034,9 @@ 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; - // Threshold arming owns the over-threshold session: no TTL double-fold. + // 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; expect( governor.interceptIdleContinuation(emptyMessage(), capabilities), ).toBeNull(); @@ -1153,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, @@ -1165,68 +1056,30 @@ describe("provider-aware idle recompress (CL-8745)", () => { [], () => nowMs, ); - governor.noteInferenceDone(inferenceDone(overThreshold), tenTurns); - expect( - governor.interceptActions(toolDone(), inferAction, capabilities), - ).not.toBeNull(); - expect(governor.resumeAfterCompact(emptyMessage())).toBe("infer"); - - governor.noteInferenceDone(inferenceDone(overThreshold), tenTurns); - expect( - governor.interceptActions(toolDone(), inferAction, capabilities), - ).toBeNull(); - - nowMs += 11 * MINUTE_MS; - expect( - governor.interceptIdleContinuation(emptyMessage(), capabilities), - ).toEqual(ttlCompact); - }); - - 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, ); - governor.noteInferenceDone(inferenceDone(overThreshold), tenTurns); expect( governor.interceptActions(toolDone(), inferAction, capabilities), ).not.toBeNull(); expect(governor.resumeAfterCompact(emptyMessage())).toBe("infer"); - governor.noteInferenceDone(inferenceDone(overThreshold), tenTurns); - nowMs += 11 * MINUTE_MS; - 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), + inferenceDone(overThreshold, "", "anthropic"), tenTurns, ); expect( governor.interceptActions(toolDone(), inferAction, capabilities), ).toBeNull(); - governor.noteInferenceDone( - inferenceDoneWithTools(overThreshold + 2 * resumeDelta), - tenTurns, - ); + nowMs += 5 * MINUTE_MS + 1; expect( - governor.interceptActions(toolDone(), inferAction, capabilities), - ).not.toBeNull(); + 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( @@ -1243,59 +1096,16 @@ 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(); - expect( - governor.interceptIdleContinuation(emptyMessage(), capabilities), - ).toEqual(ttlCompact); - expect(continuations).toBe(1); - }); - - 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), - ).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", () => { + test("a raced operator message past the window does not fold", () => { let continuations = 0; let nowMs = 60_000_000; const governor = createCompactionGovernor( @@ -1311,11 +1121,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 4cff59a80..90bf14ee2 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, @@ -103,6 +102,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, @@ -114,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; @@ -140,30 +161,8 @@ 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 - // 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 - // 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. @@ -203,7 +202,6 @@ export function createCompactionGovernor( function noteCompactIssued(): void { awaitingPostCompactMeasurement = true; - lastCompactAt = now(); } function atThresholdCompactCap(): boolean { @@ -262,18 +260,7 @@ 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 - // 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; @@ -313,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 && @@ -368,37 +354,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" @@ -450,19 +405,12 @@ 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; - if (!isTtlRecompressDue(now())) return null; - return issueIdleFold( - inboundText(event), - capabilities, - "cache-ttl-recompress", - ); + // 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; } // A context-overflow inference error would otherwise terminate the loop @@ -516,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. @@ -561,6 +505,19 @@ export function createCompactionGovernor( extraInstructions = trimmed; } + // 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); + 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 +579,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/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 267036e18..8757ae19c 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 } : {}), }; @@ -842,6 +862,7 @@ export async function runExec(config: Config): Promise { }, getDefaultSource: () => liveDefaultSource.length > 0 ? liveDefaultSource : liveSource.id, + anthropicCachePrompt: () => config.anthropicCachePrompt, getCompactor: () => createSessionPruningCompactor({ summarize: summarizeForCompaction, @@ -866,6 +887,13 @@ export async function runExec(config: Config): Promise { } }, }), + getCacheWriteSeed: () => + resumeCacheWriteSeed({ + at: activeRunHandle.lastCacheWriteAt, + storedModel: activeRunHandle.cacheWriteModel, + liveProvider: config.providerName, + liveProtocol: liveSource.provider, + }), onBuilt: (agent, storage) => { currentAgent = agent; currentStorage = storage; @@ -913,7 +941,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 0bf58fbf4..9582ce93e 100644 --- a/src/provider/cache-ttl.test.ts +++ b/src/provider/cache-ttl.test.ts @@ -1,27 +1,33 @@ import { describe, expect, test } from "bun:test"; -import { cacheTtlMsFor } from "./cache-ttl.js"; +import { + buildAnthropicSource, + buildGoSource, + buildZenSource, +} from "../config/index.js"; +import { + anthropicCacheWriteAt, + cacheTtlMsFor, + resumeCacheWriteSeed, +} 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 +43,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 +55,117 @@ describe("cacheTtlMsFor", () => { provider: "codex-responses", model: "gpt-5.6-luna", }), - ).toBe(10 * MINUTE_MS); + ).toBeUndefined(); }); - 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, - ); + 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(); + }); +}); + +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", + liveProtocol: "anthropic", + }), + ).toEqual({ at: 10, model: "anthropic:claude-opus-4-6" }); + expect( + resumeCacheWriteSeed({ + at: 10, + storedModel: "openai:gpt-5.6", + liveProvider: "openai", + liveProtocol: "openai", + }), + ).toBeUndefined(); + expect( + resumeCacheWriteSeed({ + at: 10, + storedModel: "anthropic:claude-opus-4-6", + liveProvider: "openai", + liveProtocol: "openai", + }), + ).toBeUndefined(); + expect( + resumeCacheWriteSeed({ + at: undefined, + storedModel: "anthropic:claude-opus-4-6", + liveProvider: "anthropic", + liveProtocol: "anthropic", + }), + ).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("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 a3803b52f..8cf8621e8 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,96 @@ 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]; + if (identity.provider !== undefined && identity.provider.length > 0) + 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 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 || 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}` }; +} - return ttlForModelString(identity.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/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/anthropic-cache-prompt.test.ts b/src/session/anthropic-cache-prompt.test.ts new file mode 100644 index 000000000..55dad43da --- /dev/null +++ b/src/session/anthropic-cache-prompt.test.ts @@ -0,0 +1,160 @@ +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("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); + 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..2d588073a --- /dev/null +++ b/src/session/anthropic-cache-prompt.ts @@ -0,0 +1,171 @@ +// 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; + /** When omitted, the shrink is on. Settings pass false unless opted in. */ + enabled?: () => boolean; +}; + +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) { + 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(); + 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 1e3d88e51..a40296a48 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,13 +51,18 @@ 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"; 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, @@ -503,6 +509,15 @@ 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 + * 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; /** @@ -690,6 +705,17 @@ export function assembleChatAgent(wiring: ChatAgentWiring): AssembledChatAgent { createAttachmentRehydrateTransform((key) => storageForAgent.readBlob(key), ), + createAnthropicCachePromptTransform({ + nowMs: () => Date.now(), + cacheWriteAt: () => getActiveRun()?.lastCacheWriteAt, + enabled: () => wiring.anthropicCachePrompt?.() === true, + protocol: () => { + const sources = wiring.getSources(); + const preferred = wiring.getDefaultSource(); + const match = sources.find((source) => source.id === preferred); + return (match ?? sources[0])?.provider; + }, + }), ], }, audit, @@ -718,6 +744,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 +779,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..1e1c36685 --- /dev/null +++ b/src/session/cache-ttl-resume.test.ts @@ -0,0 +1,304 @@ +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, + InferenceSource, + ReactorAction, +} from "@intx/types/runtime"; +import { createChatDirector } from "../agent/director.js"; +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 { 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 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; + source: InferenceSource; +}): 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 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", [], {}); + 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) => { + 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: args.source, + 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: args.source.id, + provider: args.source.provider, + model: args.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 does not compact", () => { + const native = buildAnthropicSource({ + id: "anthropic", + baseURL: "https://example.invalid", + model: "claude-opus-4-6", + }); + + 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.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 () => { + const result = await resumeAndInfer({ + at: Date.now() - 2 * MINUTE_MS, + source: native, + }); + + 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, + source: buildOpenAISource({ + id: "openai", + baseURL: "https://example.invalid/v1", + model: "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(); + }); + + test("expired Zen, OpenCode Go, and custom Anthropic catalog ids do not compact", 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.some((action) => action.type === "compact")).toBe( + false, + ); + expect(texts(result.stored)).toEqual(texts(result.seed)); + } + }); +}); 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/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( 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..5cfb73f61 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, @@ -695,6 +706,7 @@ export async function assembleTUISession( state.liveDefaultSource.length > 0 ? state.liveDefaultSource : state.liveSource.id, + anthropicCachePrompt: () => config.anthropicCachePrompt, getCompactor: () => compactionLifecycle.wrapCompactor( createSessionPruningCompactor({ @@ -722,6 +734,13 @@ export async function assembleTUISession( }, }), ), + getCacheWriteSeed: () => + resumeCacheWriteSeed({ + at: start.activeRunHandle.lastCacheWriteAt, + storedModel: start.activeRunHandle.cacheWriteModel, + liveProvider: state.config.providerName, + liveProtocol: state.liveSource.provider, + }), 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); 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);