Skip to content

Commit 1bdfa13

Browse files
committed
Stage compaction rewrites until the store publishes
Persist blobs and stage writeTurns first. Replace reactor memory only after commit so an interrupt resumes the old generation.
1 parent ba1d031 commit 1bdfa13

6 files changed

Lines changed: 485 additions & 50 deletions

File tree

‎src/session/optimized-context-store.test.ts‎

Lines changed: 107 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -606,3 +606,110 @@ describe("createOptimizedContextStore checkpoint", () => {
606606
expect(await headIdent(nameOnlyWhitespaceEmail)).toEqual(HARNESS_IDENT);
607607
});
608608
});
609+
610+
function turnTexts(turns: ConversationTurn[]): string[] {
611+
return turns.map((t) => (t.content[0] as { text: string }).text);
612+
}
613+
614+
async function gitLsTree(dir: string): Promise<string[]> {
615+
const proc = Bun.spawn(["git", "-C", dir, "ls-tree", "-r", "--name-only", "HEAD"], {
616+
stdout: "pipe",
617+
stderr: "pipe",
618+
});
619+
const [exitCode, stdout, stderr] = await Promise.all([
620+
proc.exited,
621+
new Response(proc.stdout).text(),
622+
new Response(proc.stderr).text(),
623+
]);
624+
if (exitCode !== 0) {
625+
throw new Error(`git ls-tree failed: ${stderr.trim() || stdout.trim()}`);
626+
}
627+
return stdout.split("\n").filter((line) => line.length > 0);
628+
}
629+
630+
describe("createOptimizedContextStore unpublished rewrite", () => {
631+
test("rewrite writeTurns stays off the live generation until commit", async () => {
632+
const dir = tempDir();
633+
const store = await createOptimizedContextStore(dir);
634+
const original = [turn("keep-a"), turn("keep-b"), turn("drop-me")];
635+
await store.writeTurns(original);
636+
await store.writeMetadata(EMPTY_CHECKPOINT_METADATA);
637+
await store.commit({ message: "published original" });
638+
639+
const compacted = [turn("[Compacted prior context]"), turn("keep-b")];
640+
await store.writeTurns(compacted);
641+
642+
const loaded = await store.load();
643+
expect(turnTexts(loaded.turns)).toEqual(["keep-a", "keep-b", "drop-me"]);
644+
645+
await store.commit({ message: "publish compact" });
646+
const published = await store.load();
647+
expect(turnTexts(published.turns)).toEqual(["[Compacted prior context]", "keep-b"]);
648+
});
649+
650+
test("omitting commit leaves a new store on the old generation", async () => {
651+
const dir = tempDir();
652+
const store = await createOptimizedContextStore(dir);
653+
await store.writeTurns([turn("old-a"), turn("old-b")]);
654+
await store.writeMetadata(EMPTY_CHECKPOINT_METADATA);
655+
await store.commit({ message: "old" });
656+
657+
await store.writeBlob(
658+
"new-blob",
659+
new TextEncoder().encode("needed-by-new-turns"),
660+
"text/plain",
661+
);
662+
await store.writeTurns([turn("[Compacted prior context]")]);
663+
664+
const crashed = await createOptimizedContextStore(dir);
665+
const loaded = await crashed.load();
666+
expect(turnTexts(loaded.turns)).toEqual(["old-a", "old-b"]);
667+
});
668+
669+
test("append writeTurns is still visible before commit", async () => {
670+
const dir = tempDir();
671+
const store = await createOptimizedContextStore(dir);
672+
const first = turn("one");
673+
await store.writeTurns([first]);
674+
await store.writeMetadata(EMPTY_CHECKPOINT_METADATA);
675+
await store.commit({ message: "one" });
676+
677+
await store.writeTurns([first, turn("two")]);
678+
const loaded = await store.load();
679+
expect(turnTexts(loaded.turns)).toEqual(["one", "two"]);
680+
});
681+
682+
test("folds evidence-archive into the compact commit tree", async () => {
683+
const dir = tempDir();
684+
const store = await createOptimizedContextStore(dir);
685+
await store.writeTurns([turn("old")]);
686+
await store.writeMetadata(EMPTY_CHECKPOINT_METADATA);
687+
await store.commit({ message: "old" });
688+
689+
const archiveDir = path.join(dir, "evidence-archive");
690+
fs.mkdirSync(archiveDir, { recursive: true });
691+
fs.writeFileSync(path.join(archiveDir, "index.jsonl"), "{}\n");
692+
await store.writeTurns([turn("compacted")]);
693+
await store.writeMetadata(EMPTY_CHECKPOINT_METADATA);
694+
await store.commit({ message: "compact" });
695+
696+
expect(await gitLsTree(dir)).toContain("evidence-archive/index.jsonl");
697+
const loaded = await store.load();
698+
expect(turnTexts(loaded.turns)).toEqual(["compacted"]);
699+
});
700+
701+
test("readAt of the old hash is not the load completeness path", async () => {
702+
const dir = tempDir();
703+
const store = await createOptimizedContextStore(dir);
704+
await store.writeTurns([turn("era-1")]);
705+
await store.writeMetadata(EMPTY_CHECKPOINT_METADATA);
706+
const first = await store.commit({ message: "era-1" });
707+
708+
await store.writeTurns([turn("era-2")]);
709+
await store.writeMetadata(EMPTY_CHECKPOINT_METADATA);
710+
await store.commit({ message: "era-2" });
711+
712+
expect(turnTexts(await store.readAt(first.hash))).toEqual(["era-1"]);
713+
expect(turnTexts((await store.load()).turns)).toEqual(["era-2"]);
714+
});
715+
});

