Skip to content

Commit 42b70f3

Browse files
committed
Restore archive leaf entries after failed compaction
1 parent 3039a71 commit 42b70f3

2 files changed

Lines changed: 230 additions & 21 deletions

File tree

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

Lines changed: 183 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
11
import { defined } from "../../tests/helpers/defined.js";
2-
import { describe, test, expect } from "bun:test";
2+
import { describe, test, expect, spyOn } from "bun:test";
33
import fs from "node:fs";
44
import os from "node:os";
55
import path from "node:path";
@@ -744,6 +744,188 @@ async function gitLsTree(dir: string): Promise<string[]> {
744744
}
745745

746746
describe("createOptimizedContextStore unpublished rewrite", () => {
747+
test("partial archive staging and failed rollback fence shared audit publication until reconciliation succeeds", async () => {
748+
const dir = tempDir();
749+
const { storage } = await createSessionStores(dir);
750+
const archiveDir = path.join(dir, "evidence-archive");
751+
fs.mkdirSync(archiveDir, { recursive: true });
752+
fs.writeFileSync(path.join(archiveDir, "index.jsonl"), "original");
753+
await storage.writeTurns([turn("old-a"), turn("old-b")]);
754+
await storage.writeMetadata(EMPTY_CHECKPOINT_METADATA);
755+
const published = await storage.commit({ message: "original" });
756+
const { audit } = await createSessionStores(dir);
757+
fs.writeFileSync(path.join(archiveDir, "index.jsonl"), "changed");
758+
fs.writeFileSync(path.join(archiveDir, "added.jsonl"), "new");
759+
await storage.writeTurns([turn("[Compacted prior context]")]);
760+
const originalAdd = git.add;
761+
const originalReset = git.resetIndex;
762+
let resetFails = true;
763+
const staged: string[] = [];
764+
const addSpy = spyOn(git, "add").mockImplementation(async (args) => {
765+
if (args.dir === dir && args.filepath === "evidence-archive/index.jsonl")
766+
throw new Error("partial staging failure");
767+
await originalAdd(args);
768+
if (args.dir === dir)
769+
staged.push(...(Array.isArray(args.filepath) ? args.filepath : [args.filepath]));
770+
});
771+
const resetSpy = spyOn(git, "resetIndex").mockImplementation(async (args) => {
772+
if (args.dir === dir && args.filepath.startsWith("evidence-archive/") && resetFails)
773+
throw new Error("injected reset failure");
774+
return originalReset(args);
775+
});
776+
const errorRecord = {
777+
source: "reactor" as const,
778+
category: "test",
779+
message: "error",
780+
fatal: false,
781+
timestamp: new Date().toISOString(),
782+
sessionId: "s",
783+
seq: 1,
784+
};
785+
try {
786+
await expect(storage.commit({ message: "partial" })).rejects.toThrow(
787+
"rollback remains pending",
788+
);
789+
expect(staged).toContain("evidence-archive/added.jsonl");
790+
await expect(audit.commitErrors([errorRecord])).rejects.toThrow("injected reset failure");
791+
await expect(
792+
audit.commitAudit([
793+
{
794+
callId: "a",
795+
tool: "read_file",
796+
arguments: {},
797+
authz: null,
798+
result: { content: "ok", isError: false },
799+
timestamp: new Date().toISOString(),
800+
sessionId: "s",
801+
seq: 2,
802+
},
803+
]),
804+
).rejects.toThrow("injected reset failure");
805+
expect(await git.resolveRef({ fs, dir, ref: "HEAD" })).toBe(published.hash);
806+
resetFails = false;
807+
await audit.commitErrors([errorRecord]);
808+
expect(
809+
(await git.listFiles({ fs, dir, ref: "HEAD" })).filter((name) =>
810+
name.startsWith("evidence-archive"),
811+
),
812+
).toEqual(["evidence-archive/index.jsonl"]);
813+
const { blob } = await git.readBlob({
814+
fs,
815+
dir,
816+
oid: await git.resolveRef({ fs, dir, ref: "HEAD" }),
817+
filepath: "evidence-archive/index.jsonl",
818+
});
819+
expect(new TextDecoder().decode(blob)).toBe("original");
820+
expect(fs.readFileSync(path.join(archiveDir, "index.jsonl"), "utf8")).toBe("changed");
821+
} finally {
822+
addSpy.mockRestore();
823+
resetSpy.mockRestore();
824+
}
825+
await storage.commit({ message: "retry" });
826+
expect(
827+
(await git.listFiles({ fs, dir, ref: "HEAD" })).filter((name) =>
828+
name.startsWith("evidence-archive"),
829+
),
830+
).toEqual(["evidence-archive/added.jsonl", "evidence-archive/index.jsonl"]);
831+
});
832+
test.each([true, false])(
833+
"failed compact restores concrete archive index leaves (already tracked: %s)",
834+
async (tracked) => {
835+
const dir = tempDir();
836+
const { loadOrCreateCommitSigner } = await import("./commit-signer.js");
837+
const sign = await loadOrCreateCommitSigner(dir);
838+
let fail = false;
839+
const { storage, audit } = await createSessionStores(dir, {
840+
signer: (payload) => {
841+
if (fail) throw new Error("injected signing failure");
842+
return sign(payload);
843+
},
844+
});
845+
const archiveDir = path.join(dir, "evidence-archive");
846+
fs.mkdirSync(archiveDir, { recursive: true });
847+
if (tracked) {
848+
fs.writeFileSync(path.join(archiveDir, "index.jsonl"), "original evidence");
849+
fs.writeFileSync(path.join(archiveDir, "deleted.jsonl"), "old evidence");
850+
}
851+
await storage.writeTurns([turn("old-a"), turn("old-b")]);
852+
await storage.writeMetadata(EMPTY_CHECKPOINT_METADATA);
853+
const original = await storage.commit({ message: "original" });
854+
const originalLeaves = (await git.listFiles({ fs, dir })).filter((name) =>
855+
name.startsWith("evidence-archive"),
856+
);
857+
fs.writeFileSync(path.join(archiveDir, "index.jsonl"), "new evidence");
858+
fs.writeFileSync(path.join(archiveDir, "added.jsonl"), "added evidence");
859+
if (tracked) fs.unlinkSync(path.join(archiveDir, "deleted.jsonl"));
860+
await storage.writeTurns([turn("[Compacted prior context]")]);
861+
fail = true;
862+
await expect(storage.commit({ message: "failed compact" })).rejects.toThrow(
863+
"injected signing failure",
864+
);
865+
expect(
866+
(await git.listFiles({ fs, dir })).filter((name) => name.startsWith("evidence-archive")),
867+
).toEqual(originalLeaves);
868+
expect(await git.resolveRef({ fs, dir, ref: "HEAD" })).toBe(original.hash);
869+
for (const [, head, , stage] of await git.statusMatrix({
870+
fs,
871+
dir,
872+
filepaths: ["evidence-archive"],
873+
}))
874+
expect(stage).toBe(head);
875+
expect(fs.readFileSync(path.join(archiveDir, "index.jsonl"), "utf8")).toBe("new evidence");
876+
fail = false;
877+
await audit.commitAudit([
878+
{
879+
callId: "audit",
880+
tool: "read_file",
881+
arguments: {},
882+
authz: null,
883+
result: { content: "ok", isError: false },
884+
timestamp: new Date().toISOString(),
885+
sessionId: "s",
886+
seq: 1,
887+
},
888+
]);
889+
await audit.commitErrors([
890+
{
891+
source: "reactor",
892+
category: "test",
893+
message: "test error",
894+
fatal: false,
895+
timestamp: new Date().toISOString(),
896+
sessionId: "s",
897+
seq: 2,
898+
},
899+
]);
900+
expect(
901+
(await git.listFiles({ fs, dir, ref: "HEAD" })).filter((name) =>
902+
name.startsWith("evidence-archive"),
903+
),
904+
).toEqual(originalLeaves);
905+
if (tracked) {
906+
const { blob } = await git.readBlob({
907+
fs,
908+
dir,
909+
oid: await git.resolveRef({ fs, dir, ref: "HEAD" }),
910+
filepath: "evidence-archive/index.jsonl",
911+
});
912+
expect(new TextDecoder().decode(blob)).toBe("original evidence");
913+
}
914+
await storage.commit({ message: "retry compact" });
915+
expect(
916+
(await git.listFiles({ fs, dir, ref: "HEAD" })).filter((name) =>
917+
name.startsWith("evidence-archive"),
918+
),
919+
).toEqual(["evidence-archive/added.jsonl", "evidence-archive/index.jsonl"]);
920+
const { blob } = await git.readBlob({
921+
fs,
922+
dir,
923+
oid: await git.resolveRef({ fs, dir, ref: "HEAD" }),
924+
filepath: "evidence-archive/index.jsonl",
925+
});
926+
expect(new TextDecoder().decode(blob)).toBe("new evidence");
927+
},
928+
);
747929
test("rewrite writeTurns stays off the live generation until commit", async () => {
748930
const dir = tempDir();
749931
const store = await createOptimizedContextStore(dir);

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

Lines changed: 47 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -369,14 +369,29 @@ async function resetIndexPaths(
369369
filepaths: readonly string[],
370370
): Promise<void> {
371371
for (const filepath of filepaths) {
372-
try {
373-
await git.resetIndex({ fs, dir, filepath });
374-
} catch {
375-
// Not in the index; vendor restore already covers its own paths.
376-
}
372+
await git.resetIndex({ fs, dir, filepath });
377373
}
378374
}
379375

376+
const pendingCheckpointRollbacks = new Map<string, () => Promise<void>>();
377+
378+
async function reconcilePendingCheckpoint(dir: string): Promise<void> {
379+
const key = path.resolve(dir);
380+
const rollback = pendingCheckpointRollbacks.get(key);
381+
if (rollback === undefined) return;
382+
await rollback();
383+
pendingCheckpointRollbacks.delete(key);
384+
}
385+
386+
function withReconciledDirLock<T>(dir: string, operation: () => Promise<T>): Promise<T> {
387+
return withResolvedDirLock(dir, async () => {
388+
// All store instances sharing this index must finish a failed rollback
389+
// before an audit, error, or checkpoint commit can publish it.
390+
await reconcilePendingCheckpoint(dir);
391+
return operation();
392+
});
393+
}
394+
380395
async function restorePublishedTurnFiles(dir: string): Promise<void> {
381396
const hash = await git.resolveRef({ fs, dir, ref: "HEAD" });
382397
const present = new Set(await git.listFiles({ fs, dir, ref: hash }));
@@ -633,7 +648,7 @@ export async function createSessionStores(
633648
pendingBlobFilepaths.add(`${TOOL_OUTPUT_DIR}/${filename}`);
634649
},
635650
async commit(options, signal) {
636-
return withResolvedDirLock(dir, async () => {
651+
return withReconciledDirLock(dir, async () => {
637652
const stagedRewrite = unpublishedRewrite;
638653
const extraPaths: string[] = [];
639654
try {
@@ -646,8 +661,14 @@ export async function createSessionStores(
646661
if (await pathExists(path.join(dir, filepath))) toAdd.push(filepath);
647662
else toRemove.push(filepath);
648663
}
649-
if (await pathExists(path.join(dir, EVIDENCE_ARCHIVE_DIR))) {
650-
toAdd.push(EVIDENCE_ARCHIVE_DIR);
664+
const archiveRows = await git.statusMatrix({
665+
fs,
666+
dir,
667+
filepaths: [EVIDENCE_ARCHIVE_DIR],
668+
});
669+
for (const [filepath, , worktree] of archiveRows) {
670+
if (worktree === 0) toRemove.push(filepath);
671+
else toAdd.push(filepath);
651672
}
652673
// Reconcile disk segments even if a prior process lost its pending set.
653674
await reconcileSegmentStaging(dir, TURNS_FILE, toAdd, toRemove);
@@ -659,11 +680,7 @@ export async function createSessionStores(
659680
await git.add({ fs, dir, filepath });
660681
}
661682
for (const filepath of remove) {
662-
try {
663-
await git.remove({ fs, dir, filepath });
664-
} catch {
665-
// Already absent from the index.
666-
}
683+
await git.remove({ fs, dir, filepath });
667684
}
668685
const committed = await base.commit(options, signal);
669686
pendingBlobFilepaths.clear();
@@ -674,20 +691,30 @@ export async function createSessionStores(
674691
}
675692
return committed;
676693
} catch (cause) {
677-
await resetIndexPaths(dir, extraPaths);
678-
if (stagedRewrite !== null) {
679-
await restorePublishedTurnFiles(dir);
680-
writeTurnsSegmented = createSegmentedJSONLWriter(dir, TURNS_FILE);
681-
liveTurnRefs = null;
694+
pendingCheckpointRollbacks.set(path.resolve(dir), async () => {
695+
await resetIndexPaths(dir, extraPaths);
696+
if (stagedRewrite !== null) {
697+
await restorePublishedTurnFiles(dir);
698+
writeTurnsSegmented = createSegmentedJSONLWriter(dir, TURNS_FILE);
699+
liveTurnRefs = null;
700+
}
701+
});
702+
try {
703+
await reconcilePendingCheckpoint(dir);
704+
} catch (rollbackCause) {
705+
throw new AggregateError(
706+
[cause, rollbackCause],
707+
"Checkpoint failed; index rollback remains pending",
708+
);
682709
}
683710
throw cause;
684711
}
685712
});
686713
},
687714
commitAudit: (records, signal) =>
688-
withResolvedDirLock(dir, () => base.commitAudit(records, signal)),
715+
withReconciledDirLock(dir, () => base.commitAudit(records, signal)),
689716
commitErrors: (records, signal) =>
690-
withResolvedDirLock(dir, () => base.commitErrors(records, signal)),
717+
withReconciledDirLock(dir, () => base.commitErrors(records, signal)),
691718
loadAudit: (sessionId, signal) => base.loadAudit(sessionId, signal),
692719
};
693720

0 commit comments

Comments
 (0)