feat(durable): optimize document replay

This commit is contained in:
Mario Zechner
2026-09-25 00:09:50 +02:00
parent 9e70c3d505
commit 5fd446ca18
4 changed files with 191 additions and 11 deletions
+20 -9
View File
@@ -1,5 +1,5 @@
import type { Context, JsonValue } from "@earendil-works/chord";
import { applyImmutable } from "@earendil-works/chord/delta";
import { applyImmutableBatches, type Op } from "@earendil-works/chord/delta";
import type {
ConversationRecord,
Cursor,
@@ -37,6 +37,21 @@ type DocumentAction = {
retire: boolean;
};
function* documentDeltaBatches(
id: Id,
version: number,
revisions: readonly DocumentRevision[],
start: number,
): Generator<readonly Op[]> {
for (let index = start; index < revisions.length; index++) {
const revision = revisions[index]!;
if (revision.kind !== "delta" || revision.version !== version) {
throw new Error(`Document ${id} crosses a stored version boundary without a base`);
}
yield revision.ops;
}
}
type DocumentAddressIndex = {
ids: Id[];
currentId?: Id;
@@ -472,14 +487,10 @@ export class MemoryStorage implements Storage {
while (baseIndex >= 0 && revisions[baseIndex]!.kind !== "base") baseIndex--;
const base = revisions[baseIndex];
if (base?.kind !== "base") throw new Error(`Document ${id} is missing a required base`);
let value = base.value;
for (let index = baseIndex + 1; index < revisions.length; index++) {
const revision = revisions[index]!;
if (revision.kind !== "delta" || revision.version !== base.version) {
throw new Error(`Document ${id} crosses a stored version boundary without a base`);
}
value = applyImmutable(value, revision.ops) as JsonObject;
}
const value = applyImmutableBatches(
base.value,
documentDeltaBatches(id, base.version, revisions, baseIndex + 1),
) as JsonObject;
return { record: clone(stored.record), version: base.version, value: clone(value) };
}
@@ -1,5 +1,5 @@
import type { Context, JsonValue } from "@earendil-works/chord";
import { applyImmutable, type Op } from "@earendil-works/chord/delta";
import { apply, type Op } from "@earendil-works/chord/delta";
import type {
ConversationRecord,
Cursor,
@@ -449,7 +449,7 @@ export class SqliteStorage implements Storage {
if (revision.kind !== "delta" || revision.version !== base.version) {
throw new Error(`Document ${id} crosses a stored version boundary without a base`);
}
value = applyImmutable(value, parseJson<readonly Op[]>(revision.content)) as JsonObject;
value = apply(value, parseJson<readonly Op[]>(revision.content)) as JsonObject;
}
return { record, version: base.version, value };
}
@@ -697,6 +697,99 @@ export function createStorageConformance(options: StorageConformanceOptions): re
expect((await storage.document(secondId, "current", context))?.value).toEqual({ items: ["new"] });
}),
createCase(options, "streams long document tails across root replacement deltas", async (storage) => {
const rootId = await createRoot(storage);
const id = await storage.mintId();
const record = {
id,
kind: "conversation.long-tail",
scope: { kind: "conversation", conversationId: rootId },
history: "rewindable",
fork: "asOf",
} satisfies DocumentCreate;
const initial = {
revision: 0,
rows: Array.from({ length: 512 }, (_, value) => ({ value, stable: `row-${value}` })),
};
const createdAt = await storage.commit(
[{ type: "document.create", record, content: { kind: "base", version: 1, value: initial } }],
context,
);
const beforeReplacement = structuredClone(initial);
let beforeReplacementAt = createdAt;
for (let revision = 1; revision <= 24; revision++) {
const index = (revision * 17) % beforeReplacement.rows.length;
beforeReplacement.rows[index]!.value = -revision;
beforeReplacement.revision = revision;
beforeReplacementAt = await storage.commit(
[
{
type: "document.change",
id,
content: {
kind: "delta",
version: 1,
ops: [
["s", ["rows", index, "value"], -revision],
["s", ["revision"], revision],
],
},
},
],
context,
);
}
const replacement = {
revision: 100,
rows: Array.from({ length: 512 }, (_, value) => ({ value: 10_000 + value, stable: `new-${value}` })),
};
const replacementSnapshot = structuredClone(replacement);
const replacementAt = await storage.commit(
[
{
type: "document.change",
id,
content: { kind: "delta", version: 1, ops: [["r", replacement]] },
},
],
context,
);
replacement.rows[0]!.value = -999;
const current = structuredClone(replacementSnapshot);
for (let revision = 101; revision <= 124; revision++) {
const index = (revision * 19) % current.rows.length;
current.rows[index]!.value = -revision;
current.revision = revision;
await storage.commit(
[
{
type: "document.change",
id,
content: {
kind: "delta",
version: 1,
ops: [
["s", ["rows", index, "value"], -revision],
["s", ["revision"], revision],
],
},
},
],
context,
);
}
expect((await storage.document(id, createdAt, context))?.value).toEqual(initial);
expect((await storage.document(id, beforeReplacementAt, context))?.value).toEqual(beforeReplacement);
expect((await storage.document(id, replacementAt, context))?.value).toEqual(replacementSnapshot);
const read = (await storage.document(id, "current", context))!;
expect(read.value).toEqual(current);
(read.value.rows as Array<{ value: number }>)[0]!.value = -1_000;
expect((await storage.document(id, "current", context))?.value).toEqual(current);
}),
createCase(
options,
"uses bases for version transitions and rejects historical reads of current-only documents",
@@ -201,6 +201,82 @@ describe("Pico SqliteStorage", () => {
);
});
it("replays detached root replacements and follow-up edits while rejecting corrupt operations", async () => {
const { storage, path } = await createSqliteStorage();
await createRoot(storage);
const id = await storage.mintId();
await storage.commit(
[
{
type: "document.create",
record: { id, kind: "replay", scope: { kind: "session" } },
content: { kind: "base", version: 1, value: { nested: { value: 1 }, rows: [] } },
},
],
context,
);
await storage.commit(
[
{
type: "document.change",
id,
content: {
kind: "delta",
version: 1,
ops: [["r", { nested: { value: 2 }, rows: [{ id: 1 }] }]],
},
},
],
context,
);
await storage.commit(
[
{
type: "document.change",
id,
content: {
kind: "delta",
version: 1,
ops: [
["s", ["nested", "value"], 3],
["p", ["rows"], 1, 0, [{ id: 2 }]],
["m", ["rows"], [1, 0]],
],
},
},
],
context,
);
await storage.commit(
[
{
type: "document.change",
id,
content: { kind: "delta", version: 1, ops: [["s", ["nested", "value"], 4]] },
},
],
context,
);
const expected = { nested: { value: 4 }, rows: [{ id: 2 }, { id: 1 }] };
const first = (await storage.document(id, "current", context))!;
expect(first.value).toEqual(expected);
(first.value.nested as { value: number }).value = 99;
(first.value.rows as Array<{ id: number }>)[0]!.id = 99;
expect((await storage.document(id, "current", context))?.value).toEqual(expected);
const database = new DatabaseSync(path);
try {
database
.prepare(`UPDATE document_revisions SET content = ? WHERE document_id = ? AND seq =
(SELECT max(seq) FROM document_revisions WHERE document_id = ?)`)
.run('[["unknown"]]', id, id);
} finally {
database.close();
}
await expect(storage.document(id, "current", context)).rejects.toThrow("unknown op verb");
});
it("rolls SQL rows and sequence allocation back as one transaction", async () => {
const { storage, path } = await createSqliteStorage();
await createRoot(storage);