fix(durable): harden async SQLite adapter queue and close

- drop AsyncLocalStorage misuse guards from the Node adapter
- run calls immediately when the queue is idle; transactions publish their barrier before the callback starts
- SqliteStorage.close() waits for admitted multi-query reads and shares one close promise
This commit is contained in:
Mario Zechner
2026-09-30 22:32:34 +02:00
parent 8f89293a26
commit ed391c4f0a
6 changed files with 180 additions and 54 deletions
+4
View File
@@ -2,6 +2,10 @@
## [Unreleased]
### Breaking Changes
- The portable SQLite facade in `@earendil-works/pi-durable/storage/sqlite` is asynchronous: `SqliteDatabase` extends the new `SqliteExecutor` (`exec`, `run`, `get`, `all` by SQL text), `prepare` and `SqliteStatement` are removed, `transaction` takes an async callback that receives a transaction handle, and `close()` returns a promise. Custom adapters must be rewritten ([#10232](https://github.com/earendil-works/pi/pull/10232) by [@christianklotz](https://github.com/christianklotz)).
## [0.99.2] - 2026-09-30
### Breaking Changes
+1 -1
View File
@@ -419,7 +419,7 @@ const usage = await harness.usage(context); // { models: { "openai/gpt-6-sol": U
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:
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, so calling `database` itself inside the callback never settles:
```typescript
await database.transaction(async (transaction) => {
@@ -23,7 +23,9 @@ export interface SqliteExecutor {
* `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.
* returned promise settles after commit or rollback. Calling the database itself (including
* `transaction` or `close`) from inside a callback therefore waits for that transaction and
* never settles.
*
* 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
+47 -30
View File
@@ -1,4 +1,3 @@
import { AsyncLocalStorage } from "node:async_hooks";
import { mkdir } from "node:fs/promises";
import { dirname } from "node:path";
import type { StatementSync } from "node:sqlite";
@@ -19,20 +18,55 @@ const DEFAULT_BUSY_TIMEOUT_MS = 5_000;
type TransactionScope = { active: boolean };
class SerialOperationQueue {
private tail = Promise.resolve();
const ignore = (): void => {};
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;
/**
* Runs operations in call order. An operation starts immediately when nothing is running or waiting;
* otherwise it waits for everything before it. An asynchronous operation holds the queue until it settles.
*/
class SerialOperationQueue {
private tail: Promise<void> = Promise.resolve();
private pending = 0;
run<T>(operation: () => T): Promise<T> {
if (this.pending > 0) return this.enqueue(operation);
try {
return await operation();
} finally {
release();
return Promise.resolve(operation());
} catch (error) {
return Promise.reject(error);
}
}
runAsync<T>(operation: () => Promise<T>): Promise<T> {
if (this.pending > 0) return this.enqueue(operation);
this.pending++;
// Publish the barrier before the operation starts, so calls it makes synchronously wait behind it.
const { promise: barrier, resolve: releaseBarrier } = Promise.withResolvers<void>();
this.tail = barrier;
let started: Promise<T>;
try {
started = operation();
} catch (error) {
started = Promise.reject(error);
}
return started.finally(() => {
this.pending--;
releaseBarrier();
});
}
private enqueue<T>(operation: () => T | Promise<T>): Promise<T> {
this.pending++;
return this.release(this.tail.then(operation));
}
private release<T>(operation: Promise<T>): Promise<T> {
const settled = operation.finally(() => {
this.pending--;
});
this.tail = settled.then(ignore, ignore);
return settled;
}
}
/**
@@ -97,8 +131,6 @@ class NodeSqliteTransaction extends NodeSqliteExecutor {
/** `SqliteDatabase` adapter backed by Node's built-in `node:sqlite`. */
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) {
@@ -106,16 +138,11 @@ export class NodeSqliteDatabase extends NodeSqliteExecutor implements SqliteData
}
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 () => {
return this.access.runAsync(async () => {
this.database.exec("BEGIN IMMEDIATE");
const scope = { active: true };
try {
const result = await this.transactionScope.run(scope, () =>
callback(new NodeSqliteTransaction(this.database, this.statements, scope)),
);
const result = await callback(new NodeSqliteTransaction(this.database, this.statements, scope));
scope.active = false;
this.database.exec("COMMIT");
return result;
@@ -132,9 +159,6 @@ export class NodeSqliteDatabase extends NodeSqliteExecutor implements SqliteData
}
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;
@@ -148,15 +172,8 @@ export class NodeSqliteDatabase extends NodeSqliteExecutor implements SqliteData
}
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;
}
}
/** Open and configure a Node-backed SQLite database facade. */
+59 -11
View File
@@ -131,6 +131,9 @@ export class SqliteStorage implements Storage {
private readonly db: SqliteDatabase;
private nextId: number;
private closed = false;
private closing: Promise<void> | undefined;
private admittedReads = 0;
private readsDrained: (() => void) | undefined;
private constructor(db: SqliteDatabase, nextId: number) {
this.db = db;
@@ -227,12 +230,36 @@ export class SqliteStorage implements Storage {
id: EntryId,
context: Context,
): Promise<{ readonly entry: EntryRecord; readonly commitSeq: Seq } | undefined>;
async entry(
entry(
idOrConversationId: EntryId | ConversationId,
idOrContext: EntryId | Context,
context?: Context,
): Promise<{ readonly entry: EntryRecord; readonly commitSeq: Seq } | undefined> {
return this.admitRead(() => this.readEntry(idOrConversationId, idOrContext, context));
}
findLatestHeadMarker(
conversationId: ConversationId,
atOrBeforeEntryId: EntryId | undefined,
_context: Context,
): Promise<(EntryRecord & { readonly head: EntryId }) | undefined> {
return this.admitRead(() => this.readLatestHeadMarker(conversationId, atOrBeforeEntryId));
}
scanEntries(
query: EntryQuery,
limit: number,
cursor: Cursor | undefined,
_context: Context,
): Promise<Page<EntryRecord, Cursor>> {
return this.admitRead(() => this.readEntries(query, limit, cursor));
}
private async readEntry(
idOrConversationId: EntryId | ConversationId,
idOrContext: EntryId | Context,
context?: Context,
): Promise<{ readonly entry: EntryRecord; readonly commitSeq: Seq } | undefined> {
this.assertOpen();
const id =
context === undefined
? idFromNumber<EntryId>(idOrConversationId)
@@ -261,12 +288,10 @@ export class SqliteStorage implements Storage {
return { entry, commitSeq: seqFromNumber(row.commit_seq) };
}
async findLatestHeadMarker(
private async readLatestHeadMarker(
conversationId: ConversationId,
atOrBeforeEntryId: EntryId | undefined,
_context: Context,
): Promise<(EntryRecord & { readonly head: EntryId }) | undefined> {
this.assertOpen();
let conversation = await this.readConversation(conversationId);
if (conversation === undefined) throw new Error(`Unknown conversation: ${conversationId}`);
let upper: number | undefined = atOrBeforeEntryId;
@@ -289,13 +314,11 @@ export class SqliteStorage implements Storage {
}
}
async scanEntries(
private async readEntries(
query: EntryQuery,
limit: number,
cursor: Cursor | undefined,
_context: Context,
): Promise<Page<EntryRecord, Cursor>> {
this.assertOpen();
let conversation = await this.readConversation(query.conversationId);
if (conversation === undefined) throw new Error(`Unknown conversation: ${query.conversationId}`);
const after = cursorId(cursor);
@@ -480,12 +503,37 @@ export class SqliteStorage implements Storage {
);
}
async close(_context: Context): Promise<void> {
if (this.closed) return;
this.closed = true;
close(_context: Context): Promise<void> {
if (this.closing === undefined) {
this.closed = true;
this.closing = this.closeDatabase();
}
return this.closing;
}
private async closeDatabase(): Promise<void> {
if (this.admittedReads > 0) {
await new Promise<void>((resolve) => {
this.readsDrained = resolve;
});
}
await this.db.close();
}
/**
* Run a read that issues several queries. Close waits for admitted reads, so their later queries never reach a
* closed database. Single-query reads and transactions are already ordered before close by the database.
*/
private async admitRead<T>(read: () => Promise<T>): Promise<T> {
this.assertOpen();
this.admittedReads++;
try {
return await read();
} finally {
if (--this.admittedReads === 0) this.readsDrained?.();
}
}
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);
+66 -11
View File
@@ -156,6 +156,34 @@ describe("portable SQLite facade settlement", () => {
await database.close();
});
it("runs operations in call order whether they start immediately or wait", async () => {
const database = await openNodeSqliteDatabase(":memory:");
await database.exec("CREATE TABLE call_order (value INTEGER)");
// Operations called during a transaction must neither see its uncommitted rows nor join its rollback.
const transaction = database.transaction(async (handle) => {
await handle.run("INSERT INTO call_order (value) VALUES (?)", 1);
await Promise.resolve();
throw new Error("roll back");
});
const beforeWrite = database.all("SELECT value FROM call_order ORDER BY value");
const write = database.run("INSERT INTO call_order (value) VALUES (?)", 2);
const afterWrite = database.all("SELECT value FROM call_order ORDER BY value");
await expect(transaction).rejects.toThrow("roll back");
await write;
expect(await beforeWrite).toEqual([]);
expect(await afterWrite).toEqual([{ value: 2 }]);
const storage = await SqliteStorage.open(database);
const commit = storage.commit(
[{ type: "conversation", value: { id: ROOT_CONVERSATION_ID } }],
BACKGROUND_CONTEXT,
);
const read = storage.conversation(ROOT_CONVERSATION_ID, BACKGROUND_CONTEXT);
await commit;
expect(await read).toEqual({ id: ROOT_CONVERSATION_ID });
await storage.close(BACKGROUND_CONTEXT);
});
it("queues ordinary operations behind an active transaction", async () => {
const database = await openNodeSqliteDatabase(":memory:");
await database.exec("CREATE TABLE operation_queue (value INTEGER)");
@@ -193,21 +221,48 @@ describe("portable SQLite facade settlement", () => {
await database.close();
});
it("rejects database operations from inside a transaction callback instead of waiting forever", async () => {
it("queues database calls made synchronously by a transaction that started immediately", 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",
);
await expect(database.transaction(() => database.close())).rejects.toThrow(
"Cannot close SQLite during an active transaction",
);
await database.exec("CREATE TABLE barrier_probe (value INTEGER)");
let outside!: Promise<void>;
const transaction = database.transaction(async (handle) => {
// Misuse: this call must wait for the transaction instead of joining it.
outside = database.run("INSERT INTO barrier_probe (value) VALUES (?)", 2);
await handle.run("INSERT INTO barrier_probe (value) VALUES (?)", 1);
await Promise.resolve();
throw new Error("roll back");
});
await expect(transaction).rejects.toThrow("roll back");
await outside;
expect(await database.all("SELECT value FROM barrier_probe")).toEqual([{ value: 2 }]);
await database.close();
});
it("lets admitted multi-query reads finish before storage closes", async () => {
const storage = await SqliteStorage.open(await openNodeSqliteDatabase(":memory:"));
const entryId = idFromNumber<EntryId>(2);
await storage.commit(
[
{ type: "conversation", value: { id: ROOT_CONVERSATION_ID } },
{ type: "entry", value: { id: entryId, conversationId: ROOT_CONVERSATION_ID, kind: "probe" } },
],
BACKGROUND_CONTEXT,
);
const scan = storage.scanEntries({ conversationId: ROOT_CONVERSATION_ID }, 10, undefined, BACKGROUND_CONTEXT);
const entry = storage.entry(ROOT_CONVERSATION_ID, entryId, BACKGROUND_CONTEXT);
const head = storage.findLatestHeadMarker(ROOT_CONVERSATION_ID, undefined, BACKGROUND_CONTEXT);
const closed = storage.close(BACKGROUND_CONTEXT);
// A repeated close settles only when the database is closed.
expect(storage.close(BACKGROUND_CONTEXT)).toBe(closed);
expect((await scan).items.map((item) => item.id)).toEqual([entryId]);
expect((await entry)?.entry.kind).toBe("probe");
expect(await head).toBeUndefined();
await closed;
await expect(
storage.scanEntries({ conversationId: ROOT_CONVERSATION_ID }, 10, undefined, BACKGROUND_CONTEXT),
).rejects.toThrow("SqliteStorage is closed");
});
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)");