Merge PR #10232: feat(durable): make SQLite storage asynchronous

This commit is contained in:
Mario Zechner
2026-09-30 21:36:10 +02:00
8 changed files with 548 additions and 343 deletions
+10 -1
View File
@@ -417,7 +417,16 @@ const usage = await harness.usage(context); // { models: { "openai/gpt-6-sol": U
| SQLite | `openNodeSqliteStorage(file)` from `@earendil-works/pi-durable/storage/sqlite/node` | One database file. WAL mode with `synchronous = NORMAL`: commits survive process crashes; the newest may be lost on power or host failure. |
| JSONL | `openNodeJsonlStorage(directory, context)` from `@earendil-works/pi-durable/storage/jsonl/node` | Append-only files in one directory. Pass `{ fsync: true }` to flush before each commit marker. |
One process owns a storage at a time; there is no cross-process locking. The portable SQLite and JSONL cores (`/storage/sqlite`, `/storage/jsonl`) run without Node APIs, for example on Bun or in Cloudflare Durable Objects, given a synchronous SQLite database or a `FileSystem` from `@earendil-works/pi-durable/env`.
One process owns a storage at a time; there is no cross-process locking. The portable SQLite and JSONL cores (`/storage/sqlite`, `/storage/jsonl`) run without Node APIs, for example on Bun or in Cloudflare Durable Objects, given an asynchronous `SqliteDatabase` facade or a `FileSystem` from `@earendil-works/pi-durable/env`.
SQLite adapters implement promise-based `exec`, `run`, `get`, `all`, `transaction`, and `close`. `run`, `get`, and `all` take SQL text plus positional bindings; adapters may cache prepared statements by SQL text. A transaction callback receives a transaction handle; all work in the transaction must use it, and the handle expires when the callback settles. Adapters must queue unrelated operations and other transactions until the transaction finishes:
```typescript
await database.transaction(async (transaction) => {
await transaction.exec("CREATE TABLE example (value TEXT)");
await transaction.run("INSERT INTO example (value) VALUES (?)", "stored atomically");
});
```
Custom backends can run the shared conformance suite with any Vitest- or Jest-compatible runner:
+24 -19
View File
@@ -2,30 +2,35 @@
export type SqliteValue = null | number | bigint | string | Uint8Array;
/**
* A prepared synchronous SQLite statement.
* Implementations must support repeated execution with new bindings across transactions.
* Asynchronous SQL operations shared by a database and its transaction handles.
*
* `exec` runs SQL text without bindings and may contain several statements. `run`, `get`, and `all`
* execute one statement with positional bindings. Adapters may cache prepared statements by SQL text,
* so callers pass values as bindings instead of interpolating them.
*/
export interface SqliteStatement {
run(...params: SqliteValue[]): void;
get<T extends object>(...params: SqliteValue[]): T | undefined;
all<T extends object>(...params: SqliteValue[]): T[];
export interface SqliteExecutor {
exec(sql: string): Promise<void>;
run(sql: string, ...params: SqliteValue[]): Promise<void>;
get<T extends object>(sql: string, ...params: SqliteValue[]): Promise<T | undefined>;
all<T extends object>(sql: string, ...params: SqliteValue[]): Promise<T[]>;
}
/**
* Minimal database facade required by `SqliteStorage`.
*
* Queries and transaction callbacks are synchronous so the same storage core can
* run on Node, Bun, and Cloudflare Durable Object SQLite. An adapter may return a
* promise from `transaction` while it waits for the transaction to settle.
* When the callback throws, the adapter must roll the transaction back before
* rethrowing that same error. If rollback fails, it must throw a different error
* (for example an `AggregateError`) so callers cannot mistake the callback error
* for a guaranteed rollback. Callers must not close the database while a
* returned settlement is pending.
* All operations are asynchronous so adapters may execute outside the harness runtime.
*
* `transaction` passes the callback a transaction handle. All work in the transaction
* must use that handle; the handle is invalid after the callback settles. Adapters must
* queue unrelated operations and other transactions until the transaction finishes. The
* returned promise settles after commit or rollback.
*
* When the callback rejects, the adapter must roll the transaction back before rejecting
* with that same error. If rollback fails, it must reject with a different error (for
* example an `AggregateError`) so callers cannot mistake the callback error for a
* guaranteed rollback.
*/
export interface SqliteDatabase {
exec(sql: string): void;
prepare(sql: string): SqliteStatement;
transaction<T>(callback: () => T): T | Promise<T>;
close(): void | Promise<void>;
export interface SqliteDatabase extends SqliteExecutor {
transaction<T>(callback: (transaction: SqliteExecutor) => Promise<T>): Promise<T>;
close(): Promise<void>;
}
+1 -1
View File
@@ -1,4 +1,4 @@
export type { SqliteDatabase, SqliteStatement, SqliteValue } from "./database.ts";
export type { SqliteDatabase, SqliteExecutor, SqliteValue } from "./database.ts";
export {
applySqliteMigrations,
CURRENT_SQLITE_SCHEMA_VERSION,
@@ -102,13 +102,13 @@ export async function applySqliteMigrations(
}
}
await database.transaction(() => {
database.exec(`CREATE TABLE IF NOT EXISTS durable_schema (
await database.transaction(async (transaction) => {
await transaction.exec(`CREATE TABLE IF NOT EXISTS durable_schema (
singleton INTEGER PRIMARY KEY CHECK (singleton = 1),
version INTEGER NOT NULL CHECK (version >= 0)
) STRICT`);
database.prepare("INSERT OR IGNORE INTO durable_schema (singleton, version) VALUES (1, 0)").run();
const row = database.prepare("SELECT version FROM durable_schema WHERE singleton = 1").get<SchemaRow>();
await transaction.run("INSERT OR IGNORE INTO durable_schema (singleton, version) VALUES (1, 0)");
const row = await transaction.get<SchemaRow>("SELECT version FROM durable_schema WHERE singleton = 1");
if (row === undefined) throw new Error("Durable SQLite schema metadata is missing");
const currentVersion = migrations.at(-1)?.version ?? 0;
if (row.version > currentVersion) {
@@ -118,8 +118,8 @@ export async function applySqliteMigrations(
}
for (const migration of migrations) {
if (migration.version <= row.version) continue;
for (const statement of migration.statements) database.exec(statement);
database.prepare("UPDATE durable_schema SET version = ? WHERE singleton = 1").run(migration.version);
for (const statement of migration.statements) await transaction.exec(statement);
await transaction.run("UPDATE durable_schema SET version = ? WHERE singleton = 1", migration.version);
}
});
}
+125 -53
View File
@@ -1,8 +1,9 @@
import { AsyncLocalStorage } from "node:async_hooks";
import { mkdir } from "node:fs/promises";
import { dirname } from "node:path";
import type { SQLInputValue, StatementSync } from "node:sqlite";
import type { StatementSync } from "node:sqlite";
import { DatabaseSync } from "node:sqlite";
import type { SqliteDatabase, SqliteStatement, SqliteValue } from "./database.ts";
import type { SqliteDatabase, SqliteExecutor, SqliteValue } from "./database.ts";
import { SqliteStorage } from "./storage.ts";
/** Node SQLite connection settings for a durable storage file. */
@@ -16,74 +17,145 @@ export type NodeSqliteStorageOptions = {
const DEFAULT_WAL_AUTO_CHECKPOINT_PAGES = 1_000;
const DEFAULT_BUSY_TIMEOUT_MS = 5_000;
class NodeSqliteStatement implements SqliteStatement {
private readonly statement: StatementSync;
type TransactionScope = { active: boolean };
constructor(statement: StatementSync) {
this.statement = statement;
class SerialOperationQueue {
private tail = Promise.resolve();
async run<T>(operation: () => T | Promise<T>): Promise<T> {
const previous = this.tail;
const { promise, resolve: release } = Promise.withResolvers<void>();
this.tail = promise;
await previous;
try {
return await operation();
} finally {
release();
}
}
}
/**
* Executes SQL on one connection. Prepared statements are cached per connection by SQL text, so the
* database and its transaction handles share them across transactions.
*/
abstract class NodeSqliteExecutor implements SqliteExecutor {
protected readonly database: DatabaseSync;
protected readonly statements: Map<string, StatementSync>;
constructor(database: DatabaseSync, statements: Map<string, StatementSync>) {
this.database = database;
this.statements = statements;
}
run(...params: SqliteValue[]): void {
this.statement.run(...(params as SQLInputValue[]));
exec(sql: string): Promise<void> {
return this.runOperation(() => {
this.database.exec(sql);
});
}
get<T extends object>(...params: SqliteValue[]): T | undefined {
return this.statement.get(...(params as SQLInputValue[])) as T | undefined;
run(sql: string, ...params: SqliteValue[]): Promise<void> {
return this.runOperation(() => {
this.statement(sql).run(...params);
});
}
all<T extends object>(...params: SqliteValue[]): T[] {
return this.statement.all(...(params as SQLInputValue[])) as T[];
get<T extends object>(sql: string, ...params: SqliteValue[]): Promise<T | undefined> {
return this.runOperation(() => this.statement(sql).get(...params) as T | undefined);
}
all<T extends object>(sql: string, ...params: SqliteValue[]): Promise<T[]> {
return this.runOperation(() => this.statement(sql).all(...params) as T[]);
}
protected abstract runOperation<T>(operation: () => T): Promise<T>;
private statement(sql: string): StatementSync {
let statement = this.statements.get(sql);
if (statement === undefined) {
statement = this.database.prepare(sql);
this.statements.set(sql, statement);
}
return statement;
}
}
class NodeSqliteTransaction extends NodeSqliteExecutor {
private readonly scope: TransactionScope;
constructor(database: DatabaseSync, statements: Map<string, StatementSync>, scope: TransactionScope) {
super(database, statements);
this.scope = scope;
}
protected async runOperation<T>(operation: () => T): Promise<T> {
if (!this.scope.active) throw new Error("SQLite transaction handle is no longer active");
return operation();
}
}
/** `SqliteDatabase` adapter backed by Node's built-in `node:sqlite`. */
export class NodeSqliteDatabase implements SqliteDatabase {
private readonly database: DatabaseSync;
export class NodeSqliteDatabase extends NodeSqliteExecutor implements SqliteDatabase {
private readonly access = new SerialOperationQueue();
/** Detects database calls from inside a transaction callback, which would otherwise wait for that transaction forever. */
private readonly transactionScope = new AsyncLocalStorage<TransactionScope>();
private closed = false;
constructor(database: DatabaseSync) {
this.database = database;
super(database, new Map());
}
exec(sql: string): void {
this.database.exec(sql);
}
prepare(sql: string): SqliteStatement {
return new NodeSqliteStatement(this.database.prepare(sql));
}
transaction<T>(callback: () => T): T {
this.database.exec("BEGIN IMMEDIATE");
try {
const result = callback();
if (
result !== null &&
(typeof result === "object" || typeof result === "function") &&
typeof Reflect.get(result, "then") === "function"
) {
throw new TypeError("SQLite transaction callbacks must be synchronous");
}
this.database.exec("COMMIT");
return result;
} catch (error) {
transaction<T>(callback: (transaction: SqliteExecutor) => Promise<T>): Promise<T> {
if (this.insideTransaction()) {
return Promise.reject(new Error("Nested SQLite transactions are not supported"));
}
return this.access.run(async () => {
this.database.exec("BEGIN IMMEDIATE");
const scope = { active: true };
try {
this.database.exec("ROLLBACK");
} catch (rollbackError) {
throw new AggregateError([error, rollbackError], "SQLite transaction failed and rollback failed");
const result = await this.transactionScope.run(scope, () =>
callback(new NodeSqliteTransaction(this.database, this.statements, scope)),
);
scope.active = false;
this.database.exec("COMMIT");
return result;
} catch (error) {
scope.active = false;
try {
this.database.exec("ROLLBACK");
} catch (rollbackError) {
throw new AggregateError([error, rollbackError], "SQLite transaction failed and rollback failed");
}
throw error;
}
throw error;
}
});
}
close(): void {
if (this.closed) return;
this.closed = true;
try {
this.database.exec("PRAGMA wal_checkpoint(TRUNCATE)");
} finally {
this.database.close();
close(): Promise<void> {
if (this.insideTransaction()) {
return Promise.reject(new Error("Cannot close SQLite during an active transaction"));
}
return this.access.run(() => {
if (this.closed) return;
this.closed = true;
this.statements.clear();
try {
this.database.exec("PRAGMA wal_checkpoint(TRUNCATE)");
} finally {
this.database.close();
}
});
}
protected runOperation<T>(operation: () => T): Promise<T> {
if (this.insideTransaction()) {
return Promise.reject(new Error("Use the transaction handle inside a transaction callback"));
}
return this.access.run(operation);
}
private insideTransaction(): boolean {
return this.transactionScope.getStore()?.active === true;
}
}
@@ -98,13 +170,13 @@ export async function openNodeSqliteDatabase(
const database = new DatabaseSync(path, { timeout });
const adapter = new NodeSqliteDatabase(database);
try {
adapter.exec("PRAGMA journal_mode = WAL");
adapter.exec("PRAGMA synchronous = NORMAL");
adapter.exec(`PRAGMA wal_autocheckpoint = ${checkpointPages}`);
await adapter.exec("PRAGMA journal_mode = WAL");
await adapter.exec("PRAGMA synchronous = NORMAL");
await adapter.exec(`PRAGMA wal_autocheckpoint = ${checkpointPages}`);
return adapter;
} catch (error) {
try {
adapter.close();
await adapter.close();
} catch {
// Preserve the configuration failure.
}
+170 -187
View File
@@ -31,7 +31,7 @@ import type {
TaskQuery,
TaskRecord,
} from "../../types.ts";
import type { SqliteDatabase, SqliteStatement, SqliteValue } from "./database.ts";
import type { SqliteDatabase, SqliteExecutor, SqliteValue } from "./database.ts";
import { applySqliteMigrations } from "./migrations.ts";
type StoredTask = TaskRecord<JsonValue, JsonValue, JsonValue>;
@@ -63,12 +63,6 @@ const encodeJson = (value: unknown): string => JSON.stringify(value) as string;
// Some SQLite bindings replace lone UTF-16 surrogates. JSON encoding keeps indexed identities lossless.
const encodeIndexedString = (value: string): string => JSON.stringify(value);
const getRow = <T extends object>(statement: SqliteStatement, ...params: SqliteValue[]): T | undefined =>
statement.get<T>(...params);
const allRows = <T extends object>(statement: SqliteStatement, ...params: SqliteValue[]): T[] =>
statement.all<T>(...params);
const cursorId = <I extends Id<string>>(cursor: Cursor | undefined): I | undefined => {
const after = cursor?.after;
if (after === undefined) return undefined;
@@ -132,37 +126,6 @@ const writeId = (write: StorageWrite): Id<string> | undefined => {
}
};
class StatementCachingDatabase implements SqliteDatabase {
private readonly database: SqliteDatabase;
private readonly statements = new Map<string, SqliteStatement>();
constructor(database: SqliteDatabase) {
this.database = database;
}
exec(sql: string): void {
this.database.exec(sql);
}
prepare(sql: string): SqliteStatement {
let statement = this.statements.get(sql);
if (statement === undefined) {
statement = this.database.prepare(sql);
this.statements.set(sql, statement);
}
return statement;
}
transaction<T>(callback: () => T): T | Promise<T> {
return this.database.transaction(callback);
}
close(): void | Promise<void> {
this.statements.clear();
return this.database.close();
}
}
/** Portable SQLite implementation of the Pico storage contract. */
export class SqliteStorage implements Storage {
private readonly db: SqliteDatabase;
@@ -170,7 +133,7 @@ export class SqliteStorage implements Storage {
private closed = false;
private constructor(db: SqliteDatabase, nextId: number) {
this.db = new StatementCachingDatabase(db);
this.db = db;
this.nextId = nextId;
}
@@ -178,8 +141,8 @@ export class SqliteStorage implements Storage {
static async open(db: SqliteDatabase): Promise<SqliteStorage> {
try {
await applySqliteMigrations(db);
const metadata = getRow<MetadataRow>(
db.prepare("SELECT next_id, next_seq FROM durable_metadata WHERE singleton = 1"),
const metadata = await db.get<MetadataRow>(
"SELECT next_id, next_seq FROM durable_metadata WHERE singleton = 1",
);
if (metadata === undefined) throw new Error("Durable SQLite metadata is missing");
return new SqliteStorage(db, Number(metadata.next_id));
@@ -197,19 +160,21 @@ export class SqliteStorage implements Storage {
this.assertOpen();
const documentActions = this.prepareDocumentActions(writes);
const candidateNextId = this.candidateNextId(writes);
const seq = await this.db.transaction(() => {
const metadata = getRow<MetadataRow>(
this.db.prepare("SELECT next_id, next_seq FROM durable_metadata WHERE singleton = 1"),
const seq = await this.db.transaction(async (transaction) => {
const metadata = await transaction.get<MetadataRow>(
"SELECT next_id, next_seq FROM durable_metadata WHERE singleton = 1",
);
if (metadata === undefined) throw new Error("Durable SQLite metadata is missing");
const committedSeq = seqFromNumber(metadata.next_seq);
this.checkGlobalIds(writes);
this.checkDocumentActions(documentActions);
for (const write of writes) this.applyTableWrite(write, committedSeq);
this.applyDocumentActions(documentActions, committedSeq);
this.db
.prepare("UPDATE durable_metadata SET next_id = ?, next_seq = ? WHERE singleton = 1")
.run(String(Math.max(Number(metadata.next_id), candidateNextId)), committedSeq + 1);
await this.checkGlobalIds(transaction, writes);
await this.checkDocumentActions(transaction, documentActions);
for (const write of writes) await this.applyTableWrite(transaction, write, committedSeq);
await this.applyDocumentActions(transaction, documentActions, committedSeq);
await transaction.run(
"UPDATE durable_metadata SET next_id = ?, next_seq = ? WHERE singleton = 1",
String(Math.max(Number(metadata.next_id), candidateNextId)),
committedSeq + 1,
);
return committedSeq;
});
this.nextId = Math.max(this.nextId, candidateNextId);
@@ -224,7 +189,7 @@ export class SqliteStorage implements Storage {
async conversation(id: ConversationId, _context: Context): Promise<ConversationRecord | undefined> {
this.assertOpen();
const row = getRow<JsonRow>(this.db.prepare("SELECT record FROM conversations WHERE id = ?"), id);
const row = await this.db.get<JsonRow>("SELECT record FROM conversations WHERE id = ?", id);
return row === undefined ? undefined : parseJson<ConversationRecord>(row.record);
}
@@ -246,8 +211,8 @@ export class SqliteStorage implements Storage {
params.push(query.ownerTaskId);
}
params.push(limit + 1);
const rows = allRows<JsonRow>(
this.db.prepare(`SELECT record FROM conversations WHERE ${clauses.join(" AND ")} ORDER BY id LIMIT ?`),
const rows = await this.db.all<JsonRow>(
`SELECT record FROM conversations WHERE ${clauses.join(" AND ")} ORDER BY id LIMIT ?`,
...params,
);
return page(
@@ -278,10 +243,10 @@ export class SqliteStorage implements Storage {
let conversation: ConversationRecord | undefined;
if (context !== undefined) {
const conversationId = idFromNumber<ConversationId>(idOrConversationId);
conversation = this.readConversation(conversationId);
conversation = await this.readConversation(conversationId);
if (conversation === undefined) throw new Error(`Unknown conversation: ${conversationId}`);
}
const row = getRow<EntryJsonRow>(this.db.prepare("SELECT record, commit_seq FROM entries WHERE id = ?"), id);
const row = await this.db.get<EntryJsonRow>("SELECT record, commit_seq FROM entries WHERE id = ?", id);
if (row === undefined) return undefined;
const entry = parseJson<EntryRecord>(row.record);
if (conversation !== undefined) {
@@ -289,7 +254,7 @@ export class SqliteStorage implements Storage {
while (conversation.id !== entry.conversationId) {
if (conversation.parent === undefined) return undefined;
upperEntryId = Math.min(upperEntryId, conversation.parent.at);
conversation = this.readConversation(conversation.parent.conversationId)!;
conversation = (await this.readConversation(conversation.parent.conversationId))!;
}
if (entry.id > upperEntryId) return undefined;
}
@@ -302,29 +267,25 @@ export class SqliteStorage implements Storage {
_context: Context,
): Promise<(EntryRecord & { readonly head: EntryId }) | undefined> {
this.assertOpen();
let conversation = this.readConversation(conversationId);
let conversation = await this.readConversation(conversationId);
if (conversation === undefined) throw new Error(`Unknown conversation: ${conversationId}`);
let upper: number | undefined = atOrBeforeEntryId;
while (true) {
const row =
upper === undefined
? getRow<JsonRow>(
this.db.prepare(
"SELECT record FROM entries WHERE conversation_id = ? AND head IS NOT NULL ORDER BY id DESC LIMIT 1",
),
? await this.db.get<JsonRow>(
"SELECT record FROM entries WHERE conversation_id = ? AND head IS NOT NULL ORDER BY id DESC LIMIT 1",
conversation.id,
)
: getRow<JsonRow>(
this.db.prepare(
"SELECT record FROM entries WHERE conversation_id = ? AND head IS NOT NULL AND id <= ? ORDER BY id DESC LIMIT 1",
),
: await this.db.get<JsonRow>(
"SELECT record FROM entries WHERE conversation_id = ? AND head IS NOT NULL AND id <= ? ORDER BY id DESC LIMIT 1",
conversation.id,
upper,
);
if (row !== undefined) return parseJson<EntryRecord & { readonly head: EntryId }>(row.record);
if (conversation.parent === undefined) return undefined;
upper = upper === undefined ? conversation.parent.at : Math.min(upper, conversation.parent.at);
conversation = this.readConversation(conversation.parent.conversationId)!;
conversation = (await this.readConversation(conversation.parent.conversationId))!;
}
}
@@ -335,7 +296,7 @@ export class SqliteStorage implements Storage {
_context: Context,
): Promise<Page<EntryRecord, Cursor>> {
this.assertOpen();
let conversation = this.readConversation(query.conversationId);
let conversation = await this.readConversation(query.conversationId);
if (conversation === undefined) throw new Error(`Unknown conversation: ${query.conversationId}`);
const after = cursorId(cursor);
let upper: number | undefined = query.maxEntryId;
@@ -353,22 +314,22 @@ export class SqliteStorage implements Storage {
params.push(upper);
}
params.push(limit + 1 - values.length);
const rows = allRows<JsonRow>(
this.db.prepare(`SELECT record FROM entries WHERE ${clauses.join(" AND ")} ORDER BY id DESC LIMIT ?`),
const rows = await this.db.all<JsonRow>(
`SELECT record FROM entries WHERE ${clauses.join(" AND ")} ORDER BY id DESC LIMIT ?`,
...params,
);
values.push(...rows.map((row) => parseJson<EntryRecord>(row.record)));
if (values.length > limit || conversation.parent === undefined) break;
upper = upper === undefined ? conversation.parent.at : Math.min(upper, conversation.parent.at);
if (query.minEntryId !== undefined && upper < query.minEntryId) break;
conversation = this.readConversation(conversation.parent.conversationId)!;
conversation = (await this.readConversation(conversation.parent.conversationId))!;
}
return page(values, limit);
}
async task(id: TaskId, _context: Context): Promise<StoredTask | undefined> {
this.assertOpen();
const row = getRow<JsonRow>(this.db.prepare("SELECT record FROM tasks WHERE id = ?"), id);
const row = await this.db.get<JsonRow>("SELECT record FROM tasks WHERE id = ?", id);
return row === undefined ? undefined : parseJson<StoredTask>(row.record);
}
@@ -402,8 +363,8 @@ export class SqliteStorage implements Storage {
params.push(query.background ? 1 : 0);
}
params.push(limit + 1);
const rows = allRows<JsonRow>(
this.db.prepare(`SELECT record FROM tasks WHERE ${clauses.join(" AND ")} ORDER BY id LIMIT ?`),
const rows = await this.db.all<JsonRow>(
`SELECT record FROM tasks WHERE ${clauses.join(" AND ")} ORDER BY id LIMIT ?`,
...params,
);
return page(
@@ -414,7 +375,7 @@ export class SqliteStorage implements Storage {
async submission(id: SubmissionId, _context: Context): Promise<SubmissionRecord | undefined> {
this.assertOpen();
const row = getRow<JsonRow>(this.db.prepare("SELECT record FROM submissions WHERE id = ?"), id);
const row = await this.db.get<JsonRow>("SELECT record FROM submissions WHERE id = ?", id);
return row === undefined ? undefined : parseJson<SubmissionRecord>(row.record);
}
@@ -436,8 +397,8 @@ export class SqliteStorage implements Storage {
params.push(query.status);
}
params.push(limit + 1);
const rows = allRows<JsonRow>(
this.db.prepare(`SELECT record FROM submissions WHERE ${clauses.join(" AND ")} ORDER BY id LIMIT ?`),
const rows = await this.db.all<JsonRow>(
`SELECT record FROM submissions WHERE ${clauses.join(" AND ")} ORDER BY id LIMIT ?`,
...params,
);
return page(
@@ -452,8 +413,8 @@ export class SqliteStorage implements Storage {
_context: Context,
): Promise<SubmissionRecord | undefined> {
this.assertOpen();
const row = getRow<JsonRow>(
this.db.prepare("SELECT record FROM submissions WHERE conversation_id = ? AND request_id = ?"),
const row = await this.db.get<JsonRow>(
"SELECT record FROM submissions WHERE conversation_id = ? AND request_id = ?",
conversationId,
encodeIndexedString(requestId),
);
@@ -467,24 +428,25 @@ export class SqliteStorage implements Storage {
): Promise<DocumentRecord | undefined> {
this.assertOpen();
const parts = addressParts(address);
const statement =
const sql =
at === "current"
? this.db.prepare(`SELECT record FROM documents
? `SELECT record FROM documents
WHERE kind = ? AND scope_kind = ? AND owner_id = ? AND family = ? AND key_value = ?
AND retired_at IS NULL ORDER BY created_at DESC LIMIT 1`)
: this.db.prepare(`SELECT record FROM documents
AND retired_at IS NULL ORDER BY created_at DESC LIMIT 1`
: `SELECT record FROM documents
WHERE kind = ? AND scope_kind = ? AND owner_id = ? AND family = ? AND key_value = ?
AND created_at <= ? AND (retired_at IS NULL OR retired_at > ?)
ORDER BY created_at DESC LIMIT 1`);
ORDER BY created_at DESC LIMIT 1`;
const params: SqliteValue[] = [parts.kind, parts.scopeKind, parts.ownerId, parts.family, parts.keyValue];
if (at !== "current") params.push(at, at);
const row = getRow<JsonRow>(statement, ...params);
const row = await this.db.get<JsonRow>(sql, ...params);
return row === undefined ? undefined : parseJson<DocumentRecord>(row.record);
}
async document(id: DocumentId, at: DocumentPoint, _context: Context): Promise<StoredDocument | undefined> {
this.assertOpen();
return this.materializeDocument(id, at);
// The record and revision queries must observe one committed state; a commit between them can replace the base.
return this.db.transaction((transaction) => this.materializeDocument(transaction, id, at));
}
async scanDocuments(
@@ -508,8 +470,8 @@ export class SqliteStorage implements Storage {
params.push(query.at, query.at);
}
params.push(limit + 1);
const rows = allRows<JsonRow>(
this.db.prepare(`SELECT record FROM documents WHERE ${clauses.join(" AND ")} ORDER BY id LIMIT ?`),
const rows = await this.db.all<JsonRow>(
`SELECT record FROM documents WHERE ${clauses.join(" AND ")} ORDER BY id LIMIT ?`,
...params,
);
return page(
@@ -524,13 +486,17 @@ export class SqliteStorage implements Storage {
await this.db.close();
}
private readConversation(id: ConversationId): ConversationRecord | undefined {
const row = getRow<JsonRow>(this.db.prepare("SELECT record FROM conversations WHERE id = ?"), id);
private async readConversation(id: ConversationId): Promise<ConversationRecord | undefined> {
const row = await this.db.get<JsonRow>("SELECT record FROM conversations WHERE id = ?", id);
return row === undefined ? undefined : parseJson<ConversationRecord>(row.record);
}
private materializeDocument(id: DocumentId, at: DocumentPoint): StoredDocument | undefined {
const row = getRow<JsonRow>(this.db.prepare("SELECT record FROM documents WHERE id = ?"), id);
private async materializeDocument(
executor: SqliteExecutor,
id: DocumentId,
at: DocumentPoint,
): Promise<StoredDocument | undefined> {
const row = await executor.get<JsonRow>("SELECT record FROM documents WHERE id = ?", id);
if (row === undefined) return undefined;
const record = parseJson<DocumentRecord>(row.record);
if (at !== "current" && isCurrentOnly(record)) {
@@ -538,17 +504,17 @@ export class SqliteStorage implements Storage {
}
if (!isAliveAt(record, at)) return undefined;
const upper = at === "current" ? Number.MAX_SAFE_INTEGER : at;
const base = getRow<RevisionRow>(
this.db.prepare(`SELECT seq, kind, version, content FROM document_revisions
WHERE document_id = ? AND kind = 'base' AND seq <= ? ORDER BY seq DESC LIMIT 1`),
const base = await executor.get<RevisionRow>(
`SELECT seq, kind, version, content FROM document_revisions
WHERE document_id = ? AND kind = 'base' AND seq <= ? ORDER BY seq DESC LIMIT 1`,
id,
upper,
);
if (base === undefined) throw new Error(`Document ${id} is missing a required base`);
let value = parseJson<JsonObject>(base.content);
const tail = allRows<RevisionRow>(
this.db.prepare(`SELECT seq, kind, version, content FROM document_revisions
WHERE document_id = ? AND seq > ? AND seq <= ? ORDER BY seq`),
const tail = await executor.all<RevisionRow>(
`SELECT seq, kind, version, content FROM document_revisions
WHERE document_id = ? AND seq > ? AND seq <= ? ORDER BY seq`,
id,
base.seq,
upper,
@@ -571,15 +537,15 @@ export class SqliteStorage implements Storage {
return nextId;
}
private checkGlobalIds(writes: readonly StorageWrite[]): void {
private async checkGlobalIds(executor: SqliteExecutor, writes: readonly StorageWrite[]): Promise<void> {
const claimed = new Map<Id<string>, TableName>();
const lookup = this.db.prepare("SELECT record_type FROM record_ids WHERE id = ?");
for (const write of writes) {
if (write.type === "document.change" || write.type === "document.retire") continue;
const document = write.type === "document.create" || write.type === "document.copy";
const table: TableName = document ? "document" : write.type;
const id = document ? write.record.id : write.value.id;
const existing = getRow<RecordIdRow>(lookup, id)?.record_type;
const existing = (await executor.get<RecordIdRow>("SELECT record_type FROM record_ids WHERE id = ?", id))
?.record_type;
const earlier = claimed.get(id);
if (table === "conversation" || table === "entry" || table === "document") {
if (existing !== undefined) throw new Error(`ID ${id} already belongs to ${existing}`);
@@ -640,22 +606,23 @@ export class SqliteStorage implements Storage {
return actions;
}
private checkDocumentActions(actions: ReadonlyMap<DocumentId, DocumentAction>): void {
private async checkDocumentActions(
executor: SqliteExecutor,
actions: ReadonlyMap<DocumentId, DocumentAction>,
): Promise<void> {
const liveCounts = new Map<string, number>();
for (const [id, action] of actions) {
if (action.copy !== undefined && actions.has(action.copy.id)) {
throw new StorageRejected(`Document copy ${id} source is changed in the copy batch`);
}
const row = getRow<JsonRow>(this.db.prepare("SELECT record FROM documents WHERE id = ?"), id);
const row = await executor.get<JsonRow>("SELECT record FROM documents WHERE id = ?", id);
const existing = row === undefined ? undefined : parseJson<DocumentRecord>(row.record);
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.content?.kind === "delta") {
const previous = getRow<{ readonly version: number }>(
this.db.prepare(
"SELECT version FROM document_revisions WHERE document_id = ? ORDER BY seq DESC LIMIT 1",
),
const previous = await executor.get<{ readonly version: number }>(
"SELECT version FROM document_revisions WHERE document_id = ? ORDER BY seq DESC LIMIT 1",
id,
);
if (previous === undefined) throw new Error(`Document ${id} delta has no base`);
@@ -666,7 +633,7 @@ export class SqliteStorage implements Storage {
const record = action.create ?? existing!;
const key = addressKey(record);
let live = liveCounts.get(key);
if (live === undefined) live = this.currentDocumentId(record) === undefined ? 0 : 1;
if (live === undefined) live = (await this.currentDocumentId(executor, record)) === undefined ? 0 : 1;
if (action.retire && existing !== undefined) live--;
if (action.create !== undefined && !action.retire) live++;
liveCounts.set(key, live);
@@ -676,73 +643,78 @@ export class SqliteStorage implements Storage {
}
}
private currentDocumentId(address: DocumentAddress | DocumentCreate | DocumentRecord): DocumentId | undefined {
private async currentDocumentId(
executor: SqliteExecutor,
address: DocumentAddress | DocumentCreate | DocumentRecord,
): Promise<DocumentId | undefined> {
const parts = addressParts(address);
const id = getRow<IdRow>(
this.db.prepare(`SELECT id FROM documents
const id = (
await executor.get<IdRow>(
`SELECT id FROM documents
WHERE kind = ? AND scope_kind = ? AND owner_id = ? AND family = ? AND key_value = ? AND retired_at IS NULL
LIMIT 1`),
parts.kind,
parts.scopeKind,
parts.ownerId,
parts.family,
parts.keyValue,
LIMIT 1`,
parts.kind,
parts.scopeKind,
parts.ownerId,
parts.family,
parts.keyValue,
)
)?.id;
return id === undefined ? undefined : idFromNumber<DocumentId>(id);
}
private applyTableWrite(write: StorageWrite, seq: Seq): void {
private async applyTableWrite(executor: SqliteExecutor, write: StorageWrite, seq: Seq): Promise<void> {
switch (write.type) {
case "conversation":
this.claimId(write.value.id, "conversation");
this.db
.prepare(
"INSERT INTO conversations (id, owner_conversation_id, owner_task_id, record) VALUES (?, ?, ?, ?)",
)
.run(
write.value.id,
write.value.owner?.conversationId ?? null,
write.value.owner?.taskId ?? null,
encodeJson(write.value),
);
await this.claimId(executor, write.value.id, "conversation");
await executor.run(
"INSERT INTO conversations (id, owner_conversation_id, owner_task_id, record) VALUES (?, ?, ?, ?)",
write.value.id,
write.value.owner?.conversationId ?? null,
write.value.owner?.taskId ?? null,
encodeJson(write.value),
);
break;
case "entry":
this.claimId(write.value.id, "entry");
this.db
.prepare("INSERT INTO entries (id, conversation_id, head, commit_seq, record) VALUES (?, ?, ?, ?, ?)")
.run(write.value.id, write.value.conversationId, write.value.head ?? null, seq, encodeJson(write.value));
await this.claimId(executor, write.value.id, "entry");
await executor.run(
"INSERT INTO entries (id, conversation_id, head, commit_seq, record) VALUES (?, ?, ?, ?, ?)",
write.value.id,
write.value.conversationId,
write.value.head ?? null,
seq,
encodeJson(write.value),
);
break;
case "task":
this.claimId(write.value.id, "task");
this.db
.prepare(`INSERT INTO tasks (id, conversation_id, kind, status, abort_requested, background, record)
await this.claimId(executor, write.value.id, "task");
await executor.run(
`INSERT INTO tasks (id, conversation_id, kind, status, abort_requested, background, record)
VALUES (?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(id) DO UPDATE SET conversation_id = excluded.conversation_id, kind = excluded.kind,
status = excluded.status, abort_requested = excluded.abort_requested,
background = excluded.background, record = excluded.record`)
.run(
write.value.id,
write.value.conversationId,
encodeIndexedString(write.value.kind),
write.value.state.status,
write.value.abortRequested ? 1 : 0,
write.value.background ? 1 : 0,
encodeJson(write.value),
);
background = excluded.background, record = excluded.record`,
write.value.id,
write.value.conversationId,
encodeIndexedString(write.value.kind),
write.value.state.status,
write.value.abortRequested ? 1 : 0,
write.value.background ? 1 : 0,
encodeJson(write.value),
);
break;
case "submission":
this.claimId(write.value.id, "submission");
this.db
.prepare(`INSERT INTO submissions (id, conversation_id, request_id, status, record) VALUES (?, ?, ?, ?, ?)
await this.claimId(executor, write.value.id, "submission");
await executor.run(
`INSERT INTO submissions (id, conversation_id, request_id, status, record) VALUES (?, ?, ?, ?, ?)
ON CONFLICT(id) DO UPDATE SET conversation_id = excluded.conversation_id,
request_id = excluded.request_id, status = excluded.status, record = excluded.record`)
.run(
write.value.id,
write.value.conversationId,
write.value.requestId === undefined ? null : encodeIndexedString(write.value.requestId),
write.value.status,
encodeJson(write.value),
);
request_id = excluded.request_id, status = excluded.status, record = excluded.record`,
write.value.id,
write.value.conversationId,
write.value.requestId === undefined ? null : encodeIndexedString(write.value.requestId),
write.value.status,
encodeJson(write.value),
);
break;
case "document.create":
case "document.copy":
@@ -752,16 +724,20 @@ export class SqliteStorage implements Storage {
}
}
private claimId(id: Id<string>, table: TableName): void {
this.db.prepare("INSERT OR IGNORE INTO record_ids (id, record_type) VALUES (?, ?)").run(id, table);
private async claimId(executor: SqliteExecutor, id: Id<string>, table: TableName): Promise<void> {
await executor.run("INSERT OR IGNORE INTO record_ids (id, record_type) VALUES (?, ?)", id, table);
}
private applyDocumentActions(actions: ReadonlyMap<DocumentId, DocumentAction>, seq: Seq): void {
private async applyDocumentActions(
executor: SqliteExecutor,
actions: ReadonlyMap<DocumentId, DocumentAction>,
seq: Seq,
): Promise<void> {
for (const [id, action] of actions) {
let content = action.content;
if (action.copy !== undefined) {
try {
const stored = this.materializeDocument(action.copy.id, action.copy.at);
const stored = await this.materializeDocument(executor, action.copy.id, action.copy.at);
if (stored === undefined) throw new Error(`Fork source document ${action.copy.id} cannot be read`);
const create = action.create!;
if (
@@ -788,47 +764,54 @@ export class SqliteStorage implements Storage {
...(action.retire ? { retiredAt: seq } : {}),
};
const parts = addressParts(record);
this.claimId(id, "document");
this.db
.prepare(`INSERT INTO documents
await this.claimId(executor, id, "document");
await executor.run(
`INSERT INTO documents
(id, kind, family, key_value, scope_kind, owner_id, created_at, retired_at, record)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)`)
.run(
id,
parts.kind,
parts.family,
parts.keyValue,
parts.scopeKind,
parts.ownerId,
seq,
action.retire ? seq : null,
encodeJson(record),
);
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)`,
id,
parts.kind,
parts.family,
parts.keyValue,
parts.scopeKind,
parts.ownerId,
seq,
action.retire ? seq : null,
encodeJson(record),
);
} else {
const row = getRow<JsonRow>(this.db.prepare("SELECT record FROM documents WHERE id = ?"), id)!;
const row = (await executor.get<JsonRow>("SELECT record FROM documents WHERE id = ?", id))!;
record = parseJson<DocumentRecord>(row.record);
}
if (content !== undefined) {
if (content.kind === "base" && isCurrentOnly(record)) {
this.db.prepare("DELETE FROM document_revisions WHERE document_id = ?").run(id);
await executor.run("DELETE FROM document_revisions WHERE document_id = ?", id);
}
const encodedContent = content.kind === "base" ? encodeJson(content.value) : encodeJson(content.ops);
this.db
.prepare(
"INSERT INTO document_revisions (document_id, seq, kind, version, content) VALUES (?, ?, ?, ?, ?)",
)
.run(id, seq, content.kind, content.version, encodedContent);
await executor.run(
"INSERT INTO document_revisions (document_id, seq, kind, version, content) VALUES (?, ?, ?, ?, ?)",
id,
seq,
content.kind,
content.version,
encodedContent,
);
}
if (action.retire) {
if (action.create === undefined) {
record = { ...record, retiredAt: seq };
this.db
.prepare("UPDATE documents SET retired_at = ?, record = ? WHERE id = ?")
.run(seq, encodeJson(record), id);
await executor.run(
"UPDATE documents SET retired_at = ?, record = ? WHERE id = ?",
seq,
encodeJson(record),
id,
);
}
if (isCurrentOnly(record)) {
await executor.run("DELETE FROM document_revisions WHERE document_id = ?", id);
}
if (isCurrentOnly(record)) this.db.prepare("DELETE FROM document_revisions WHERE document_id = ?").run(id);
}
}
}
+188 -53
View File
@@ -1,11 +1,12 @@
import { DatabaseSync, type StatementSync } from "node:sqlite";
import { BACKGROUND_CONTEXT } from "@earendil-works/chord/context";
import { describe, expect, it } from "vitest";
import { StorageRejected } from "../src/errors.ts";
import { idFromNumber } from "../src/ids.ts";
import type { SqliteDatabase, SqliteStatement } from "../src/storage/sqlite/index.ts";
import type { SqliteDatabase, SqliteExecutor, SqliteValue } from "../src/storage/sqlite/index.ts";
import { SqliteStorage } from "../src/storage/sqlite/index.ts";
import { type NodeSqliteDatabase, openNodeSqliteDatabase } from "../src/storage/sqlite/node.ts";
import { type EntryId, ROOT_CONVERSATION_ID } from "../src/types.ts";
import { NodeSqliteDatabase, openNodeSqliteDatabase } from "../src/storage/sqlite/node.ts";
import { type DocumentId, type EntryId, ROOT_CONVERSATION_ID } from "../src/types.ts";
type SettlementMode = "immediate" | "delay" | "reject";
@@ -13,47 +14,43 @@ class ControlledSettlementDatabase implements SqliteDatabase {
private readonly delegate: NodeSqliteDatabase;
private mode: SettlementMode = "immediate";
private pendingSettlement: (() => void) | undefined;
private readonly prepareCounts = new Map<string, number>();
constructor(delegate: NodeSqliteDatabase) {
this.delegate = delegate;
}
exec(sql: string): void {
this.delegate.exec(sql);
exec(sql: string): Promise<void> {
return this.delegate.exec(sql);
}
prepare(sql: string): SqliteStatement {
this.prepareCounts.set(sql, (this.prepareCounts.get(sql) ?? 0) + 1);
return this.delegate.prepare(sql);
run(sql: string, ...params: SqliteValue[]): Promise<void> {
return this.delegate.run(sql, ...params);
}
transaction<T>(callback: () => T): T | Promise<T> {
get<T extends object>(sql: string, ...params: SqliteValue[]): Promise<T | undefined> {
return this.delegate.get<T>(sql, ...params);
}
all<T extends object>(sql: string, ...params: SqliteValue[]): Promise<T[]> {
return this.delegate.all<T>(sql, ...params);
}
transaction<T>(callback: (transaction: SqliteExecutor) => Promise<T>): Promise<T> {
const mode = this.mode;
this.mode = "immediate";
if (mode === "immediate") return this.delegate.transaction(callback);
try {
const result = this.delegate.transaction(() => {
const value = callback();
if (mode === "reject") throw new Error("controlled settlement rejection");
return value;
});
return new Promise<T>((resolve) => {
this.pendingSettlement = () => resolve(result);
});
} catch (error) {
return new Promise<T>((_resolve, reject) => {
this.pendingSettlement = () => reject(error);
});
}
const settlement = this.delegate.transaction(async (transaction) => {
const value = await callback(transaction);
if (mode === "reject") throw new Error("controlled settlement rejection");
return value;
});
return new Promise<T>((resolve, reject) => {
this.pendingSettlement = () => void settlement.then(resolve, reject);
});
}
close(): void {
this.delegate.close();
}
prepareCount(sql: string): number {
return this.prepareCounts.get(sql) ?? 0;
close(): Promise<void> {
return this.delegate.close();
}
controlNextSettlement(mode: Exclude<SettlementMode, "immediate">): void {
@@ -69,10 +66,23 @@ class ControlledSettlementDatabase implements SqliteDatabase {
}
}
class PrepareCountingDatabaseSync extends DatabaseSync {
private readonly prepareCounts = new Map<string, number>();
override prepare(sql: string): StatementSync {
this.prepareCounts.set(sql, (this.prepareCounts.get(sql) ?? 0) + 1);
return super.prepare(sql);
}
repeatedPrepares(): string[] {
return [...this.prepareCounts].filter(([, count]) => count > 1).map(([sql]) => sql);
}
}
describe("portable SQLite facade settlement", () => {
it("prepares each storage statement once and rebinds it across commits", async () => {
const database = new ControlledSettlementDatabase(await openNodeSqliteDatabase(":memory:"));
const storage = await SqliteStorage.open(database);
it("prepares each storage statement once per connection and reuses it across transactions", async () => {
const connection = new PrepareCountingDatabaseSync(":memory:");
const storage = await SqliteStorage.open(new NodeSqliteDatabase(connection));
await storage.commit([{ type: "conversation", value: { id: ROOT_CONVERSATION_ID } }], BACKGROUND_CONTEXT);
await storage.commit(
Array.from({ length: 100 }, (_, index) => ({
@@ -96,37 +106,135 @@ describe("portable SQLite facade settlement", () => {
storage.commit([{ type: "conversation", value: { id: ROOT_CONVERSATION_ID } }], BACKGROUND_CONTEXT),
).rejects.toThrow("ID 1 already belongs to conversation");
expect((await storage.entry(idFromNumber<EntryId>(2), BACKGROUND_CONTEXT))?.entry.kind).toBe("cached");
expect(database.prepareCount("SELECT record, commit_seq FROM entries WHERE id = ?")).toBe(1);
expect(database.prepareCount("INSERT OR IGNORE INTO record_ids (id, record_type) VALUES (?, ?)")).toBe(1);
expect(
database.prepareCount(
"INSERT INTO entries (id, conversation_id, head, commit_seq, record) VALUES (?, ?, ?, ?, ?)",
),
).toBe(1);
expect(connection.repeatedPrepares()).toEqual([]);
await storage.close(BACKGROUND_CONTEXT);
});
it("rejects asynchronous Node transaction callbacks and closes idempotently", async () => {
it("commits work done through the transaction handle and closes idempotently", async () => {
const database = await openNodeSqliteDatabase(":memory:");
expect(() => database.transaction(() => Promise.resolve())).toThrow(
"SQLite transaction callbacks must be synchronous",
await database.transaction(async (transaction) => {
await transaction.exec("CREATE TABLE async_probe (value INTEGER)");
await transaction.run("INSERT INTO async_probe (value) VALUES (?)", 1);
});
expect(await database.get("SELECT value FROM async_probe")).toEqual({ value: 1 });
await database.close();
await database.close();
});
it("serializes concurrent transactions", async () => {
const database = await openNodeSqliteDatabase(":memory:");
await database.exec("CREATE TABLE transaction_queue (value INTEGER)");
let markFirstStarted!: () => void;
const firstStarted = new Promise<void>((resolve) => {
markFirstStarted = resolve;
});
let releaseFirst!: () => void;
const firstGate = new Promise<void>((resolve) => {
releaseFirst = resolve;
});
const first = database.transaction(async (transaction) => {
await transaction.exec("INSERT INTO transaction_queue (value) VALUES (1)");
markFirstStarted();
await firstGate;
});
await firstStarted;
let secondStarted = false;
const second = database.transaction(async (transaction) => {
secondStarted = true;
await transaction.exec("INSERT INTO transaction_queue (value) VALUES (2)");
});
await Promise.resolve();
expect(secondStarted).toBe(false);
releaseFirst();
await Promise.all([first, second]);
expect(await database.all("SELECT value FROM transaction_queue ORDER BY value")).toEqual([
{ value: 1 },
{ value: 2 },
]);
await database.close();
});
it("queues ordinary operations behind an active transaction", async () => {
const database = await openNodeSqliteDatabase(":memory:");
await database.exec("CREATE TABLE operation_queue (value INTEGER)");
let markTransactionStarted!: () => void;
const transactionStarted = new Promise<void>((resolve) => {
markTransactionStarted = resolve;
});
let releaseTransaction!: () => void;
const transactionGate = new Promise<void>((resolve) => {
releaseTransaction = resolve;
});
const pending = database.transaction(async (transaction) => {
await transaction.exec("INSERT INTO operation_queue (value) VALUES (1)");
markTransactionStarted();
await transactionGate;
});
await transactionStarted;
let writeSettled = false;
const write = database.exec("INSERT INTO operation_queue (value) VALUES (2)").finally(() => {
writeSettled = true;
});
let readSettled = false;
const read = database.all("SELECT value FROM operation_queue ORDER BY value").finally(() => {
readSettled = true;
});
await Promise.resolve();
expect(writeSettled).toBe(false);
expect(readSettled).toBe(false);
releaseTransaction();
await pending;
await write;
await expect(read).resolves.toEqual([{ value: 1 }, { value: 2 }]);
await database.close();
});
it("rejects database operations from inside a transaction callback instead of waiting forever", async () => {
const database = await openNodeSqliteDatabase(":memory:");
await database.exec("CREATE TABLE misuse_probe (value INTEGER)");
const misuse = "Use the transaction handle inside a transaction callback";
await expect(database.transaction(() => database.exec("SELECT 1"))).rejects.toThrow(misuse);
await expect(database.transaction(() => database.all("SELECT value FROM misuse_probe"))).rejects.toThrow(misuse);
await expect(database.transaction(() => database.transaction(async () => undefined))).rejects.toThrow(
"Nested SQLite transactions are not supported",
);
database.close();
database.close();
await expect(database.transaction(() => database.close())).rejects.toThrow(
"Cannot close SQLite during an active transaction",
);
await database.close();
});
it("rejects a transaction handle used after its transaction settles", async () => {
const database = await openNodeSqliteDatabase(":memory:");
await database.exec("CREATE TABLE stale_probe (value INTEGER)");
let handle!: SqliteExecutor;
await database.transaction(async (transaction) => {
handle = transaction;
await transaction.run("INSERT INTO stale_probe (value) VALUES (?)", 1);
});
const stale = "SQLite transaction handle is no longer active";
await expect(handle.exec("INSERT INTO stale_probe (value) VALUES (2)")).rejects.toThrow(stale);
await expect(handle.run("INSERT INTO stale_probe (value) VALUES (?)", 3)).rejects.toThrow(stale);
expect(await database.all("SELECT value FROM stale_probe")).toEqual([{ value: 1 }]);
await database.close();
});
it("does not preserve a guaranteed rejection when rollback itself fails", async () => {
const database = await openNodeSqliteDatabase(":memory:");
database.exec("CREATE TABLE rollback_probe (value INTEGER)");
expect(() =>
database.transaction(() => {
database.exec("INSERT INTO rollback_probe (value) VALUES (1)");
database.exec("COMMIT");
await database.exec("CREATE TABLE rollback_probe (value INTEGER)");
await expect(
database.transaction(async (transaction) => {
await transaction.exec("INSERT INTO rollback_probe (value) VALUES (1)");
await transaction.exec("COMMIT");
throw new StorageRejected("rejected after an escaped commit");
}),
).toThrow(AggregateError);
expect(database.prepare("SELECT value FROM rollback_probe").get()).toEqual({ value: 1 });
database.close();
).rejects.toThrow(AggregateError);
expect(await database.get("SELECT value FROM rollback_probe")).toEqual({ value: 1 });
await database.close();
});
it("awaits async transaction settlement and adopts IDs only after success", async () => {
@@ -180,4 +288,31 @@ describe("portable SQLite facade settlement", () => {
expect(await storage.entry(idFromNumber<EntryId>(200), BACKGROUND_CONTEXT)).toBeUndefined();
await storage.close(BACKGROUND_CONTEXT);
});
it("reads a document from one committed state while a commit replaces its base", async () => {
const id = idFromNumber<DocumentId>(5);
// Each yield count starts the commit at a different point of the read's record and revision queries.
for (let yields = 0; yields < 16; yields++) {
const storage = await SqliteStorage.open(await openNodeSqliteDatabase(":memory:"));
await storage.commit(
[
{
type: "document.create",
record: { id, kind: "replaced", scope: { kind: "session" } },
content: { kind: "base", version: 1, value: { value: 1 } },
},
],
BACKGROUND_CONTEXT,
);
const read = storage.document(id, "current", BACKGROUND_CONTEXT);
for (let index = 0; index < yields; index++) await Promise.resolve();
const replace = storage.commit(
[{ type: "document.change", id, content: { kind: "base", version: 1, value: { value: 2 } } }],
BACKGROUND_CONTEXT,
);
const [stored] = await Promise.all([read, replace]);
expect([{ value: 1 }, { value: 2 }]).toContainEqual(stored?.value);
await storage.close(BACKGROUND_CONTEXT);
}
});
});
+24 -23
View File
@@ -32,15 +32,15 @@ describe("durable SQLite migrations", () => {
try {
await applySqliteMigrations(database);
await applySqliteMigrations(database);
expect(database.prepare("SELECT version FROM durable_schema WHERE singleton = 1").get()).toEqual({
expect(await database.get("SELECT version FROM durable_schema WHERE singleton = 1")).toEqual({
version: CURRENT_SQLITE_SCHEMA_VERSION,
});
expect(database.prepare("SELECT next_id, next_seq FROM durable_metadata WHERE singleton = 1").get()).toEqual({
expect(await database.get("SELECT next_id, next_seq FROM durable_metadata WHERE singleton = 1")).toEqual({
next_id: "2",
next_seq: 1,
});
} finally {
database.close();
await database.close();
}
});
@@ -48,10 +48,11 @@ describe("durable SQLite migrations", () => {
const path = await databasePath();
const database = await openNodeSqliteDatabase(path);
await applySqliteMigrations(database);
database
.prepare("UPDATE durable_schema SET version = ? WHERE singleton = 1")
.run(CURRENT_SQLITE_SCHEMA_VERSION + 1);
database.close();
await database.run(
"UPDATE durable_schema SET version = ? WHERE singleton = 1",
CURRENT_SQLITE_SCHEMA_VERSION + 1,
);
await database.close();
await expect(openNodeSqliteStorage(path)).rejects.toThrow("is newer than supported version");
});
@@ -74,23 +75,23 @@ describe("durable SQLite migrations", () => {
];
await expect(applySqliteMigrations(database, failed)).rejects.toThrow();
expect(
database
.prepare(
"SELECT count(*) AS count FROM sqlite_schema WHERE name IN ('durable_schema', 'migration_first', 'migration_second')",
)
.get(),
await database.get(
"SELECT count(*) AS count FROM sqlite_schema WHERE name IN ('durable_schema', 'migration_first', 'migration_second')",
),
).toEqual({ count: 0 });
await applySqliteMigrations(database, [
failed[0],
{ version: 2, statements: ["CREATE TABLE migration_second (value TEXT) STRICT"] },
]);
expect(database.prepare("SELECT version FROM durable_schema WHERE singleton = 1").get()).toEqual({
expect(await database.get("SELECT version FROM durable_schema WHERE singleton = 1")).toEqual({
version: 2,
});
expect(database.prepare("SELECT value FROM migration_first").get()).toEqual({ value: "retained" });
expect(await database.get("SELECT value FROM migration_first")).toEqual({
value: "retained",
});
} finally {
database.close();
await database.close();
}
});
@@ -124,13 +125,13 @@ describe("durable SQLite migrations", () => {
},
];
await expect(applySqliteMigrations(database, failedMigrations)).rejects.toThrow();
expect(database.prepare("SELECT version FROM durable_schema WHERE singleton = 1").get()).toEqual({
expect(await database.get("SELECT version FROM durable_schema WHERE singleton = 1")).toEqual({
version: CURRENT_SQLITE_SCHEMA_VERSION,
});
expect(
database
.prepare("SELECT count(*) AS count FROM sqlite_schema WHERE type = 'table' AND name = 'migration_probe'")
.get(),
await database.get(
"SELECT count(*) AS count FROM sqlite_schema WHERE type = 'table' AND name = 'migration_probe'",
),
).toEqual({ count: 0 });
const successfulMigrations: readonly SqliteMigration[] = [
@@ -138,10 +139,10 @@ describe("durable SQLite migrations", () => {
{ version: nextVersion, statements: ["CREATE TABLE migration_probe (value TEXT) STRICT"] },
];
await applySqliteMigrations(database, successfulMigrations);
expect(database.prepare("SELECT version FROM durable_schema WHERE singleton = 1").get()).toEqual({
expect(await database.get("SELECT version FROM durable_schema WHERE singleton = 1")).toEqual({
version: nextVersion,
});
expect(database.prepare("SELECT record, commit_seq FROM entries WHERE id = 2").get()).toEqual({
expect(await database.get("SELECT record, commit_seq FROM entries WHERE id = 2")).toEqual({
record: JSON.stringify({
id: 2,
conversationId: ROOT_CONVERSATION_ID,
@@ -150,10 +151,10 @@ describe("durable SQLite migrations", () => {
}),
commit_seq: 1,
});
expect(database.prepare("SELECT next_id, next_seq FROM durable_metadata WHERE singleton = 1").get()).toEqual({
expect(await database.get("SELECT next_id, next_seq FROM durable_metadata WHERE singleton = 1")).toEqual({
next_id: "3",
next_seq: 2,
});
database.close();
await database.close();
});
});