Skip to content

Commit e4744a0

Browse files
committed
feat(session): shrink anthropic prompts once the cache write expires
1 parent 4ab98d8 commit e4744a0

3 files changed

Lines changed: 325 additions & 0 deletions

File tree

Lines changed: 147 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,147 @@
1+
import { mkdtemp, readFile } from "node:fs/promises";
2+
import { tmpdir } from "node:os";
3+
import { join } from "node:path";
4+
import { describe, expect, test } from "bun:test";
5+
import type { ConversationTurn } from "@intx/types/runtime";
6+
7+
import { createAnthropicCachePromptTransform } from "./anthropic-cache-prompt.js";
8+
import { createOptimizedContextStore } from "./optimized-context-store.js";
9+
10+
const MINUTE_MS = 60_000;
11+
const NOW = 1_000_000_000_000;
12+
const BODY = "file-body-".repeat(2_000);
13+
14+
const CTX = { state: {} as never, trigger: "test" };
15+
16+
function transformFor(args: {
17+
protocol: string;
18+
cacheWriteAt: number | undefined;
19+
nowMs?: number;
20+
}) {
21+
return createAnthropicCachePromptTransform({
22+
nowMs: () => args.nowMs ?? NOW,
23+
cacheWriteAt: () => args.cacheWriteAt,
24+
protocol: () => args.protocol,
25+
});
26+
}
27+
28+
function history(): ConversationTurn[] {
29+
return [
30+
{
31+
role: "user",
32+
content: [{ type: "text", text: "read the sources" }],
33+
timestamp: 1,
34+
},
35+
{
36+
role: "assistant",
37+
content: [
38+
{
39+
type: "tool_call",
40+
id: "read-1",
41+
name: "read_file",
42+
arguments: { path: "src/big.ts" },
43+
},
44+
],
45+
timestamp: 2,
46+
},
47+
{
48+
role: "user",
49+
content: [
50+
{
51+
type: "tool_result",
52+
callId: "read-1",
53+
content: [{ type: "text", text: BODY }],
54+
},
55+
],
56+
timestamp: 3,
57+
},
58+
{
59+
role: "assistant",
60+
content: [
61+
{
62+
type: "tool_call",
63+
id: "bash-1",
64+
name: "bash",
65+
arguments: { command: "wc -l src/big.ts" },
66+
},
67+
],
68+
timestamp: 4,
69+
},
70+
{
71+
role: "user",
72+
content: [
73+
{
74+
type: "tool_result",
75+
callId: "bash-1",
76+
content: [{ type: "text", text: BODY }],
77+
},
78+
],
79+
timestamp: 5,
80+
},
81+
{
82+
role: "assistant",
83+
content: [{ type: "text", text: "newest assistant" }],
84+
timestamp: 6,
85+
},
86+
{
87+
role: "user",
88+
content: [{ type: "text", text: "newest user" }],
89+
timestamp: 7,
90+
},
91+
];
92+
}
93+
94+
describe("anthropic cache prompt transform", () => {
95+
test("expired or missing Anthropic stamp stubs old tool bodies and leaves stored turns", async () => {
96+
const dir = await mkdtemp(join(tmpdir(), "anthropic-cache-prompt-"));
97+
const store = await createOptimizedContextStore(dir);
98+
const turns = history();
99+
await store.writeTurns(turns);
100+
const stored = (await store.load()).turns;
101+
const diskBefore = await readFile(join(dir, "turns.jsonl"), "utf8");
102+
const storedJson = JSON.stringify(stored);
103+
104+
for (const cacheWriteAt of [undefined, NOW - 5 * MINUTE_MS]) {
105+
const result = await transformFor({
106+
protocol: "anthropic",
107+
cacheWriteAt,
108+
}).apply(stored, CTX);
109+
expect(JSON.stringify(result.output).length).toBeLessThan(
110+
storedJson.length,
111+
);
112+
expect(result.output).toHaveLength(stored.length);
113+
expect(JSON.stringify(result.output)).not.toContain(BODY);
114+
expect(JSON.stringify(result.output)).toContain("newest assistant");
115+
expect(JSON.stringify(result.output)).toContain("newest user");
116+
expect(JSON.stringify(result.output)).toContain(
117+
"[read_file result omitted]",
118+
);
119+
expect(JSON.stringify(result.output)).toContain("[bash result omitted]");
120+
expect(result.record.reason).toBe("stubbed-tool-results");
121+
}
122+
123+
expect(JSON.stringify((await store.load()).turns)).toBe(storedJson);
124+
expect(await readFile(join(dir, "turns.jsonl"), "utf8")).toBe(diskBefore);
125+
expect(storedJson).toContain(BODY);
126+
});
127+
128+
test("a stamp two minutes old returns the same turns", async () => {
129+
const turns = history();
130+
const result = await transformFor({
131+
protocol: "anthropic",
132+
cacheWriteAt: NOW - 2 * MINUTE_MS,
133+
}).apply(turns, CTX);
134+
expect(result.output).toBe(turns);
135+
expect(result.record.reason).toBe("cache-warm");
136+
});
137+
138+
test("openai leaves the prompt unchanged when the stamp is old", async () => {
139+
const turns = history();
140+
const result = await transformFor({
141+
protocol: "openai",
142+
cacheWriteAt: NOW - 60 * MINUTE_MS,
143+
}).apply(turns, CTX);
144+
expect(result.output).toBe(turns);
145+
expect(result.record.reason).toBe("non-anthropic");
146+
});
147+
});
Lines changed: 166 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,166 @@
1+
// Prompt-only shrink after an Anthropic ephemeral cache write expires.
2+
// Compaction rewrites turns.jsonl and keeps a raw tail, so the next infer
3+
// still cache-writes those tool bodies. This transform runs inside
4+
// executeInfer: its output is what gets written to prompt.jsonl. It never
5+
// writes turns and never calls the compaction governor.
6+
7+
import type {
8+
ContentBlock,
9+
ContextTransform,
10+
ConversationTurn,
11+
} from "@intx/types/runtime";
12+
13+
import { cacheTtlMsFor } from "../provider/cache-ttl.js";
14+
15+
export type AnthropicCachePromptDeps = {
16+
nowMs: () => number;
17+
/** Wall time of the last Anthropic-protocol cache write, if one was stamped. */
18+
cacheWriteAt: () => number | undefined;
19+
/** Live adapter protocol (`InferenceSource.provider`). */
20+
protocol: () => string | undefined;
21+
};
22+
23+
const STRATEGY = "anthropic-cache-prompt";
24+
25+
type ToolResultBlock = Extract<ContentBlock, { type: "tool_result" }>;
26+
27+
function passthrough(
28+
turns: ConversationTurn[],
29+
reason: string,
30+
): Awaited<ReturnType<ContextTransform["apply"]>> {
31+
return {
32+
output: turns,
33+
record: {
34+
strategy: STRATEGY,
35+
version: "1",
36+
parameters: {},
37+
reason,
38+
decisions: { stubbed: 0 },
39+
},
40+
};
41+
}
42+
43+
function newestAssistantTextIndex(turns: readonly ConversationTurn[]): number {
44+
for (let i = turns.length - 1; i >= 0; i--) {
45+
const turn = turns[i];
46+
if (turn === undefined || turn.role !== "assistant") continue;
47+
if (
48+
turn.content.some(
49+
(block) => block.type === "text" && block.text.length > 0,
50+
)
51+
) {
52+
return i;
53+
}
54+
}
55+
return -1;
56+
}
57+
58+
// Tool results before the newest assistant text are already consumed.
59+
// When the model has not written that text yet, the newest result turn is
60+
// the unconsumed suffix and stays whole.
61+
function stubBoundary(turns: readonly ConversationTurn[]): number {
62+
const assistantText = newestAssistantTextIndex(turns);
63+
if (assistantText >= 0) return assistantText;
64+
for (let i = turns.length - 1; i >= 0; i--) {
65+
const turn = turns[i];
66+
if (turn?.content.some((block) => block.type === "tool_result")) return i;
67+
}
68+
return turns.length;
69+
}
70+
71+
function toolNames(turns: readonly ConversationTurn[]): Map<string, string> {
72+
const names = new Map<string, string>();
73+
for (const turn of turns) {
74+
for (const block of turn.content) {
75+
if (block.type === "tool_call") names.set(block.id, block.name);
76+
}
77+
}
78+
return names;
79+
}
80+
81+
function stubText(name: string | undefined): string {
82+
return name === undefined
83+
? "[tool result omitted]"
84+
: `[${name} result omitted]`;
85+
}
86+
87+
function stubToolResult(
88+
block: ToolResultBlock,
89+
names: ReadonlyMap<string, string>,
90+
): { block: ToolResultBlock; stubbed: boolean } {
91+
const text = stubText(names.get(block.callId));
92+
if (JSON.stringify(block.content).length <= text.length) {
93+
return { block, stubbed: false };
94+
}
95+
const next: ToolResultBlock = {
96+
type: "tool_result",
97+
callId: block.callId,
98+
content: [{ type: "text", text }],
99+
};
100+
if (block.isError === true) next.isError = true;
101+
return { block: next, stubbed: true };
102+
}
103+
104+
function stubTurn(
105+
turn: ConversationTurn,
106+
names: ReadonlyMap<string, string>,
107+
): { turn: ConversationTurn; stubbed: number } {
108+
let stubbed = 0;
109+
let changed = false;
110+
const content = turn.content.map((block) => {
111+
if (block.type !== "tool_result") return block;
112+
const next = stubToolResult(block, names);
113+
if (!next.stubbed) return block;
114+
changed = true;
115+
stubbed += 1;
116+
return next.block;
117+
});
118+
return { turn: changed ? { ...turn, content } : turn, stubbed };
119+
}
120+
121+
function shrinkPrompt(turns: ConversationTurn[]): {
122+
output: ConversationTurn[];
123+
stubbed: number;
124+
} {
125+
const boundary = stubBoundary(turns);
126+
const names = toolNames(turns);
127+
let stubbed = 0;
128+
let changed = false;
129+
const output = turns.map((turn, index) => {
130+
if (index >= boundary) return turn;
131+
const next = stubTurn(turn, names);
132+
stubbed += next.stubbed;
133+
if (next.turn !== turn) changed = true;
134+
return next.turn;
135+
});
136+
if (!changed) return { output: turns, stubbed: 0 };
137+
return { output, stubbed };
138+
}
139+
140+
export function createAnthropicCachePromptTransform(
141+
deps: AnthropicCachePromptDeps,
142+
): ContextTransform {
143+
return {
144+
name: STRATEGY,
145+
version: "1",
146+
async apply(turns, _ctx) {
147+
const ttl = cacheTtlMsFor(deps.protocol());
148+
if (ttl === undefined) return passthrough(turns, "non-anthropic");
149+
const at = deps.cacheWriteAt();
150+
if (at !== undefined && deps.nowMs() - at < ttl) {
151+
return passthrough(turns, "cache-warm");
152+
}
153+
const shrunk = shrinkPrompt(turns);
154+
return {
155+
output: shrunk.output,
156+
record: {
157+
strategy: STRATEGY,
158+
version: "1",
159+
parameters: {},
160+
reason: shrunk.stubbed > 0 ? "stubbed-tool-results" : "noop",
161+
decisions: { stubbed: shrunk.stubbed },
162+
},
163+
};
164+
},
165+
};
166+
}

