feat: add JSONL sidecar reclamation

This commit is contained in:
Mario Zechner
2026-09-23 09:54:05 +02:00
parent 4bc1a2fe53
commit b313731b80
4 changed files with 742 additions and 73 deletions
+11 -4
View File
@@ -79,9 +79,10 @@ arguments, and reads never expose backend-owned cached objects.
JSONL creation has `fsync?: boolean`, defaulting to `false`. With `false`, append
sidecars and then the marker without an explicit flush. With `true`, append all
affected sidecars, flush each affected sidecar, and then append the main marker.
Do not explicitly flush `main.jsonl`. A main-only commit has no sidecars to
flush. Any uncertain append or flush failure poisons the open backend and
publishes no prepared in-memory mutation.
Do not explicitly flush `main.jsonl` for ordinary publication. A main-only commit
has no sidecars to flush. Any uncertain publication append or flush failure
poisons the open backend and publishes no prepared in-memory mutation. Package 5
adds the separate post-publication flush required to authorize reclamation.
Fault-test torn/short sidecar writes, failures between sidecars, every marker
boundary, unconfirmed tails, missing confirmed data, poisoned writes, exact-byte
@@ -92,7 +93,13 @@ conformance suite directly and after reopen.
## 5. JSONL reclamation
Implement task-document retirement and current-only base reclamation using
committed markers, temporary replacement, rename, and descriptor invalidation.
committed markers and descriptor invalidation. Remove a sidecar directly when no
records remain; otherwise use temporary replacement and rename. With `fsync:
true`, flush `main.jsonl` once before destructive reclamation so the authorizing
marker cannot disappear while cleanup survives; if that flush fails, skip
reclamation without failing the already-published commit. Flush a non-empty
temporary replacement before rename. This is not publication flushing or main
compaction.
Crash-test every rewrite/rename boundary. Verify that rewindable history is
never reclaimed and default no-fsync behavior matches the specification.
+14 -8
View File
@@ -2120,10 +2120,14 @@ JSONL creation accepts an `fsync` option that defaults to `false`. Without
`fsync`, it guarantees ordinary process-crash consistency, not survival of
power, host, kernel, or filesystem failure. With `fsync: true`, the backend
appends all affected sidecar records, flushes each affected sidecar, and only
then appends the main marker. It does not explicitly flush `main.jsonl`; an
acknowledged tail commit may therefore still disappear, but a marker that
survives should not overtake its sidecar data. A main-only commit has no
sidecars to flush.
then appends the main marker. Ordinary publication does not explicitly flush
`main.jsonl`; an acknowledged tail commit may therefore still disappear, but a
marker that survives should not overtake its sidecar data. A main-only commit has
no sidecars to flush. Before destructive reclamation with `fsync: true`, the
backend flushes `main.jsonl` once so the authorizing marker cannot disappear
while its replacement or removal survives. If that flush fails, the committed
state remains published and reclamation is deferred. A non-empty temporary
replacement is also flushed before rename.
Recovery:
@@ -2135,10 +2139,12 @@ Recovery:
record unnecessary.
- Any uncertain append failure poisons the open backend.
Reclamation starts only after the authorizing base/retirement commits. It writes
a temporary replacement, renames it, and invalidates cached file descriptors so
future appends cannot target an unlinked inode. `main.jsonl` is not compacted in
the initial implementation.
Reclamation starts only after the authorizing base/retirement commits. When no
sidecar records remain, it removes the sidecar directly. Otherwise, it writes a
temporary replacement, renames it, and invalidates cached file descriptors so
future appends cannot target an unlinked inode. Flushing `main.jsonl` to
authorize reclamation does not compact it. `main.jsonl` is not compacted in the
initial implementation.
## 12. API footguns
+230 -50
View File
@@ -19,6 +19,7 @@ import { MemoryStorage } from "../memory.ts";
const FORMAT_VERSION = 1;
const MAIN_FILE = "main.jsonl";
const RECLAIM_SUFFIX = ".reclaim";
const textDecoder = new TextDecoder("utf-8", { fatal: true });
type StoredTask = TaskRecord<JsonValue, JsonValue, JsonValue>;
@@ -92,6 +93,11 @@ const sidecarFileName = (kind: "doc" | "task", id: Id): string => `${kind}-${id}
const isSidecarFileName = (name: string): boolean => /^(?:doc|task)-(?:0|[1-9]\d*)\.jsonl$/.test(name);
const isReclaimFileName = (name: string): boolean => /^(?:doc|task)-(?:0|[1-9]\d*)\.jsonl\.reclaim$/.test(name);
const isCurrentOnly = (record: DocumentCreate): boolean =>
record.scope.kind !== "conversation" || record.history === "latest";
const sidecarKey = (file: string, seq: Seq, ordinal: number): string => JSON.stringify([file, seq, ordinal]);
const jsonLine = (value: unknown): string => {
@@ -233,21 +239,16 @@ export class JsonlStorage implements Storage {
private readonly directory: string;
private readonly mainPath: string;
private readonly fsync: boolean;
private readonly memory: MemoryStorage;
private readonly memory = new MemoryStorage();
private readonly currentOnlyDocuments = new Set<Id>();
private readonly liveTaskSidecars = new Set<Id>();
private closed = false;
private poisonError: JsonlStoragePoisonedError | undefined;
private constructor(
fs: FileSystem,
directory: string,
mainPath: string,
memory: MemoryStorage,
options: JsonlStorageOptions,
) {
private constructor(fs: FileSystem, directory: string, mainPath: string, options: JsonlStorageOptions) {
this.fs = fs;
this.directory = directory;
this.mainPath = mainPath;
this.memory = memory;
this.fsync = options.fsync ?? false;
}
@@ -264,15 +265,16 @@ export class JsonlStorage implements Storage {
if (!created.ok) throw errorFromFile("directory creation", created.error);
const mainPathResult = await fs.joinPath([absolute.value, MAIN_FILE], context);
if (!mainPathResult.ok) throw errorFromFile("path join", mainPathResult.error);
const memory = new MemoryStorage();
await JsonlStorage.recover(fs, absolute.value, mainPathResult.value, memory, context);
return new JsonlStorage(fs, absolute.value, mainPathResult.value, memory, options);
const storage = new JsonlStorage(fs, absolute.value, mainPathResult.value, options);
await storage.recover(context);
return storage;
}
async commit(writes: readonly StorageWrite[], context: Context): Promise<Seq> {
this.assertUsable();
const prepared = this.memory.prepareCommit(writes);
const encoded = this.encodeCommit(prepared.seq, prepared.writes);
const reclamations = this.planReclamations(prepared.writes, encoded);
const sidecars = await Promise.all(
[...encoded.sidecars].map(async ([file, content]) => ({
file,
@@ -293,7 +295,10 @@ export class JsonlStorage implements Storage {
}
const marker = await this.fs.appendFile(this.mainPath, encoded.marker, context);
if (!marker.ok) throw this.poison(errorFromFile(`append to ${MAIN_FILE}`, marker.error));
return prepared.apply();
const seq = prepared.apply();
this.adoptSidecarState(prepared.writes);
await this.reclaimSidecars(reclamations, context);
return seq;
}
async mintId() {
@@ -421,14 +426,95 @@ export class JsonlStorage implements Storage {
return { marker: jsonLine(marker), sidecars };
}
private static async recover(
fs: FileSystem,
directory: string,
mainPath: string,
memory: MemoryStorage,
context: Context,
): Promise<void> {
const main = await JsonlStorage.readLines(fs, mainPath, MAIN_FILE, context, (text, line) =>
private planReclamations(writes: readonly StorageWrite[], encoded: EncodedCommit): ReadonlyMap<string, string> {
const createdCurrentOnlyDocuments = new Set<Id>();
const retiredDocuments = new Set<Id>();
const baseDocuments = new Set<Id>();
const finalTasks = new Map<Id, StoredTask>();
for (const write of writes) {
switch (write.type) {
case "document.create":
if (isCurrentOnly(write.record)) createdCurrentOnlyDocuments.add(write.record.id);
break;
case "document.change":
if (write.content.kind === "base") baseDocuments.add(write.id);
break;
case "document.retire":
retiredDocuments.add(write.id);
break;
case "task":
finalTasks.set(write.value.id, write.value);
break;
}
}
const replacements = new Map<string, string>();
const isCurrentOnlyDocument = (id: Id): boolean =>
this.currentOnlyDocuments.has(id) || createdCurrentOnlyDocuments.has(id);
for (const id of retiredDocuments) {
if (isCurrentOnlyDocument(id)) replacements.set(sidecarFileName("doc", id), "");
}
for (const id of baseDocuments) {
if (!isCurrentOnlyDocument(id) || retiredDocuments.has(id)) continue;
const file = sidecarFileName("doc", id);
const content = encoded.sidecars.get(file);
if (content !== undefined) replacements.set(file, content);
}
for (const [id, task] of finalTasks) {
if (
task.state.status === "terminal" &&
(this.liveTaskSidecars.has(id) || encoded.sidecars.has(sidecarFileName("task", id)))
) {
replacements.set(sidecarFileName("task", id), "");
}
}
return replacements;
}
private adoptSidecarState(writes: readonly StorageWrite[]): void {
for (const write of writes) {
if (write.type === "document.create") {
if (isCurrentOnly(write.record)) this.currentOnlyDocuments.add(write.record.id);
} else if (write.type === "task") {
if (write.value.state.status === "terminal") this.liveTaskSidecars.delete(write.value.id);
else this.liveTaskSidecars.add(write.value.id);
}
}
}
/** The marker already published this state, so reclamation is retryable best-effort maintenance. */
private async reclaimSidecars(replacements: ReadonlyMap<string, string>, context: Context): Promise<void> {
if (replacements.size === 0) return;
if (this.fsync) {
const flushed = await this.fs.flushFile(this.mainPath, context);
if (!flushed.ok) return;
}
for (const [file, content] of replacements) await this.replaceSidecar(file, content, context);
}
private async replaceSidecar(file: string, content: string, context: Context): Promise<void> {
const path = await this.fs.joinPath([this.directory, file], context);
if (!path.ok) return;
if (content === "") {
await this.fs.remove(path.value, { force: true }, context);
return;
}
const temporaryPath = await this.fs.joinPath([this.directory, `${file}${RECLAIM_SUFFIX}`], context);
if (!temporaryPath.ok) return;
const written = await this.fs.writeFile(temporaryPath.value, content, context);
if (!written.ok) return;
if (this.fsync) {
const flushed = await this.fs.flushFile(temporaryPath.value, context);
if (!flushed.ok) return;
}
await this.fs.renameFile(temporaryPath.value, path.value, context);
}
private async recover(context: Context): Promise<void> {
const fs = this.fs;
const directory = this.directory;
const memory = this.memory;
const main = await JsonlStorage.readLines(fs, this.mainPath, MAIN_FILE, context, (text, line) =>
parseMainMarker(text, line),
);
let previousSeq = 0;
@@ -441,6 +527,11 @@ export class JsonlStorage implements Storage {
const listed = await fs.listDir(directory, context);
if (!listed.ok) throw errorFromFile("directory listing", listed.error);
for (const info of listed.value) {
if (info.kind === "file" && isReclaimFileName(info.name)) {
await fs.remove(info.path, { force: true }, context);
}
}
const sidecarFiles = listed.value
.filter((info) => info.kind === "file" && isSidecarFileName(info.name))
.map((info) => info.name)
@@ -468,6 +559,56 @@ export class JsonlStorage implements Storage {
}
}
const currentOnlyDocuments = new Set<Id>();
const retiredDocuments = new Set<Id>();
const finalTaskIsLive = new Map<Id, boolean>();
for (const { value: marker } of main.lines) {
for (const operation of marker.writes) {
if (operation.type === "document.create") {
if (isCurrentOnly(operation.record)) currentOnlyDocuments.add(operation.record.id);
} else if (operation.type === "document.retire") {
retiredDocuments.add(operation.id);
} else if (operation.type === "task") {
finalTaskIsLive.set(operation.value.id, false);
} else if (operation.type === "task.sidecar") {
finalTaskIsLive.set(operation.id, true);
}
}
}
const retiredCurrentOnlyDocuments = new Set([...retiredDocuments].filter((id) => currentOnlyDocuments.has(id)));
const latestBases = new Map<Id, SidecarRecord>();
for (const { value: marker } of main.lines) {
for (const operation of marker.writes) {
if (operation.type !== "document.create" && operation.type !== "document.change") continue;
const id = operation.type === "document.create" ? operation.record.id : operation.id;
if (!currentOnlyDocuments.has(id)) continue;
const record = recordByKey.get(
sidecarKey(sidecarFileName("doc", id), marker.seq, operation.ordinal),
)?.value;
if (
record?.payload.type !== "document" ||
record.payload.id !== id ||
record.payload.content.kind !== "base"
) {
continue;
}
const previous = latestBases.get(id);
if (
previous === undefined ||
record.seq > previous.seq ||
(record.seq === previous.seq && record.ordinal > previous.ordinal)
) {
latestBases.set(id, record);
}
}
}
const isBeforeLatestBase = (id: Id, seq: Seq, ordinal: number): boolean => {
const base = latestBases.get(id);
return base !== undefined && (seq < base.seq || (seq === base.seq && ordinal < base.ordinal));
};
const terminalTasks = new Set([...finalTaskIsLive].filter(([, live]) => !live).map(([id]) => id));
const confirmed = new Set<string>();
for (const line of main.lines) {
const marker = line.value;
@@ -482,52 +623,61 @@ export class JsonlStorage implements Storage {
writes.push(operation);
break;
case "task.sidecar": {
const optional = terminalTasks.has(operation.id);
const record = JsonlStorage.confirmRecord(
marker,
operation.ordinal,
sidecarFileName("task", operation.id),
recordByKey,
confirmed,
optional,
);
if (record.payload.type !== "task" || record.payload.value.id !== operation.id) {
throw new JsonlCorruptionError(`Confirmed task sidecar data does not match commit ${marker.seq}`);
if (record !== undefined) {
if (record.payload.type !== "task" || record.payload.value.id !== operation.id) {
throw new JsonlCorruptionError(
`Confirmed task sidecar data does not match commit ${marker.seq}`,
);
}
if (!optional) writes.push({ type: "task", value: record.payload.value });
}
writes.push({ type: "task", value: record.payload.value });
break;
}
case "document.create": {
const record = JsonlStorage.confirmRecord(
marker,
operation.ordinal,
sidecarFileName("doc", operation.record.id),
recordByKey,
confirmed,
);
if (record.payload.type !== "document" || record.payload.id !== operation.record.id) {
throw new JsonlCorruptionError(
`Confirmed document sidecar data does not match commit ${marker.seq}`,
);
}
if (record.payload.content.kind !== "base") {
throw new JsonlCorruptionError(`Document creation lacks a confirmed base in commit ${marker.seq}`);
}
writes.push({ type: "document.create", record: operation.record, content: record.payload.content });
break;
}
case "document.create":
case "document.change": {
const id = operation.type === "document.create" ? operation.record.id : operation.id;
const reclaimed =
retiredCurrentOnlyDocuments.has(id) || isBeforeLatestBase(id, marker.seq, operation.ordinal);
const record = JsonlStorage.confirmRecord(
marker,
operation.ordinal,
sidecarFileName("doc", operation.id),
sidecarFileName("doc", id),
recordByKey,
confirmed,
reclaimed,
);
if (record.payload.type !== "document" || record.payload.id !== operation.id) {
throw new JsonlCorruptionError(
`Confirmed document sidecar data does not match commit ${marker.seq}`,
);
let content: DocumentContent | undefined;
if (record !== undefined) {
if (record.payload.type !== "document" || record.payload.id !== id) {
throw new JsonlCorruptionError(
`Confirmed document sidecar data does not match commit ${marker.seq}`,
);
}
content = record.payload.content;
}
if (operation.type === "document.create") {
if (content !== undefined && content.kind !== "base") {
throw new JsonlCorruptionError(
`Document creation lacks a confirmed base in commit ${marker.seq}`,
);
}
writes.push({
type: "document.create",
record: operation.record,
content: reclaimed || content === undefined ? { kind: "base", version: 1, value: {} } : content,
});
} else if (!reclaimed && content !== undefined) {
writes.push({ type: "document.change", id, content });
}
writes.push({ type: "document.change", id: operation.id, content: record.payload.content });
break;
}
}
@@ -542,6 +692,7 @@ export class JsonlStorage implements Storage {
}
}
const reclamations = new Map<string, string>();
for (const [file, parsed] of parsedFiles) {
let unconfirmedAt: number | undefined;
for (const line of parsed.lines) {
@@ -558,6 +709,33 @@ export class JsonlStorage implements Storage {
const truncated = await fs.truncateFile(parsed.path, unconfirmedAt, context);
if (!truncated.ok) throw errorFromFile(`tail truncation of ${file}`, truncated.error);
}
const id = Number(file.slice(file.indexOf("-") + 1, -".jsonl".length));
const confirmedLines = parsed.lines.filter((line) =>
confirmed.has(sidecarKey(file, line.value.seq, line.value.ordinal)),
);
let retainedLines: readonly ParsedLine<SidecarRecord>[] | undefined;
if (file.startsWith("task-") && terminalTasks.has(id)) {
retainedLines = [];
} else if (file.startsWith("doc-") && retiredCurrentOnlyDocuments.has(id)) {
retainedLines = [];
} else if (file.startsWith("doc-") && latestBases.has(id)) {
retainedLines = confirmedLines.filter(
(line) => !isBeforeLatestBase(id, line.value.seq, line.value.ordinal),
);
}
if (
retainedLines !== undefined &&
(retainedLines.length < confirmedLines.length || retainedLines.length === 0)
) {
reclamations.set(file, retainedLines.map((line) => jsonLine(line.value)).join(""));
}
}
await this.reclaimSidecars(reclamations, context);
for (const id of currentOnlyDocuments) this.currentOnlyDocuments.add(id);
for (const [id, live] of finalTaskIsLive) {
if (live) this.liveTaskSidecars.add(id);
}
}
@@ -567,11 +745,13 @@ export class JsonlStorage implements Storage {
file: string,
recordByKey: ReadonlyMap<string, ParsedLine<SidecarRecord>>,
confirmed: Set<string>,
): SidecarRecord {
optional: boolean,
): SidecarRecord | undefined {
const key = sidecarKey(file, marker.seq, ordinal);
if (confirmed.has(key)) throw new JsonlCorruptionError(`Sidecar record is confirmed more than once`);
const line = recordByKey.get(key);
if (line === undefined) {
if (optional) return undefined;
throw new JsonlCorruptionError(`Missing confirmed sidecar record ${file} at sequence ${marker.seq}`);
}
confirmed.add(key);
+487 -11
View File
@@ -1,4 +1,4 @@
import { mkdtemp, readFile, rm, stat, writeFile } from "node:fs/promises";
import { mkdtemp, readdir, readFile, rm, stat, writeFile } from "node:fs/promises";
import { tmpdir } from "node:os";
import { basename, join } from "node:path";
import type { Context, JsonValue } from "@earendil-works/chord";
@@ -139,7 +139,7 @@ function sessionDocument(id: number, kind = "test.document"): DocumentCreate {
}
type Failure = {
readonly operation: "append" | "flush";
readonly operation: "append" | "flush" | "write" | "rename" | "remove";
readonly call: number;
readonly mode: "before" | "after" | "short";
};
@@ -149,18 +149,26 @@ class InstrumentedEnv extends NodeExecutionEnv {
private failure: Failure | undefined;
private appendCalls = 0;
private flushCalls = 0;
private writeCalls = 0;
private renameCalls = 0;
private removeCalls = 0;
fail(failure: Failure): void {
this.failure = failure;
this.appendCalls = 0;
this.flushCalls = 0;
this.operations.length = 0;
this.resetObservations();
}
clear(): void {
this.failure = undefined;
this.resetObservations();
}
private resetObservations(): void {
this.appendCalls = 0;
this.flushCalls = 0;
this.writeCalls = 0;
this.renameCalls = 0;
this.removeCalls = 0;
this.operations.length = 0;
}
@@ -201,6 +209,66 @@ class InstrumentedEnv extends NodeExecutionEnv {
}
return err(new FileError("unknown", "injected flush failure", path));
}
override async writeFile(
path: string,
content: string | Uint8Array,
writeContext: Context,
): Promise<Result<void, FileError>> {
this.writeCalls++;
this.operations.push(`write:${basename(path)}`);
const failure = this.failure;
if (failure?.operation !== "write" || failure.call !== this.writeCalls) {
return super.writeFile(path, content, writeContext);
}
if (failure.mode === "before") return err(new FileError("unknown", "injected write failure", path));
if (failure.mode === "short") {
const bytes = typeof content === "string" ? new TextEncoder().encode(content) : content;
const partial = bytes.subarray(0, Math.max(1, Math.floor(bytes.length / 2)));
const written = await super.writeFile(path, partial, writeContext);
if (!written.ok) return written;
return err(new FileError("unknown", "injected short write", path));
}
const written = await super.writeFile(path, content, writeContext);
if (!written.ok) return written;
return err(new FileError("unknown", "injected post-write failure", path));
}
override async renameFile(
sourcePath: string,
destinationPath: string,
renameContext: Context,
): Promise<Result<void, FileError>> {
this.renameCalls++;
this.operations.push(`rename:${basename(sourcePath)}->${basename(destinationPath)}`);
const failure = this.failure;
if (failure?.operation !== "rename" || failure.call !== this.renameCalls) {
return super.renameFile(sourcePath, destinationPath, renameContext);
}
if (failure.mode === "after") {
const renamed = await super.renameFile(sourcePath, destinationPath, renameContext);
if (!renamed.ok) return renamed;
}
return err(new FileError("unknown", "injected rename failure", sourcePath));
}
override async remove(
path: string,
options: { recursive?: boolean; force?: boolean } | undefined,
removeContext: Context,
): Promise<Result<void, FileError>> {
this.removeCalls++;
this.operations.push(`remove:${basename(path)}`);
const failure = this.failure;
if (failure?.operation !== "remove" || failure.call !== this.removeCalls) {
return super.remove(path, options, removeContext);
}
if (failure.mode === "after") {
const removed = await super.remove(path, options, removeContext);
if (!removed.ok) return removed;
}
return err(new FileError("unknown", "injected remove failure", path));
}
}
async function readLines(path: string): Promise<string[]> {
@@ -208,6 +276,16 @@ async function readLines(path: string): Promise<string[]> {
return text === "" ? [] : text.trimEnd().split("\n");
}
async function fileExists(path: string): Promise<boolean> {
try {
await stat(path);
return true;
} catch (error) {
if ((error as NodeJS.ErrnoException).code === "ENOENT") return false;
throw error;
}
}
describe("Pico JsonlStorage publication and recovery", () => {
it("opens through the Node adapter", async () => {
const storage = await openNodeJsonlStorage(await tempDirectory(), context);
@@ -216,7 +294,7 @@ describe("Pico JsonlStorage publication and recovery", () => {
expect(await storage.conversation(ROOT_CONVERSATION_ID, context)).toEqual({ id: ROOT_CONVERSATION_ID });
});
it("persists live tasks and document revisions in sidecars with one marker for every commit", async () => {
it("publishes every commit before reclaiming current-only document and terminal-task sidecars", async () => {
const directory = await tempDirectory();
const storage = await openStorage(directory);
await createRoot(storage);
@@ -247,12 +325,15 @@ describe("Pico JsonlStorage publication and recovery", () => {
],
context,
);
expect(await readLines(join(directory, `task-${taskId}.jsonl`))).toHaveLength(1);
expect(await readLines(join(directory, `doc-${documentId}.jsonl`))).toHaveLength(1);
await storage.commit([{ type: "document.retire", id: documentId }], context);
await storage.commit([{ type: "task", value: terminalTask(taskId) }], context);
expect(await readLines(join(directory, "main.jsonl"))).toHaveLength(7);
expect(await readLines(join(directory, `task-${taskId}.jsonl`))).toHaveLength(1);
expect(await readLines(join(directory, `doc-${documentId}.jsonl`))).toHaveLength(3);
expect(await fileExists(join(directory, `task-${taskId}.jsonl`))).toBe(false);
expect(await fileExists(join(directory, `doc-${documentId}.jsonl`))).toBe(false);
const markerTypes = (await readLines(join(directory, "main.jsonl"))).map(
(line) => (JSON.parse(line) as { readonly type: string }).type,
);
@@ -420,7 +501,7 @@ describe("Pico JsonlStorage publication and recovery", () => {
});
}
it("orders appends and optional flushes exactly and never flushes main.jsonl", async () => {
it("orders publication flushes exactly and flushes main only to authorize reclamation", async () => {
for (const fsync of [false, true]) {
const directory = await tempDirectory();
const env = new InstrumentedEnv({ cwd: directory });
@@ -452,6 +533,27 @@ describe("Pico JsonlStorage publication and recovery", () => {
];
expect(env.operations).toEqual(expected);
env.clear();
await storage.commit(
[
{
type: "document.change",
id: firstId,
content: { kind: "base", version: 1, value: { checkpoint: true } },
},
],
context,
);
expect(env.operations).toEqual([
`append:doc-${firstId}.jsonl`,
...(fsync ? [`flush:doc-${firstId}.jsonl`] : []),
"append:main.jsonl",
...(fsync ? ["flush:main.jsonl"] : []),
`write:doc-${firstId}.jsonl.reclaim`,
...(fsync ? [`flush:doc-${firstId}.jsonl.reclaim`] : []),
`rename:doc-${firstId}.jsonl.reclaim->doc-${firstId}.jsonl`,
]);
env.clear();
await storage.commit(
[
@@ -467,9 +569,383 @@ describe("Pico JsonlStorage publication and recovery", () => {
context,
);
expect(env.operations).toEqual(["append:main.jsonl"]);
const taskId = await storage.mintId();
await storage.commit([{ type: "task", value: pendingTask(taskId) }], context);
env.clear();
await storage.commit([{ type: "task", value: terminalTask(taskId) }], context);
expect(env.operations).toEqual([
"append:main.jsonl",
...(fsync ? ["flush:main.jsonl"] : []),
`remove:task-${taskId}.jsonl`,
]);
}
});
for (const failure of [
{ operation: "write", call: 1, mode: "before" },
{ operation: "write", call: 1, mode: "short" },
{ operation: "write", call: 1, mode: "after" },
{ operation: "rename", call: 1, mode: "before" },
{ operation: "rename", call: 1, mode: "after" },
] as const satisfies readonly Failure[]) {
it(`recovers a committed base across reclaim ${failure.operation} ${failure.mode}`, async () => {
const directory = await tempDirectory();
const env = new InstrumentedEnv({ cwd: directory });
const storage = await openStorage(directory, env);
await createRoot(storage);
const id = await storage.mintId();
await storage.commit(
[
{
type: "document.create",
record: sessionDocument(id),
content: { kind: "base", version: 1, value: { count: 0 } },
},
],
context,
);
await storage.commit(
[
{
type: "document.change",
id,
content: { kind: "delta", version: 1, ops: [["s", ["count"], 1]] },
},
],
context,
);
env.fail(failure);
await expect(
storage.commit(
[
{
type: "document.change",
id,
content: { kind: "base", version: 1, value: { count: 2 } },
},
],
context,
),
).resolves.toBe(4);
expect((await storage.document(id, "current", context))?.value).toEqual({ count: 2 });
await storage.close(context);
const reopened = await openStorage(directory);
expect((await reopened.document(id, "current", context))?.value).toEqual({ count: 2 });
expect(await readLines(join(directory, `doc-${id}.jsonl`))).toHaveLength(1);
expect((await readdir(directory)).filter((name) => name.endsWith(".reclaim"))).toEqual([]);
});
}
for (const mode of ["before", "after"] as const) {
it(`recovers document-retirement reclamation across remove ${mode}`, async () => {
const directory = await tempDirectory();
const env = new InstrumentedEnv({ cwd: directory });
const storage = await openStorage(directory, env);
await createRoot(storage);
const taskId = await storage.mintId();
const id = await storage.mintId();
await storage.commit(
[
{ type: "task", value: pendingTask(taskId) },
{
type: "document.create",
record: { id, kind: "task.document", scope: { kind: "task", taskId } },
content: { kind: "base", version: 1, value: { count: 1 } },
},
],
context,
);
env.fail({ operation: "remove", call: 1, mode });
await expect(storage.commit([{ type: "document.retire", id }], context)).resolves.toBe(3);
expect(await storage.document(id, "current", context)).toBeUndefined();
await storage.close(context);
const reopened = await openStorage(directory);
expect(await reopened.document(id, "current", context)).toBeUndefined();
expect(await reopened.task(taskId, context)).toEqual(pendingTask(taskId));
expect(await fileExists(join(directory, `doc-${id}.jsonl`))).toBe(false);
expect((await readdir(directory)).filter((name) => name.endsWith(".reclaim"))).toEqual([]);
});
}
for (const mode of ["before", "after"] as const) {
it(`defers reclamation after authorizing-main flush ${mode} failure`, async () => {
const directory = await tempDirectory();
const env = new InstrumentedEnv({ cwd: directory });
const storage = await openStorage(directory, env, { fsync: true });
await createRoot(storage);
const id = await storage.mintId();
await storage.commit(
[
{
type: "document.create",
record: sessionDocument(id),
content: { kind: "base", version: 1, value: { count: 0 } },
},
],
context,
);
env.fail({ operation: "flush", call: 2, mode });
await expect(
storage.commit(
[
{
type: "document.change",
id,
content: { kind: "base", version: 1, value: { count: 2 } },
},
],
context,
),
).resolves.toBe(3);
expect(env.operations).toEqual([
`append:doc-${id}.jsonl`,
`flush:doc-${id}.jsonl`,
"append:main.jsonl",
"flush:main.jsonl",
]);
expect((await storage.document(id, "current", context))?.value).toEqual({ count: 2 });
expect(await readLines(join(directory, `doc-${id}.jsonl`))).toHaveLength(2);
await storage.close(context);
const recoveryEnv = new InstrumentedEnv({ cwd: directory });
recoveryEnv.fail({ operation: "flush", call: 1, mode });
const deferred = await openStorage(directory, recoveryEnv, { fsync: true });
expect((await deferred.document(id, "current", context))?.value).toEqual({ count: 2 });
expect(recoveryEnv.operations).toEqual(["flush:main.jsonl"]);
expect(await readLines(join(directory, `doc-${id}.jsonl`))).toHaveLength(2);
await deferred.close(context);
const reclaimed = await openStorage(directory, new NodeExecutionEnv({ cwd: directory }), { fsync: true });
expect((await reclaimed.document(id, "current", context))?.value).toEqual({ count: 2 });
expect(await readLines(join(directory, `doc-${id}.jsonl`))).toHaveLength(1);
});
}
for (const mode of ["before", "after"] as const) {
it(`keeps a committed base usable after reclaim-temp flush ${mode} failure`, async () => {
const directory = await tempDirectory();
const env = new InstrumentedEnv({ cwd: directory });
const storage = await openStorage(directory, env, { fsync: true });
await createRoot(storage);
const id = await storage.mintId();
await storage.commit(
[
{
type: "document.create",
record: sessionDocument(id),
content: { kind: "base", version: 1, value: { count: 0 } },
},
],
context,
);
env.fail({ operation: "flush", call: 3, mode });
await expect(
storage.commit(
[
{
type: "document.change",
id,
content: { kind: "base", version: 1, value: { count: 2 } },
},
],
context,
),
).resolves.toBe(3);
env.clear();
await storage.commit(
[
{
type: "document.change",
id,
content: { kind: "delta", version: 1, ops: [["s", ["count"], 3]] },
},
],
context,
);
await storage.close(context);
const reopened = await openStorage(directory, new NodeExecutionEnv({ cwd: directory }), { fsync: true });
expect((await reopened.document(id, "current", context))?.value).toEqual({ count: 3 });
expect(await readLines(join(directory, `doc-${id}.jsonl`))).toHaveLength(2);
});
}
for (const mode of ["before", "after"] as const) {
it(`recovers terminal-task reclamation across remove ${mode}`, async () => {
const directory = await tempDirectory();
const env = new InstrumentedEnv({ cwd: directory });
const storage = await openStorage(directory, env);
await createRoot(storage);
const id = await storage.mintId();
await storage.commit([{ type: "task", value: pendingTask(id) }], context);
env.fail({ operation: "remove", call: 1, mode });
await expect(storage.commit([{ type: "task", value: terminalTask(id) }], context)).resolves.toBe(3);
expect(await storage.task(id, context)).toEqual(terminalTask(id));
await storage.close(context);
const reopened = await openStorage(directory);
expect(await reopened.task(id, context)).toEqual(terminalTask(id));
expect(await fileExists(join(directory, `task-${id}.jsonl`))).toBe(false);
expect((await readdir(directory)).filter((name) => name.endsWith(".reclaim"))).toEqual([]);
});
}
it("appends later deltas to the replacement sidecar after a current-only base", async () => {
const directory = await tempDirectory();
const storage = await openStorage(directory);
await createRoot(storage);
const id = await storage.mintId();
await storage.commit(
[
{
type: "document.create",
record: sessionDocument(id),
content: { kind: "base", version: 1, value: { count: 0 } },
},
],
context,
);
await storage.commit(
[
{
type: "document.change",
id,
content: { kind: "base", version: 1, value: { count: 10 } },
},
],
context,
);
await storage.commit(
[
{
type: "document.change",
id,
content: { kind: "delta", version: 1, ops: [["s", ["count"], 11]] },
},
],
context,
);
expect(await readLines(join(directory, `doc-${id}.jsonl`))).toHaveLength(2);
const reopened = await openStorage(directory);
expect((await reopened.document(id, "current", context))?.value).toEqual({ count: 11 });
});
it("never reclaims rewindable document history, including after a base and retirement", async () => {
const directory = await tempDirectory();
const storage = await openStorage(directory);
await createRoot(storage);
const id = await storage.mintId();
const record = {
id,
kind: "rewindable",
scope: { kind: "conversation" as const, conversationId: ROOT_CONVERSATION_ID },
history: "rewindable" as const,
fork: "asOf" as const,
};
const createdAt = await storage.commit(
[{ type: "document.create", record, content: { kind: "base", version: 1, value: { count: 0 } } }],
context,
);
const changedAt = await storage.commit(
[
{
type: "document.change",
id,
content: { kind: "delta", version: 1, ops: [["s", ["count"], 1]] },
},
],
context,
);
await storage.commit(
[
{
type: "document.change",
id,
content: { kind: "base", version: 1, value: { count: 2 } },
},
],
context,
);
await storage.commit([{ type: "document.retire", id }], context);
expect(await readLines(join(directory, `doc-${id}.jsonl`))).toHaveLength(3);
const reopened = await openStorage(directory);
expect((await reopened.document(id, createdAt, context))?.value).toEqual({ count: 0 });
expect((await reopened.document(id, changedAt, context))?.value).toEqual({ count: 1 });
expect(await reopened.document(id, "current", context)).toBeUndefined();
expect(await readLines(join(directory, `doc-${id}.jsonl`))).toHaveLength(3);
});
it("reclaims retired task-, session-, and latest-conversation document sidecars", async () => {
const directory = await tempDirectory();
const storage = await openStorage(directory);
await createRoot(storage);
const taskId = await storage.mintId();
const sessionId = await storage.mintId();
const latestId = await storage.mintId();
const taskDocumentId = await storage.mintId();
const createdAt = await storage.commit(
[
{ type: "task", value: pendingTask(taskId) },
{
type: "document.create",
record: sessionDocument(sessionId, "session"),
content: { kind: "base", version: 1, value: {} },
},
{
type: "document.create",
record: {
id: latestId,
kind: "latest",
scope: { kind: "conversation", conversationId: ROOT_CONVERSATION_ID },
history: "latest",
fork: "current",
},
content: { kind: "base", version: 1, value: {} },
},
{
type: "document.create",
record: { id: taskDocumentId, kind: "task", scope: { kind: "task", taskId } },
content: { kind: "base", version: 1, value: {} },
},
],
context,
);
const retiredAt = await storage.commit(
[
{ type: "document.retire", id: sessionId },
{ type: "document.retire", id: latestId },
{ type: "document.retire", id: taskDocumentId },
{ type: "task", value: terminalTask(taskId) },
],
context,
);
for (const file of [
`doc-${sessionId}.jsonl`,
`doc-${latestId}.jsonl`,
`doc-${taskDocumentId}.jsonl`,
`task-${taskId}.jsonl`,
]) {
expect(await fileExists(join(directory, file))).toBe(false);
}
const reopened = await openStorage(directory);
expect(await reopened.task(taskId, context)).toEqual(terminalTask(taskId));
expect(await reopened.document(sessionId, "current", context)).toBeUndefined();
expect(await reopened.document(latestId, "current", context)).toBeUndefined();
expect(await reopened.document(taskDocumentId, "current", context)).toBeUndefined();
expect(
await reopened.findDocument({ kind: "session", scope: { kind: "session" } }, createdAt, context),
).toMatchObject({ id: sessionId, createdAt, retiredAt });
expect(
await reopened.findDocument({ kind: "session", scope: { kind: "session" } }, retiredAt, context),
).toBeUndefined();
});
it("truncates torn UTF-8 tails at exact byte offsets and reuses the unconfirmed sequence", async () => {
const directory = await tempDirectory();
const storage = await openStorage(directory);
@@ -512,7 +988,7 @@ describe("Pico JsonlStorage publication and recovery", () => {
await storage.commit([{ type: "task", value: pendingTask(taskId) }], context);
await storage.commit([{ type: "task", value: terminalTask(taskId) }], context);
const sidecarPath = join(directory, `task-${taskId}.jsonl`);
const confirmedSize = (await stat(sidecarPath)).size;
expect(await fileExists(sidecarPath)).toBe(false);
await writeFile(
sidecarPath,
`${JSON.stringify({
@@ -526,7 +1002,7 @@ describe("Pico JsonlStorage publication and recovery", () => {
);
const reopened = await openStorage(directory);
expect((await stat(sidecarPath)).size).toBe(confirmedSize);
expect(await fileExists(sidecarPath)).toBe(false);
expect(await reopened.task(taskId, context)).toEqual(terminalTask(taskId));
expect(await reopened.commit([], context)).toBe(4);
});