‎src/session/optimized-context-store.ts‎

Lines changed: 101 additions & 42 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@ const RESPONSE_FILE = "response.jsonl";
2626
const MANIFEST_FILE = "manifest.jsonl";
2727
const METADATA_FILE = "metadata.json";
2828
const TOOL_OUTPUT_DIR = "tool-output";
29+
const EVIDENCE_ARCHIVE_DIR = "evidence-archive";
2930

3031
const log = getLogger([LOG_NAMESPACE_ROOT, "session", "context-store"]);
3132

@@ -447,6 +448,30 @@ export async function createOptimizedContextStore(
447448
const pendingSegmentPaths = new Set<string>();
448449
const writeTurnsSegmented = createSegmentedJSONLWriter(dir, TURNS_FILE);
449450
const writePromptSegmented = createSegmentedJSONLWriter(dir, PROMPT_FILE);
451+
let liveTurnRefs: readonly ConversationTurn[] | null = null;
452+
let unpublishedRewrite: ConversationTurn[] | null = null;
453+
454+
function refPrefixLength(
455+
prev: readonly ConversationTurn[],
456+
next: readonly ConversationTurn[],
457+
): number {
458+
const max = Math.min(prev.length, next.length);
459+
let prefix = 0;
460+
while (prefix < max && prev[prefix] === next[prefix]) prefix++;
461+
return prefix;
462+
}
463+
464+
function contentPrefixLength(
465+
prev: readonly ConversationTurn[],
466+
next: readonly ConversationTurn[],
467+
): number {
468+
const max = Math.min(prev.length, next.length);
469+
let prefix = 0;
470+
while (prefix < max && JSON.stringify(prev[prefix]) === JSON.stringify(next[prefix])) {
471+
prefix++;
472+
}
473+
return prefix;
474+
}
450475

451476
async function writeSegmented(
452477
writer: ReturnType<typeof createSegmentedJSONLWriter>,
@@ -488,6 +513,70 @@ export async function createOptimizedContextStore(
488513
return [...baseTurns, ...parsedExtras.slice(0, keepExtras).flat()];
489514
}
490515

516+
async function loadLive(signal?: AbortSignal) {
517+
try {
518+
const baseResult = await base.load(signal);
519+
const extraTexts = await readExtraSegmentTexts(dir, TURNS_FILE);
520+
if (extraTexts.length === 0) return baseResult;
521+
const turns = await loadTurnsWithoutMalformedToolSequence(baseResult.turns, extraTexts);
522+
return { ...baseResult, turns };
523+
} catch (cause) {
524+
log.warn("base context store load failed; recovering turns from disk segments", {
525+
cause: cause instanceof Error ? cause.message : String(cause),
526+
});
527+
let baseTurns: ConversationTurn[];
528+
try {
529+
// Prefer resilient parse of segment 0 alone so orphan-tail heal still runs.
530+
// skipMalformed: mid-file garbage/interleaved records must not kill resume
531+
// (CL-7052); null-pad stripping and torn-tail drop still apply.
532+
const basePath = path.join(dir, TURNS_FILE);
533+
if (await pathExists(basePath)) {
534+
const text = await fs.promises.readFile(basePath, "utf-8");
535+
baseTurns = parseSegmentTurns(text, true, TURNS_FILE, true);
536+
} else {
537+
baseTurns = [];
538+
}
539+
} catch (parseCause) {
540+
throw new Error(
541+
`failed to load ${TURNS_FILE}: ${
542+
parseCause instanceof Error ? parseCause.message : String(parseCause)
543+
}`,
544+
{ cause: parseCause },
545+
);
546+
}
547+
const extraTexts = await readExtraSegmentTexts(dir, TURNS_FILE);
548+
const turns =
549+
extraTexts.length === 0
550+
? baseTurns
551+
: await loadTurnsWithoutMalformedToolSequence(baseTurns, extraTexts);
552+
const metadata = await loadMetadataSoft(() => base.loadMetadata());
553+
return { turns, ...metadata };
554+
}
555+
}
556+
557+
async function writeTurnsLiveOrStage(turns: readonly ConversationTurn[]): Promise<void> {
558+
if (unpublishedRewrite !== null) {
559+
unpublishedRewrite = [...turns];
560+
return;
561+
}
562+
if (liveTurnRefs !== null) {
563+
if (refPrefixLength(liveTurnRefs, turns) < liveTurnRefs.length) {
564+
unpublishedRewrite = [...turns];
565+
return;
566+
}
567+
await writeSegmented(writeTurnsSegmented, turns);
568+
liveTurnRefs = [...turns];
569+
return;
570+
}
571+
const live = (await loadLive()).turns;
572+
if (live.length > 0 && contentPrefixLength(live, turns) < live.length) {
573+
unpublishedRewrite = [...turns];
574+
return;
575+
}
576+
await writeSegmented(writeTurnsSegmented, turns);
577+
liveTurnRefs = [...turns];
578+
}
579+
491580
return {
492581
// Full-history read. Called by the reactor during initialization, where
493582
// the complete turn history is the actual live conversation state, not an
@@ -498,47 +587,7 @@ export async function createOptimizedContextStore(
498587
// garbage turns.jsonl), recover usable turns via resilient segment parse
499588
// and re-read metadata via the base schema (soft-empty only if that fails
500589
// too) so resume does not die on a bare Bun JSON token or wipe pending ops.
501-
async load(signal) {
502-
try {
503-
const baseResult = await base.load(signal);
504-
const extraTexts = await readExtraSegmentTexts(dir, TURNS_FILE);
505-
if (extraTexts.length === 0) return baseResult;
506-
const turns = await loadTurnsWithoutMalformedToolSequence(baseResult.turns, extraTexts);
507-
return { ...baseResult, turns };
508-
} catch (cause) {
509-
log.warn("base context store load failed; recovering turns from disk segments", {
510-
cause: cause instanceof Error ? cause.message : String(cause),
511-
});
512-
let baseTurns: ConversationTurn[];
513-
try {
514-
// Prefer resilient parse of segment 0 alone so orphan-tail heal still runs.
515-
// skipMalformed: mid-file garbage/interleaved records must not kill resume
516-
// (CL-7052); null-pad stripping and torn-tail drop still apply.
517-
const basePath = path.join(dir, TURNS_FILE);
518-
if (await pathExists(basePath)) {
519-
const text = await fs.promises.readFile(basePath, "utf-8");
520-
baseTurns = parseSegmentTurns(text, true, TURNS_FILE, true);
521-
} else {
522-
baseTurns = [];
523-
}
524-
} catch (parseCause) {
525-
// Unrecoverable: rethrow with the file name in the message.
526-
throw new Error(
527-
`failed to load ${TURNS_FILE}: ${
528-
parseCause instanceof Error ? parseCause.message : String(parseCause)
529-
}`,
530-
{ cause: parseCause },
531-
);
532-
}
533-
const extraTexts = await readExtraSegmentTexts(dir, TURNS_FILE);
534-
const turns =
535-
extraTexts.length === 0
536-
? baseTurns
537-
: await loadTurnsWithoutMalformedToolSequence(baseTurns, extraTexts);
538-
const metadata = await loadMetadataSoft(() => base.loadMetadata());
539-
return { turns, ...metadata };
540-
}
541-
},
590+
load: (signal) => loadLive(signal),
542591
setConnectorState: (state) => base.setConnectorState(state),
543592
branch: (name, signal) => base.branch(name, signal),
544593
log: (limit, signal) => base.log(limit, signal),
@@ -562,7 +611,7 @@ export async function createOptimizedContextStore(
562611
writePrompt: (turns) => writeSegmented(writePromptSegmented, turns),
563612
writeResponse: (turn, signal) => base.writeResponse(turn, signal),
564613
writeManifest: (records, signal) => base.writeManifest(records, signal),
565-
writeTurns: (turns) => writeSegmented(writeTurnsSegmented, turns),
614+
writeTurns: (turns) => writeTurnsLiveOrStage(turns),
566615
writeMetadata: (metadata, signal) => base.writeMetadata(metadata, signal),
567616
readManifestHistory: (limit, signal) => base.readManifestHistory(limit, signal),
568617
async writeBlob(key, bytes, contentType, signal) {
@@ -571,6 +620,12 @@ export async function createOptimizedContextStore(
571620
pendingBlobFilepaths.add(`${TOOL_OUTPUT_DIR}/${filename}`);
572621
},
573622
async commit(options, _signal) {
623+
if (unpublishedRewrite !== null) {
624+
await writeSegmented(writeTurnsSegmented, unpublishedRewrite);
625+
liveTurnRefs = unpublishedRewrite;
626+
unpublishedRewrite = null;
627+
}
628+
574629
const toAdd: string[] = [];
575630
const toRemove: string[] = [];
576631

@@ -584,6 +639,10 @@ export async function createOptimizedContextStore(
584639
else toRemove.push(filepath);
585640
}
586641

642+
if (await pathExists(path.join(dir, EVIDENCE_ARCHIVE_DIR))) {
643+
toAdd.push(EVIDENCE_ARCHIVE_DIR);
644+
}
645+
587646
// Disk is source of truth for which turn/prompt segments should remain
588647
// tracked after a rewrite or heal, even if pendingSegmentPaths was lost.
589648
await reconcileSegmentStaging(dir, TURNS_FILE, toAdd, toRemove);

0 commit comments

Comments
 (0)