‎src/session/assemble-runtime.ts‎

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -60,7 +60,9 @@ import { createDoomLoopCorrectiveNote } from "../agent/doom-loop-note.js";
6060
import type { AgentToolset } from "../agent/tools.js";
6161
import { createAgentWithLiveToolDispatch } from "../agent/live-tool-dispatch.js";
6262
import { createSessionStores } from "./optimized-context-store.js";
63+
import { getActiveRun } from "./active-run.js";
6364
import { createAttachmentRehydrateTransform } from "./attachment-store.js";
65+
import { createAnthropicCachePromptTransform } from "./anthropic-cache-prompt.js";
6466
import {
6567
applyRecordingPolicyToText,
6668
createCompactionArchive,
@@ -701,6 +703,16 @@ export function assembleChatAgent(wiring: ChatAgentWiring): AssembledChatAgent {
701703
createAttachmentRehydrateTransform((key) =>
702704
storageForAgent.readBlob(key),
703705
),
706+
createAnthropicCachePromptTransform({
707+
nowMs: () => Date.now(),
708+
cacheWriteAt: () => getActiveRun()?.lastCacheWriteAt,
709+
protocol: () => {
710+
const sources = wiring.getSources();
711+
const preferred = wiring.getDefaultSource();
712+
const match = sources.find((source) => source.id === preferred);
713+
return (match ?? sources[0])?.provider;
714+
},
715+
}),
704716
],
705717
},
706718
audit,

0 commit comments

Comments
 (0)