mirror of
https://github.com/earendil-works/pi.git
synced 2026-10-02 00:35:27 +08:00
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:
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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. */
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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)");
|
||||
|
||||
Reference in New Issue
Block a user