Skip to content

Commit 31b8041

Browse files
committed
fix(compaction): preserve fresh usage tail pressure
1 parent 7a99462 commit 31b8041

2 files changed

Lines changed: 142 additions & 0 deletions

File tree

src/auto-reply/reply/agent-runner-memory.test.ts

Lines changed: 78 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -746,6 +746,11 @@ describe("runMemoryFlushIfNeeded", () => {
746746
usage: { input: 240_000, output: 1_000 },
747747
},
748748
}),
749+
JSON.stringify({
750+
type: "compaction",
751+
timestamp: new Date(1_700_000_000_000).toISOString(),
752+
result: { summary: "compacted" },
753+
}),
749754
JSON.stringify({
750755
message: {
751756
role: "assistant",
@@ -955,6 +960,79 @@ describe("runMemoryFlushIfNeeded", () => {
955960
expect(compactCall.currentTokenCount).toBeGreaterThanOrEqual(96_000);
956961
});
957962

963+
it("preserves latest usage plus tail pressure when no post-compaction marker exists", async () => {
964+
const sessionFile = path.join(rootDir, "fresh-usage-tail-pressure-session.jsonl");
965+
await fs.writeFile(
966+
sessionFile,
967+
[
968+
JSON.stringify({
969+
type: "session",
970+
id: "session",
971+
}),
972+
JSON.stringify({
973+
id: "u1",
974+
parentId: null,
975+
message: {
976+
role: "user",
977+
content: "continue the current long context turn",
978+
},
979+
}),
980+
JSON.stringify({
981+
id: "a1",
982+
parentId: "u1",
983+
message: {
984+
role: "assistant",
985+
content: "small latest answer",
986+
usage: { input: 94_000, output: 500 },
987+
},
988+
}),
989+
JSON.stringify({
990+
id: "t1",
991+
parentId: "a1",
992+
message: {
993+
role: "tool",
994+
content: `post-usage tool tail ${"x".repeat(36_000)}`,
995+
},
996+
}),
997+
].join("\n"),
998+
"utf8",
999+
);
1000+
registerMemoryFlushPlanResolverForTest(() => ({
1001+
softThresholdTokens: 4_000,
1002+
forceFlushTranscriptBytes: 1_000_000_000,
1003+
reserveTokensFloor: 0,
1004+
prompt: "Pre-compaction memory flush.\nNO_REPLY",
1005+
systemPrompt: "Write memory to memory/YYYY-MM-DD.md.",
1006+
relativePath: "memory/2023-11-14.md",
1007+
}));
1008+
const sessionEntry: SessionEntry = {
1009+
sessionId: "session",
1010+
sessionFile,
1011+
updatedAt: Date.now(),
1012+
totalTokensFresh: false,
1013+
};
1014+
1015+
await runPreflightCompactionIfNeeded({
1016+
cfg: { agents: { defaults: { compaction: { memoryFlush: {} } } } },
1017+
followupRun: createTestFollowupRun({
1018+
sessionId: "session",
1019+
sessionFile,
1020+
sessionKey: "main",
1021+
}),
1022+
defaultModel: "anthropic/claude-opus-4-6",
1023+
agentCfgContextTokens: 100_000,
1024+
sessionEntry,
1025+
sessionStore: { main: sessionEntry },
1026+
sessionKey: "main",
1027+
storePath: path.join(rootDir, "sessions.json"),
1028+
isHeartbeat: false,
1029+
replyOperation: createReplyOperation(),
1030+
});
1031+
1032+
const compactCall = requireCompactEmbeddedPiSessionCall();
1033+
expect(compactCall.currentTokenCount).toBeGreaterThanOrEqual(100_000);
1034+
});
1035+
9581036
it("does not count bytes from a large latest usage record as post-usage tail pressure", async () => {
9591037
const sessionFile = path.join(rootDir, "large-usage-record-session.jsonl");
9601038
await fs.writeFile(

src/auto-reply/reply/agent-runner-memory.ts

Lines changed: 64 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -191,6 +191,7 @@ export type SessionTranscriptUsageSnapshot = {
191191
promptTokens?: number;
192192
outputTokens?: number;
193193
trailingBytesTokens?: number;
194+
hasPostUsageCompactionMarker?: boolean;
194195
};
195196

196197
// Keep a generous near-threshold window so large assistant outputs still trigger
@@ -220,6 +221,58 @@ function parseUsageFromTranscriptLine(line: string): ReturnType<typeof normalize
220221
return undefined;
221222
}
222223

224+
function collectTranscriptText(value: unknown): string {
225+
if (typeof value === "string") {
226+
return value;
227+
}
228+
if (!Array.isArray(value)) {
229+
return "";
230+
}
231+
return value
232+
.map((item) => {
233+
if (!item || typeof item !== "object" || Array.isArray(item)) {
234+
return "";
235+
}
236+
const text = (item as { text?: unknown }).text;
237+
return typeof text === "string" ? text : "";
238+
})
239+
.filter(Boolean)
240+
.join("\n");
241+
}
242+
243+
function transcriptLineHasPostUsageCompactionMarker(line: string): boolean {
244+
const trimmed = line.trim();
245+
if (!trimmed) {
246+
return false;
247+
}
248+
try {
249+
const parsed = JSON.parse(trimmed) as {
250+
type?: unknown;
251+
message?: { content?: unknown };
252+
payload?: { type?: unknown; text?: unknown };
253+
};
254+
if (parsed.type === "compaction" || parsed.type === "session.compacted") {
255+
return true;
256+
}
257+
const payloadType = parsed.payload?.type;
258+
if (payloadType === "compaction" || payloadType === "session.compacted") {
259+
return true;
260+
}
261+
const text = [
262+
collectTranscriptText(parsed.message?.content),
263+
collectTranscriptText(parsed.payload?.text),
264+
]
265+
.filter(Boolean)
266+
.join("\n");
267+
return (
268+
text.includes("[Post-compaction context refresh]") ||
269+
text.includes("Session was just compacted.")
270+
);
271+
} catch {
272+
return false;
273+
}
274+
}
275+
223276
function resolveSessionLogPath(
224277
sessionId?: string,
225278
sessionEntry?: SessionEntry,
@@ -257,6 +310,7 @@ function deriveTranscriptUsageSnapshot(
257310
| {
258311
usage: ReturnType<typeof normalizeUsage> | undefined;
259312
trailingBytes?: number;
313+
hasPostUsageCompactionMarker?: boolean;
260314
}
261315
| undefined,
262316
): SessionTranscriptUsageSnapshot | undefined {
@@ -276,6 +330,7 @@ function deriveTranscriptUsageSnapshot(
276330
return {
277331
promptTokens,
278332
outputTokens,
333+
hasPostUsageCompactionMarker: snapshot.hasPostUsageCompactionMarker === true ? true : undefined,
279334
trailingBytesTokens:
280335
typeof snapshot.trailingBytes === "number" &&
281336
Number.isFinite(snapshot.trailingBytes) &&
@@ -360,6 +415,7 @@ async function readLastNonzeroUsageFromSessionLog(logPath: string) {
360415
const stat = await handle.stat();
361416
let position = stat.size;
362417
let leadingPartial = "";
418+
let postUsageCompactionMarkerSeen = false;
363419
while (position > 0) {
364420
const chunkSize = Math.min(TRANSCRIPT_TAIL_CHUNK_BYTES, position);
365421
const start = position - chunkSize;
@@ -384,16 +440,23 @@ async function readLastNonzeroUsageFromSessionLog(logPath: string) {
384440
return {
385441
usage,
386442
trailingBytes: suffixBytesOutsideCombined + trailingBytesInChunk,
443+
hasPostUsageCompactionMarker:
444+
postUsageCompactionMarkerSeen ||
445+
trailingLines.some(transcriptLineHasPostUsageCompactionMarker),
387446
};
388447
}
389448
}
449+
if (!postUsageCompactionMarkerSeen) {
450+
postUsageCompactionMarkerSeen = lines.some(transcriptLineHasPostUsageCompactionMarker);
451+
}
390452
position = start;
391453
}
392454
const usage = parseUsageFromTranscriptLine(leadingPartial);
393455
return usage
394456
? {
395457
usage,
396458
trailingBytes: Math.max(0, stat.size - Buffer.byteLength(leadingPartial, "utf8")),
459+
hasPostUsageCompactionMarker: postUsageCompactionMarkerSeen,
397460
}
398461
: undefined;
399462
} finally {
@@ -461,6 +524,7 @@ async function estimatePromptTokensFromSessionTranscript(params: {
461524
const outputTokens = snapshot.usage?.outputTokens;
462525
const rawUsagePromptTokens = Math.ceil(promptTokens);
463526
const boundedUsagePromptTokens =
527+
snapshot.usage?.hasPostUsageCompactionMarker === true &&
464528
typeof estimatedMessageTokens === "number" &&
465529
rawUsagePromptTokens > estimatedMessageTokens * 2 + 10_000
466530
? estimatedMessageTokens

0 commit comments

Comments
 (0)