From 5d4de953ce48ebe1fbdd72e9819b8bccabaeb08d Mon Sep 17 00:00:00 2001 From: Mario Zechner Date: Thu, 24 Sep 2026 23:09:15 +0200 Subject: [PATCH] feat(durable): add checkpoints and document migrations --- packages/durable/CHANGELOG.md | 2 + packages/durable/docs/pico-v5.md | 7 +- packages/durable/src/documents.ts | 18 +- packages/durable/src/session/session.ts | 72 +- packages/durable/src/session/transaction.ts | 46 +- packages/durable/src/storage/jsonl/storage.ts | 8 +- packages/durable/src/storage/memory.ts | 17 +- .../durable/src/storage/sqlite/storage.ts | 29 +- .../src/testing/storage-conformance.ts | 14 + packages/durable/src/types.ts | 20 + packages/durable/test/jsonl-storage.test.ts | 8 +- .../session-checkpoints-migrations.test.ts | 684 ++++++++++++++++++ packages/durable/test/sqlite-storage.test.ts | 10 +- 13 files changed, 897 insertions(+), 38 deletions(-) create mode 100644 packages/durable/test/session-checkpoints-migrations.test.ts diff --git a/packages/durable/CHANGELOG.md b/packages/durable/CHANGELOG.md index 80e4ca78d..aeb7d3185 100644 --- a/packages/durable/CHANGELOG.md +++ b/packages/durable/CHANGELOG.md @@ -5,10 +5,12 @@ ### Breaking Changes - Reordered Storage scan arguments so the limit precedes the cursor. +- Added the required conversation-visible `Storage.entry(conversationId, id, context)` overload. ### Added - Added transactional Sessions with typed durable documents, task creation, snapshots, retirement, and commit publications. +- Added document checkpoint selection, lazy version migration, and `Session.snapshotAsOf()` for rewindable conversation documents. ## [0.87.1] - 2026-09-22 diff --git a/packages/durable/docs/pico-v5.md b/packages/durable/docs/pico-v5.md index 6796be810..83276768b 100644 --- a/packages/durable/docs/pico-v5.md +++ b/packages/durable/docs/pico-v5.md @@ -2043,6 +2043,7 @@ interface Storage { scanConversations(limit: number, cursor: Cursor | undefined, context: Context): Promise>; entry(id: Id, context: Context): Promise<{ readonly entry: EntryRecord; readonly commitSeq: Seq } | undefined>; + entry(conversationId: Id, id: Id, context: Context): Promise<{ readonly entry: EntryRecord; readonly commitSeq: Seq } | undefined>; findLatestHeadMarker(conversationId: Id, atOrBeforeEntryId: Id | undefined, context: Context): Promise<(EntryRecord & { readonly head: Id }) | undefined>; scanEntries(query: EntryQuery, limit: number, cursor: Cursor | undefined, context: Context): Promise>; @@ -2072,8 +2073,10 @@ inclusive ID range in newest-first order while applying every conversation ancestry cap. With no bounds it pages complete visible history. To read context through entry `E`, find the marker at or before `E`, then scan from `marker?.head` through `E`. For current context the upper bound is omitted. -`entry()` combines exact global lookup with the commit sequence required by -historical document reads. `limit` is always the maximum page size. +`entry(id)` combines exact global lookup with the commit sequence required by +historical document reads. `entry(conversationId, id)` returns that pair only +when the entry is visible through the requested conversation's ancestry. `limit` +is always the maximum page size. `findDocument()` resolves one exact logical kind/scope/key address at current or historical membership. A missing key means the singleton, not every family member. `scanDocuments()` enumerates only the incarnations alive in one exact diff --git a/packages/durable/src/documents.ts b/packages/durable/src/documents.ts index 6bb13689b..91c1dd252 100644 --- a/packages/durable/src/documents.ts +++ b/packages/durable/src/documents.ts @@ -1,4 +1,5 @@ -import type { JsonValue } from "@earendil-works/chord"; +import { copyJson, type JsonValue } from "@earendil-works/chord"; +import type { Op } from "@earendil-works/chord/delta"; import type { CommonDocDefinition, ConversationDocFamilyToken, @@ -19,6 +20,7 @@ import type { RewindableConversationSemantics, SessionDocFamilyToken, SessionDocToken, + StoredDocument, TaskDocFamilyToken, TaskDocToken, } from "./types.ts"; @@ -78,6 +80,8 @@ export type AnyDocDefinition = DocumentSemantics & { readonly version: number; readonly family?: true; initial(seed?: JsonValue): JsonObject; + migrate?(value: JsonObject, fromVersion: number): JsonObject; + checkpointWhen?(value: Readonly, ops: readonly Op[]): boolean; }; /** Erased singleton or family token. */ @@ -166,12 +170,20 @@ export function checkRecordScope(definition: AnyDocDefinition, record: DocumentR } } -/** Reject typed access to a stored version the token cannot use without migration. */ +/** Reject typed access to a stored version the supplied definition cannot use. */ export function checkRecordVersion(definition: AnyDocDefinition, record: DocumentRecord, version: number): void { if (version > definition.version) { throw new Error(`Document ${record.id} (${record.kind}) has newer version ${version} than ${definition.version}`); } - if (version < definition.version) { + if (version < definition.version && definition.migrate === undefined) { throw new Error(`Document ${record.id} (${record.kind}) requires migration from version ${version}`); } } + +/** Validate and materialize a detached stored value for typed access. */ +export function materializeDocument(definition: AnyDocDefinition, stored: StoredDocument): JsonObject { + checkRecordScope(definition, stored.record); + checkRecordVersion(definition, stored.record, stored.version); + if (stored.version === definition.version) return stored.value; + return copyJson(definition.migrate!(stored.value, stored.version)) as JsonObject; +} diff --git a/packages/durable/src/session/session.ts b/packages/durable/src/session/session.ts index 62215376f..f55186ee7 100644 --- a/packages/durable/src/session/session.ts +++ b/packages/durable/src/session/session.ts @@ -1,13 +1,21 @@ import type { Context, JsonValue } from "@earendil-works/chord"; import { awaitWithContext, withoutAbortSignal } from "@earendil-works/chord/context"; import { track } from "@earendil-works/chord/delta"; -import { type AnyDocToken, checkRecordScope, checkRecordVersion, resolveAddress } from "../documents.ts"; +import { + type AnyDocToken, + checkRecordScope, + checkRecordVersion, + materializeDocument, + resolveAddress, +} from "../documents.ts"; import type { ConversationDocFamilyToken, ConversationDocToken, DocumentAddress, Id, JsonObject, + RewindableConversationDocFamilyToken, + RewindableConversationDocToken, Session, SessionDocFamilyToken, SessionDocToken, @@ -45,7 +53,7 @@ export class SessionKernel implements Session { this.#host = { storage, cached: (id) => this.#documents.get(id), - load: (addressId, address, context) => this.#load(addressId, address, context), + load: (definition, addressId, address, context) => this.#load(definition, addressId, address, context), install: (document) => { this.#documents.set(document.addressId, document); }, @@ -101,14 +109,57 @@ export class SessionKernel implements Session { this.#documents.get(resolved.id) ?? (await this.#enqueue(async () => { this.#assertHealthy(); - return this.#load(resolved.id, resolved.address, context); + return this.#load(definition, resolved.id, resolved.address, context); })); if (loaded === undefined) return undefined; checkRecordScope(definition, loaded.record); - checkRecordVersion(definition, loaded.record, loaded.version); + checkRecordVersion(definition, loaded.record, loaded.storedVersion); return loaded.tracker.value; } + snapshotAsOf( + token: RewindableConversationDocToken, + conversationId: Id, + at: Id, + context: Context, + ): Promise | undefined>; + snapshotAsOf( + token: RewindableConversationDocFamilyToken, + conversationId: Id, + key: string, + at: Id, + context: Context, + ): Promise | undefined>; + async snapshotAsOf(token: AnyDocToken, ...args: readonly unknown[]): Promise { + this.#assertUsable(); + const definition = token.definition; + const resolved = resolveAddress(definition, args); + if (resolved.address.scope.kind !== "conversation") { + throw new TypeError("Session.snapshotAsOf() requires a conversation document"); + } + const conversationId = resolved.address.scope.conversationId; + const at = args[resolved.nextArgument] as Id; + const context = args[resolved.nextArgument + 1] as Context; + return this.#enqueue(async () => { + this.#assertHealthy(); + const storedEntry = await this.#storage.entry(conversationId, at, context); + if (storedEntry === undefined) { + throw new Error(`Entry ${at} is not visible from conversation ${conversationId}`); + } + const address: DocumentAddress = { + ...resolved.address, + scope: { kind: "conversation", conversationId: storedEntry.entry.conversationId }, + }; + const record = await this.#storage.findDocument(address, storedEntry.commitSeq, context); + if (record === undefined) return undefined; + const stored = await this.#storage.document(record.id, storedEntry.commitSeq, context); + if (stored === undefined) { + throw new Error(`Historical document ${record.id} (${record.kind}) cannot be read`); + } + return materializeDocument(definition, stored); + }); + } + close(context: Context): Promise { if (this.#closing === undefined) { const cleanup = withoutAbortSignal(context); @@ -197,19 +248,24 @@ export class SessionKernel implements Session { }); } - async #load(addressId: string, address: DocumentAddress, context: Context): Promise { + async #load( + definition: AnyDocToken["definition"], + addressId: string, + address: DocumentAddress, + context: Context, + ): Promise { const cached = this.#documents.get(addressId); if (cached !== undefined) return cached; const record = await this.#storage.findDocument(address, "current", context); if (record === undefined) return undefined; const stored = await this.#storage.document(record.id, "current", context); if (stored === undefined) throw new Error(`Current document ${record.id} (${record.kind}) cannot be read`); - // Storage returns detached strict JSON, which the tracker owns without another copy. + const value = materializeDocument(definition, stored); const loaded: LoadedDocument = { addressId, record: stored.record, - version: stored.version, - tracker: track(stored.value), + storedVersion: stored.version, + tracker: track(value), }; this.#documents.set(addressId, loaded); return loaded; diff --git a/packages/durable/src/session/transaction.ts b/packages/durable/src/session/transaction.ts index 5f10b72b9..e5498248e 100644 --- a/packages/durable/src/session/transaction.ts +++ b/packages/durable/src/session/transaction.ts @@ -68,8 +68,8 @@ export type DocumentCommitChange = { export type LoadedDocument = { readonly addressId: string; readonly record: DocumentRecord; - /** Stored definition version of the tracked value. */ - readonly version: number; + /** Persisted definition version; older while the tracked value is migrated only in memory. */ + storedVersion: number; readonly tracker: Tracker; }; @@ -78,8 +78,13 @@ export interface TransactionHost { readonly storage: Storage; /** Return the cached current incarnation without loading. */ cached(addressId: string): LoadedDocument | undefined; - /** Return the cached current incarnation, cold-loading it on the mutation line when necessary. */ - load(addressId: string, address: DocumentAddress, context: Context): Promise; + /** Return the cached current incarnation, cold-loading and migrating it when necessary. */ + load( + definition: AnyDocDefinition, + addressId: string, + address: DocumentAddress, + context: Context, + ): Promise; /** Install a newly committed incarnation. */ install(document: LoadedDocument): void; /** Remove a retired incarnation if it is still the cached occupant of its address. */ @@ -343,11 +348,13 @@ export class Transaction implements Tx { async #acquire(entry: DocumentEntry, seed: JsonValue | undefined, skipLoad: boolean): Promise> { const definition = entry.definition!; - const loaded = skipLoad ? undefined : await this.#host.load(entry.addressId, entry.address, this.#context); + const loaded = skipLoad + ? undefined + : await this.#host.load(definition, entry.addressId, entry.address, this.#context); this.#assertOpen(); if (loaded !== undefined) { checkRecordScope(definition, loaded.record); - checkRecordVersion(definition, loaded.record, loaded.version); + checkRecordVersion(definition, loaded.record, loaded.storedVersion); entry.target = { kind: "loaded", document: loaded }; entry.change = loaded.tracker.beginChange(); return entry.change.state; @@ -441,7 +448,7 @@ export class Transaction implements Tx { this.#host.install({ addressId: document.addressId, record, - version: target.version, + storedVersion: target.version, tracker: target.tracker, }); } @@ -459,6 +466,9 @@ export class Transaction implements Tx { const changed = prepared.ops.length > 0; if (changed) target.document.tracker.adopt(prepared); else prepared.abort(); + if (target.document.storedVersion < document.definition!.version) { + target.document.storedVersion = document.definition!.version; + } if (!document.retireOnCommit && !changed) break; if (document.retireOnCommit) { this.#host.evict(target.document.addressId, target.document.record.id); @@ -573,22 +583,30 @@ export class Transaction implements Tx { }); if (document.retireOnCommit) writes.push({ type: "document.retire", id: target.record.id }); break; - case "loaded": - if (document.prepared!.ops.length > 0) { + case "loaded": { + const definition = document.definition!; + const prepared = document.prepared!; + if (target.document.storedVersion < definition.version) { writes.push({ type: "document.change", id: target.document.record.id, - content: { - version: target.document.version, - kind: "delta", - ops: document.prepared!.ops, - }, + content: { version: definition.version, kind: "base", value: prepared.value }, + }); + } else if (prepared.ops.length > 0) { + const useBase = definition.checkpointWhen?.(prepared.value, prepared.ops) ?? false; + writes.push({ + type: "document.change", + id: target.document.record.id, + content: useBase + ? { version: definition.version, kind: "base", value: prepared.value } + : { version: definition.version, kind: "delta", ops: prepared.ops }, }); } if (document.retireOnCommit) { writes.push({ type: "document.retire", id: target.document.record.id }); } break; + } case "retire-only": writes.push({ type: "document.retire", id: target.record.id }); break; diff --git a/packages/durable/src/storage/jsonl/storage.ts b/packages/durable/src/storage/jsonl/storage.ts index fc22cdc41..2aacfcaa4 100644 --- a/packages/durable/src/storage/jsonl/storage.ts +++ b/packages/durable/src/storage/jsonl/storage.ts @@ -309,8 +309,12 @@ export class JsonlStorage implements Storage { return this.store.scanConversations(limit, cursor, context); } - async entry(id: Id, context: Context) { - return this.store.entry(id, context); + entry(id: Id, context: Context): ReturnType; + entry(conversationId: Id, id: Id, context: Context): ReturnType; + async entry(idOrConversationId: Id, idOrContext: Id | Context, context?: Context) { + return context === undefined + ? this.store.entry(idOrConversationId, idOrContext as Context) + : this.store.entry(idOrConversationId, idOrContext as Id, context); } async findLatestHeadMarker(conversationId: Id, atOrBeforeEntryId: Id | undefined, context: Context) { diff --git a/packages/durable/src/storage/memory.ts b/packages/durable/src/storage/memory.ts index feba4a886..0c5d25f9e 100644 --- a/packages/durable/src/storage/memory.ts +++ b/packages/durable/src/storage/memory.ts @@ -333,12 +333,23 @@ export class MemoryStorage implements Storage { return page(values, limit); } - async entry( + entry(id: Id, context: Context): Promise<{ readonly entry: EntryRecord; readonly commitSeq: Seq } | undefined>; + entry( + conversationId: Id, id: Id, - _context: Context, + context: Context, + ): Promise<{ readonly entry: EntryRecord; readonly commitSeq: Seq } | undefined>; + async entry( + idOrConversationId: Id, + idOrContext: Id | Context, + context?: Context, ): Promise<{ readonly entry: EntryRecord; readonly commitSeq: Seq } | undefined> { this.assertOpen(); - const entry = this.state.entries.get(id); + const id = context === undefined ? idOrConversationId : (idOrContext as Id); + const entry = + context === undefined + ? this.state.entries.get(id) + : this.visibleEntries(idOrConversationId, id, id).next().value; if (entry === undefined) return undefined; return { entry: clone(entry), commitSeq: this.state.entryCommitSeqs.get(id)! }; } diff --git a/packages/durable/src/storage/sqlite/storage.ts b/packages/durable/src/storage/sqlite/storage.ts index d6c6fd4db..d1a7bab3e 100644 --- a/packages/durable/src/storage/sqlite/storage.ts +++ b/packages/durable/src/storage/sqlite/storage.ts @@ -230,13 +230,36 @@ export class SqliteStorage implements Storage { ); } - async entry( + entry(id: Id, context: Context): Promise<{ readonly entry: EntryRecord; readonly commitSeq: Seq } | undefined>; + entry( + conversationId: Id, id: Id, - _context: Context, + context: Context, + ): Promise<{ readonly entry: EntryRecord; readonly commitSeq: Seq } | undefined>; + async entry( + idOrConversationId: Id, + idOrContext: Id | Context, + context?: Context, ): Promise<{ readonly entry: EntryRecord; readonly commitSeq: Seq } | undefined> { this.assertOpen(); + let conversation = context === undefined ? undefined : this.readConversation(idOrConversationId); + if (context !== undefined && conversation === undefined) { + throw new Error(`Unknown conversation: ${idOrConversationId}`); + } + const id = context === undefined ? idOrConversationId : (idOrContext as Id); const row = getRow(this.db.prepare("SELECT record, commit_seq FROM entries WHERE id = ?"), id); - return row === undefined ? undefined : { entry: parseJson(row.record), commitSeq: row.commit_seq }; + if (row === undefined) return undefined; + const entry = parseJson(row.record); + if (conversation !== undefined) { + let upperEntryId = Number.POSITIVE_INFINITY; + while (conversation.id !== entry.conversationId) { + if (conversation.parent === undefined) return undefined; + upperEntryId = Math.min(upperEntryId, conversation.parent.at); + conversation = this.readConversation(conversation.parent.conversationId)!; + } + if (entry.id > upperEntryId) return undefined; + } + return { entry, commitSeq: row.commit_seq }; } async findLatestHeadMarker( diff --git a/packages/durable/src/testing/storage-conformance.ts b/packages/durable/src/testing/storage-conformance.ts index cc7d651b2..b569326e3 100644 --- a/packages/durable/src/testing/storage-conformance.ts +++ b/packages/durable/src/testing/storage-conformance.ts @@ -412,6 +412,20 @@ export function createStorageConformance(options: StorageConformanceOptions): re expect((await storage.entry(grandchildHead, context))?.commitSeq).toBe(grandchildEntriesSeq); expect((await storage.entry(grandchildTail, context))?.commitSeq).toBe(grandchildEntriesSeq); expect(await storage.entry(999_999, context)).toBeUndefined(); + + expect(await storage.entry(grandchildId, rootFirst, context)).toEqual({ + entry: entry(rootFirst, rootId), + commitSeq: rootEntriesSeq, + }); + expect((await storage.entry(grandchildId, childForkPoint, context))?.entry.conversationId).toBe(childId); + expect((await storage.entry(grandchildId, grandchildTail, context))?.commitSeq).toBe(grandchildEntriesSeq); + expect(await storage.entry(grandchildId, rootExcludedSameCommit, context)).toBeUndefined(); + expect(await storage.entry(grandchildId, rootExcludedLater, context)).toBeUndefined(); + expect(await storage.entry(grandchildId, childExcluded, context)).toBeUndefined(); + expect(await storage.entry(grandchildId, childExcludedLater, context)).toBeUndefined(); + expect(await storage.entry(rootId, grandchildHead, context)).toBeUndefined(); + expect(await storage.entry(grandchildId, 999_999, context)).toBeUndefined(); + await expect(storage.entry(999_999, rootFirst, context)).rejects.toThrow("Unknown conversation"); await expect(storage.scanEntries({ conversationId: 999_999 }, 10, undefined, context)).rejects.toThrow( "Unknown conversation", ); diff --git a/packages/durable/src/types.ts b/packages/durable/src/types.ts index d41d14bba..ce7a34dc4 100644 --- a/packages/durable/src/types.ts +++ b/packages/durable/src/types.ts @@ -600,6 +600,20 @@ export interface Session { key: string, context: Context, ): Promise | undefined>; + + snapshotAsOf( + token: RewindableConversationDocToken, + conversationId: Id, + at: Id, + context: Context, + ): Promise | undefined>; + snapshotAsOf( + token: RewindableConversationDocFamilyToken, + conversationId: Id, + key: string, + at: Id, + context: Context, + ): Promise | undefined>; } /** @@ -631,6 +645,12 @@ export interface Storage { /** Look up one global entry and the sequence of the commit that persisted it. */ entry(id: Id, context: Context): Promise<{ readonly entry: EntryRecord; readonly commitSeq: Seq } | undefined>; + /** Look up one entry only when it is visible through the requested conversation's ancestry. */ + entry( + conversationId: Id, + id: Id, + context: Context, + ): Promise<{ readonly entry: EntryRecord; readonly commitSeq: Seq } | undefined>; /** * Return the newest visible entry with a `head` at or below the optional inclusive cutoff. diff --git a/packages/durable/test/jsonl-storage.test.ts b/packages/durable/test/jsonl-storage.test.ts index 13d8a74b3..3947a467c 100644 --- a/packages/durable/test/jsonl-storage.test.ts +++ b/packages/durable/test/jsonl-storage.test.ts @@ -68,7 +68,13 @@ class ReopeningStorage implements Storage { conversation: Storage["conversation"] = (id, readContext) => this.current.conversation(id, readContext); scanConversations: Storage["scanConversations"] = (limit, cursor, readContext) => this.current.scanConversations(limit, cursor, readContext); - entry: Storage["entry"] = (id, readContext) => this.current.entry(id, readContext); + entry(id: number, readContext: Context): ReturnType; + entry(conversationId: number, id: number, readContext: Context): ReturnType; + entry(idOrConversationId: number, idOrContext: number | Context, readContext?: Context) { + return readContext === undefined + ? this.current.entry(idOrConversationId, idOrContext as Context) + : this.current.entry(idOrConversationId, idOrContext as number, readContext); + } findLatestHeadMarker: Storage["findLatestHeadMarker"] = (conversationId, at, readContext) => this.current.findLatestHeadMarker(conversationId, at, readContext); scanEntries: Storage["scanEntries"] = (query, limit, cursor, readContext) => diff --git a/packages/durable/test/session-checkpoints-migrations.test.ts b/packages/durable/test/session-checkpoints-migrations.test.ts new file mode 100644 index 000000000..eccfb279f --- /dev/null +++ b/packages/durable/test/session-checkpoints-migrations.test.ts @@ -0,0 +1,684 @@ +import type { Op } from "@earendil-works/chord/delta"; +import { defineDoc, defineDocFamily, type Id, type JsonObject, type StorageWrite } from "@earendil-works/pi-durable"; +import { describe, expect, it } from "vitest"; +import { context, createConversation, documentChanges, flush, openTestSession } from "./session-support.ts"; + +function documentWrites(writes: readonly StorageWrite[]): readonly StorageWrite[] { + return writes.filter( + (write) => + write.type === "document.create" || write.type === "document.change" || write.type === "document.retire", + ); +} + +describe("Session document checkpoints", () => { + it("selects bases only for nonempty ordinary batches and passes the exact prepared revision and ops", async () => { + const calls: { value: Readonly; ops: readonly Op[] }[] = []; + const falseCalls: { value: Readonly; ops: readonly Op[] }[] = []; + const BaseDoc = defineDoc<{ items: string[] }>({ + kind: "checkpoint.base", + version: 1, + scope: "session", + initial: () => ({ items: [] }), + checkpointWhen: (value, ops) => { + calls.push({ value, ops }); + return true; + }, + }); + const DeltaDoc = defineDoc<{ count: number }>({ + kind: "checkpoint.delta", + version: 1, + scope: "session", + initial: () => ({ count: 0 }), + checkpointWhen: (value, ops) => { + falseCalls.push({ value, ops }); + return false; + }, + }); + const DefaultDoc = defineDoc<{ count: number }>({ + kind: "checkpoint.default", + version: 1, + scope: "session", + initial: () => ({ count: 0 }), + }); + const { session, storage, publications } = openTestSession(); + + await session.commit(async (tx) => { + await tx.doc(BaseDoc); + await tx.doc(DeltaDoc); + await tx.doc(DefaultDoc); + }, context); + expect(calls).toHaveLength(0); + expect(falseCalls).toHaveLength(0); + expect(documentWrites(storage.commits.at(-1)!)).toSatisfy((writes: readonly StorageWrite[]) => + writes.every((write) => write.type === "document.create" && write.content.kind === "base"), + ); + + await session.commit(async (tx) => { + (await tx.doc(BaseDoc)).items.push("x"); + (await tx.doc(DeltaDoc)).count++; + (await tx.doc(DefaultDoc)).count++; + }, context); + await flush(); + const writes = storage.admittedCommits.at(-1)!.filter((write) => write.type === "document.change"); + expect(writes.map((write) => write.content.kind)).toEqual(["base", "delta", "delta"]); + expect(calls).toHaveLength(1); + expect(falseCalls).toHaveLength(1); + const snapshot = (await session.snapshot(BaseDoc, context))!; + const deltaSnapshot = (await session.snapshot(DeltaDoc, context))!; + const [published, publishedDelta] = documentChanges(publications.at(-1)!); + expect(calls[0]!.value).toBe(snapshot); + expect(calls[0]!.ops).toBe(published!.ops); + expect(falseCalls[0]!.value).toBe(deltaSnapshot); + expect(falseCalls[0]!.ops).toBe(publishedDelta!.ops); + if (writes[1]!.type !== "document.change" || writes[1]!.content.kind !== "delta") { + throw new Error("Expected ordinary delta"); + } + expect(writes[1]!.content.ops).toBe(falseCalls[0]!.ops); + expect(writes[0]!.type).toBe("document.change"); + if (writes[0]!.type !== "document.change" || writes[0]!.content.kind !== "base") { + throw new Error("Expected checkpoint base"); + } + expect(writes[0]!.content.value).toBe(snapshot); + }); + + it("skips the predicate for empty batches but calls it for nonempty structural no-ops", async () => { + let calls = 0; + const Doc = defineDoc<{ items: string[] }>({ + kind: "checkpoint.no-op", + version: 1, + scope: "session", + initial: () => ({ items: ["a", "b"] }), + checkpointWhen: () => { + calls++; + return false; + }, + }); + const { session, storage } = openTestSession(); + await session.commit((tx) => tx.doc(Doc).then(() => undefined), context); + const commits = storage.commits.length; + + await session.commit(async (tx) => { + const value = await tx.doc(Doc); + value.items.push("x"); + value.items.pop(); + }, context); + expect(calls).toBe(0); + expect(storage.commits).toHaveLength(commits); + + await session.commit(async (tx) => { + const value = await tx.doc(Doc); + const first = value.items.shift()!; + value.items.unshift(first); + }, context); + expect(calls).toBe(1); + expect(storage.commits.at(-1)![0]).toMatchObject({ + type: "document.change", + content: { kind: "delta" }, + }); + }); + + it("rolls back every prepared document when a checkpoint predicate throws", async () => { + let throwCheckpoint = true; + let firstCalls = 0; + const First = defineDoc<{ count: number }>({ + kind: "checkpoint.rollback.first", + version: 1, + scope: "session", + initial: () => ({ count: 0 }), + checkpointWhen: () => { + firstCalls++; + return false; + }, + }); + const Second = defineDoc<{ count: number }>({ + kind: "checkpoint.rollback.second", + version: 1, + scope: "session", + initial: () => ({ count: 0 }), + checkpointWhen: () => { + if (throwCheckpoint) throw new Error("checkpoint failed"); + return false; + }, + }); + const { session, storage, publications } = openTestSession(); + await session.commit(async (tx) => { + await tx.doc(First); + await tx.doc(Second); + }, context); + await flush(); + const first = (await session.snapshot(First, context))!; + const second = (await session.snapshot(Second, context))!; + const commits = storage.commits.length; + const published = publications.length; + + await expect( + session.commit(async (tx) => { + (await tx.doc(First)).count = 1; + (await tx.doc(Second)).count = 2; + }, context), + ).rejects.toThrow("checkpoint failed"); + await flush(); + expect(firstCalls).toBe(1); + expect(storage.commits).toHaveLength(commits); + expect(publications).toHaveLength(published); + expect(await session.snapshot(First, context)).toBe(first); + expect(await session.snapshot(Second, context)).toBe(second); + + throwCheckpoint = false; + await session.commit(async (tx) => { + (await tx.doc(First)).count = 3; + (await tx.doc(Second)).count = 4; + }, context); + expect(await session.snapshot(First, context)).toEqual({ count: 3 }); + expect(await session.snapshot(Second, context)).toEqual({ count: 4 }); + }); + + it("persists repeated false decisions as deltas and replays the complete tail", async () => { + let calls = 0; + const Doc = defineDoc<{ values: number[] }>({ + kind: "checkpoint.tail", + version: 1, + scope: "session", + initial: () => ({ values: [] }), + checkpointWhen: () => { + calls++; + return false; + }, + }); + const { session, storage } = openTestSession(); + await session.commit((tx) => tx.doc(Doc).then(() => undefined), context); + for (let value = 1; value <= 8; value++) { + await session.commit(async (tx) => { + (await tx.doc(Doc)).values.push(value); + }, context); + } + const changes = storage.commits.flat().filter((write) => write.type === "document.change"); + expect(changes).toHaveLength(8); + expect(changes.every((write) => write.content.kind === "delta")).toBe(true); + expect(calls).toBe(8); + await session.unloadDocuments(); + expect(await session.snapshot(Doc, context)).toEqual({ values: [1, 2, 3, 4, 5, 6, 7, 8] }); + }); + + it("keeps a prepared root replacement as a delta when the predicate is false", async () => { + const initial: Record = {}; + for (let index = 0; index < 4_100; index++) initial[`field${index}`] = 0; + const Doc = defineDoc>({ + kind: "checkpoint.root-replacement", + version: 1, + scope: "session", + initial: () => initial, + checkpointWhen: () => false, + }); + const { session, storage } = openTestSession(); + await session.commit((tx) => tx.doc(Doc).then(() => undefined), context); + await session.commit(async (tx) => { + const value = await tx.doc(Doc); + for (let index = 0; index < 4_100; index++) value[`field${index}`] = 1; + }, context); + const write = storage.commits.at(-1)![0]!; + expect(write.type).toBe("document.change"); + if (write.type !== "document.change" || write.content.kind !== "delta") { + throw new Error("Expected root replacement delta"); + } + expect(write.content.ops).toHaveLength(1); + expect(write.content.ops[0]![0]).toBe("r"); + await session.unloadDocuments(); + expect((await session.snapshot(Doc, context))!.field4099).toBe(1); + }); + + it("uses ordinary checkpoint selection before retirement", async () => { + let calls = 0; + const Doc = defineDoc<{ count: number }>({ + kind: "checkpoint.retire", + version: 1, + scope: "session", + initial: () => ({ count: 0 }), + checkpointWhen: () => { + calls++; + return true; + }, + }); + const { session, storage } = openTestSession(); + await session.commit((tx) => tx.doc(Doc).then(() => undefined), context); + await session.commit(async (tx) => { + (await tx.doc(Doc)).count = 1; + await tx.retireDoc(Doc); + }, context); + expect(calls).toBe(1); + expect(documentWrites(storage.commits.at(-1)!)).toEqual([ + expect.objectContaining({ type: "document.change", content: expect.objectContaining({ kind: "base" }) }), + expect.objectContaining({ type: "document.retire" }), + ]); + }); +}); + +describe("Session document migrations", () => { + it("migrates read-only once per cold load, copies the callback result, and writes nothing", async () => { + type Current = { count: number; labels: string[] }; + const Old = defineDoc<{ count: number }>({ + kind: "migration.read-only", + version: 1, + scope: "session", + initial: () => ({ count: 2 }), + }); + let calls = 0; + let retained: Current | undefined; + const Current = defineDoc({ + kind: "migration.read-only", + version: 3, + scope: "session", + initial: () => ({ count: 0, labels: [] }), + migrate: (value, fromVersion) => { + expect(fromVersion).toBe(1); + calls++; + retained = { count: value.count as number, labels: ["migrated"] }; + return retained; + }, + }); + const { session, storage } = openTestSession(); + await session.commit((tx) => tx.doc(Old).then(() => undefined), context); + await session.unloadDocuments(); + const commits = storage.commits.length; + + const first = (await session.snapshot(Current, context))!; + expect(await session.snapshot(Current, context)).toBe(first); + expect(calls).toBe(1); + expect(storage.commits).toHaveLength(commits); + retained!.count = 99; + retained!.labels.push("mutated"); + expect(first).toEqual({ count: 2, labels: ["migrated"] }); + + await session.unloadDocuments(); + const second = (await session.snapshot(Current, context))!; + expect(second).not.toBe(first); + expect(second).toEqual(first); + expect(calls).toBe(2); + expect(storage.commits).toHaveLength(commits); + }); + + it("writes the required base on the first successful transaction, then writes deltas", async () => { + const Old = defineDoc<{ count: number }>({ + kind: "migration.transition", + version: 1, + scope: "session", + initial: () => ({ count: 4 }), + }); + let checkpoints = 0; + const Current = defineDoc<{ count: number }>({ + kind: "migration.transition", + version: 3, + scope: "session", + initial: () => ({ count: 0 }), + migrate: (value, fromVersion) => ({ count: (value.count as number) + fromVersion - 1 }), + checkpointWhen: () => { + checkpoints++; + return false; + }, + }); + const { session, storage, publications } = openTestSession(); + await session.commit((tx) => tx.doc(Old).then(() => undefined), context); + await session.unloadDocuments(); + expect(await session.snapshot(Current, context)).toEqual({ count: 4 }); + const snapshot = await session.snapshot(Current, context); + + await session.commit((tx) => tx.doc(Current).then(() => undefined), context); + await flush(); + expect(storage.commits.at(-1)![0]).toMatchObject({ + type: "document.change", + content: { kind: "base", version: 3, value: { count: 4 } }, + }); + expect(checkpoints).toBe(0); + expect(documentChanges(publications.at(-1)!)).toHaveLength(0); + expect(await session.snapshot(Current, context)).toBe(snapshot); + + await session.commit(async (tx) => { + (await tx.doc(Current)).count = 7; + }, context); + expect(storage.commits.at(-1)![0]).toMatchObject({ + type: "document.change", + content: { kind: "delta", version: 3 }, + }); + expect(checkpoints).toBe(1); + await session.unloadDocuments(); + expect(await session.snapshot(Current, context)).toEqual({ count: 7 }); + }); + + it("rolls migration and edits back with the callback, then coalesces later edits into one base", async () => { + const Old = defineDoc<{ count: number }>({ + kind: "migration.rollback", + version: 1, + scope: "session", + initial: () => ({ count: 1 }), + }); + let migrations = 0; + const Current = defineDoc<{ count: number; migrated: boolean }>({ + kind: "migration.rollback", + version: 2, + scope: "session", + initial: () => ({ count: 0, migrated: false }), + migrate: (value) => { + migrations++; + return { count: value.count as number, migrated: true }; + }, + }); + const { session, storage, publications } = openTestSession(); + await session.commit((tx) => tx.doc(Old).then(() => undefined), context); + await session.unloadDocuments(); + const commits = storage.commits.length; + + await expect( + session.commit(async (tx) => { + (await tx.doc(Current)).count = 8; + throw new Error("rollback"); + }, context), + ).rejects.toThrow("rollback"); + expect(storage.commits).toHaveLength(commits); + expect(await session.snapshot(Current, context)).toEqual({ count: 1, migrated: true }); + expect(migrations).toBe(1); + + await session.commit(async (tx) => { + const value = await tx.doc(Current); + value.count = 9; + value.migrated = false; + }, context); + expect(documentWrites(storage.commits.at(-1)!)).toEqual([ + expect.objectContaining({ + type: "document.change", + content: { kind: "base", version: 2, value: { count: 9, migrated: false } }, + }), + ]); + await flush(); + const published = documentChanges(publications.at(-1)!)[0]!; + const admitted = storage.admittedCommits.at(-1)![0]!; + if (admitted.type !== "document.change" || admitted.content.kind !== "base") { + throw new Error("Expected migration base"); + } + expect(published.value).toBe(await session.snapshot(Current, context)); + expect(published.value).toBe(admitted.content.value); + expect(published.ops.length).toBeGreaterThan(0); + expect(migrations).toBe(1); + }); + + it("strict-checks migration results before tracker ownership and remains usable after rejection", async () => { + const Old = defineDoc<{ count: number }>({ + kind: "migration.invalid", + version: 1, + scope: "session", + initial: () => ({ count: 1 }), + }); + const Invalid = defineDoc({ + kind: "migration.invalid", + version: 2, + scope: "session", + initial: () => ({}), + migrate: () => ({ invalid: new Date() }) as unknown as JsonObject, + }); + const Other = defineDoc<{ ok: boolean }>({ + kind: "migration.invalid.other", + version: 1, + scope: "session", + initial: () => ({ ok: true }), + }); + const { session, storage } = openTestSession(); + await session.commit((tx) => tx.doc(Old).then(() => undefined), context); + await session.unloadDocuments(); + const commits = storage.commits.length; + await expect(session.snapshot(Invalid, context)).rejects.toThrow("strict JSON"); + await expect(session.commit((tx) => tx.doc(Invalid).then(() => undefined), context)).rejects.toThrow( + "strict JSON", + ); + expect(storage.commits).toHaveLength(commits); + await session.commit((tx) => tx.doc(Other).then(() => undefined), context); + expect(await session.snapshot(Other, context)).toEqual({ ok: true }); + }); + + it("rejects newer stored versions and older versions without migration for snapshots and transactions", async () => { + const V2 = defineDoc<{ count: number }>({ + kind: "migration.compatibility", + version: 2, + scope: "session", + initial: () => ({ count: 2 }), + }); + const V1 = defineDoc<{ count: number }>({ + kind: "migration.compatibility", + version: 1, + scope: "session", + initial: () => ({ count: 1 }), + }); + const V3WithoutMigration = defineDoc<{ count: number }>({ + kind: "migration.compatibility", + version: 3, + scope: "session", + initial: () => ({ count: 3 }), + }); + const { session, storage } = openTestSession(); + await session.commit((tx) => tx.doc(V2).then(() => undefined), context); + await session.unloadDocuments(); + const commits = storage.commits.length; + + await expect(session.snapshot(V1, context)).rejects.toThrow("newer version 2 than 1"); + await expect(session.commit((tx) => tx.doc(V1).then(() => undefined), context)).rejects.toThrow( + "newer version 2 than 1", + ); + await expect(session.snapshot(V3WithoutMigration, context)).rejects.toThrow("requires migration from version 2"); + await expect(session.commit((tx) => tx.doc(V3WithoutMigration).then(() => undefined), context)).rejects.toThrow( + "requires migration from version 2", + ); + expect(storage.commits).toHaveLength(commits); + }); + + it("persists a required migration base before retirement without consulting the checkpoint predicate", async () => { + const Old = defineDoc<{ count: number }>({ + kind: "migration.retire", + version: 1, + scope: "session", + initial: () => ({ count: 1 }), + }); + const Current = defineDoc<{ count: number }>({ + kind: "migration.retire", + version: 2, + scope: "session", + initial: () => ({ count: 0 }), + migrate: (value) => ({ count: value.count as number }), + checkpointWhen: () => { + throw new Error("must not run"); + }, + }); + const { session, storage } = openTestSession(); + await session.commit((tx) => tx.doc(Old).then(() => undefined), context); + await session.unloadDocuments(); + await session.commit(async (tx) => { + await tx.doc(Current); + await tx.retireDoc(Current); + }, context); + expect(documentWrites(storage.commits.at(-1)!)).toEqual([ + expect.objectContaining({ + type: "document.change", + content: { kind: "base", version: 2, value: { count: 1 } }, + }), + expect.objectContaining({ type: "document.retire" }), + ]); + }); + + it("leaves unaccessed older documents and unavailable definitions untouched", async () => { + const FirstV1 = defineDoc<{ count: number }>({ + kind: "migration.lazy.first", + version: 1, + scope: "session", + initial: () => ({ count: 1 }), + }); + const SecondV1 = defineDoc<{ count: number }>({ + kind: "migration.lazy.second", + version: 1, + scope: "session", + initial: () => ({ count: 2 }), + }); + let secondMigrations = 0; + const FirstV2 = defineDoc<{ count: number }>({ + kind: "migration.lazy.first", + version: 2, + scope: "session", + initial: () => ({ count: 0 }), + migrate: (value) => ({ count: value.count as number }), + }); + defineDoc<{ count: number }>({ + kind: "migration.lazy.second", + version: 2, + scope: "session", + initial: () => ({ count: 0 }), + migrate: (value) => { + secondMigrations++; + return { count: value.count as number }; + }, + }); + const { session, storage } = openTestSession(); + await session.commit(async (tx) => { + await tx.doc(FirstV1); + await tx.doc(SecondV1); + }, context); + await session.unloadDocuments(); + const commits = storage.commits.length; + expect(await session.snapshot(FirstV2, context)).toEqual({ count: 1 }); + expect(storage.commits).toHaveLength(commits); + expect(secondMigrations).toBe(0); + const secondRecord = await storage.findDocument( + { kind: "migration.lazy.second", scope: { kind: "session" } }, + "current", + context, + ); + expect((await storage.document(secondRecord!.id, "current", context))!.version).toBe(1); + }); +}); + +describe("Session historical document snapshots", () => { + it("migrates current and historical rewindable values independently and follows fork ancestry", async () => { + const V1 = defineDoc<{ count: number }>({ + kind: "history.migration", + version: 1, + scope: "conversation", + history: "rewindable", + fork: "asOf", + initial: () => ({ count: 0 }), + }); + const Family = defineDocFamily<{ seed: string; count: number }, string>({ + kind: "history.family", + version: 1, + family: true, + scope: "conversation", + history: "rewindable", + fork: "asOf", + initial: (seed) => ({ seed, count: 0 }), + }); + const migrations: number[] = []; + const V3 = defineDoc<{ count: number; version: number }>({ + kind: "history.migration", + version: 3, + scope: "conversation", + history: "rewindable", + fork: "asOf", + initial: () => ({ count: 0, version: 3 }), + migrate: (value, fromVersion) => { + migrations.push(fromVersion); + return { count: value.count as number, version: 3 }; + }, + }); + const { session, storage } = openTestSession(); + const conversationId = await createConversation(session); + let firstEntry!: Id; + let secondEntry!: Id; + let thirdEntry!: Id; + await session.commit(async (tx) => { + firstEntry = (await tx.appendEntry(conversationId, { kind: "first" })).id; + (await tx.doc(V1, conversationId)).count = 1; + (await tx.doc(Family, conversationId, "member", "seed")).count = 1; + }, context); + await session.commit(async (tx) => { + secondEntry = (await tx.appendEntry(conversationId, { kind: "second" })).id; + (await tx.doc(V1, conversationId)).count = 2; + }, context); + await session.unloadDocuments(); + const commits = storage.commits.length; + expect(await session.snapshot(V3, conversationId, context)).toEqual({ count: 2, version: 3 }); + expect(migrations).toEqual([1]); + expect(storage.commits).toHaveLength(commits); + + await session.commit(async (tx) => { + thirdEntry = (await tx.appendEntry(conversationId, { kind: "third" })).id; + await tx.doc(V3, conversationId); + }, context); + expect(storage.commits.at(-1)!).toContainEqual( + expect.objectContaining({ + type: "document.change", + content: expect.objectContaining({ kind: "base", version: 3 }), + }), + ); + + expect(await session.snapshotAsOf(V3, conversationId, firstEntry, context)).toEqual({ count: 1, version: 3 }); + expect(await session.snapshotAsOf(V3, conversationId, secondEntry, context)).toEqual({ count: 2, version: 3 }); + expect(await session.snapshotAsOf(V3, conversationId, thirdEntry, context)).toEqual({ count: 2, version: 3 }); + expect(migrations).toEqual([1, 1, 1]); + expect(await session.snapshotAsOf(Family, conversationId, "member", firstEntry, context)).toEqual({ + seed: "seed", + count: 1, + }); + + const childId = await session.commit( + async (tx) => + ( + await tx.createConversation({ + parent: { conversationId, at: secondEntry }, + }) + ).id, + context, + ); + expect(await session.snapshotAsOf(V3, childId, firstEntry, context)).toEqual({ count: 1, version: 3 }); + expect(await session.snapshotAsOf(V3, childId, secondEntry, context)).toEqual({ count: 2, version: 3 }); + expect(await session.snapshotAsOf(Family, childId, "member", firstEntry, context)).toEqual({ + seed: "seed", + count: 1, + }); + await expect(session.snapshotAsOf(V3, childId, thirdEntry, context)).rejects.toThrow( + `Entry ${thirdEntry} is not visible`, + ); + }); + + it("selects the incarnation alive at the entry commit across retirement and recreation", async () => { + const Doc = defineDoc<{ value: string }>({ + kind: "history.incarnation", + version: 1, + scope: "conversation", + history: "rewindable", + fork: "asOf", + initial: () => ({ value: "initial" }), + }); + const { session } = openTestSession(); + const conversationId = await createConversation(session); + const beforeCreation = await session.commit( + async (tx) => (await tx.appendEntry(conversationId, { kind: "before" })).id, + context, + ); + let createdAt!: Id; + let retiredAt!: Id; + let recreatedAt!: Id; + await session.commit(async (tx) => { + createdAt = (await tx.appendEntry(conversationId, { kind: "create" })).id; + (await tx.doc(Doc, conversationId)).value = "old"; + }, context); + await session.commit(async (tx) => { + retiredAt = (await tx.appendEntry(conversationId, { kind: "retire" })).id; + await tx.retireDoc(Doc, conversationId); + }, context); + await session.commit(async (tx) => { + recreatedAt = (await tx.appendEntry(conversationId, { kind: "recreate" })).id; + (await tx.doc(Doc, conversationId)).value = "new"; + }, context); + + expect(await session.snapshotAsOf(Doc, conversationId, beforeCreation, context)).toBeUndefined(); + expect(await session.snapshotAsOf(Doc, conversationId, createdAt, context)).toEqual({ value: "old" }); + expect(await session.snapshotAsOf(Doc, conversationId, retiredAt, context)).toBeUndefined(); + expect(await session.snapshotAsOf(Doc, conversationId, recreatedAt, context)).toEqual({ value: "new" }); + await session.close(context); + await expect(session.snapshotAsOf(Doc, conversationId, recreatedAt, context)).rejects.toThrow("closed"); + }); +}); diff --git a/packages/durable/test/sqlite-storage.test.ts b/packages/durable/test/sqlite-storage.test.ts index 73b45d0fa..24f5ce034 100644 --- a/packages/durable/test/sqlite-storage.test.ts +++ b/packages/durable/test/sqlite-storage.test.ts @@ -2,7 +2,7 @@ import { mkdtemp, rm, stat } from "node:fs/promises"; import { tmpdir } from "node:os"; import { join } from "node:path"; import { DatabaseSync } from "node:sqlite"; -import type { JsonValue } from "@earendil-works/chord"; +import type { Context, JsonValue } from "@earendil-works/chord"; import { BACKGROUND_CONTEXT } from "@earendil-works/chord/context"; import { registerStorageConformance } from "@earendil-works/pi-durable/testing"; import { afterEach, describe, expect, it } from "vitest"; @@ -57,7 +57,13 @@ class ReopeningStorage implements Storage { conversation: Storage["conversation"] = (id, readContext) => this.current.conversation(id, readContext); scanConversations: Storage["scanConversations"] = (limit, cursor, readContext) => this.current.scanConversations(limit, cursor, readContext); - entry: Storage["entry"] = (id, readContext) => this.current.entry(id, readContext); + entry(id: number, readContext: Context): ReturnType; + entry(conversationId: number, id: number, readContext: Context): ReturnType; + entry(idOrConversationId: number, idOrContext: number | Context, readContext?: Context) { + return readContext === undefined + ? this.current.entry(idOrConversationId, idOrContext as Context) + : this.current.entry(idOrConversationId, idOrContext as number, readContext); + } findLatestHeadMarker: Storage["findLatestHeadMarker"] = (conversationId, at, readContext) => this.current.findLatestHeadMarker(conversationId, at, readContext); scanEntries: Storage["scanEntries"] = (query, limit, cursor, readContext) =>