feat(durable): add transactional sessions and documents

Add typed document definitions, atomic Session transactions, complete commit publication, task/document lifecycle handling, and paginated Tx scans. Share Chord's strict JSON copier across draft placements and durable ownership boundaries.
This commit is contained in:
Mario Zechner
2026-09-24 22:24:04 +02:00
parent 7cf037c218
commit 19a0361be8
24 changed files with 2797 additions and 143 deletions
+4 -48
View File
@@ -1,3 +1,4 @@
import { copyJson } from "../json.ts";
import type { JsonValue } from "../types.ts";
import { applyImmutable } from "./apply-immutable-batch.ts";
import { applyImmutableTrusted } from "./apply-immutable-trusted.ts";
@@ -1540,52 +1541,7 @@ function singletonPiece(piece: Piece, sourceIndex: number): Piece {
function clonePlacement(value: unknown): Stored {
const proxyNode = isContainer(value) ? (Reflect.get(value, NODE) as OverlayNode | undefined) : undefined;
return proxyNode === undefined ? cloneJson(value) : clonePlacementNode(proxyNode);
}
function cloneJson(value: unknown, ancestors?: Set<object>): Stored {
if (!isContainer(value)) return strictPrimitive(value);
const active = ancestors ?? new Set<object>();
if (active.has(value)) throw new TypeError("Draft placements cannot contain cycles");
active.add(value);
try {
if (Array.isArray(value)) {
if (Object.getPrototypeOf(value) !== Array.prototype || Reflect.ownKeys(value).length !== value.length + 1) {
throw new TypeError("Draft placements must contain dense plain arrays");
}
const result = new Array<JsonValue>(value.length);
for (let index = 0; index < value.length; index++) {
const descriptor = Object.getOwnPropertyDescriptor(value, String(index));
if (descriptor === undefined || !descriptor.enumerable || !("value" in descriptor)) {
throw new TypeError("Draft placements must contain enumerable indexed data properties");
}
defineData(result, String(index), cloneJson(descriptor.value, active));
}
return result;
}
const prototype = Object.getPrototypeOf(value);
if (prototype !== Object.prototype && prototype !== null) {
throw new TypeError("Draft placements must contain plain objects or arrays");
}
const result = Object.create(prototype) as Record<string, JsonValue>;
for (const key of Reflect.ownKeys(value)) {
if (typeof key === "symbol") throw new TypeError("Draft placements cannot contain symbol properties");
const descriptor = Object.getOwnPropertyDescriptor(value, key);
if (descriptor === undefined || !descriptor.enumerable || !("value" in descriptor)) {
throw new TypeError("Draft placements must contain enumerable data properties");
}
defineData(result, key, cloneJson(descriptor.value, active));
}
return result;
} finally {
active.delete(value);
}
}
function strictPrimitive(value: unknown): Primitive {
if (value === null || typeof value === "string" || typeof value === "boolean") return value;
if (typeof value === "number" && Number.isFinite(value)) return value;
throw new TypeError("Draft placements must contain strict JSON values");
return proxyNode === undefined ? copyJson(value) : clonePlacementNode(proxyNode);
}
function cloneNode(node: OverlayNode): Container {
@@ -1633,9 +1589,9 @@ function clonePlacementNode(node: OverlayNode): Container {
}
function clonePlacementStored(value: Stored, context: OverlayContext): Stored {
if (!isContainer(value)) return strictPrimitive(value);
if (!isContainer(value)) return copyJson(value);
const node = context.rawNodes!.get(value);
return node === undefined ? cloneJson(value) : clonePlacementNode(node);
return node === undefined ? copyJson(value) : clonePlacementNode(node);
}
function defineData(target: object, key: PropertyKey, value: JsonValue): void {
+1 -1
View File
@@ -11,7 +11,7 @@ export {
replicatedState,
} from "./api.ts";
export type { Draft } from "./delta/index.ts";
export { isJsonValue } from "./json.ts";
export { type CopyJsonOptions, copyJson, isJsonValue } from "./json.ts";
export {
isRemoteServiceErrorCode,
REMOTE_SERVICE_ERROR_CODES,
+75 -5
View File
@@ -1,15 +1,85 @@
import type { JsonValue } from "./types.ts";
const DATA_DESCRIPTOR: PropertyDescriptor = {
value: undefined,
writable: true,
enumerable: true,
configurable: true,
};
export type CopyJsonOptions = {
/** Omit undefined object properties while preserving strict array semantics. */
readonly omitUndefinedProperties?: boolean;
};
/** Copy a value into an alias-free strict-JSON tree owned by the caller. */
export function copyJson(value: unknown, options?: CopyJsonOptions): JsonValue {
return copy(value, undefined, options?.omitUndefinedProperties === true);
}
function copy(value: unknown, ancestors: Set<object> | undefined, omitUndefinedProperties: boolean): JsonValue {
if (value === null || typeof value === "string" || typeof value === "boolean") return value;
if (typeof value === "number") {
if (Number.isFinite(value)) return value;
throw new TypeError("Value contains a non-finite number and is not strict JSON");
}
if (typeof value !== "object")
throw new TypeError(`Value contains a non-JSON ${typeof value}; expected strict JSON`);
const active = ancestors ?? new Set<object>();
if (active.has(value)) throw new TypeError("Value contains cycles and is not strict JSON");
active.add(value);
try {
if (Array.isArray(value)) {
// Indices plus `length` only: extra, symbol, and missing keys all change the count.
if (Object.getPrototypeOf(value) !== Array.prototype || Reflect.ownKeys(value).length !== value.length + 1) {
throw new TypeError("Value must contain strict JSON dense plain arrays");
}
const result = new Array<JsonValue>(value.length);
for (let index = 0; index < value.length; index++) {
const descriptor = Object.getOwnPropertyDescriptor(value, index);
if (descriptor === undefined || !descriptor.enumerable || !("value" in descriptor)) {
throw new TypeError("Value must contain strict JSON enumerable indexed data properties");
}
defineData(result, String(index), copy(descriptor.value, active, omitUndefinedProperties));
}
return result;
}
const prototype = Object.getPrototypeOf(value);
if (prototype !== Object.prototype && prototype !== null) {
throw new TypeError("Value must contain strict JSON plain objects or arrays");
}
const result = Object.create(prototype) as Record<string, JsonValue>;
for (const key of Reflect.ownKeys(value)) {
if (typeof key === "symbol") throw new TypeError("Value contains a symbol key and is not strict JSON");
const descriptor = Object.getOwnPropertyDescriptor(value, key)!;
if (!descriptor.enumerable || !("value" in descriptor)) {
throw new TypeError("Value must contain strict JSON enumerable data properties");
}
if (descriptor.value === undefined && omitUndefinedProperties) continue;
defineData(result, key, copy(descriptor.value, active, omitUndefinedProperties));
}
return result;
} finally {
active.delete(value);
}
}
function defineData(target: object, key: PropertyKey, value: JsonValue): void {
DATA_DESCRIPTOR.value = value;
Object.defineProperty(target, key, DATA_DESCRIPTOR);
DATA_DESCRIPTOR.value = undefined;
}
/** Return whether a value is finite strict JSON with plain objects and no cycles. */
export function isJsonValue(value: unknown): value is JsonValue {
return check(value, new Set<object>(), 0);
return check(value, new Set<object>());
}
function check(value: unknown, ancestors: Set<object>, depth: number): boolean {
if (depth > 512) return false;
function check(value: unknown, ancestors: Set<object>): boolean {
if (value === null || typeof value === "string" || typeof value === "boolean") return true;
if (typeof value === "number") return Number.isFinite(value);
if (Array.isArray(value)) {
if (Object.getPrototypeOf(value) !== Array.prototype) return false;
const keys = Reflect.ownKeys(value);
if (keys.length !== value.length + 1 || keys.some((key) => typeof key !== "string")) return false;
if (ancestors.has(value)) return false;
@@ -21,7 +91,7 @@ function check(value: unknown, ancestors: Set<object>, depth: number): boolean {
descriptor === undefined ||
!descriptor.enumerable ||
!("value" in descriptor) ||
!check(descriptor.value, ancestors, depth + 1)
!check(descriptor.value, ancestors)
) {
return false;
}
@@ -39,7 +109,7 @@ function check(value: unknown, ancestors: Set<object>, depth: number): boolean {
ancestors.add(value);
try {
for (const descriptor of Object.values(Object.getOwnPropertyDescriptors(value))) {
if (!descriptor.enumerable || !("value" in descriptor) || !check(descriptor.value, ancestors, depth + 1)) {
if (!descriptor.enumerable || !("value" in descriptor) || !check(descriptor.value, ancestors)) {
return false;
}
}
+52 -1
View File
@@ -1,5 +1,5 @@
import { describe, expect, test } from "vitest";
import { isJsonValue } from "../src/index.ts";
import { copyJson, isJsonValue } from "../src/index.ts";
describe("isJsonValue", () => {
test("checks strict JSON without normalizing it", () => {
@@ -12,3 +12,54 @@ describe("isJsonValue", () => {
expect(isJsonValue(cyclic)).toBe(false);
});
});
describe("copyJson", () => {
test("copies strict JSON without retaining aliases", () => {
const shared = { value: 1 };
const input = { left: shared, right: shared };
const copied = copyJson(input) as typeof input;
expect(copied).toEqual(input);
expect(copied).not.toBe(input);
expect(copied.left).not.toBe(shared);
expect(copied.right).not.toBe(shared);
expect(copied.left).not.toBe(copied.right);
});
test("optionally omits undefined object properties without normalizing arrays", () => {
const input = { kept: 1, omitted: undefined, nested: { omitted: undefined, kept: true } };
expect(() => copyJson(input)).toThrow(/strict JSON/);
expect(copyJson(input, { omitUndefinedProperties: true })).toEqual({ kept: 1, nested: { kept: true } });
expect(() => copyJson([undefined], { omitUndefinedProperties: true })).toThrow(/strict JSON/);
});
test("preserves null prototypes and own __proto__ data properties", () => {
const input = Object.create(null) as Record<string, unknown>;
Object.defineProperty(input, "__proto__", {
value: { safe: true },
enumerable: true,
writable: true,
configurable: true,
});
const copied = copyJson(input) as Record<string, unknown>;
expect(Object.getPrototypeOf(copied)).toBeNull();
expect(Object.hasOwn(copied, "__proto__")).toBe(true);
expect(copied.__proto__).toEqual({ safe: true });
});
test("rejects cycles and non-strict container properties", () => {
const cyclic: { self?: unknown } = {};
cyclic.self = cyclic;
const sparse: unknown[] = [];
sparse[1] = 1;
const extra = Object.assign([1], { extra: 2 });
class ArraySubclass extends Array<unknown> {}
const subclass = new ArraySubclass(1);
const accessor = Object.defineProperty({}, "value", { enumerable: true, get: () => 1 });
const hidden = Object.defineProperty({}, "value", { enumerable: false, value: 1 });
const symbol = { [Symbol("value")]: 1 };
for (const invalid of [cyclic, sparse, extra, subclass, accessor, hidden, symbol]) {
expect(isJsonValue(invalid)).toBe(false);
expect(() => copyJson(invalid)).toThrow(/strict JSON/);
}
});
});
+8
View File
@@ -2,6 +2,14 @@
## [Unreleased]
### Breaking Changes
- Reordered Storage scan arguments so the limit precedes the cursor.
### Added
- Added transactional Sessions with typed durable documents, task creation, snapshots, retirement, and commit publications.
## [0.87.1] - 2026-09-22
## [0.87.0] - 2026-09-21
+5 -3
View File
@@ -11,7 +11,7 @@ facades, membranes, document routing, view projection, events, or clone chains.
- Obsolete `pico` and `pico4` prototypes were removed.
- `pico3` remains.
- Packages 1–5 are implemented in `packages/durable`; later Pico5 runtime packages remain.
- Packages 1–7 are implemented in `packages/durable`; later Pico5 runtime packages remain.
## 1. Records, cursors, and memory tables
@@ -113,8 +113,10 @@ Experimental variants under other Delta directories are not Pico APIs.
Implement these packages as one milestone. Keep the implementation layers
separate, but do not build a temporary untyped document-acquisition seam.
Implement the generic `Tx` table surface these tests require: table reads,
`ReadAfterWrite`, ID-creating writes, and full task replacement. Semantic
Implement the generic `Tx` table surface these tests require: exact table reads,
paginated conversation/entry/task scans, `ReadAfterWrite`, ID-creating writes,
and full task replacement. Scans expose the Storage cursor and caller-selected
limit; they never hide an unbounded full scan. Semantic
conversation, entry, task, and scheduler behavior remains in Packages 13–17.
Keep one Astra-immutable tracker per loaded document. Its trusted immutable
+8 -7
View File
@@ -366,8 +366,8 @@ interface Conversation {
context(context: Context): Promise<ContextView>;
entries(
query: Omit<EntryQuery, "conversationId">,
cursor: Cursor | undefined,
limit: number,
cursor: Cursor | undefined,
context: Context,
): Promise<Page<EntryRecord, Cursor>>;
fork(
@@ -804,8 +804,9 @@ interface Tx {
conversation(id: Id): Promise<ConversationRecord | undefined>;
entry(id: Id): Promise<EntryRecord | undefined>;
task(id: Id): Promise<TaskRecord<JsonValue, JsonValue, JsonValue> | undefined>;
scanEntries(query: EntryQuery): Promise<readonly EntryRecord[]>;
scanTasks(query: TaskQuery): Promise<readonly TaskRecord<JsonValue, JsonValue, JsonValue>[]>;
scanConversations(limit: number, cursor?: Cursor): Promise<Page<ConversationRecord, Cursor>>;
scanEntries(query: EntryQuery, limit: number, cursor?: Cursor): Promise<Page<EntryRecord, Cursor>>;
scanTasks(query: TaskQuery, limit: number, cursor?: Cursor): Promise<Page<TaskRecord<JsonValue, JsonValue, JsonValue>, Cursor>>;
createConversation(value: Omit<ConversationRecord, "id">): Promise<ConversationRecord>;
appendEntry(conversationId: Id, value: EntryDraft): Promise<EntryRecord>;
@@ -2039,21 +2040,21 @@ interface Storage {
mintId(): Promise<Id>;
conversation(id: Id, context: Context): Promise<ConversationRecord | undefined>;
scanConversations(cursor: Cursor | undefined, limit: number, context: Context): Promise<Page<ConversationRecord, Cursor>>;
scanConversations(limit: number, cursor: Cursor | undefined, context: Context): Promise<Page<ConversationRecord, Cursor>>;
entry(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, cursor: Cursor | undefined, limit: number, context: Context): Promise<Page<EntryRecord, Cursor>>;
scanEntries(query: EntryQuery, limit: number, cursor: Cursor | undefined, context: Context): Promise<Page<EntryRecord, Cursor>>;
task(id: Id, context: Context): Promise<TaskRecord<JsonValue, JsonValue, JsonValue> | undefined>;
scanTasks(query: TaskQuery, cursor: Cursor | undefined, limit: number, context: Context): Promise<Page<TaskRecord<JsonValue, JsonValue, JsonValue>, Cursor>>;
scanTasks(query: TaskQuery, limit: number, cursor: Cursor | undefined, context: Context): Promise<Page<TaskRecord<JsonValue, JsonValue, JsonValue>, Cursor>>;
submission(id: Id, context: Context): Promise<SubmissionRecord | undefined>;
submissionByRequest(conversationId: Id, requestId: string, context: Context): Promise<SubmissionRecord | undefined>;
findDocument(address: DocumentAddress, at: DocumentPoint, context: Context): Promise<DocumentRecord | undefined>;
document(id: Id, at: DocumentPoint, context: Context): Promise<StoredDocument | undefined>;
scanDocuments(query: DocumentQuery, cursor: Cursor | undefined, limit: number, context: Context): Promise<Page<DocumentRecord, Cursor>>;
scanDocuments(query: DocumentQuery, limit: number, cursor: Cursor | undefined, context: Context): Promise<Page<DocumentRecord, Cursor>>;
close(context: Context): Promise<void>;
}
+177
View File
@@ -0,0 +1,177 @@
import type { JsonValue } from "@earendil-works/chord";
import type {
CommonDocDefinition,
ConversationDocFamilyToken,
ConversationDocToken,
DocDefinition,
DocFamilyDefinition,
DocFamilyToken,
DocToken,
DocumentAddress,
DocumentCreate,
DocumentRecord,
DocumentSemantics,
Id,
JsonObject,
LatestConversationSemantics,
RewindableConversationDocFamilyToken,
RewindableConversationDocToken,
RewindableConversationSemantics,
SessionDocFamilyToken,
SessionDocToken,
TaskDocFamilyToken,
TaskDocToken,
} from "./types.ts";
type FamilyInput<T extends JsonObject, I extends JsonValue> = Omit<CommonDocDefinition<T>, "initial"> & {
readonly family: true;
initial(seed: I): T;
};
/** Define a Session-scoped singleton document. */
export function defineDoc<T extends JsonObject>(
definition: CommonDocDefinition<T> & { readonly scope: "session" },
): SessionDocToken<T>;
/** Define a latest-only conversation singleton document. */
export function defineDoc<T extends JsonObject>(
definition: CommonDocDefinition<T> & LatestConversationSemantics,
): ConversationDocToken<T>;
/** Define a rewindable conversation singleton document. */
export function defineDoc<T extends JsonObject>(
definition: CommonDocDefinition<T> & RewindableConversationSemantics,
): RewindableConversationDocToken<T>;
/** Define a task-scoped singleton document. */
export function defineDoc<T extends JsonObject>(
definition: CommonDocDefinition<T> & { readonly scope: "task" },
): TaskDocToken<T>;
export function defineDoc<T extends JsonObject>(definition: DocDefinition<T>): DocToken<T, DocDefinition<T>> {
validateDefinition(definition);
return { definition };
}
/** Define a Session-scoped document family. */
export function defineDocFamily<T extends JsonObject, I extends JsonValue>(
definition: FamilyInput<T, I> & { readonly scope: "session" },
): SessionDocFamilyToken<T, I>;
/** Define a latest-only conversation document family. */
export function defineDocFamily<T extends JsonObject, I extends JsonValue>(
definition: FamilyInput<T, I> & LatestConversationSemantics,
): ConversationDocFamilyToken<T, I>;
/** Define a rewindable conversation document family. */
export function defineDocFamily<T extends JsonObject, I extends JsonValue>(
definition: FamilyInput<T, I> & RewindableConversationSemantics,
): RewindableConversationDocFamilyToken<T, I>;
/** Define a task-scoped document family. */
export function defineDocFamily<T extends JsonObject, I extends JsonValue>(
definition: FamilyInput<T, I> & { readonly scope: "task" },
): TaskDocFamilyToken<T, I>;
export function defineDocFamily<T extends JsonObject, I extends JsonValue>(
definition: DocFamilyDefinition<T, I>,
): DocFamilyToken<T, I, DocFamilyDefinition<T, I>> {
validateDefinition(definition);
return { definition };
}
/** Erased definition shape used by the Session after overload resolution. */
export type AnyDocDefinition = DocumentSemantics & {
readonly kind: string;
readonly version: number;
readonly family?: true;
initial(seed?: JsonValue): JsonObject;
};
/** Erased singleton or family token. */
export type AnyDocToken = { readonly definition: AnyDocDefinition };
function validateDefinition(definition: AnyDocDefinition): void {
if (!Number.isSafeInteger(definition.version) || definition.version < 1) {
throw new TypeError(`Document ${definition.kind} version must be a positive integer`);
}
}
/** Logical address plus its string identity for maps. */
export type ResolvedAddress = {
readonly address: DocumentAddress;
readonly id: string;
readonly nextArgument: number;
};
/** Resolve an overloaded argument list and return the index after the owner and family key. */
export function resolveAddress(definition: AnyDocDefinition, args: readonly unknown[]): ResolvedAddress {
let index = 0;
let scope: DocumentAddress["scope"];
switch (definition.scope) {
case "session":
scope = { kind: "session" };
break;
case "conversation":
scope = { kind: "conversation", conversationId: ownerId(args[index++], definition) };
break;
case "task":
scope = { kind: "task", taskId: ownerId(args[index++], definition) };
break;
}
const key = definition.family === true ? (args[index++] as string) : undefined;
const address: DocumentAddress =
key === undefined ? { kind: definition.kind, scope } : { kind: definition.kind, scope, key };
return { address, id: addressId(address), nextArgument: index };
}
function ownerId(value: unknown, definition: AnyDocDefinition): Id {
if (!Number.isSafeInteger(value as number)) {
throw new TypeError(`Document ${definition.kind} requires a ${definition.scope} ID`);
}
return value as Id;
}
/** Stable string identity of one logical address. */
export function addressId(address: DocumentAddress): string {
const owner =
address.scope.kind === "session"
? null
: address.scope.kind === "conversation"
? address.scope.conversationId
: address.scope.taskId;
return JSON.stringify([address.kind, address.scope.kind, owner, address.key ?? null]);
}
/** Build the storage create record for a new incarnation at an address. */
export function documentCreate(definition: AnyDocDefinition, address: DocumentAddress, id: Id): DocumentCreate {
const key = address.key === undefined ? {} : { key: address.key };
switch (address.scope.kind) {
case "session":
return { id, kind: address.kind, ...key, scope: address.scope };
case "task":
return { id, kind: address.kind, ...key, scope: address.scope };
case "conversation":
return {
id,
kind: address.kind,
...key,
scope: address.scope,
history: definition.history!,
fork: definition.fork!,
} as DocumentCreate;
}
}
/** Reject typed access whose token disagrees with the persisted scope, history, or fork semantics. */
export function checkRecordScope(definition: AnyDocDefinition, record: DocumentRecord): void {
if (
record.scope.kind !== definition.scope ||
(record.scope.kind === "conversation" &&
(record.history !== definition.history || record.fork !== definition.fork))
) {
throw new TypeError(`Document ${record.id} (${record.kind}) does not match the supplied definition semantics`);
}
}
/** Reject typed access to a stored version the token cannot use without migration. */
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) {
throw new Error(`Document ${record.id} (${record.kind}) requires migration from version ${version}`);
}
}
+25
View File
@@ -1,30 +1,55 @@
export { defineDoc, defineDocFamily } from "./documents.ts";
export { createSession } from "./session/session.ts";
export { ReadAfterWrite } from "./session/transaction.ts";
export { MemoryStorage } from "./storage/memory.ts";
export type {
CommonDocDefinition,
ContextEdit,
ConversationDocFamilyToken,
ConversationDocToken,
ConversationRecord,
Cursor,
DocDefinition,
DocFamilyDefinition,
DocFamilyToken,
DocToken,
DocumentAddress,
DocumentContent,
DocumentCreate,
DocumentPoint,
DocumentQuery,
DocumentRecord,
DocumentSemantics,
EntryDraft,
EntryQuery,
EntryRecord,
Id,
JsonObject,
LatestConversationSemantics,
Page,
RewindableConversationDocFamilyToken,
RewindableConversationDocToken,
RewindableConversationSemantics,
Seq,
Session,
SessionDocFamilyToken,
SessionDocToken,
Storage,
StorageWrite,
StoredDocument,
SubmissionCreate,
SubmissionRecord,
Task,
TaskDefinition,
TaskDocFamilyToken,
TaskDocToken,
TaskOptions,
TaskOutcome,
TaskOutcomeError,
TaskQuery,
TaskRecord,
TaskRef,
TaskState,
Tx,
} from "./types.ts";
export { ROOT_CONVERSATION_ID } from "./types.ts";
@@ -0,0 +1,16 @@
import type { Seq, StorageWrite } from "../types.ts";
import type { DocumentCommitChange } from "./transaction.ts";
/** Complete table record committed without another publication copy. */
export type TableCommitChange = Extract<
StorageWrite,
{ readonly type: "conversation" | "entry" | "task" | "submission" }
>;
export type CommitChange = TableCommitChange | DocumentCommitChange;
/** Every immutable change from one successful Session commit. Change order is unspecified. */
export type CommitPublication = {
readonly seq: Seq;
readonly changes: readonly CommitChange[];
};
+239
View File
@@ -0,0 +1,239 @@
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 {
ConversationDocFamilyToken,
ConversationDocToken,
DocumentAddress,
Id,
JsonObject,
Session,
SessionDocFamilyToken,
SessionDocToken,
Storage,
StorageWrite,
TaskDocFamilyToken,
TaskDocToken,
Tx,
} from "../types.ts";
import type { CommitChange, CommitPublication } from "./publications.ts";
import { type DocumentCommitChange, type LoadedDocument, Transaction, type TransactionHost } from "./transaction.ts";
/** Open a Session kernel over one storage backend. */
export function createSession(storage: Storage): Session {
return new SessionKernel(storage);
}
/**
* Session kernel: one mutation line, the loaded document tracker cache, and committed publication.
*
* Only committed state is observable. Every commit callback, preparation, Storage settlement, adoption, and
* publication enqueue runs while the line is held; listeners run later.
*/
export class SessionKernel implements Session {
readonly #storage: Storage;
readonly #documents = new Map<string, LoadedDocument>();
readonly #listeners = new Set<(publication: CommitPublication, context: Context) => void>();
readonly #host: TransactionHost;
#tail: Promise<void> = Promise.resolve();
#closing: Promise<void> | undefined;
#poison: { readonly error: unknown } | undefined;
constructor(storage: Storage) {
this.#storage = storage;
this.#host = {
storage,
cached: (id) => this.#documents.get(id),
load: (addressId, address, context) => this.#load(addressId, address, context),
install: (document) => {
this.#documents.set(document.addressId, document);
},
evict: (id, recordId) => {
if (this.#documents.get(id)?.record.id === recordId) this.#documents.delete(id);
},
};
}
commit<T>(change: (tx: Tx) => T | Promise<T>, context: Context): Promise<T> {
try {
this.#assertUsable();
} catch (error) {
return Promise.reject(error);
}
return this.#enqueue(() => this.#runCommit(change, context));
}
snapshot<T extends JsonObject>(token: SessionDocToken<T>, context: Context): Promise<Readonly<T> | undefined>;
snapshot<T extends JsonObject>(
token: ConversationDocToken<T>,
conversationId: Id,
context: Context,
): Promise<Readonly<T> | undefined>;
snapshot<T extends JsonObject>(
token: TaskDocToken<T>,
taskId: Id,
context: Context,
): Promise<Readonly<T> | undefined>;
snapshot<T extends JsonObject, I extends JsonValue>(
token: SessionDocFamilyToken<T, I>,
key: string,
context: Context,
): Promise<Readonly<T> | undefined>;
snapshot<T extends JsonObject, I extends JsonValue>(
token: ConversationDocFamilyToken<T, I>,
conversationId: Id,
key: string,
context: Context,
): Promise<Readonly<T> | undefined>;
snapshot<T extends JsonObject, I extends JsonValue>(
token: TaskDocFamilyToken<T, I>,
taskId: Id,
key: string,
context: Context,
): Promise<Readonly<T> | undefined>;
async snapshot(token: AnyDocToken, ...args: readonly unknown[]): Promise<JsonObject | undefined> {
this.#assertUsable();
const definition = token.definition;
const resolved = resolveAddress(definition, args);
const context = args[resolved.nextArgument] as Context;
const loaded =
this.#documents.get(resolved.id) ??
(await this.#enqueue(async () => {
this.#assertHealthy();
return this.#load(resolved.id, resolved.address, context);
}));
if (loaded === undefined) return undefined;
checkRecordScope(definition, loaded.record);
checkRecordVersion(definition, loaded.record, loaded.version);
return loaded.tracker.value;
}
close(context: Context): Promise<void> {
if (this.#closing === undefined) {
const cleanup = withoutAbortSignal(context);
this.#closing = this.#enqueue(async () => {
this.#documents.clear();
await this.#storage.close(cleanup);
});
}
return awaitWithContext(this.#closing, context);
}
/** Register an internal commit listener. Callbacks run later in commit order and must not throw. */
subscribeCommits(listener: (publication: CommitPublication, context: Context) => void): () => void {
this.#listeners.add(listener);
return () => {
this.#listeners.delete(listener);
};
}
/** Drop every loaded tracker on the mutation line; later access cold-loads from Storage. */
unloadDocuments(): Promise<void> {
return this.#enqueue(async () => {
this.#documents.clear();
});
}
async #runCommit<T>(change: (tx: Tx) => T | Promise<T>, context: Context): Promise<T> {
this.#assertHealthy();
context.abortSignal?.throwIfAborted();
const tx = new Transaction(this.#host, context);
let result: T;
try {
result = await change(tx);
} catch (error) {
await tx.settleFailure();
throw error;
}
const writes = await tx.settleSuccess();
if (writes.length === 0) {
tx.discard();
return result;
}
let seq: number;
try {
// Once admitted, caller cancellation does not interrupt Storage settlement.
seq = await this.#storage.commit(writes, withoutAbortSignal(context));
} catch (error) {
tx.discard();
this.#poison = { error };
throw error;
}
let documents: DocumentCommitChange[];
try {
documents = tx.adopt(seq);
} catch (error) {
// Storage already committed; a failed adoption leaves memory behind durable state.
this.#poison = { error };
throw error;
}
this.#publish(seq, writes, documents, context);
return result;
}
#publish(
seq: number,
writes: readonly StorageWrite[],
documents: readonly DocumentCommitChange[],
context: Context,
): void {
if (this.#listeners.size === 0) return;
const changes: CommitChange[] = [];
for (const write of writes) {
switch (write.type) {
case "conversation":
case "entry":
case "task":
case "submission":
changes.push(write);
}
}
for (const document of documents) changes.push(document);
const publication: CommitPublication = { seq, changes };
const listeners = [...this.#listeners];
queueMicrotask(() => {
for (const listener of listeners) listener(publication, context);
});
}
async #load(addressId: string, address: DocumentAddress, context: Context): Promise<LoadedDocument | undefined> {
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 loaded: LoadedDocument = {
addressId,
record: stored.record,
version: stored.version,
tracker: track(stored.value),
};
this.#documents.set(addressId, loaded);
return loaded;
}
#enqueue<T>(job: () => Promise<T>): Promise<T> {
const run = this.#tail.then(job);
this.#tail = run.then(
() => undefined,
() => undefined,
);
return run;
}
#assertUsable(): void {
if (this.#closing !== undefined) throw new Error("Session is closed");
this.#assertHealthy();
}
#assertHealthy(): void {
if (this.#poison !== undefined) {
throw new Error("Session is poisoned by a failed commit after storage admission; reopen it", {
cause: this.#poison.error,
});
}
}
}
+677
View File
@@ -0,0 +1,677 @@
import { type Context, copyJson, type Draft, type JsonValue } from "@earendil-works/chord";
import { type Change, type Op, type Prepared, type Tracker, track } from "@earendil-works/chord/delta";
import {
type AnyDocDefinition,
type AnyDocToken,
addressId,
checkRecordScope,
checkRecordVersion,
documentCreate,
type ResolvedAddress,
resolveAddress,
} from "../documents.ts";
import type {
ConversationDocFamilyToken,
ConversationDocToken,
ConversationRecord,
Cursor,
DocumentAddress,
DocumentCreate,
DocumentRecord,
EntryDraft,
EntryQuery,
EntryRecord,
Id,
JsonObject,
Seq,
SessionDocFamilyToken,
SessionDocToken,
Storage,
StorageWrite,
Task,
TaskDocFamilyToken,
TaskDocToken,
TaskOptions,
TaskQuery,
TaskRecord,
TaskRef,
Tx,
} from "../types.ts";
type AnyTaskRecord = TaskRecord<JsonValue, JsonValue, JsonValue>;
const INTERNAL_SCAN_PAGE_SIZE = 256;
const EMPTY_OPERATIONS: readonly Op[] = [];
const TABLE_JSON_COPY_OPTIONS = { omitUndefinedProperties: true } as const;
/** A transaction read a table after its first table write. Read every required row before writing. */
export class ReadAfterWrite extends Error {
constructor(method: string) {
super(`Tx.${method}() cannot read tables after the first table write`);
this.name = "ReadAfterWrite";
}
}
/** Committed change of one document incarnation. */
export type DocumentCommitChange = {
readonly type: "document";
readonly record: DocumentRecord;
/** Conversation owning the document; task documents derive it from their task record. */
readonly conversationId: Id | undefined;
/** Exact adopted immutable revision, or `null` when this commit retired the incarnation. */
readonly value: JsonObject | null;
/** Exact adopted operations for an ordinary update; empty for creation and retirement. */
readonly ops: readonly Op[];
};
/** One committed document incarnation owned by the Session tracker cache. */
export type LoadedDocument = {
readonly addressId: string;
readonly record: DocumentRecord;
/** Stored definition version of the tracked value. */
readonly version: number;
readonly tracker: Tracker<JsonObject>;
};
/** Session services used by a transaction while it holds the mutation line. */
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<LoadedDocument | undefined>;
/** Install a newly committed incarnation. */
install(document: LoadedDocument): void;
/** Remove a retired incarnation if it is still the cached occupant of its address. */
evict(addressId: string, recordId: Id): void;
}
/** Committed and candidate state for one task touched by this transaction. */
type TransactionTask = {
committedRead?: Promise<AnyTaskRecord | undefined>;
write?: { readonly kind: "create" | "replace"; readonly record: AnyTaskRecord };
publicationConversationId?: Id;
};
/** Storage/cache provenance of one staged document incarnation. */
type DocumentTarget =
| { readonly kind: "loaded"; readonly document: LoadedDocument }
| {
readonly kind: "created";
readonly record: DocumentCreate;
readonly version: number;
readonly tracker: Tracker<JsonObject>;
}
| { readonly kind: "retire-only"; readonly record: DocumentRecord };
/** One document incarnation acquired, created, or retired by this transaction. */
type DocumentEntry = {
readonly addressId: string;
readonly address: DocumentAddress;
/** Absent only for retirement entries discovered by a terminal-task scan. */
readonly definition?: AnyDocDefinition;
/** Memoized public acquisition; absent for metadata-only retirement. */
draftPromise?: Promise<Draft<JsonObject>>;
/** Set after acquisition or retirement lookup finds the affected incarnation. */
target?: DocumentTarget;
change?: Change<JsonObject>;
prepared?: Prepared<JsonObject>;
retireOnCommit: boolean;
/** Resolved before Storage admission so adoption performs no reads. */
conversationId?: Id;
};
/**
* Transaction for one Session commit callback.
*
* Every asynchronous operation is tracked so callback settlement can reject and drain unfinished work. Session calls
* one settlement method, then either discards prepared changes or adopts them once after Storage succeeds.
*/
export class Transaction implements Tx {
readonly #host: TransactionHost;
readonly #context: Context;
readonly #pendingOperations = new Set<Promise<unknown>>();
#sealed = false;
#hasTableWrite = false;
/** Atomic batch; conversation and entry writes stage eagerly, while task and document writes assemble later. */
readonly #writes: StorageWrite[] = [];
readonly #createdConversationIds = new Set<Id>();
/** One entry per task touched by a public read, candidate write, or document-owner lookup. */
readonly #tasksById = new Map<Id, TransactionTask>();
/** Every document acquisition or retirement marker in staging order. */
readonly #documents: DocumentEntry[] = [];
/** Latest transaction-local incarnation or retirement marker at each logical address. */
readonly #latestDocumentByAddress = new Map<string, DocumentEntry>();
constructor(host: TransactionHost, context: Context) {
this.#host = host;
this.#context = context;
}
// ─── Table reads ────────────────────────────────────────────────────────
conversation(id: Id): Promise<ConversationRecord | undefined> {
return this.#read("conversation", () => this.#host.storage.conversation(id, this.#context));
}
entry(id: Id): Promise<EntryRecord | undefined> {
return this.#read("entry", async () => (await this.#host.storage.entry(id, this.#context))?.entry);
}
task(id: Id): Promise<AnyTaskRecord | undefined> {
return this.#read("task", () => this.#committedTask(id));
}
scanConversations(limit: number, cursor?: Cursor) {
return this.#read("scanConversations", () => this.#host.storage.scanConversations(limit, cursor, this.#context));
}
scanEntries(query: EntryQuery, limit: number, cursor?: Cursor) {
return this.#read("scanEntries", () => this.#host.storage.scanEntries(query, limit, cursor, this.#context));
}
scanTasks(query: TaskQuery, limit: number, cursor?: Cursor) {
return this.#read("scanTasks", () => this.#host.storage.scanTasks(query, limit, cursor, this.#context));
}
// ─── Table writes ───────────────────────────────────────────────────────
createConversation(value: Omit<ConversationRecord, "id">): Promise<ConversationRecord> {
return this.#write(async () => {
const id = await this.#host.storage.mintId();
this.#assertOpen();
const record = copyJson({ ...value, id }, TABLE_JSON_COPY_OPTIONS) as ConversationRecord;
this.#createdConversationIds.add(id);
this.#writes.push({ type: "conversation", value: record });
return record;
});
}
appendEntry(conversationId: Id, value: EntryDraft): Promise<EntryRecord> {
return this.#write(async () => {
await this.#requireConversation(conversationId);
this.#assertOpen();
const id = await this.#host.storage.mintId();
this.#assertOpen();
const { head, ...rest } = value;
const record = copyJson(
head === undefined
? { ...rest, id, conversationId }
: { ...rest, id, conversationId, head: head === "self" ? id : head },
TABLE_JSON_COPY_OPTIONS,
) as unknown as EntryRecord;
this.#writes.push({ type: "entry", value: record });
return record;
});
}
createTask<I, S extends { phase: string }, R, H extends object>(
task: Task<I, S, R, H>,
input: I,
options?: TaskOptions,
): Promise<TaskRef<R>> {
return this.#write(async () => {
const conversationId = options?.conversationId;
if (conversationId === undefined) throw new TypeError("Tx.createTask() requires options.conversationId");
await this.#requireConversation(conversationId);
this.#assertOpen();
const definition = task.definition;
const checkpoint = definition.initial(input);
const id = await this.#host.storage.mintId();
this.#assertOpen();
const record = copyJson(
{
id,
conversationId,
kind: definition.name,
version: definition.version,
input,
after: options?.after ?? [],
background: options?.background ?? false,
abortRequested: false,
state: { status: "pending", checkpoint },
},
TABLE_JSON_COPY_OPTIONS,
) as unknown as AnyTaskRecord;
this.#tasksById.set(id, { write: { kind: "create", record } });
return { id };
});
}
setTask(value: AnyTaskRecord): void {
this.#assertOpen();
this.#hasTableWrite = true;
const task = this.#taskEntry(value.id);
if (task.write?.record.state.status === "terminal") {
throw new Error(`Task ${value.id} already has a terminal candidate`);
}
task.write = {
kind: task.write?.kind === "create" ? "create" : "replace",
record: copyJson(value, TABLE_JSON_COPY_OPTIONS) as unknown as AnyTaskRecord,
};
}
// ─── Documents ──────────────────────────────────────────────────────────
doc<T extends JsonObject>(token: SessionDocToken<T>): Promise<Draft<T>>;
doc<T extends JsonObject>(token: ConversationDocToken<T>, conversationId: Id): Promise<Draft<T>>;
doc<T extends JsonObject>(token: TaskDocToken<T>, taskId: Id): Promise<Draft<T>>;
doc<T extends JsonObject, I extends JsonValue>(
token: SessionDocFamilyToken<T, I>,
key: string,
seed: I,
): Promise<Draft<T>>;
doc<T extends JsonObject, I extends JsonValue>(
token: ConversationDocFamilyToken<T, I>,
conversationId: Id,
key: string,
seed: I,
): Promise<Draft<T>>;
doc<T extends JsonObject, I extends JsonValue>(
token: TaskDocFamilyToken<T, I>,
taskId: Id,
key: string,
seed: I,
): Promise<Draft<T>>;
doc(token: AnyDocToken, ...args: readonly unknown[]): Promise<Draft<JsonObject>> {
try {
this.#assertOpen();
const definition = token.definition;
const resolved = resolveAddress(definition, args);
this.#assertTaskDocumentsOpen(resolved);
const latest = this.#latestDocumentByAddress.get(resolved.id);
if (latest !== undefined && !latest.retireOnCommit && latest.draftPromise !== undefined) {
return latest.draftPromise;
}
const seed = definition.family === true ? copyJson(args[resolved.nextArgument]) : undefined;
const docEntry: DocumentEntry = {
addressId: resolved.id,
address: resolved.address,
definition,
retireOnCommit: false,
};
this.#documents.push(docEntry);
this.#latestDocumentByAddress.set(docEntry.addressId, docEntry);
// Capture retirement before awaiting so a pending old acquisition and its replacement stay distinct.
docEntry.draftPromise = this.#track(this.#acquire(docEntry, seed, latest?.retireOnCommit === true));
return docEntry.draftPromise;
} catch (error) {
return Promise.reject(error);
}
}
retireDoc<T extends JsonObject>(token: SessionDocToken<T>): Promise<void>;
retireDoc<T extends JsonObject>(token: ConversationDocToken<T>, conversationId: Id): Promise<void>;
retireDoc<T extends JsonObject>(token: TaskDocToken<T>, taskId: Id): Promise<void>;
retireDoc<T extends JsonObject, I extends JsonValue>(token: SessionDocFamilyToken<T, I>, key: string): Promise<void>;
retireDoc<T extends JsonObject, I extends JsonValue>(
token: ConversationDocFamilyToken<T, I>,
conversationId: Id,
key: string,
): Promise<void>;
retireDoc<T extends JsonObject, I extends JsonValue>(
token: TaskDocFamilyToken<T, I>,
taskId: Id,
key: string,
): Promise<void>;
retireDoc(token: AnyDocToken, ...args: readonly unknown[]): Promise<void> {
try {
this.#assertOpen();
const definition = token.definition;
const resolved = resolveAddress(definition, args);
const latest = this.#latestDocumentByAddress.get(resolved.id);
if (latest?.retireOnCommit) return Promise.resolve();
if (latest?.draftPromise !== undefined) {
// Retirement of an acquired draft persists its final content before retirement.
latest.retireOnCommit = true;
return this.#track(latest.draftPromise.then(() => undefined));
}
const entry: DocumentEntry = {
addressId: resolved.id,
address: resolved.address,
definition,
retireOnCommit: true,
};
this.#documents.push(entry);
this.#latestDocumentByAddress.set(entry.addressId, entry);
return this.#track(this.#findRetirement(entry));
} catch (error) {
return Promise.reject(error);
}
}
async #acquire(entry: DocumentEntry, seed: JsonValue | undefined, skipLoad: boolean): Promise<Draft<JsonObject>> {
const definition = entry.definition!;
const loaded = skipLoad ? undefined : await this.#host.load(entry.addressId, entry.address, this.#context);
this.#assertOpen();
if (loaded !== undefined) {
checkRecordScope(definition, loaded.record);
checkRecordVersion(definition, loaded.record, loaded.version);
entry.target = { kind: "loaded", document: loaded };
entry.change = loaded.tracker.beginChange();
return entry.change.state;
}
const scope = entry.address.scope;
if (scope.kind === "conversation") await this.#requireConversation(scope.conversationId);
if (scope.kind === "task") {
const task = await this.#currentTask(scope.taskId);
if (task === undefined) throw new Error(`Task ${scope.taskId} does not exist`);
if (task.state.status === "terminal") throw new Error(`Task ${scope.taskId} is terminal`);
}
this.#assertOpen();
const value = copyJson(
definition.family === true ? definition.initial(seed) : definition.initial(),
) as JsonObject;
const id = await this.#host.storage.mintId();
this.#assertOpen();
const tracker = track(value);
entry.target = {
kind: "created",
record: documentCreate(definition, entry.address, id),
version: definition.version,
tracker,
};
entry.change = tracker.beginChange();
return entry.change.state;
}
async #findRetirement(entry: DocumentEntry): Promise<void> {
const record =
this.#host.cached(entry.addressId)?.record ??
(await this.#host.storage.findDocument(entry.address, "current", this.#context));
this.#assertOpen();
if (record === undefined) return;
checkRecordScope(entry.definition!, record);
entry.target = { kind: "retire-only", record };
}
// ─── Settlement ─────────────────────────────────────────────────────────
/** Seal after callback failure: abort every change and observe every pending operation. */
async settleFailure(): Promise<void> {
this.#sealed = true;
this.#abortChanges();
await this.#drain();
}
/**
* Seal after callback success, prepare every open change, and assemble the atomic batch.
* Any failure aborts every change before Storage admission.
*/
async settleSuccess(): Promise<readonly StorageWrite[]> {
this.#sealed = true;
if (this.#pendingOperations.size > 0) {
this.#abortChanges();
await this.#drain();
throw new Error("Session commit callback settled before its pending Tx operations");
}
try {
// Synchronously prepare every open change; this revokes every draft.
for (const document of this.#documents) {
if (document.change !== undefined) document.prepared = document.change.prepare();
}
return await this.#assemble();
} catch (error) {
this.#abortChanges();
throw error;
}
}
/** Abort every prepared change after Storage failure or when no write is required. */
discard(): void {
this.#abortChanges();
}
/** Adopt every prepared change by pointer swap after Storage success and describe the publication. */
adopt(seq: Seq): DocumentCommitChange[] {
const publications: DocumentCommitChange[] = [];
for (const document of this.#documents) {
const target = document.target;
if (target === undefined) continue;
switch (target.kind) {
case "created": {
const prepared = document.prepared!;
const record: DocumentRecord = document.retireOnCommit
? { ...target.record, createdAt: seq, retiredAt: seq }
: { ...target.record, createdAt: seq };
if (document.retireOnCommit) prepared.abort();
else {
target.tracker.adopt(prepared);
this.#host.install({
addressId: document.addressId,
record,
version: target.version,
tracker: target.tracker,
});
}
publications.push({
type: "document",
record,
conversationId: document.conversationId,
value: document.retireOnCommit ? null : prepared.value,
ops: EMPTY_OPERATIONS,
});
break;
}
case "loaded": {
const prepared = document.prepared!;
const changed = prepared.ops.length > 0;
if (changed) target.document.tracker.adopt(prepared);
else prepared.abort();
if (!document.retireOnCommit && !changed) break;
if (document.retireOnCommit) {
this.#host.evict(target.document.addressId, target.document.record.id);
}
publications.push({
type: "document",
record: document.retireOnCommit
? { ...target.document.record, retiredAt: seq }
: target.document.record,
conversationId: document.conversationId,
value: document.retireOnCommit ? null : prepared.value,
ops: document.retireOnCommit ? EMPTY_OPERATIONS : prepared.ops,
});
break;
}
case "retire-only":
this.#host.evict(document.addressId, target.record.id);
publications.push({
type: "document",
record: { ...target.record, retiredAt: seq },
conversationId: document.conversationId,
value: null,
ops: EMPTY_OPERATIONS,
});
break;
}
}
return publications;
}
async #assemble(): Promise<StorageWrite[]> {
const storage = this.#host.storage;
for (const [id, task] of this.#tasksById) {
if (task.write?.kind !== "replace") continue;
const committed = await this.#committedTask(id);
if (committed === undefined) throw new Error(`Task ${id} does not exist`);
if (committed.state.status === "terminal") throw new Error(`Task ${id} is already terminal`);
}
// Terminal settlement retires every task document, including ones created by this transaction.
let terminalTaskIds: Set<Id> | undefined;
for (const task of this.#tasksById.values()) {
const candidate = task.write?.record;
if (candidate?.state.status !== "terminal") continue;
terminalTaskIds ??= new Set();
terminalTaskIds.add(candidate.id);
}
if (terminalTaskIds !== undefined) {
const targetedDocumentIds = new Set<Id>();
for (const document of this.#documents) {
const scope = document.address.scope;
if (scope.kind !== "task" || !terminalTaskIds.has(scope.taskId)) continue;
document.retireOnCommit = true;
const target = document.target;
if (target?.kind === "loaded") targetedDocumentIds.add(target.document.record.id);
if (target?.kind === "created" || target?.kind === "retire-only") targetedDocumentIds.add(target.record.id);
}
for (const taskId of terminalTaskIds) {
if (this.#tasksById.get(taskId)?.write?.kind === "create") continue;
let cursor: Cursor | undefined;
do {
const page = await storage.scanDocuments(
{ scope: { kind: "task", taskId }, at: "current" },
INTERNAL_SCAN_PAGE_SIZE,
cursor,
this.#context,
);
for (const record of page.items) {
if (targetedDocumentIds.has(record.id)) continue;
this.#documents.push({
addressId: addressId(record),
address: record,
target: { kind: "retire-only", record },
retireOnCommit: true,
});
targetedDocumentIds.add(record.id);
}
cursor = page.next;
} while (cursor !== undefined);
}
}
// Resolve task-document publication ownership before Storage admission so adoption remains synchronous.
for (const document of this.#documents) {
const target = document.target;
if (target === undefined) continue;
if (target.kind === "loaded" && !document.retireOnCommit && document.prepared!.ops.length === 0) continue;
const scope = document.address.scope;
if (scope.kind === "conversation") document.conversationId = scope.conversationId;
if (scope.kind !== "task") continue;
const task = this.#taskEntry(scope.taskId);
if (task.publicationConversationId === undefined) {
const current = await this.#currentTask(scope.taskId);
if (current !== undefined) task.publicationConversationId = current.conversationId;
}
document.conversationId = task.publicationConversationId;
}
const writes = this.#writes;
for (const task of this.#tasksById.values()) {
if (task.write !== undefined) writes.push({ type: "task", value: task.write.record });
}
for (const document of this.#documents) {
const target = document.target;
if (target === undefined) continue;
switch (target.kind) {
case "created":
writes.push({
type: "document.create",
record: target.record,
content: { version: target.version, kind: "base", value: document.prepared!.value },
});
if (document.retireOnCommit) writes.push({ type: "document.retire", id: target.record.id });
break;
case "loaded":
if (document.prepared!.ops.length > 0) {
writes.push({
type: "document.change",
id: target.document.record.id,
content: {
version: target.document.version,
kind: "delta",
ops: document.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;
}
}
return writes;
}
// ─── Helpers ────────────────────────────────────────────────────────────
#abortChanges(): void {
for (const document of this.#documents) document.change?.abort();
}
async #drain(): Promise<void> {
await Promise.allSettled(this.#pendingOperations);
}
#assertOpen(): void {
if (this.#sealed) throw new Error("Transaction has settled");
}
#assertTaskDocumentsOpen(resolved: ResolvedAddress): void {
const scope = resolved.address.scope;
if (scope.kind === "task" && this.#tasksById.get(scope.taskId)?.write?.record.state.status === "terminal") {
throw new Error(`Task ${scope.taskId} is terminal`);
}
}
/** Register an operation so callback settlement can reject and drain it. */
#track<T>(operation: Promise<T>): Promise<T> {
this.#pendingOperations.add(operation);
const settle = (): void => {
this.#pendingOperations.delete(operation);
};
operation.then(settle, settle);
return operation;
}
#read<T>(method: string, read: () => Promise<T>): Promise<T> {
try {
this.#assertOpen();
if (this.#hasTableWrite) throw new ReadAfterWrite(method);
return this.#track(read());
} catch (error) {
return Promise.reject(error);
}
}
#write<T>(write: () => Promise<T>): Promise<T> {
try {
this.#assertOpen();
this.#hasTableWrite = true;
return this.#track(write());
} catch (error) {
return Promise.reject(error);
}
}
async #requireConversation(id: Id): Promise<void> {
if (this.#createdConversationIds.has(id)) return;
if ((await this.#host.storage.conversation(id, this.#context)) === undefined) {
throw new Error(`Conversation ${id} does not exist`);
}
}
#taskEntry(id: Id): TransactionTask {
let task = this.#tasksById.get(id);
if (task === undefined) {
task = {};
this.#tasksById.set(id, task);
}
return task;
}
/** Latest candidate task record, falling back to committed state; not a caller table read. */
async #currentTask(id: Id): Promise<AnyTaskRecord | undefined> {
return this.#tasksById.get(id)?.write?.record ?? (await this.#committedTask(id));
}
#committedTask(id: Id): Promise<AnyTaskRecord | undefined> {
const task = this.#taskEntry(id);
task.committedRead ??= this.#host.storage.task(id, this.#context);
return task.committedRead;
}
}
+9 -13
View File
@@ -100,11 +100,7 @@ const isCurrentOnly = (record: DocumentCreate): boolean =>
const sidecarKey = (file: string, seq: Seq, ordinal: number): string => JSON.stringify([file, seq, ordinal]);
const jsonLine = (value: unknown): string => {
const encoded = JSON.stringify(value);
if (encoded === undefined) throw new TypeError("JSONL record is not serializable");
return `${encoded}\n`;
};
const jsonLine = (value: MainMarker | SidecarRecord): string => `${JSON.stringify(value)}\n`;
const errorFromFile = (action: string, error: FileError): Error =>
new Error(`JSONL ${action} failed: ${error.message}`, { cause: error });
@@ -309,8 +305,8 @@ export class JsonlStorage implements Storage {
return this.store.conversation(id, context);
}
async scanConversations(cursor: Cursor | undefined, limit: number, context: Context) {
return this.store.scanConversations(cursor, limit, context);
async scanConversations(limit: number, cursor: Cursor | undefined, context: Context) {
return this.store.scanConversations(limit, cursor, context);
}
async entry(id: Id, context: Context) {
@@ -321,16 +317,16 @@ export class JsonlStorage implements Storage {
return this.store.findLatestHeadMarker(conversationId, atOrBeforeEntryId, context);
}
async scanEntries(query: EntryQuery, cursor: Cursor | undefined, limit: number, context: Context) {
return this.store.scanEntries(query, cursor, limit, context);
async scanEntries(query: EntryQuery, limit: number, cursor: Cursor | undefined, context: Context) {
return this.store.scanEntries(query, limit, cursor, context);
}
async task(id: Id, context: Context) {
return this.store.task(id, context);
}
async scanTasks(query: TaskQuery, cursor: Cursor | undefined, limit: number, context: Context) {
return this.store.scanTasks(query, cursor, limit, context);
async scanTasks(query: TaskQuery, limit: number, cursor: Cursor | undefined, context: Context) {
return this.store.scanTasks(query, limit, cursor, context);
}
async submission(id: Id, context: Context) {
@@ -349,8 +345,8 @@ export class JsonlStorage implements Storage {
return this.store.document(id, at, context);
}
async scanDocuments(query: DocumentQuery, cursor: Cursor | undefined, limit: number, context: Context) {
return this.store.scanDocuments(query, cursor, limit, context);
async scanDocuments(query: DocumentQuery, limit: number, cursor: Cursor | undefined, context: Context) {
return this.store.scanDocuments(query, limit, cursor, context);
}
async close(context: Context): Promise<void> {
+4 -7
View File
@@ -320,8 +320,8 @@ export class MemoryStorage implements Storage {
}
async scanConversations(
cursor: Cursor | undefined,
limit: number,
cursor: Cursor | undefined,
_context: Context,
): Promise<Page<ConversationRecord, Cursor>> {
this.assertOpen();
@@ -370,8 +370,8 @@ export class MemoryStorage implements Storage {
async scanEntries(
query: EntryQuery,
cursor: Cursor | undefined,
limit: number,
cursor: Cursor | undefined,
_context: Context,
): Promise<Page<EntryRecord, Cursor>> {
this.assertOpen();
@@ -394,8 +394,8 @@ export class MemoryStorage implements Storage {
async scanTasks(
query: TaskQuery,
cursor: Cursor | undefined,
limit: number,
cursor: Cursor | undefined,
_context: Context,
): Promise<Page<StoredTask, Cursor>> {
this.assertOpen();
@@ -474,8 +474,8 @@ export class MemoryStorage implements Storage {
async scanDocuments(
query: DocumentQuery,
cursor: Cursor | undefined,
limit: number,
cursor: Cursor | undefined,
_context: Context,
): Promise<Page<DocumentRecord, Cursor>> {
this.assertOpen();
@@ -581,9 +581,6 @@ export class MemoryStorage implements Storage {
if (action.create === undefined && existing === undefined) throw new Error(`Unknown document: ${id}`);
if (action.create !== undefined && existing !== undefined) throw new Error(`Document ${id} already exists`);
if (existing?.record.retiredAt !== undefined) throw new Error(`Document ${id} is retired`);
if (action.create !== undefined && action.content?.kind !== "base") {
throw new Error(`Document ${id} creation requires a base`);
}
const previous = existing?.revisions.at(-1);
if (action.content?.kind === "delta") {
if (previous === undefined) throw new Error(`Document ${id} delta has no base`);
@@ -213,8 +213,8 @@ export class SqliteStorage implements Storage {
}
async scanConversations(
cursor: Cursor | undefined,
limit: number,
cursor: Cursor | undefined,
_context: Context,
): Promise<Page<ConversationRecord, Cursor>> {
this.assertOpen();
@@ -273,8 +273,8 @@ export class SqliteStorage implements Storage {
async scanEntries(
query: EntryQuery,
cursor: Cursor | undefined,
limit: number,
cursor: Cursor | undefined,
_context: Context,
): Promise<Page<EntryRecord, Cursor>> {
this.assertOpen();
@@ -317,8 +317,8 @@ export class SqliteStorage implements Storage {
async scanTasks(
query: TaskQuery,
cursor: Cursor | undefined,
limit: number,
cursor: Cursor | undefined,
_context: Context,
): Promise<Page<StoredTask, Cursor>> {
this.assertOpen();
@@ -433,8 +433,8 @@ export class SqliteStorage implements Storage {
async scanDocuments(
query: DocumentQuery,
cursor: Cursor | undefined,
limit: number,
cursor: Cursor | undefined,
_context: Context,
): Promise<Page<DocumentRecord, Cursor>> {
this.assertOpen();
@@ -544,9 +544,6 @@ export class SqliteStorage implements Storage {
if (action.create === undefined && existing === undefined) throw new Error(`Unknown document: ${id}`);
if (action.create !== undefined && existing !== undefined) throw new Error(`Document ${id} already exists`);
if (existing?.retiredAt !== undefined) throw new Error(`Document ${id} is retired`);
if (action.create !== undefined && action.content?.kind !== "base") {
throw new Error(`Document ${id} creation requires a base`);
}
if (action.content?.kind === "delta") {
const previous = getRow<{ readonly version: number }>(
this.db.prepare(
@@ -280,7 +280,7 @@ export const STORAGE_READ_BENCHMARKS: readonly StorageReadBenchmark[] = [
name: "entry page scan (100)",
async run(storage) {
return (
await storage.scanEntries({ conversationId: ROOT_CONVERSATION_ID }, undefined, 100, BACKGROUND_CONTEXT)
await storage.scanEntries({ conversationId: ROOT_CONVERSATION_ID }, 100, undefined, BACKGROUND_CONTEXT)
).items.length;
},
expected: () => 100,
@@ -291,8 +291,8 @@ export const STORAGE_READ_BENCHMARKS: readonly StorageReadBenchmark[] = [
return (
await storage.scanTasks(
{ kind: "benchmark.filtered", status: "pending", background: true },
undefined,
50,
undefined,
BACKGROUND_CONTEXT,
)
).items.length;
@@ -349,8 +349,8 @@ export const STORAGE_READ_BENCHMARKS: readonly StorageReadBenchmark[] = [
return (
await storage.scanEntries(
{ conversationId: dataset.deepestConversationId },
undefined,
100,
undefined,
BACKGROUND_CONTEXT,
)
).items.length;
@@ -245,7 +245,7 @@ export function createStorageConformance(options: StorageConformanceOptions): re
);
expect(
(await storage.scanEntries({ conversationId: rootId }, undefined, 10, context)).items.map(({ id }) => id),
(await storage.scanEntries({ conversationId: rootId }, 10, undefined, context)).items.map(({ id }) => id),
).toEqual([30, 20, 10]);
expect((await storage.findLatestHeadMarker(rootId, undefined, context))?.id).toBe(20);
}),
@@ -264,11 +264,11 @@ export function createStorageConformance(options: StorageConformanceOptions): re
context,
);
const first = await storage.scanEntries({ conversationId: rootId }, undefined, 2, context);
const first = await storage.scanEntries({ conversationId: rootId }, 2, undefined, context);
expect(first.items.map(({ id }) => id)).toEqual([newestId, middleId]);
const appendedId = await storage.mintId();
await storage.commit([{ type: "entry", value: entry(appendedId, rootId) }], context);
const second = await storage.scanEntries({ conversationId: rootId }, first.next, 2, context);
const second = await storage.scanEntries({ conversationId: rootId }, 2, first.next, context);
expect(second.items.map(({ id }) => id)).toEqual([oldestId]);
expect(second.next).toBeUndefined();
}),
@@ -285,11 +285,11 @@ export function createStorageConformance(options: StorageConformanceOptions): re
context,
);
const first = await storage.scanConversations(undefined, 2, context);
const first = await storage.scanConversations(2, undefined, context);
expect(first.items.map(({ id }) => id)).toEqual([rootId, secondId]);
expect(first.next).toBeDefined();
const roundTrippedCursor = JSON.parse(JSON.stringify(first.next)) as NonNullable<typeof first.next>;
const second = await storage.scanConversations(roundTrippedCursor, 2, context);
const second = await storage.scanConversations(2, roundTrippedCursor, context);
expect(second.items.map(({ id }) => id)).toEqual([thirdId]);
expect(second.next).toBeUndefined();
}),
@@ -356,11 +356,11 @@ export function createStorageConformance(options: StorageConformanceOptions): re
const childExcludedLater = await storage.mintId();
await storage.commit([{ type: "entry", value: entry(childExcludedLater, childId) }], context);
const first = await storage.scanEntries({ conversationId: grandchildId }, undefined, 2, context);
const first = await storage.scanEntries({ conversationId: grandchildId }, 2, undefined, context);
expect(first.items.map(({ id }) => id)).toEqual([grandchildTail, grandchildHead]);
const second = await storage.scanEntries({ conversationId: grandchildId }, first.next, 2, context);
const second = await storage.scanEntries({ conversationId: grandchildId }, 2, first.next, context);
expect(second.items.map(({ id }) => id)).toEqual([childForkPoint, rootForkPoint]);
const third = await storage.scanEntries({ conversationId: grandchildId }, second.next, 2, context);
const third = await storage.scanEntries({ conversationId: grandchildId }, 2, second.next, context);
expect(third.items.map(({ id }) => id)).toEqual([rootFirst]);
expect(third.next).toBeUndefined();
@@ -374,16 +374,16 @@ export function createStorageConformance(options: StorageConformanceOptions): re
const activeFirst = await storage.scanEntries(
{ conversationId: grandchildId, minEntryId: currentMarker?.head },
undefined,
1,
undefined,
context,
);
expect(activeFirst.items.map(({ id }) => id)).toEqual([grandchildTail]);
expect(activeFirst.next).toBeDefined();
const activeSecond = await storage.scanEntries(
{ conversationId: grandchildId, minEntryId: currentMarker?.head },
activeFirst.next,
1,
activeFirst.next,
context,
);
expect(activeSecond.items.map(({ id }) => id)).toEqual([grandchildHead]);
@@ -397,8 +397,8 @@ export function createStorageConformance(options: StorageConformanceOptions): re
minEntryId: historicalMarker?.head,
maxEntryId: childForkPoint,
},
undefined,
10,
undefined,
context,
)
).items.map(({ id }) => id),
@@ -412,7 +412,7 @@ 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();
await expect(storage.scanEntries({ conversationId: 999_999 }, undefined, 10, context)).rejects.toThrow(
await expect(storage.scanEntries({ conversationId: 999_999 }, 10, undefined, context)).rejects.toThrow(
"Unknown conversation",
);
}),
@@ -455,17 +455,17 @@ export function createStorageConformance(options: StorageConformanceOptions): re
await storage.commit([{ type: "task", value: terminal }], context);
expect(await storage.task(firstId, context)).toEqual(terminal);
const pendingPage = await storage.scanTasks({ status: "pending" }, undefined, 1, context);
const pendingPage = await storage.scanTasks({ status: "pending" }, 1, undefined, context);
expect(pendingPage.items.map(({ id }) => id)).toEqual([secondId]);
expect(pendingPage.next).toBeDefined();
expect(
(await storage.scanTasks({ status: "pending" }, pendingPage.next, 1, context)).items.map(({ id }) => id),
(await storage.scanTasks({ status: "pending" }, 1, pendingPage.next, context)).items.map(({ id }) => id),
).toEqual([thirdId]);
expect(
(await storage.scanTasks({ status: "terminal", abortRequested: true }, undefined, 10, context)).items,
(await storage.scanTasks({ status: "terminal", abortRequested: true }, 10, undefined, context)).items,
).toEqual([terminal]);
expect(
(await storage.scanTasks({ background: true }, undefined, 10, context)).items.map(({ id }) => id),
(await storage.scanTasks({ background: true }, 10, undefined, context)).items.map(({ id }) => id),
).toEqual([secondId]);
}),
@@ -671,12 +671,12 @@ export function createStorageConformance(options: StorageConformanceOptions): re
});
expect(
(
await storage.scanDocuments({ scope: firstRecord.scope, at: changedAt }, undefined, 10, context)
await storage.scanDocuments({ scope: firstRecord.scope, at: changedAt }, 10, undefined, context)
).items.map(({ id }) => id),
).toEqual([firstId]);
expect(
(
await storage.scanDocuments({ scope: firstRecord.scope, at: retiredAt }, undefined, 10, context)
await storage.scanDocuments({ scope: firstRecord.scope, at: retiredAt }, 10, undefined, context)
).items.map(({ id }) => id),
).toEqual([secondId]);
expect(await storage.document(firstId, retiredAt, context)).toBeUndefined();
@@ -809,18 +809,18 @@ export function createStorageConformance(options: StorageConformanceOptions): re
)?.id,
).toBe(firstId);
expect(
(await storage.scanDocuments({ scope: { kind: "session" }, at: "current" }, undefined, 1, context)).items,
(await storage.scanDocuments({ scope: { kind: "session" }, at: "current" }, 1, undefined, context)).items,
).toHaveLength(1);
const first = await storage.scanDocuments(
{ scope: { kind: "session" }, at: "current" },
undefined,
1,
undefined,
context,
);
const second = await storage.scanDocuments(
{ scope: { kind: "session" }, at: "current" },
first.next,
1,
first.next,
context,
);
expect([...first.items, ...second.items].map(({ id }) => id)).toEqual([firstId, secondId]);
@@ -828,8 +828,8 @@ export function createStorageConformance(options: StorageConformanceOptions): re
(
await storage.scanDocuments(
{ scope: { kind: "conversation", conversationId: rootId }, at: "current" },
undefined,
10,
undefined,
context,
)
).items.map(({ id }) => id),
@@ -851,8 +851,8 @@ export function createStorageConformance(options: StorageConformanceOptions): re
(
await storage.scanDocuments(
{ scope: { kind: "task", taskId }, at: "current", kind: "task.cache" },
undefined,
10,
undefined,
context,
)
).items.map(({ id }) => id),
@@ -984,7 +984,7 @@ export function createStorageConformance(options: StorageConformanceOptions): re
).rejects.toThrow("already has a current incarnation");
expect(await storage.task(taskId, context)).toEqual(task);
expect((await storage.scanTasks({ status: "pending" }, undefined, 10, context)).items).toEqual([task]);
expect((await storage.scanTasks({ status: "pending" }, 10, undefined, context)).items).toEqual([task]);
expect(await storage.submissionByRequest(rootId, "atomic", context)).toEqual(submission);
expect(await storage.entry(entryId, context)).toBeUndefined();
expect(await storage.document(conflictingDocumentId, "current", context)).toBeUndefined();
@@ -1065,10 +1065,10 @@ export function createStorageConformance(options: StorageConformanceOptions): re
context,
);
expect((await storage.scanTasks({ kind: first }, undefined, 10, context)).items.map(({ id }) => id)).toEqual([
expect((await storage.scanTasks({ kind: first }, 10, undefined, context)).items.map(({ id }) => id)).toEqual([
firstTaskId,
]);
expect((await storage.scanTasks({ kind: second }, undefined, 10, context)).items.map(({ id }) => id)).toEqual([
expect((await storage.scanTasks({ kind: second }, 10, undefined, context)).items.map(({ id }) => id)).toEqual([
secondTaskId,
]);
expect((await storage.task(firstTaskId, context))?.kind).toBe(first);
@@ -1099,8 +1099,8 @@ export function createStorageConformance(options: StorageConformanceOptions): re
(
await storage.scanDocuments(
{ scope: { kind: "session" }, at: "current", kind: first },
undefined,
10,
undefined,
context,
)
).items.map(({ id }) => id),
+228 -5
View File
@@ -1,4 +1,4 @@
import type { Context, JsonValue } from "@earendil-works/chord";
import type { Context, Draft, JsonValue } from "@earendil-works/chord";
import type { Op } from "@earendil-works/chord/delta";
import type { Message } from "@earendil-works/pi-ai";
@@ -14,6 +14,127 @@ export type Seq = number;
/** The root conversation always uses this reserved ID. */
export const ROOT_CONVERSATION_ID: Id = 1;
/** Conversation document that retains only its current state. */
export type LatestConversationSemantics = {
readonly scope: "conversation";
readonly history: "latest";
readonly fork: "current" | "initial";
};
/** Conversation document whose history remains addressable for as-of reads. */
export type RewindableConversationSemantics = {
readonly scope: "conversation";
readonly history: "rewindable";
readonly fork: "asOf" | "current" | "initial";
};
/** Ownership and lifetime of a document; only conversation documents declare history and fork behavior. */
export type DocumentSemantics =
| { readonly scope: "session"; readonly history?: never; readonly fork?: never }
| LatestConversationSemantics
| RewindableConversationSemantics
| { readonly scope: "task"; readonly history?: never; readonly fork?: never };
/** Definition fields shared by singleton documents and document families. */
export type CommonDocDefinition<T extends JsonObject> = {
/** Stable persisted kind; part of the public protocol. */
readonly kind: string;
/** Positive integer version of the stored value shape. */
readonly version: number;
initial(): T;
migrate?(value: JsonObject, fromVersion: number): T;
checkpointWhen?(value: Readonly<T>, ops: readonly Op[]): boolean;
};
/** Singleton document definition. */
export type DocDefinition<T extends JsonObject> = CommonDocDefinition<T> & DocumentSemantics;
/** Keyed document family definition; `initial(seed)` runs only when a member is absent. */
export type DocFamilyDefinition<T extends JsonObject, I extends JsonValue> = Omit<CommonDocDefinition<T>, "initial"> &
DocumentSemantics & {
readonly family: true;
initial(seed: I): T;
};
declare const docType: unique symbol;
/** Typed singleton document token passed explicitly to typed access. */
export interface DocToken<T extends JsonObject, D extends DocDefinition<T>> {
readonly definition: D;
readonly [docType]?: T;
}
/** Typed document family token passed explicitly to typed access. */
export interface DocFamilyToken<T extends JsonObject, I extends JsonValue, D extends DocFamilyDefinition<T, I>> {
readonly definition: D;
readonly [docType]?: T;
}
export type SessionDocToken<T extends JsonObject> = DocToken<T, CommonDocDefinition<T> & { readonly scope: "session" }>;
export type ConversationDocToken<T extends JsonObject> = DocToken<
T,
CommonDocDefinition<T> & (LatestConversationSemantics | RewindableConversationSemantics)
>;
export type RewindableConversationDocToken<T extends JsonObject> = DocToken<
T,
CommonDocDefinition<T> & RewindableConversationSemantics
>;
export type TaskDocToken<T extends JsonObject> = DocToken<T, CommonDocDefinition<T> & { readonly scope: "task" }>;
export type SessionDocFamilyToken<T extends JsonObject, I extends JsonValue> = DocFamilyToken<
T,
I,
DocFamilyDefinition<T, I> & { readonly scope: "session" }
>;
export type ConversationDocFamilyToken<T extends JsonObject, I extends JsonValue> = DocFamilyToken<
T,
I,
DocFamilyDefinition<T, I> & (LatestConversationSemantics | RewindableConversationSemantics)
>;
export type RewindableConversationDocFamilyToken<T extends JsonObject, I extends JsonValue> = DocFamilyToken<
T,
I,
DocFamilyDefinition<T, I> & RewindableConversationSemantics
>;
export type TaskDocFamilyToken<T extends JsonObject, I extends JsonValue> = DocFamilyToken<
T,
I,
DocFamilyDefinition<T, I> & { readonly scope: "task" }
>;
declare const taskResultType: unique symbol;
/** Task definition fields currently supported by Session task creation. */
export type TaskDefinition<I, S extends { phase: string }, R, H extends object> = {
/** Registered task kind persisted in `TaskRecord.kind`. */
readonly name: string;
/** Definition version persisted with live input and checkpoints. */
readonly version: number;
/** First durable checkpoint for a newly created task. */
initial(input: I): S;
readonly hooks?: H;
/** Type-only result marker until phase handlers commit typed outcomes. */
readonly [taskResultType]?: R;
};
/** Typed executable task definition. */
export interface Task<I, S extends { phase: string }, R, H extends object> {
readonly definition: TaskDefinition<I, S, R, H>;
}
/** Task ID carrying its result type for typed waits. */
export type TaskRef<R> = { readonly id: Id; readonly [taskResultType]?: R };
/** Creation options for a durable task. */
export type TaskOptions = {
/** Owning conversation; required for Session commits that are not bound to a conversation. */
readonly conversationId?: Id;
/** Tasks that must be terminal before ordinary execution may begin. */
readonly after?: readonly Id[];
/** Excluded from ordinary idle waits and conversation aborts. */
readonly background?: boolean;
};
/** Immutable identity, history ancestry, and task ownership of a transcript scope. */
export type ConversationRecord = {
readonly id: Id;
@@ -379,6 +500,108 @@ export type StorageWrite =
}
| { readonly type: "document.retire"; readonly id: Id };
/**
* Transaction surface of one Session commit callback.
* Table reads and creation results are trusted immutable values and may be shared with internal commit state.
*/
export interface Tx {
conversation(id: Id): Promise<ConversationRecord | undefined>;
entry(id: Id): Promise<EntryRecord | undefined>;
task(id: Id): Promise<TaskRecord<JsonValue, JsonValue, JsonValue> | undefined>;
scanConversations(limit: number, cursor?: Cursor): Promise<Page<ConversationRecord, Cursor>>;
scanEntries(query: EntryQuery, limit: number, cursor?: Cursor): Promise<Page<EntryRecord, Cursor>>;
scanTasks(
query: TaskQuery,
limit: number,
cursor?: Cursor,
): Promise<Page<TaskRecord<JsonValue, JsonValue, JsonValue>, Cursor>>;
/** Returned records are Session-owned immutable values and may be shared with commit listeners. */
createConversation(value: Omit<ConversationRecord, "id">): Promise<ConversationRecord>;
/** Returned records are Session-owned immutable values and may be shared with commit listeners. */
appendEntry(conversationId: Id, value: EntryDraft): Promise<EntryRecord>;
createTask<I, S extends { phase: string }, R, H extends object>(
task: Task<I, S, R, H>,
input: I,
options?: TaskOptions,
): Promise<TaskRef<R>>;
/** Replace one task record completely. */
setTask(value: TaskRecord<JsonValue, JsonValue, JsonValue>): void;
doc<T extends JsonObject>(token: SessionDocToken<T>): Promise<Draft<T>>;
doc<T extends JsonObject>(token: ConversationDocToken<T>, conversationId: Id): Promise<Draft<T>>;
doc<T extends JsonObject>(token: TaskDocToken<T>, taskId: Id): Promise<Draft<T>>;
doc<T extends JsonObject, I extends JsonValue>(
token: SessionDocFamilyToken<T, I>,
key: string,
seed: I,
): Promise<Draft<T>>;
doc<T extends JsonObject, I extends JsonValue>(
token: ConversationDocFamilyToken<T, I>,
conversationId: Id,
key: string,
seed: I,
): Promise<Draft<T>>;
doc<T extends JsonObject, I extends JsonValue>(
token: TaskDocFamilyToken<T, I>,
taskId: Id,
key: string,
seed: I,
): Promise<Draft<T>>;
retireDoc<T extends JsonObject>(token: SessionDocToken<T>): Promise<void>;
retireDoc<T extends JsonObject>(token: ConversationDocToken<T>, conversationId: Id): Promise<void>;
retireDoc<T extends JsonObject>(token: TaskDocToken<T>, taskId: Id): Promise<void>;
retireDoc<T extends JsonObject, I extends JsonValue>(token: SessionDocFamilyToken<T, I>, key: string): Promise<void>;
retireDoc<T extends JsonObject, I extends JsonValue>(
token: ConversationDocFamilyToken<T, I>,
conversationId: Id,
key: string,
): Promise<void>;
retireDoc<T extends JsonObject, I extends JsonValue>(
token: TaskDocFamilyToken<T, I>,
taskId: Id,
key: string,
): Promise<void>;
}
/** Owner of one mutation line, its records, and its tracked documents. */
export interface Session {
/** Run one atomic transaction on the Session mutation line. */
commit<T>(change: (tx: Tx) => T | Promise<T>, context: Context): Promise<T>;
/** Seal admission, settle admitted commits, then close storage. */
close(context: Context): Promise<void>;
snapshot<T extends JsonObject>(token: SessionDocToken<T>, context: Context): Promise<Readonly<T> | undefined>;
snapshot<T extends JsonObject>(
token: ConversationDocToken<T>,
conversationId: Id,
context: Context,
): Promise<Readonly<T> | undefined>;
snapshot<T extends JsonObject>(
token: TaskDocToken<T>,
taskId: Id,
context: Context,
): Promise<Readonly<T> | undefined>;
snapshot<T extends JsonObject, I extends JsonValue>(
token: SessionDocFamilyToken<T, I>,
key: string,
context: Context,
): Promise<Readonly<T> | undefined>;
snapshot<T extends JsonObject, I extends JsonValue>(
token: ConversationDocFamilyToken<T, I>,
conversationId: Id,
key: string,
context: Context,
): Promise<Readonly<T> | undefined>;
snapshot<T extends JsonObject, I extends JsonValue>(
token: TaskDocFamilyToken<T, I>,
taskId: Id,
key: string,
context: Context,
): Promise<Readonly<T> | undefined>;
}
/**
* Atomic persistence boundary for Session records.
*
@@ -401,8 +624,8 @@ export interface Storage {
/** Scan conversations in ascending ID order. */
scanConversations(
cursor: Cursor | undefined,
limit: number,
cursor: Cursor | undefined,
context: Context,
): Promise<Page<ConversationRecord, Cursor>>;
@@ -422,8 +645,8 @@ export interface Storage {
/** Scan the inclusive visible range newest-first, returning at most `limit` entries. */
scanEntries(
query: EntryQuery,
cursor: Cursor | undefined,
limit: number,
cursor: Cursor | undefined,
context: Context,
): Promise<Page<EntryRecord, Cursor>>;
@@ -433,8 +656,8 @@ export interface Storage {
/** Scan task records matching every supplied filter. */
scanTasks(
query: TaskQuery,
cursor: Cursor | undefined,
limit: number,
cursor: Cursor | undefined,
context: Context,
): Promise<Page<TaskRecord<JsonValue, JsonValue, JsonValue>, Cursor>>;
@@ -453,8 +676,8 @@ export interface Storage {
/** Scan incarnations alive in one exact scope at the selected point. */
scanDocuments(
query: DocumentQuery,
cursor: Cursor | undefined,
limit: number,
cursor: Cursor | undefined,
context: Context,
): Promise<Page<DocumentRecord, Cursor>>;
+8 -8
View File
@@ -66,24 +66,24 @@ class ReopeningStorage implements Storage {
mintId: Storage["mintId"] = () => this.current.mintId();
conversation: Storage["conversation"] = (id, readContext) => this.current.conversation(id, readContext);
scanConversations: Storage["scanConversations"] = (cursor, limit, readContext) =>
this.current.scanConversations(cursor, limit, readContext);
scanConversations: Storage["scanConversations"] = (limit, cursor, readContext) =>
this.current.scanConversations(limit, cursor, readContext);
entry: Storage["entry"] = (id, readContext) => this.current.entry(id, readContext);
findLatestHeadMarker: Storage["findLatestHeadMarker"] = (conversationId, at, readContext) =>
this.current.findLatestHeadMarker(conversationId, at, readContext);
scanEntries: Storage["scanEntries"] = (query, cursor, limit, readContext) =>
this.current.scanEntries(query, cursor, limit, readContext);
scanEntries: Storage["scanEntries"] = (query, limit, cursor, readContext) =>
this.current.scanEntries(query, limit, cursor, readContext);
task: Storage["task"] = (id, readContext) => this.current.task(id, readContext);
scanTasks: Storage["scanTasks"] = (query, cursor, limit, readContext) =>
this.current.scanTasks(query, cursor, limit, readContext);
scanTasks: Storage["scanTasks"] = (query, limit, cursor, readContext) =>
this.current.scanTasks(query, limit, cursor, readContext);
submission: Storage["submission"] = (id, readContext) => this.current.submission(id, readContext);
submissionByRequest: Storage["submissionByRequest"] = (conversationId, requestId, readContext) =>
this.current.submissionByRequest(conversationId, requestId, readContext);
findDocument: Storage["findDocument"] = (address, at, readContext) =>
this.current.findDocument(address, at, readContext);
document: Storage["document"] = (id, at, readContext) => this.current.document(id, at, readContext);
scanDocuments: Storage["scanDocuments"] = (query, cursor, limit, readContext) =>
this.current.scanDocuments(query, cursor, limit, readContext);
scanDocuments: Storage["scanDocuments"] = (query, limit, cursor, readContext) =>
this.current.scanDocuments(query, limit, cursor, readContext);
async close(closeContext: Context): Promise<void> {
if (this.closed) return;
@@ -0,0 +1,111 @@
import type { Draft } from "@earendil-works/chord";
import {
createSession,
defineDoc,
defineDocFamily,
type Id,
MemoryStorage,
type Session,
type Tx,
} from "@earendil-works/pi-durable";
import { describe, expect, expectTypeOf, it } from "vitest";
import { context } from "./session-support.ts";
type State = { value: number };
const initial = (): State => ({ value: 0 });
const SessionDoc = defineDoc<State>({ kind: "t.session", version: 1, scope: "session", initial });
const LatestDoc = defineDoc<State>({
kind: "t.latest",
version: 1,
scope: "conversation",
history: "latest",
fork: "current",
initial,
});
const RewindableDoc = defineDoc<State>({
kind: "t.rewindable",
version: 1,
scope: "conversation",
history: "rewindable",
fork: "asOf",
initial,
});
const TaskDoc = defineDoc<State>({ kind: "t.task", version: 1, scope: "task", initial });
const SessionFamily = defineDocFamily<State, number>({
kind: "t.session-family",
version: 1,
family: true,
scope: "session",
initial: (seed) => ({ value: seed }),
});
const ConversationFamily = defineDocFamily<State, number>({
kind: "t.conversation-family",
version: 1,
family: true,
scope: "conversation",
history: "rewindable",
fork: "initial",
initial: (seed) => ({ value: seed }),
});
const TaskFamily = defineDocFamily<State, number>({
kind: "t.task-family",
version: 1,
family: true,
scope: "task",
initial: (seed) => ({ value: seed }),
});
describe("document definitions", () => {
it("validates persisted version semantics", () => {
expect(() => defineDoc<State>({ kind: "k", version: 0, scope: "session", initial })).toThrow("positive integer");
expect(() => defineDoc<State>({ kind: "k", version: 1.5, scope: "session", initial })).toThrow(
"positive integer",
);
expect(defineDoc<State>({ kind: "", version: 1, scope: "session", initial }).definition.kind).toBe("");
expect(SessionDoc.definition.kind).toBe("t.session");
});
it("types every owner, key, and seed overload", async () => {
const session: Session = createSession(new MemoryStorage());
const check = async (tx: Tx, conversationId: Id, taskId: Id): Promise<void> => {
expectTypeOf(await tx.doc(SessionDoc)).toEqualTypeOf<Draft<State>>();
expectTypeOf(await tx.doc(LatestDoc, conversationId)).toEqualTypeOf<Draft<State>>();
expectTypeOf(await tx.doc(RewindableDoc, conversationId)).toEqualTypeOf<Draft<State>>();
expectTypeOf(await tx.doc(TaskDoc, taskId)).toEqualTypeOf<Draft<State>>();
expectTypeOf(await tx.doc(SessionFamily, "k", 1)).toEqualTypeOf<Draft<State>>();
expectTypeOf(await tx.doc(ConversationFamily, conversationId, "k", 1)).toEqualTypeOf<Draft<State>>();
expectTypeOf(await tx.doc(TaskFamily, taskId, "k", 1)).toEqualTypeOf<Draft<State>>();
await tx.retireDoc(SessionDoc);
await tx.retireDoc(LatestDoc, conversationId);
await tx.retireDoc(TaskDoc, taskId);
await tx.retireDoc(SessionFamily, "k");
await tx.retireDoc(ConversationFamily, conversationId, "k");
await tx.retireDoc(TaskFamily, taskId, "k");
// @ts-expect-error Session documents take no owner
await tx.doc(SessionDoc, conversationId);
// @ts-expect-error conversation documents require a conversation ID
await tx.doc(LatestDoc);
// @ts-expect-error family access requires a creation seed
await tx.doc(SessionFamily, "k");
// @ts-expect-error family seeds are typed
await tx.doc(SessionFamily, "k", "seed");
// @ts-expect-error family retirement takes no seed
await tx.retireDoc(TaskFamily, taskId, "k", 1);
};
expect(check).toBeTypeOf("function");
expectTypeOf(await session.snapshot(SessionDoc, context)).toEqualTypeOf<Readonly<State> | undefined>();
expectTypeOf(await session.snapshot(LatestDoc, 1, context)).toEqualTypeOf<Readonly<State> | undefined>();
expectTypeOf(await session.snapshot(TaskDoc, 1, context)).toEqualTypeOf<Readonly<State> | undefined>();
expectTypeOf(await session.snapshot(SessionFamily, "k", context)).toEqualTypeOf<Readonly<State> | undefined>();
expectTypeOf(await session.snapshot(ConversationFamily, 1, "k", context)).toEqualTypeOf<
Readonly<State> | undefined
>();
expectTypeOf(await session.snapshot(TaskFamily, 1, "k", context)).toEqualTypeOf<Readonly<State> | undefined>();
// @ts-expect-error snapshots never take a creation seed
await session.snapshot(SessionFamily, "k", 1, context).catch(() => undefined);
await session.close(context);
});
});
@@ -0,0 +1,651 @@
import type { Draft } from "@earendil-works/chord";
import { defineDoc, defineDocFamily, type Id, type JsonObject } from "@earendil-works/pi-durable";
import { describe, expect, it } from "vitest";
import { context, createConversation, documentChanges, flush, openTestSession } from "./session-support.ts";
type Live = { message?: string; items: string[]; nested: { count: number }; other: { label: string } };
let liveInitCount = 0;
const LiveDoc = defineDoc<Live>({
kind: "test.live",
version: 1,
scope: "conversation",
history: "latest",
fork: "initial",
initial: () => {
liveInitCount++;
return { items: [], nested: { count: 0 }, other: { label: "x" } };
},
});
const RewindableLiveDoc = defineDoc<Live>({
kind: "test.live",
version: 1,
scope: "conversation",
history: "rewindable",
fork: "asOf",
initial: () => ({ items: [], nested: { count: 0 }, other: { label: "x" } }),
});
const LiveDocV2 = defineDoc<Live>({
kind: "test.live",
version: 2,
scope: "conversation",
history: "latest",
fork: "initial",
initial: () => ({ items: [], nested: { count: 0 }, other: { label: "x" } }),
});
type Counter = { count: number };
const CounterDoc = defineDoc<Counter>({
kind: "test.counter",
version: 1,
scope: "session",
initial: () => ({ count: 0 }),
});
type Member = { seed: string; hits: number };
const seeds: string[] = [];
const MemberDoc = defineDocFamily<Member, string>({
kind: "test.member",
version: 1,
family: true,
scope: "session",
initial: (seed) => {
seeds.push(seed);
return { seed, hits: 0 };
},
});
async function setupLive(): Promise<ReturnType<typeof openTestSession> & { readonly conversationId: Id }> {
const harness = openTestSession();
const conversationId = await createConversation(harness.session);
await harness.session.commit(async (tx) => {
const live = await tx.doc(LiveDoc, conversationId);
live.items.push("a", "b");
}, context);
return { ...harness, conversationId };
}
describe("Session document transactions", () => {
it("creates an initial base on first access and adopts it after Storage success", async () => {
const { session, storage, publications } = openTestSession();
const conversationId = await createConversation(session);
await session.commit(async (tx) => {
const live = await tx.doc(LiveDoc, conversationId);
live.message = "hello";
}, context);
const writes = storage.commits.at(-1)!;
expect(writes).toHaveLength(1);
expect(writes[0]).toMatchObject({
type: "document.create",
record: {
kind: "test.live",
scope: { kind: "conversation", conversationId },
history: "latest",
fork: "initial",
},
content: { kind: "base", version: 1, value: { message: "hello", items: [], nested: { count: 0 } } },
});
const snapshot = await session.snapshot(LiveDoc, conversationId, context);
expect(snapshot).toEqual({ message: "hello", items: [], nested: { count: 0 }, other: { label: "x" } });
await flush();
const publication = publications.at(-1)!;
const published = documentChanges(publication)[0]!;
expect(published.record.createdAt).toBe(publication.seq);
expect(published.value).toBe(snapshot);
expect(published.conversationId).toBe(conversationId);
const create = storage.admittedCommits.at(-1)!.find((write) => write.type === "document.create")!;
expect(create.content.value).toBe(published.value);
});
it("never creates on snapshot and returns undefined when absent", async () => {
const { session, storage } = openTestSession();
const conversationId = await createConversation(session);
const before = storage.commits.length;
expect(await session.snapshot(LiveDoc, conversationId, context)).toBeUndefined();
expect(await session.snapshot(CounterDoc, context)).toBeUndefined();
expect(await session.snapshot(MemberDoc, "k", context)).toBeUndefined();
expect(storage.commits.length).toBe(before);
expect(storage.mintCount).toBe(1);
});
it("returns shared immutable snapshots and keeps prior revisions stable", async () => {
const { session, conversationId } = await setupLive();
const first = (await session.snapshot(LiveDoc, conversationId, context))!;
expect(await session.snapshot(LiveDoc, conversationId, context)).toBe(first);
await session.commit(async (tx) => {
const live = await tx.doc(LiveDoc, conversationId);
live.nested.count = 1;
}, context);
const second = (await session.snapshot(LiveDoc, conversationId, context))!;
expect(second).not.toBe(first);
expect(first.nested.count).toBe(0);
expect(second.nested.count).toBe(1);
// Unchanged subtrees are structurally shared between immutable revisions.
expect(second.items).toBe(first.items);
expect(second.other).toBe(first.other);
});
it("adopts by pointer swap and shares operation payloads with the published revision", async () => {
const { session, storage, publications, conversationId } = await setupLive();
await session.commit(async (tx) => {
const live = await tx.doc(LiveDoc, conversationId);
live.other = { label: "y" };
live.items.push("c");
}, context);
await flush();
const snapshot = (await session.snapshot(LiveDoc, conversationId, context))!;
const published = documentChanges(publications.at(-1)!)[0]!;
expect(published.value).toBe(snapshot);
const admitted = storage.admittedCommits.at(-1)!.find((write) => write.type === "document.change")!;
expect(admitted.content.kind).toBe("delta");
if (admitted.content.kind !== "delta") throw new Error("Expected delta");
expect(published.ops).toBe(admitted.content.ops);
const set = published.ops.find((op) => op[0] === "s");
expect(set).toEqual(["s", ["other"], { label: "y" }]);
// Trusted immutability: the Session makes no second copy of operation payloads.
expect(set![2]).toBe(snapshot.other);
});
it("copies assigned values per placement", async () => {
const { session, conversationId } = await setupLive();
const value = { label: "shared" };
await session.commit(async (tx) => {
const live = await tx.doc(LiveDoc, conversationId);
live.other = value;
(live as Draft<Live> & { copy?: { label: string } }).copy = value;
value.label = "mutated";
}, context);
const snapshot = (await session.snapshot(LiveDoc, conversationId, context)) as Live & { copy: { label: string } };
expect(snapshot.other).toEqual({ label: "shared" });
expect(snapshot.copy).toEqual({ label: "shared" });
expect(snapshot.copy).not.toBe(snapshot.other);
expect(snapshot.other).not.toBe(value);
});
it("suppresses writes and publications for empty batches", async () => {
const { session, storage, publications, conversationId } = await setupLive();
await flush();
const commits = storage.commits.length;
const published = publications.length;
await session.commit(async (tx) => {
const live = await tx.doc(LiveDoc, conversationId);
live.nested.count = 0;
live.items.push("z");
live.items.pop();
}, context);
await flush();
expect(storage.commits.length).toBe(commits);
expect(publications.length).toBe(published);
});
it("writes and publishes replayable nonempty structural no-ops", async () => {
const { session, storage, publications, conversationId } = await setupLive();
const before = (await session.snapshot(LiveDoc, conversationId, context))!;
await session.commit(async (tx) => {
const live = await tx.doc(LiveDoc, conversationId);
const first = live.items.shift()!;
live.items.unshift(first);
}, context);
await flush();
const writes = storage.commits.at(-1)!;
expect(writes[0]).toMatchObject({ type: "document.change", content: { kind: "delta" } });
const after = (await session.snapshot(LiveDoc, conversationId, context))!;
expect(after).toEqual(before);
expect(after).not.toBe(before);
expect(documentChanges(publications.at(-1)!)[0]!.ops.length).toBeGreaterThan(0);
});
it("revokes escaped drafts when the callback settles", async () => {
const { session, conversationId } = await setupLive();
let escaped: Draft<Live> | undefined;
let items: Draft<string[]> | undefined;
await session.commit(async (tx) => {
escaped = await tx.doc(LiveDoc, conversationId);
items = escaped.items;
escaped.message = "inside";
}, context);
expect(() => escaped!.message).toThrow();
expect(() => {
escaped!.message = "outside";
}).toThrow();
expect(() => items!.length).toThrow();
expect(() => items!.push("outside")).toThrow();
expect((await session.snapshot(LiveDoc, conversationId, context))!.message).toBe("inside");
const returned = await session.commit((tx) => tx.doc(LiveDoc, conversationId), context);
expect(() => returned.message).toThrow();
});
it("aborts every change when the callback fails", async () => {
const { session, storage, conversationId } = await setupLive();
const before = (await session.snapshot(LiveDoc, conversationId, context))!;
const commits = storage.commits.length;
await expect(
session.commit(async (tx) => {
const live = await tx.doc(LiveDoc, conversationId);
const counter = await tx.doc(CounterDoc);
live.message = "lost";
counter.count = 5;
throw new Error("callback failed");
}, context),
).rejects.toThrow("callback failed");
expect(storage.commits.length).toBe(commits);
expect(await session.snapshot(LiveDoc, conversationId, context)).toBe(before);
expect(await session.snapshot(CounterDoc, context)).toBeUndefined();
await session.commit(async (tx) => {
(await tx.doc(LiveDoc, conversationId)).message = "kept";
}, context);
expect((await session.snapshot(LiveDoc, conversationId, context))!.message).toBe("kept");
});
it("memoizes concurrent duplicate acquisition and initializes once", async () => {
const { session, storage } = openTestSession();
const conversationId = await createConversation(session);
const initCount = liveInitCount;
const mints = storage.mintCount;
await session.commit(async (tx) => {
const [first, second] = await Promise.all([tx.doc(LiveDoc, conversationId), tx.doc(LiveDoc, conversationId)]);
expect(first).toBe(second);
expect(await tx.doc(LiveDoc, conversationId)).toBe(first);
first.message = "once";
}, context);
expect(liveInitCount).toBe(initCount + 1);
expect(storage.mintCount).toBe(mints + 1);
expect(storage.commits.at(-1)!.filter((write) => write.type === "document.create")).toHaveLength(1);
});
it("uses the first family seed and ignores seeds for existing members", async () => {
const { session } = openTestSession();
seeds.length = 0;
await session.commit(async (tx) => {
const first = await tx.doc(MemberDoc, "k", "first");
const second = await tx.doc(MemberDoc, "k", "second");
expect(second).toBe(first);
first.hits++;
}, context);
await session.commit(async (tx) => {
(await tx.doc(MemberDoc, "k", "third")).hits++;
(await tx.doc(MemberDoc, "other", "fourth")).hits++;
}, context);
expect(seeds).toEqual(["first", "fourth"]);
expect(await session.snapshot(MemberDoc, "k", context)).toEqual({ seed: "first", hits: 2 });
expect(await session.snapshot(MemberDoc, "other", context)).toEqual({ seed: "fourth", hits: 1 });
});
it("rejects a callback that succeeds with a pending acquisition and drains it", async () => {
const { session, storage, conversationId } = await setupLive();
await session.unloadDocuments();
const gate = storage.holdFindDocument();
const commits = storage.commits.length;
let pending: Promise<unknown> | undefined;
const commit = session.commit((tx) => {
pending = tx.doc(LiveDoc, conversationId);
}, context);
await gate.entered;
let settled = false;
void commit.catch(() => {
settled = true;
});
await flush();
// The line stays held until the late acquisition settles.
expect(settled).toBe(false);
gate.release();
await expect(commit).rejects.toThrow("pending Tx operations");
await expect(pending).rejects.toThrow("Transaction has settled");
expect(storage.commits.length).toBe(commits);
await session.commit(async (tx) => {
(await tx.doc(LiveDoc, conversationId)).message = "after";
}, context);
expect((await session.snapshot(LiveDoc, conversationId, context))!.message).toBe("after");
});
it("does not initialize or mint for an absent acquisition that finishes after settlement", async () => {
const { session, storage } = openTestSession();
let initialized = 0;
const LateDoc = defineDoc<Counter>({
kind: "test.late",
version: 1,
scope: "session",
initial: () => {
initialized++;
return { count: 0 };
},
});
const gate = storage.holdFindDocument();
const mints = storage.mintCount;
let pending: Promise<unknown> | undefined;
const commit = session.commit((tx) => {
pending = tx.doc(LateDoc);
}, context);
await gate.entered;
gate.release();
await expect(commit).rejects.toThrow("pending Tx operations");
await expect(pending).rejects.toThrow("Transaction has settled");
expect(initialized).toBe(0);
expect(storage.mintCount).toBe(mints);
});
it("rejects with the callback error when it fails with a pending acquisition", async () => {
const { session, storage, conversationId } = await setupLive();
await session.unloadDocuments();
const gate = storage.holdFindDocument();
let pending: Promise<unknown> | undefined;
const commit = session.commit((tx) => {
pending = tx.doc(LiveDoc, conversationId);
throw new Error("callback failed");
}, context);
await gate.entered;
gate.release();
await expect(commit).rejects.toThrow("callback failed");
await expect(pending).rejects.toThrow("Transaction has settled");
});
it("rejects Tx use after the callback settles", async () => {
const { session, conversationId } = await setupLive();
let captured: Parameters<Parameters<typeof session.commit>[0]>[0] | undefined;
await session.commit((tx) => {
captured = tx;
}, context);
await expect(captured!.doc(LiveDoc, conversationId)).rejects.toThrow("Transaction has settled");
await expect(captured!.conversation(conversationId)).rejects.toThrow("Transaction has settled");
expect(() => captured!.setTask({} as never)).toThrow("Transaction has settled");
});
it("rejects tokens whose semantics or version disagree with the stored incarnation", async () => {
const { session, conversationId } = await setupLive();
await expect(session.snapshot(RewindableLiveDoc, conversationId, context)).rejects.toThrow(
"does not match the supplied definition semantics",
);
await expect(
session.commit(async (tx) => {
await tx.doc(RewindableLiveDoc, conversationId);
}, context),
).rejects.toThrow("does not match the supplied definition semantics");
await expect(
session.commit(async (tx) => {
await tx.doc(LiveDocV2, conversationId);
}, context),
).rejects.toThrow("requires migration from version 1");
});
it("rejects non-JSON initializer values and draft placements before Storage admission", async () => {
const { session, storage } = openTestSession();
const conversationId = await createConversation(session);
const DateDoc = defineDoc<JsonObject>({
kind: "test.date",
version: 1,
scope: "session",
initial: () => ({ at: new Date() }) as unknown as JsonObject,
});
await expect(
session.commit(async (tx) => {
await tx.doc(DateDoc);
}, context),
).rejects.toThrow("strict JSON");
const commits = storage.commits.length;
await expect(
session.commit(async (tx) => {
const live = await tx.doc(LiveDoc, conversationId);
live.items.push(undefined as unknown as string);
}, context),
).rejects.toThrow("strict JSON");
expect(storage.commits.length).toBe(commits);
expect(await session.snapshot(LiveDoc, conversationId, context)).toBeUndefined();
});
it("rolls back prepared documents when batch assembly fails", async () => {
const { session, storage, conversationId } = await setupLive();
await session.commit(async (tx) => {
(await tx.doc(CounterDoc)).count = 1;
}, context);
const live = await session.snapshot(LiveDoc, conversationId, context);
const counter = await session.snapshot(CounterDoc, context);
const commits = storage.commits.length;
await expect(
session.commit(async (tx) => {
(await tx.doc(LiveDoc, conversationId)).message = "lost";
(await tx.doc(CounterDoc)).count = 2;
// Replacing a missing task fails during assembly, after every change was prepared.
tx.setTask({
id: 999,
conversationId,
kind: "missing",
version: 1,
input: null,
after: [],
background: false,
abortRequested: false,
state: { status: "pending", checkpoint: { phase: "start" } },
});
}, context),
).rejects.toThrow("Task 999 does not exist");
expect(storage.commits.length).toBe(commits);
expect(await session.snapshot(LiveDoc, conversationId, context)).toBe(live);
expect(await session.snapshot(CounterDoc, context)).toBe(counter);
await session.commit(async (tx) => {
(await tx.doc(LiveDoc, conversationId)).message = "next";
(await tx.doc(CounterDoc)).count = 3;
}, context);
expect((await session.snapshot(CounterDoc, context))!.count).toBe(3);
});
it("poisons the Session after an uncertain Storage failure and publishes nothing", async () => {
const { session, storage, publications, conversationId } = await setupLive();
await flush();
const before = (await session.snapshot(LiveDoc, conversationId, context))!;
const published = publications.length;
storage.failNextCommit(new Error("disk vanished"));
await expect(
session.commit(async (tx) => {
(await tx.doc(LiveDoc, conversationId)).message = "uncertain";
}, context),
).rejects.toThrow("disk vanished");
await flush();
expect(publications.length).toBe(published);
expect(before.message).toBeUndefined();
await expect(session.snapshot(LiveDoc, conversationId, context)).rejects.toThrow("poisoned");
await expect(session.commit(() => undefined, context)).rejects.toThrow("poisoned");
await session.close(context);
});
it("keeps the previous revision unchanged through Storage settlement", async () => {
const { session, storage, conversationId } = await setupLive();
const before = (await session.snapshot(LiveDoc, conversationId, context))!;
const copy = structuredClone(before);
const gate = storage.holdCommits();
const commit = session.commit(async (tx) => {
const live = await tx.doc(LiveDoc, conversationId);
live.items.push("c");
live.nested.count = 9;
}, context);
await gate.entered;
expect(await session.snapshot(LiveDoc, conversationId, context)).toBe(before);
expect(before).toEqual(copy);
gate.release();
await commit;
const after = (await session.snapshot(LiveDoc, conversationId, context))!;
expect(after.items).toEqual(["a", "b", "c"]);
expect(before).toEqual(copy);
});
it("retires documents and creates a new incarnation at the same address", async () => {
const { session, storage, publications, conversationId } = await setupLive();
await flush();
const oldId = documentChanges(publications.at(-1)!)[0]!.record.id;
let replacement: Draft<Live> | undefined;
await session.commit(async (tx) => {
const live = await tx.doc(LiveDoc, conversationId);
live.message = "final";
await tx.retireDoc(LiveDoc, conversationId);
replacement = await tx.doc(LiveDoc, conversationId);
expect(replacement).not.toBe(live);
replacement.message = "new";
}, context);
const writes = storage.commits.at(-1)!;
expect(writes).toHaveLength(3);
expect(writes).toContainEqual(expect.objectContaining({ type: "document.change", id: oldId }));
expect(writes).toContainEqual({ type: "document.retire", id: oldId });
expect(writes).toContainEqual(
expect.objectContaining({ type: "document.create", record: expect.objectContaining({ kind: "test.live" }) }),
);
await flush();
const publication = publications.at(-1)!;
const [retired, created] = documentChanges(publication);
expect(retired).toMatchObject({ record: { id: oldId }, value: null, ops: [] });
expect(retired!.record.retiredAt).toBe(publication.seq);
expect(created!.record.id).not.toBe(oldId);
expect(created!.record.createdAt).toBe(publication.seq);
expect(created).toMatchObject({ value: { message: "new", items: [] }, ops: [] });
expect(await session.snapshot(LiveDoc, conversationId, context)).toBe(created!.value);
await session.commit((tx) => tx.retireDoc(LiveDoc, conversationId), context);
expect(await session.snapshot(LiveDoc, conversationId, context)).toBeUndefined();
await session.unloadDocuments();
expect(await session.snapshot(LiveDoc, conversationId, context)).toBeUndefined();
// Retiring an absent address is a no-op.
const commits = storage.commits.length;
await session.commit((tx) => tx.retireDoc(LiveDoc, conversationId), context);
expect(storage.commits.length).toBe(commits);
});
it("retires without acquisition and recreates both existing and absent addresses", async () => {
const { session, storage, publications, conversationId } = await setupLive();
await flush();
const oldId = documentChanges(publications.at(-1)!)[0]!.record.id;
await session.unloadDocuments();
await session.commit(async (tx) => {
const retired = tx.retireDoc(LiveDoc, conversationId);
const replacement = tx.doc(LiveDoc, conversationId);
const [live] = await Promise.all([replacement, retired]);
live.message = "replacement";
}, context);
const existingWrites = storage.commits.at(-1)!;
expect(existingWrites).toHaveLength(2);
expect(existingWrites).toContainEqual({ type: "document.retire", id: oldId });
expect(existingWrites).toContainEqual(
expect.objectContaining({ type: "document.create", record: expect.objectContaining({ kind: "test.live" }) }),
);
await session.commit(async (tx) => {
const retired = tx.retireDoc(MemberDoc, "absent");
const replacement = tx.doc(MemberDoc, "absent", "seed");
const [member] = await Promise.all([replacement, retired]);
member.hits = 1;
}, context);
const absentWrites = storage.commits.at(-1)!;
expect(absentWrites.filter((write) => write.type === "document.retire")).toHaveLength(0);
expect(absentWrites.filter((write) => write.type === "document.create")).toHaveLength(1);
expect(await session.snapshot(MemberDoc, "absent", context)).toEqual({ seed: "seed", hits: 1 });
});
it("retires the existing incarnation when retirement races a pending acquisition", async () => {
const { session, storage, publications, conversationId } = await setupLive();
await flush();
const oldId = documentChanges(publications.at(-1)!)[0]!.record.id;
await session.commit(async (tx) => {
const acquired = tx.doc(LiveDoc, conversationId);
const retired = tx.retireDoc(LiveDoc, conversationId);
(await acquired).message = "final";
await retired;
}, context);
const firstWrites = storage.commits.at(-1)!;
expect(firstWrites).toHaveLength(2);
expect(firstWrites).toContainEqual({
type: "document.change",
id: oldId,
content: expect.objectContaining({ kind: "delta" }),
});
expect(firstWrites).toContainEqual({ type: "document.retire", id: oldId });
expect(await session.snapshot(LiveDoc, conversationId, context)).toBeUndefined();
await session.commit(async (tx) => {
(await tx.doc(LiveDoc, conversationId)).message = "second";
}, context);
await flush();
const secondId = documentChanges(publications.at(-1)!)[0]!.record.id;
await session.commit(async (tx) => {
const acquired = tx.doc(LiveDoc, conversationId);
const retired = tx.retireDoc(LiveDoc, conversationId);
const recreated = tx.doc(LiveDoc, conversationId);
const [old, fresh] = await Promise.all([acquired, recreated, retired]);
expect(fresh).not.toBe(old);
fresh!.message = "third";
}, context);
const writes = storage.commits.at(-1)!;
expect(writes).toHaveLength(2);
expect(writes).toContainEqual({ type: "document.retire", id: secondId });
expect(writes).toContainEqual(
expect.objectContaining({ type: "document.create", record: expect.objectContaining({ kind: "test.live" }) }),
);
expect((await session.snapshot(LiveDoc, conversationId, context))!.message).toBe("third");
});
it("reloads an unloaded document from Storage", async () => {
const { session, conversationId } = await setupLive();
const loaded = (await session.snapshot(LiveDoc, conversationId, context))!;
await session.unloadDocuments();
const reloaded = (await session.snapshot(LiveDoc, conversationId, context))!;
expect(reloaded).not.toBe(loaded);
expect(reloaded).toEqual(loaded);
await session.commit(async (tx) => {
(await tx.doc(LiveDoc, conversationId)).items.push("c");
}, context);
await session.unloadDocuments();
expect((await session.snapshot(LiveDoc, conversationId, context))!.items).toEqual(["a", "b", "c"]);
});
it("delivers publications off the mutation line so listeners can start a nested commit", async () => {
const { session, publications, conversationId } = await setupLive();
await flush();
const published = publications.length;
let listenerContext: typeof context | undefined;
const nested = new Promise<void>((resolve, reject) => {
const unsubscribe = session.subscribeCommits((_publication, deliveredContext) => {
unsubscribe();
listenerContext = deliveredContext;
void session
.commit(async (tx) => {
(await tx.doc(CounterDoc)).count = 1;
}, context)
.then(resolve, reject);
});
});
const result = await session.commit(async (tx) => {
(await tx.doc(LiveDoc, conversationId)).message = "m";
return "done";
}, context);
expect(result).toBe("done");
await nested;
expect(listenerContext).toBe(context);
expect(await session.snapshot(CounterDoc, context)).toEqual({ count: 1 });
await flush();
expect(publications.length).toBe(published + 2);
});
it("settles admitted commits before close and rejects later admission", async () => {
const { session, storage, conversationId } = await setupLive();
const gate = storage.holdCommits();
const commit = session.commit(async (tx) => {
(await tx.doc(LiveDoc, conversationId)).message = "admitted";
}, context);
await gate.entered;
const queued = session.commit(async (tx) => {
(await tx.doc(LiveDoc, conversationId)).message = "queued";
}, context);
const admittedSnapshot = session.snapshot(MemberDoc, "absent", context);
const closed = session.close(context);
await expect(session.commit(() => undefined, context)).rejects.toThrow("closed");
await expect(session.snapshot(LiveDoc, conversationId, context)).rejects.toThrow("closed");
gate.release();
await commit;
await queued;
expect(await admittedSnapshot).toBeUndefined();
await closed;
const stored = await storage
.findDocument({ kind: "test.live", scope: { kind: "conversation", conversationId } }, "current", context)
.catch((error: unknown) => error);
expect(stored).toBeInstanceOf(Error);
});
});
+129
View File
@@ -0,0 +1,129 @@
import type { Context } from "@earendil-works/chord";
import { BACKGROUND_CONTEXT } from "@earendil-works/chord/context";
import {
type DocumentAddress,
type DocumentPoint,
type DocumentRecord,
type Id,
MemoryStorage,
type Seq,
type StorageWrite,
} from "@earendil-works/pi-durable";
import type { CommitPublication } from "../src/session/publications.ts";
import { SessionKernel } from "../src/session/session.ts";
import type { DocumentCommitChange } from "../src/session/transaction.ts";
export const context: Context = BACKGROUND_CONTEXT;
type Deferred = { readonly promise: Promise<void>; readonly resolve: () => void };
function deferred(): Deferred {
let resolve!: () => void;
const promise = new Promise<void>((done) => {
resolve = done;
});
return { promise, resolve };
}
/** A gate that holds calls until released and reports when the first held call arrives. */
export type Gate = {
readonly entered: Promise<void>;
release(): void;
};
/** Memory storage with observable commits, held calls, and injected commit failures. */
export class ControlledStorage extends MemoryStorage {
/** Exact borrowed batches admitted by Session. */
readonly admittedCommits: (readonly StorageWrite[])[] = [];
/** Detached batches for value assertions. */
readonly commits: (readonly StorageWrite[])[] = [];
mintCount = 0;
#commitGate: { gate: Deferred; entered: Deferred } | undefined;
#findGate: { gate: Deferred; entered: Deferred } | undefined;
#commitFailure: Error | undefined;
holdCommits(): Gate {
const held = { gate: deferred(), entered: deferred() };
this.#commitGate = held;
return { entered: held.entered.promise, release: () => this.#release("commit", held) };
}
holdFindDocument(): Gate {
const held = { gate: deferred(), entered: deferred() };
this.#findGate = held;
return { entered: held.entered.promise, release: () => this.#release("find", held) };
}
failNextCommit(error: Error): void {
this.#commitFailure = error;
}
#release(kind: "commit" | "find", held: { gate: Deferred }): void {
if (kind === "commit" && this.#commitGate === held) this.#commitGate = undefined;
if (kind === "find" && this.#findGate === held) this.#findGate = undefined;
held.gate.resolve();
}
override async commit(writes: readonly StorageWrite[], commitContext: Context): Promise<Seq> {
this.admittedCommits.push(writes);
this.commits.push(structuredClone(writes));
const held = this.#commitGate;
if (held !== undefined) {
held.entered.resolve();
await held.gate.promise;
}
const failure = this.#commitFailure;
if (failure !== undefined) {
this.#commitFailure = undefined;
throw failure;
}
return super.commit(writes, commitContext);
}
override mintId(): Promise<Id> {
this.mintCount++;
return super.mintId();
}
override async findDocument(
address: DocumentAddress,
at: DocumentPoint,
callContext: Context,
): Promise<DocumentRecord | undefined> {
const held = this.#findGate;
if (held !== undefined) {
held.entered.resolve();
await held.gate.promise;
}
return super.findDocument(address, at, callContext);
}
}
/** Session kernel plus its controlled storage and every committed publication. */
export function openTestSession(): {
readonly storage: ControlledStorage;
readonly session: SessionKernel;
readonly publications: CommitPublication[];
} {
const storage = new ControlledStorage();
const session = new SessionKernel(storage);
const publications: CommitPublication[] = [];
session.subscribeCommits((publication) => {
publications.push(publication);
});
return { storage, session, publications };
}
export function documentChanges(publication: CommitPublication): readonly DocumentCommitChange[] {
return publication.changes.filter((change): change is DocumentCommitChange => change.type === "document");
}
/** Create one conversation and return its ID. */
export async function createConversation(session: SessionKernel): Promise<Id> {
return session.commit(async (tx) => (await tx.createConversation({})).id, context);
}
/** Resolve after pending microtasks and one macrotask turn. */
export function flush(): Promise<void> {
return new Promise((resolve) => setTimeout(resolve, 0));
}
@@ -0,0 +1,328 @@
import type { JsonValue } from "@earendil-works/chord";
import {
defineDoc,
defineDocFamily,
type Id,
ReadAfterWrite,
type Task,
type TaskRecord,
} from "@earendil-works/pi-durable";
import { describe, expect, it } from "vitest";
import { context, createConversation, documentChanges, flush, openTestSession } from "./session-support.ts";
type Checkpoint = { phase: "start" } | { phase: "next"; step: number };
const WorkTask: Task<{ path: string }, Checkpoint, { ok: boolean }, object> = {
definition: { name: "test.work", version: 1, initial: () => ({ phase: "start" }) },
};
type Progress = { lines: string[] };
const ProgressDoc = defineDoc<Progress>({
kind: "test.progress",
version: 1,
scope: "task",
initial: () => ({ lines: [] }),
});
const StepDoc = defineDocFamily<Progress, null>({
kind: "test.step",
version: 1,
family: true,
scope: "task",
initial: () => ({ lines: [] }),
});
type Notes = { text: string };
const NotesDoc = defineDoc<Notes>({
kind: "test.notes",
version: 1,
scope: "conversation",
history: "rewindable",
fork: "asOf",
initial: () => ({ text: "" }),
});
type AnyTask = TaskRecord<JsonValue, JsonValue, JsonValue>;
function terminal(task: AnyTask): AnyTask {
return {
id: task.id,
conversationId: task.conversationId,
kind: task.kind,
version: task.version,
input: task.input,
after: task.after,
background: task.background,
abortRequested: task.abortRequested,
state: { status: "terminal", outcome: { status: "completed", result: { ok: true } } },
};
}
async function createTask(
session: ReturnType<typeof openTestSession>["session"],
conversationId: Id,
withDocument = false,
): Promise<Id> {
return session.commit(async (tx) => {
const ref = await tx.createTask(WorkTask, { path: "a" }, { conversationId });
if (withDocument) (await tx.doc(ProgressDoc, ref.id)).lines.push("started");
return ref.id;
}, context);
}
describe("Session transaction tables", () => {
it("allows table reads only before the first table write", async () => {
const { session } = openTestSession();
const conversationId = await createConversation(session);
await session.commit(async (tx) => {
expect(await tx.conversation(conversationId)).toEqual({ id: conversationId });
expect(await tx.scanConversations(1)).toEqual({ items: [{ id: conversationId }] });
expect(await tx.scanTasks({ conversationId }, 10)).toEqual({ items: [] });
expect(await tx.scanEntries({ conversationId }, 10)).toEqual({ items: [] });
await tx.appendEntry(conversationId, { kind: "note" });
await expect(tx.conversation(conversationId)).rejects.toBeInstanceOf(ReadAfterWrite);
await expect(tx.task(1)).rejects.toThrow("Tx.task() cannot read tables after the first table write");
await expect(tx.entry(1)).rejects.toBeInstanceOf(ReadAfterWrite);
await expect(tx.scanConversations(10)).rejects.toBeInstanceOf(ReadAfterWrite);
await expect(tx.scanEntries({ conversationId }, 10)).rejects.toBeInstanceOf(ReadAfterWrite);
// Document access remains available after table writes.
(await tx.doc(NotesDoc, conversationId)).text = "after write";
}, context);
expect(await session.snapshot(NotesDoc, conversationId, context)).toEqual({ text: "after write" });
});
it("passes caller-selected limits and cursors through table scans", async () => {
const { session } = openTestSession();
const ids = [
await createConversation(session),
await createConversation(session),
await createConversation(session),
];
await session.commit(async (tx) => {
const first = await tx.scanConversations(2);
expect(first.items.map(({ id }) => id)).toEqual(ids.slice(0, 2));
expect(first.next).toBeDefined();
const second = await tx.scanConversations(2, first.next);
expect(second.items.map(({ id }) => id)).toEqual(ids.slice(2));
expect(second.next).toBeUndefined();
}, context);
});
it("treats synchronous setTask as the first table write", async () => {
const { session } = openTestSession();
const conversationId = await createConversation(session);
const taskId = await createTask(session, conversationId);
await session.commit(async (tx) => {
const task = (await tx.task(taskId))!;
tx.setTask(task);
await expect(tx.task(taskId)).rejects.toBeInstanceOf(ReadAfterWrite);
}, context);
});
it("creates conversations, entries, and tasks with minted IDs", async () => {
const { session, storage, publications } = openTestSession();
const created = await session.commit(async (tx) => {
const conversation = await tx.createConversation({});
const first = await tx.appendEntry(conversation.id, { kind: "note", data: "one" });
const headed = await tx.appendEntry(conversation.id, { kind: "summary", head: "self" });
const task = await tx.createTask(
WorkTask,
{ path: "x" },
{ conversationId: conversation.id, background: true },
);
return { conversation, first, headed, task };
}, context);
const ids = [created.conversation.id, created.first.id, created.headed.id, created.task.id];
expect(new Set(ids).size).toBe(4);
expect(created.headed.head).toBe(created.headed.id);
expect(created.first).toEqual({
id: created.first.id,
conversationId: created.conversation.id,
kind: "note",
data: "one",
});
expect((await storage.entry(created.headed.id, context))!.entry).toEqual(created.headed);
expect(await storage.task(created.task.id, context)).toEqual({
id: created.task.id,
conversationId: created.conversation.id,
kind: "test.work",
version: 1,
input: { path: "x" },
after: [],
background: true,
abortRequested: false,
state: { status: "pending", checkpoint: { phase: "start" } },
});
await flush();
const changes = publications.at(-1)!.changes;
expect(changes).toHaveLength(4);
const admitted = storage.admittedCommits.at(-1)!;
for (const change of changes) expect(admitted.some((write) => write === change)).toBe(true);
expect(changes.map((change) => change.type)).toEqual(
expect.arrayContaining(["conversation", "entry", "entry", "task"]),
);
expect(changes.find((change) => change.type === "conversation")!.value).toBe(created.conversation);
const entries = changes.filter((change) => change.type === "entry");
expect(entries.map((change) => change.value)).toContain(created.first);
expect(entries.map((change) => change.value)).toContain(created.headed);
await expect(session.commit((tx) => tx.createTask(WorkTask, { path: "x" }), context)).rejects.toThrow(
"requires options.conversationId",
);
await expect(session.commit((tx) => tx.appendEntry(12345, { kind: "note" }), context)).rejects.toThrow(
"Conversation 12345 does not exist",
);
});
it("takes ownership of table JSON and rejects non-strict values", async () => {
const { session, storage } = openTestSession();
const conversationId = await createConversation(session);
const payload = { nested: { value: 1 } };
const entry = await session.commit(async (tx) => {
const created = await tx.appendEntry(conversationId, { kind: "data", data: payload });
payload.nested.value = 2;
return created;
}, context);
expect((await storage.entry(entry.id, context))!.entry.data).toEqual({ nested: { value: 1 } });
const omitted = await session.commit(
(tx) => tx.appendEntry(conversationId, { kind: "omitted", data: undefined }),
context,
);
expect(Object.hasOwn(omitted, "data")).toBe(false);
expect(Object.hasOwn((await storage.entry(omitted.id, context))!.entry, "data")).toBe(false);
const commits = storage.commits.length;
await expect(
session.commit(
(tx) => tx.appendEntry(conversationId, { kind: "invalid", data: Number.NaN }).then(() => undefined),
context,
),
).rejects.toThrow("strict JSON");
expect(storage.commits.length).toBe(commits);
});
it("replaces task records completely", async () => {
const { session, storage } = openTestSession();
const conversationId = await createConversation(session);
const taskId = await createTask(session, conversationId);
await session.commit(async (tx) => {
const task = (await tx.task(taskId))!;
tx.setTask({
id: task.id,
conversationId: task.conversationId,
kind: task.kind,
version: task.version,
input: task.input,
after: task.after,
background: task.background,
abortRequested: task.abortRequested,
state: { status: "running", checkpoint: { phase: "next", step: 2 } },
memos: { choice: "b" },
});
}, context);
expect(await storage.task(taskId, context)).toMatchObject({
state: { status: "running", checkpoint: { phase: "next", step: 2 } },
memos: { choice: "b" },
});
});
it("creates a task and then its document in one transaction without ReadAfterWrite", async () => {
const { session, publications } = openTestSession();
const conversationId = await createConversation(session);
const taskId = await session.commit(async (tx) => {
const ref = await tx.createTask(WorkTask, { path: "a" }, { conversationId });
// Validation uses the candidate task record, not a caller table read.
(await tx.doc(ProgressDoc, ref.id)).lines.push("created");
(await tx.doc(StepDoc, ref.id, "one", null)).lines.push("step");
return ref.id;
}, context);
expect(await session.snapshot(ProgressDoc, taskId, context)).toEqual({ lines: ["created"] });
await flush();
const publication = publications.at(-1)!;
const documents = documentChanges(publication);
expect(publication.changes).toContainEqual(
expect.objectContaining({ type: "task", value: expect.objectContaining({ id: taskId }) }),
);
expect(documents).toHaveLength(2);
// Task documents derive their conversation from the task record.
for (const document of documents) expect(document.conversationId).toBe(conversationId);
await session.commit(async (tx) => {
await tx.createConversation({});
(await tx.doc(ProgressDoc, taskId)).lines.push("committed task");
}, context);
await flush();
expect(documentChanges(publications.at(-1)!)[0]!.conversationId).toBe(conversationId);
});
it("rejects task documents after a terminal candidate", async () => {
const { session } = openTestSession();
const conversationId = await createConversation(session);
const taskId = await createTask(session, conversationId, true);
await session.commit(async (tx) => {
const task = (await tx.task(taskId))!;
const progress = await tx.doc(ProgressDoc, taskId);
tx.setTask(terminal(task));
await expect(tx.doc(ProgressDoc, taskId)).rejects.toThrow(`Task ${taskId} is terminal`);
await expect(tx.doc(StepDoc, taskId, "late", null)).rejects.toThrow(`Task ${taskId} is terminal`);
expect(() => tx.setTask(task)).toThrow("terminal candidate");
progress.lines.push("final");
}, context);
await expect(session.commit((tx) => tx.doc(ProgressDoc, taskId).then(() => undefined), context)).rejects.toThrow(
`Task ${taskId} is terminal`,
);
});
it("retires task documents at terminal settlement, including documents created in the same transaction", async () => {
const { session, storage, publications } = openTestSession();
const conversationId = await createConversation(session);
const taskId = await createTask(session, conversationId, true);
await session.commit(async (tx) => {
(await tx.doc(StepDoc, taskId, "committed", null)).lines.push("x");
}, context);
await flush();
const published = publications.length;
await session.commit(async (tx) => {
const task = (await tx.task(taskId))!;
(await tx.doc(StepDoc, taskId, "new", null)).lines.push("created then retired");
tx.setTask(terminal(task));
}, context);
const writes = storage.commits.at(-1)!;
expect(writes).toHaveLength(5);
expect(writes.filter((write) => write.type === "task")).toHaveLength(1);
const creations = writes.filter((write) => write.type === "document.create");
const retirements = writes.filter((write) => write.type === "document.retire");
expect(creations).toHaveLength(1);
expect(retirements).toHaveLength(3);
expect(retirements.map((write) => write.id)).toContain(creations[0]!.record.id);
await flush();
expect(publications.length).toBe(published + 1);
const publication = publications.at(-1)!;
const documents = documentChanges(publication);
expect(publication.changes.filter((change) => change.type === "task")).toHaveLength(1);
expect(documents.map((document) => document.value)).toEqual([null, null, null]);
for (const document of documents) expect(document.ops).toEqual([]);
for (const document of documents) expect(document.conversationId).toBe(conversationId);
expect(await session.snapshot(ProgressDoc, taskId, context)).toBeUndefined();
expect(await session.snapshot(StepDoc, taskId, "committed", context)).toBeUndefined();
const alive = await storage.scanDocuments(
{ scope: { kind: "task", taskId }, at: "current" },
10,
undefined,
context,
);
expect(alive.items).toEqual([]);
await expect(
session.commit(async (tx) => {
tx.setTask(terminal((await tx.task(taskId))!));
}, context),
).rejects.toThrow(`Task ${taskId} is already terminal`);
});
it("validates document owners", async () => {
const { session } = openTestSession();
await expect(session.commit((tx) => tx.doc(ProgressDoc, 4242).then(() => undefined), context)).rejects.toThrow(
"Task 4242 does not exist",
);
await expect(session.commit((tx) => tx.doc(NotesDoc, 4242).then(() => undefined), context)).rejects.toThrow(
"Conversation 4242 does not exist",
);
});
});
+8 -8
View File
@@ -55,24 +55,24 @@ class ReopeningStorage implements Storage {
mintId: Storage["mintId"] = () => this.current.mintId();
conversation: Storage["conversation"] = (id, readContext) => this.current.conversation(id, readContext);
scanConversations: Storage["scanConversations"] = (cursor, limit, readContext) =>
this.current.scanConversations(cursor, limit, readContext);
scanConversations: Storage["scanConversations"] = (limit, cursor, readContext) =>
this.current.scanConversations(limit, cursor, readContext);
entry: Storage["entry"] = (id, readContext) => this.current.entry(id, readContext);
findLatestHeadMarker: Storage["findLatestHeadMarker"] = (conversationId, at, readContext) =>
this.current.findLatestHeadMarker(conversationId, at, readContext);
scanEntries: Storage["scanEntries"] = (query, cursor, limit, readContext) =>
this.current.scanEntries(query, cursor, limit, readContext);
scanEntries: Storage["scanEntries"] = (query, limit, cursor, readContext) =>
this.current.scanEntries(query, limit, cursor, readContext);
task: Storage["task"] = (id, readContext) => this.current.task(id, readContext);
scanTasks: Storage["scanTasks"] = (query, cursor, limit, readContext) =>
this.current.scanTasks(query, cursor, limit, readContext);
scanTasks: Storage["scanTasks"] = (query, limit, cursor, readContext) =>
this.current.scanTasks(query, limit, cursor, readContext);
submission: Storage["submission"] = (id, readContext) => this.current.submission(id, readContext);
submissionByRequest: Storage["submissionByRequest"] = (conversationId, requestId, readContext) =>
this.current.submissionByRequest(conversationId, requestId, readContext);
findDocument: Storage["findDocument"] = (address, at, readContext) =>
this.current.findDocument(address, at, readContext);
document: Storage["document"] = (id, at, readContext) => this.current.document(id, at, readContext);
scanDocuments: Storage["scanDocuments"] = (query, cursor, limit, readContext) =>
this.current.scanDocuments(query, cursor, limit, readContext);
scanDocuments: Storage["scanDocuments"] = (query, limit, cursor, readContext) =>
this.current.scanDocuments(query, limit, cursor, readContext);
async close(closeContext: Parameters<Storage["close"]>[0]): Promise<void> {
if (this.closed) return;