taler-typescript-core

Wallet core logic and WebUIs for various components
Log | Files | Refs | Submodules | README | LICENSE

commit a4f5256e98d70be2fa635ee5037750a1ce256598
parent 291502fd0e53624cc3839005a38f21cec8d5b7c1
Author: Florian Dold <dold@taler.net>
Date:   Sat, 29 Aug 2026 12:28:40 +0200

wallet core: finalize SQLite statements before closing

Diffstat:
Mpackages/idb-bridge/src/SqliteBackend.ts | 31+++++++++++++++++++++++++++++++
Mpackages/idb-bridge/src/browser-sqlite3-impl.test.ts | 14++++++++++++++
Mpackages/idb-bridge/src/browser-sqlite3-impl.ts | 64++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--
Mpackages/idb-bridge/src/node-helper-sqlite3-impl.test.ts | 5+++++
Mpackages/idb-bridge/src/node-helper-sqlite3-impl.ts | 51+++++++++++++++++++++++++++++++++++++++++++--------
Mpackages/idb-bridge/src/sqlite-error-recovery.test.ts | 1+
Mpackages/idb-bridge/src/sqlite3-interface.ts | 3+++
Mpackages/idb-bridge/taler-helper-sqlite3 | 9++++++++-
Mpackages/taler-wallet-core/src/db/indexeddb/handle.ts | 19++++++++++++++++---
Mpackages/taler-wallet-core/src/db/query-sqlite-error-recovery.test.ts | 1+
Mpackages/taler-wallet-core/src/db/sqlite/database.ts | 69++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---
Mpackages/taler-wallet-core/src/db/sqlite/handle.ts | 3++-
Mpackages/taler-wallet-core/src/host-impl.node.ts | 22++++++++++++++++------
Mpackages/taler-wallet-core/src/host-impl.qtart.ts | 80+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++------------
Mpackages/taler-wallet-core/src/requests.test.ts | 10+++++++++-
Mpackages/taler-wallet-core/src/shepherd.ts | 5+++--
Mpackages/taler-wallet-core/src/wallet-db-gate.test.ts | 51+++++++++++++++++++++++++++++++++++++++++++++++++++
Mpackages/taler-wallet-core/src/wallet.ts | 107++++++++++++++++++++++++++++++++++++++++++++-----------------------------------
18 files changed, 459 insertions(+), 86 deletions(-)

diff --git a/packages/idb-bridge/src/SqliteBackend.ts b/packages/idb-bridge/src/SqliteBackend.ts @@ -296,6 +296,8 @@ export class SqliteBackend implements Backend { private sqlPrepCache: Map<string, Sqlite3Statement> = new Map(); + private disposed = false; + enableTracing: boolean = false; constructor( @@ -304,6 +306,9 @@ export class SqliteBackend implements Backend { ) {} private async _prep(sql: string): Promise<Sqlite3Statement> { + if (this.disposed) { + throw Error("sqlite backend is disposed"); + } const stmt = this.sqlPrepCache.get(sql); if (stmt) { return stmt; @@ -314,12 +319,38 @@ export class SqliteBackend implements Backend { } private async _acquireTransactionLevel(level: TransactionLevel) { + if (this.disposed) { + throw Error("sqlite backend is disposed"); + } while (this.txLevel !== TransactionLevel.None) { await this.transactionDoneCond.wait(); } + if (this.disposed) { + throw Error("sqlite backend is disposed"); + } this.txLevel = level; } + /** Finalize cached statements without closing the caller-owned connection. */ + async dispose(): Promise<void> { + if (this.disposed) return; + while (this.txLevel !== TransactionLevel.None) { + await this.transactionDoneCond.wait(); + } + this.disposed = true; + const statements = [...this.sqlPrepCache.values()]; + this.sqlPrepCache.clear(); + let firstError: unknown; + for (const statement of statements) { + try { + await statement.finalize(); + } catch (error) { + firstError ??= error; + } + } + if (firstError) throw firstError; + } + private _releaseTransactionLevel() { this.txLevel = TransactionLevel.None; this.transactionDoneCond.trigger(); diff --git a/packages/idb-bridge/src/browser-sqlite3-impl.test.ts b/packages/idb-bridge/src/browser-sqlite3-impl.test.ts @@ -82,3 +82,17 @@ test("official SQLite WASM implements the shared database contract", async () => await db.close(); } }); + +test("official SQLite WASM finalizes statements before database close", async () => { + const sqlite3 = await sqlite3InitModule(); + const db = await createBrowserSqlite3Impl(sqlite3).open(":memory:"); + const statement = await db.prepare("SELECT 1 AS value"); + assert.deepEqual(await statement.getFirst(), { value: 1 }); + + await db.close(); + await db.close(); + await statement.finalize(); + await assert.rejects(statement.getFirst(), /statement is finalized/); + await assert.rejects(db.exec("SELECT 1"), /database is closed/); + await assert.rejects(db.prepare("SELECT 1"), /database is closed/); +}); diff --git a/packages/idb-bridge/src/browser-sqlite3-impl.ts b/packages/idb-bridge/src/browser-sqlite3-impl.ts @@ -124,8 +124,22 @@ function prepareStatement( sqlite3: Sqlite3Static, db: Database, statement: PreparedStatement, + databaseIsOpen: () => boolean, + onFinalize: () => void, ): Sqlite3Statement { + let finalized = false; + + const ensureUsable = (): void => { + if (finalized) { + throw Error("sqlite3 statement is finalized"); + } + if (!databaseIsOpen()) { + throw Error("sqlite3 database is closed"); + } + }; + const begin = (params: BindParams | undefined): void => { + ensureUsable(); statement.reset(true); const bindings = bindParams(params); if (Object.keys(bindings).length > 0) { @@ -134,6 +148,7 @@ function prepareStatement( }; const reset = (): void => { + if (finalized) return; try { statement.reset(true); } catch { @@ -154,6 +169,15 @@ function prepareStatement( return { internalStatement: statement, + async finalize(): Promise<void> { + if (finalized) return; + finalized = true; + try { + statement.finalize(); + } finally { + onFinalize(); + } + }, async run(params?: BindParams): Promise<RunResult> { return wrap(() => { begin(params); @@ -210,12 +234,38 @@ export function createBrowserSqlite3Impl( } throw error; } + let open = true; + let closePromise: Promise<void> | undefined; + const statements = new Set<Sqlite3Statement>(); + const ensureOpen = (): void => { + if (!open) throw Error("sqlite3 database is closed"); + }; return { internalDbHandle: db, async close(): Promise<void> { - db.close(); + if (!closePromise) { + open = false; + closePromise = (async () => { + let firstError: unknown; + for (const statement of [...statements]) { + try { + await statement.finalize(); + } catch (error) { + firstError ??= error; + } + } + try { + db.close(); + } catch (error) { + firstError ??= error; + } + if (firstError) throw firstError; + })(); + } + await closePromise; }, async exec(sqlStr: string): Promise<void> { + ensureOpen(); try { db.exec(sqlStr); } catch (error) { @@ -223,8 +273,18 @@ export function createBrowserSqlite3Impl( } }, async prepare(stmtStr: string): Promise<Sqlite3Statement> { + ensureOpen(); try { - return prepareStatement(sqlite3, db, db.prepare(stmtStr)); + let wrapped!: Sqlite3Statement; + wrapped = prepareStatement( + sqlite3, + db, + db.prepare(stmtStr), + () => open, + () => statements.delete(wrapped), + ); + statements.add(wrapped); + return wrapped; } catch (error) { throw sqliteError(sqlite3, db, error); } diff --git a/packages/idb-bridge/src/node-helper-sqlite3-impl.test.ts b/packages/idb-bridge/src/node-helper-sqlite3-impl.test.ts @@ -105,5 +105,10 @@ test("sqlite3 helper", async (t) => { ); await db.close(); + await db.close(); + await stmt1.finalize(); + await assert.rejects(stmt1.run(), /statement is finalized/); + await assert.rejects(db.exec("SELECT 1"), /database is closed/); + await assert.rejects(db.prepare("SELECT 1"), /database is closed/); await impl.shutdown(); }); diff --git a/packages/idb-bridge/src/node-helper-sqlite3-impl.ts b/packages/idb-bridge/src/node-helper-sqlite3-impl.ts @@ -468,19 +468,38 @@ export async function createNodeHelperSqlite3Impl( if (enableTracing) { console.error(`opened database ${filename}`); } + let open = true; + let closePromise: Promise<void> | undefined; + const statements = new Set<Sqlite3Statement>(); + const ensureOpen = (): void => { + if (!open) throw Error("sqlite3 database is closed"); + }; return { internalDbHandle: undefined, async close() { - if (enableTracing) { - console.error(`closing database`); + if (!closePromise) { + open = false; + closePromise = (async () => { + for (const statement of [...statements]) { + await statement.finalize(); + } + if (enableTracing) { + console.error(`closing database`); + } + const wr = new Writer(); + wr.writeUint16(myDbId); + const payload = wr.reap(); + const commRes = await helper.communicate( + HelperCmd.CLOSE, + payload, + ); + expectCommunicateSuccess(commRes); + })(); } - const wr = new Writer(); - wr.writeUint16(myDbId); - const payload = wr.reap(); - const commRes = await helper.communicate(HelperCmd.CLOSE, payload); - expectCommunicateSuccess(commRes); + await closePromise; }, async prepare(stmtStr): Promise<Sqlite3Statement> { + ensureOpen(); const myPrepId = counterPrep++; if (enableTracing) { console.error(`preparing statement ${myPrepId}`); @@ -500,9 +519,20 @@ export async function createNodeHelperSqlite3Impl( if (enableTracing) { console.error(`prepared statement ${myPrepId}`); } - return { + let finalized = false; + const ensureUsable = (): void => { + if (finalized) throw Error("sqlite3 statement is finalized"); + ensureOpen(); + }; + const statement: Sqlite3Statement = { internalStatement: undefined, + async finalize(): Promise<void> { + if (finalized) return; + finalized = true; + statements.delete(statement); + }, async getAll(params?: BindParams): Promise<ResultRow[]> { + ensureUsable(); if (enableTracing) { console.error(`running getAll`); } @@ -528,6 +558,7 @@ export async function createNodeHelperSqlite3Impl( async getFirst( params?: BindParams, ): Promise<ResultRow | undefined> { + ensureUsable(); if (enableTracing) { console.error(`running getFirst`); } @@ -551,6 +582,7 @@ export async function createNodeHelperSqlite3Impl( } }, async run(params?: BindParams): Promise<RunResult> { + ensureUsable(); if (enableTracing) { console.error(`running run`); } @@ -592,8 +624,11 @@ export async function createNodeHelperSqlite3Impl( } }, }; + statements.add(statement); + return statement; }, async exec(sqlStr: string): Promise<void> { + ensureOpen(); { if (enableTracing) { console.error(`running execute`); diff --git a/packages/idb-bridge/src/sqlite-error-recovery.test.ts b/packages/idb-bridge/src/sqlite-error-recovery.test.ts @@ -47,6 +47,7 @@ function wrapStatement( ): Sqlite3Statement { return { internalStatement: stmt.internalStatement, + finalize: () => stmt.finalize(), run(params?: BindParams) { faults.check(sql); return stmt.run(params); diff --git a/packages/idb-bridge/src/sqlite3-interface.ts b/packages/idb-bridge/src/sqlite3-interface.ts @@ -2,6 +2,7 @@ export type Sqlite3Database = { internalDbHandle: any; exec(sqlStr: string): Promise<void>; prepare(stmtStr: string): Promise<Sqlite3Statement>; + /** Finalize every owned statement and close the connection, once. */ close(): Promise<void>; }; export type Sqlite3Statement = { @@ -10,6 +11,8 @@ export type Sqlite3Statement = { run(params?: BindParams): Promise<RunResult>; getAll(params?: BindParams): Promise<ResultRow[]>; getFirst(params?: BindParams): Promise<ResultRow | undefined>; + /** Release the native statement. Repeated calls are harmless. */ + finalize(): Promise<void>; }; export interface RunResult { diff --git a/packages/idb-bridge/taler-helper-sqlite3 b/packages/idb-bridge/taler-helper-sqlite3 @@ -254,7 +254,14 @@ while True: continue if cmd == CMD_CLOSE: # close - dbconn.close() + db_id = pr.read_uint16() + closing = db_handles.pop(db_id) + for prep_id, (prep_db, _stmt) in list(prep_handles.items()): + if prep_db is closing: + del prep_handles[prep_id] + closing.close() + if dbconn is closing: + dbconn = None write_resp(req_id, RESP_OK) continue if cmd == CMD_PREPARE: diff --git a/packages/taler-wallet-core/src/db/indexeddb/handle.ts b/packages/taler-wallet-core/src/db/indexeddb/handle.ts @@ -91,6 +91,8 @@ export class IdbWalletDbHandle implements WalletDbHandle { private opening: | Promise<{ fixupsApplied: number; schemaUpgraded: boolean }> | undefined; + private closed = false; + private closePromise: Promise<void> | undefined; private notify: (n: WalletNotification) => void = () => {}; @@ -132,6 +134,7 @@ export class IdbWalletDbHandle implements WalletDbHandle { */ private rawStats?: () => AccessStats | undefined, private applyDbFixups: typeof applyFixups = applyFixups, + private disposeBackend?: () => Promise<void>, ) {} /** @@ -144,6 +147,9 @@ export class IdbWalletDbHandle implements WalletDbHandle { fixupsApplied: number; schemaUpgraded: boolean; }> { + if (this.closed) { + throw Error("wallet database handle is closed"); + } if (this.dbAccess) { return { fixupsApplied: 0, schemaUpgraded: false }; } @@ -392,8 +398,15 @@ export class IdbWalletDbHandle implements WalletDbHandle { } async close(): Promise<void> { - this.idbHandle?.close(); - this.idbHandle = undefined; - this.dbAccess = undefined; + if (!this.closePromise) { + this.closed = true; + this.closePromise = (async () => { + this.idbHandle?.close(); + this.idbHandle = undefined; + this.dbAccess = undefined; + await this.disposeBackend?.(); + })(); + } + await this.closePromise; } } diff --git a/packages/taler-wallet-core/src/db/query-sqlite-error-recovery.test.ts b/packages/taler-wallet-core/src/db/query-sqlite-error-recovery.test.ts @@ -45,6 +45,7 @@ function wrapStatement( ): Sqlite3Statement { return { internalStatement: stmt.internalStatement, + finalize: () => stmt.finalize(), run(params?: BindParams) { faults.check(sql); return stmt.run(params); diff --git a/packages/taler-wallet-core/src/db/sqlite/database.ts b/packages/taler-wallet-core/src/db/sqlite/database.ts @@ -55,6 +55,7 @@ const logger = new Logger("db/sqlite/database.ts"); * surrounding work; use prepared statements throughout instead. */ export class SqliteTxControl { + private finalized = false; private constructor( private beginStmt: Sqlite3Statement, private commitStmt: Sqlite3Statement, @@ -78,6 +79,24 @@ export class SqliteTxControl { async rollback(): Promise<void> { await this.rollbackStmt.run({}); } + + async finalize(): Promise<void> { + if (this.finalized) return; + this.finalized = true; + let firstError: unknown; + for (const statement of [ + this.beginStmt, + this.commitStmt, + this.rollbackStmt, + ]) { + try { + await statement.finalize(); + } catch (error) { + firstError ??= error; + } + } + if (firstError) throw firstError; + } } /** @@ -283,6 +302,8 @@ export interface NativeSqliteWalletDb { * for concurrent callers. */ lock: TxQueue; + state: "open" | "closing" | "closed"; + closePromise?: Promise<void>; } /** @@ -347,9 +368,45 @@ export async function openNativeSqliteWalletDb( lock: new TxQueue(), stmtCache: new Map(), stats: { rowsRead: 0 }, + state: "open", }; } +/** Finalize all wallet-owned statements and close the shared connection. */ +export async function closeNativeSqliteWalletDb( + ndb: NativeSqliteWalletDb, +): Promise<void> { + if (!ndb.closePromise) { + ndb.state = "closing"; + ndb.closePromise = ndb.lock.run(async () => { + let firstError: unknown; + try { + await ndb.txc.finalize(); + } catch (error) { + firstError ??= error; + } + const statements = [...ndb.stmtCache.values()]; + ndb.stmtCache.clear(); + for (const statement of statements) { + try { + await statement.finalize(); + } catch (error) { + firstError ??= error; + } + } + try { + await ndb.db.close(); + } catch (error) { + firstError ??= error; + } finally { + ndb.state = "closed"; + } + if (firstError) throw firstError; + }); + } + await ndb.closePromise; +} + /** * Run f in one native sqlite transaction. * @@ -361,9 +418,15 @@ export async function runNativeSqliteWalletTx<T>( notifyFn: (n: WalletNotification) => void, f: (tx: SqliteWalletTransaction) => Promise<T>, ): Promise<T> { - return await ndb.lock.run(() => - runNativeSqliteWalletTxLocked(ndb, notifyFn, f), - ); + if (ndb.state !== "open") { + throw Error("native sqlite wallet database is closed"); + } + return await ndb.lock.run(() => { + if (ndb.state !== "open") { + throw Error("native sqlite wallet database is closed"); + } + return runNativeSqliteWalletTxLocked(ndb, notifyFn, f); + }); } async function runNativeSqliteWalletTxLocked<T>( diff --git a/packages/taler-wallet-core/src/db/sqlite/handle.ts b/packages/taler-wallet-core/src/db/sqlite/handle.ts @@ -23,6 +23,7 @@ import { } from "../handle.js"; import { WalletDbTransaction } from "../transaction.js"; import { + closeNativeSqliteWalletDb, clearNativeSqliteWalletDb, exportNativeSqliteDb, importNativeSqliteDb, @@ -131,6 +132,6 @@ export class SqliteWalletDbHandle implements WalletDbHandle { } async close(): Promise<void> { - await this.ndb.db.close(); + await closeNativeSqliteWalletDb(this.ndb); } } diff --git a/packages/taler-wallet-core/src/host-impl.node.ts b/packages/taler-wallet-core/src/host-impl.node.ts @@ -173,9 +173,15 @@ async function makeSqliteDb( if (process.env.TALER_WALLET_STATS) { myBackend.trackStats = true; } + let connectionTransferred = false; const handle = new IdbWalletDbHandle( new BridgeIDBFactory(myBackend), () => myBackend.accessStats, + undefined, + async () => { + await myBackend.dispose(); + if (!connectionTransferred) await db.close(); + }, ); handle.exportToFile = async (directory, stem, forceFormat) => { if (forceFormat != null && forceFormat !== "sqlite3") { @@ -186,13 +192,17 @@ async function makeSqliteDb( return { path }; }; handle.getDiagnosticStats = () => myBackend.accessStats; - handle.migrateToNative = async (options) => - addNodeDatabaseCapabilities( - (await migrateWalletDbToNative(db, handle, options)).handle, - ); + handle.migrateToNative = async (options) => { + const result = await migrateWalletDbToNative(db, handle, options); + connectionTransferred = true; + return addNodeDatabaseCapabilities(result.handle); + }; if (kind === "empty") { - handle.openNativeIfEmpty = async () => - addNodeDatabaseCapabilities(await openNativeWalletDbForEmptyStorage(db)); + handle.openNativeIfEmpty = async () => { + const result = await openNativeWalletDbForEmptyStorage(db); + connectionTransferred = true; + return addNodeDatabaseCapabilities(result); + }; } return addNodeDatabaseCapabilities(handle); } diff --git a/packages/taler-wallet-core/src/host-impl.qtart.ts b/packages/taler-wallet-core/src/host-impl.qtart.ts @@ -120,30 +120,73 @@ export async function createQtartSqlite3Impl(): Promise<Sqlite3Interface> { return { async open(filename: string) { const internalDbHandle = tart.sqlite3Open(filename); + let open = true; + let closePromise: Promise<void> | undefined; + const statements = new Set<Sqlite3Statement>(); + const ensureOpen = (): void => { + if (!open) throw Error("sqlite3 database is closed"); + }; return { internalDbHandle, async close() { - tart.sqlite3Close(internalDbHandle); + if (!closePromise) { + open = false; + closePromise = (async () => { + let firstError: unknown; + for (const statement of [...statements]) { + try { + await statement.finalize(); + } catch (error) { + firstError ??= error; + } + } + try { + tart.sqlite3Close(internalDbHandle); + } catch (error) { + firstError ??= error; + } + if (firstError) throw firstError; + })(); + } + await closePromise; }, async prepare(stmtStr): Promise<Sqlite3Statement> { + ensureOpen(); const stmtHandle = tart.sqlite3Prepare(internalDbHandle, stmtStr); - return { + let finalized = false; + const ensureUsable = (): void => { + if (finalized) throw Error("sqlite3 statement is finalized"); + ensureOpen(); + }; + const statement: Sqlite3Statement = { internalStatement: stmtHandle, + async finalize(): Promise<void> { + if (finalized) return; + finalized = true; + statements.delete(statement); + tart.sqlite3Finalize(stmtHandle); + }, async getAll(params): Promise<ResultRow[]> { + ensureUsable(); numStmt++; return tart.sqlite3StmtGetAll(stmtHandle, params); }, async getFirst(params): Promise<ResultRow | undefined> { + ensureUsable(); numStmt++; return tart.sqlite3StmtGetFirst(stmtHandle, params); }, async run(params) { + ensureUsable(); numStmt++; return tart.sqlite3StmtRun(stmtHandle, params); }, }; + statements.add(statement); + return statement; }, async exec(sqlStr): Promise<void> { + ensureOpen(); numStmt++; tart.sqlite3Exec(internalDbHandle, sqlStr); }, @@ -191,10 +234,19 @@ async function makeSqliteDb( const myBackend = await createSqliteBackendOverDb(imp, db); myBackend.trackStats = true; myBackend.enableTracing = false; - const handle = new IdbWalletDbHandle(new BridgeIDBFactory(myBackend), () => ({ - ...myBackend.accessStats, - primitiveStatements: numStmt, - })); + let connectionTransferred = false; + const handle = new IdbWalletDbHandle( + new BridgeIDBFactory(myBackend), + () => ({ + ...myBackend.accessStats, + primitiveStatements: numStmt, + }), + undefined, + async () => { + await myBackend.dispose(); + if (!connectionTransferred) await db.close(); + }, + ); handle.getDiagnosticStats = () => ({ ...myBackend.accessStats, primitiveStatements: numStmt, @@ -207,13 +259,17 @@ async function makeSqliteDb( await myBackend.backupToFile(path); return { path }; }; - handle.migrateToNative = async (options) => - addQtartDatabaseCapabilities( - (await migrateWalletDbToNative(db, handle, options)).handle, - ); + handle.migrateToNative = async (options) => { + const result = await migrateWalletDbToNative(db, handle, options); + connectionTransferred = true; + return addQtartDatabaseCapabilities(result.handle); + }; if (kind === "empty") { - handle.openNativeIfEmpty = async () => - addQtartDatabaseCapabilities(await openNativeWalletDbForEmptyStorage(db)); + handle.openNativeIfEmpty = async () => { + const result = await openNativeWalletDbForEmptyStorage(db); + connectionTransferred = true; + return addQtartDatabaseCapabilities(result); + }; } return addQtartDatabaseCapabilities(handle); } diff --git a/packages/taler-wallet-core/src/requests.test.ts b/packages/taler-wallet-core/src/requests.test.ts @@ -141,8 +141,16 @@ for (const [expectedBackend, makeRunner] of backendCases) { { migrated: false, databaseBackend: "sqlite" }, ); } + await wallet.client.call(WalletApiOperation.Shutdown, {}); + await wallet.client.call(WalletApiOperation.Shutdown, {}); + await assert.rejects( + wallet.client.call(WalletApiOperation.GetBalances, {}), + (error: unknown) => + error instanceof TalerError && + error.errorDetail.code === TalerErrorCode.WALLET_CORE_NOT_AVAILABLE, + ); } finally { - if (initialized) { + if (initialized && closeCalls === 0) { await wallet.client.call(WalletApiOperation.Shutdown, {}); } if (closeCalls === 0) { diff --git a/packages/taler-wallet-core/src/shepherd.ts b/packages/taler-wallet-core/src/shepherd.ts @@ -233,11 +233,12 @@ export class TaskSchedulerImpl implements TaskScheduler { } async shutdown(): Promise<void> { - const tasksIds = [...this.sheps.keys()]; + const tasks = [...this.sheps.entries()]; logger.info(`Stopping task shepherd.`); - for (const taskId of tasksIds) { + for (const [taskId] of tasks) { this.stopShepherdTask(taskId); } + await Promise.all(tasks.map(([, info]) => info.latch)); } /** diff --git a/packages/taler-wallet-core/src/wallet-db-gate.test.ts b/packages/taler-wallet-core/src/wallet-db-gate.test.ts @@ -216,3 +216,54 @@ test("database gate cancels an exclusive waiter", async () => { await shared; await gate.runShared(async () => {}); }); + +test("terminal database close drains active work and rejects queued work", async () => { + const gate = new DbOperationGate(); + const active = deferred(); + const releaseActive = deferred(); + const closeStarted = deferred(); + const shared = gate.runShared(async () => { + active.resolve(); + await releaseActive.promise; + }); + await active.promise; + + const closing = gate.runTerminalExclusive(async () => { + closeStarted.resolve(); + }); + const queued = gate.runShared(async () => { + assert.fail("queued database operation was admitted after close"); + }); + await Promise.resolve(); + releaseActive.resolve(); + await shared; + await closeStarted.promise; + await closing; + await assert.rejects(queued, /wallet database is closed/); + await assert.rejects( + gate.runShared(async () => {}), + /wallet database is closed/, + ); + await assert.rejects( + gate.runExclusive(async () => {}), + /wallet database is closed/, + ); +}); + +test("admitted database close is idempotent", async () => { + const gate = new DbOperationGate(); + const handle = fakeHandle("db", []); + let closeCalls = 0; + handle.close = async () => { + closeCalls++; + }; + const admitted = new AdmittedWalletDbHandle(() => handle, gate); + + await Promise.all([admitted.close(), admitted.close()]); + await admitted.close(); + assert.strictEqual(closeCalls, 1); + await assert.rejects( + admitted.runReadWriteTx(async () => {}), + /wallet database is closed/, + ); +}); diff --git a/packages/taler-wallet-core/src/wallet.ts b/packages/taler-wallet-core/src/wallet.ts @@ -161,6 +161,7 @@ export class DbOperationGate { private exclusive = false; private waitingExclusive = 0; private waiters: Array<() => void> = []; + private closed = false; private async changed(): Promise<void> { await new Promise<void>((resolve) => this.waiters.push(resolve)); @@ -173,7 +174,9 @@ export class DbOperationGate { } async runShared<T>(f: () => Promise<T>): Promise<T> { + if (this.closed) throw Error("wallet database is closed"); while (this.exclusive || this.waitingExclusive !== 0) await this.changed(); + if (this.closed) throw Error("wallet database is closed"); this.shared++; try { return await f(); @@ -186,6 +189,7 @@ export class DbOperationGate { async acquireExclusive( cancellationToken: CancellationToken = CancellationToken.CONTINUE, ): Promise<() => void> { + if (this.closed) throw Error("wallet database is closed"); this.waitingExclusive++; try { while (this.exclusive || this.shared !== 0) { @@ -193,6 +197,7 @@ export class DbOperationGate { await cancellationToken.racePromise(this.changed()); } cancellationToken.throwIfCancelled(); + if (this.closed) throw Error("wallet database is closed"); this.exclusive = true; } finally { this.waitingExclusive--; @@ -218,10 +223,33 @@ export class DbOperationGate { release(); } } + + /** Drain admitted work, reject future work, and run the terminal close. */ + async runTerminalExclusive<T>(f: () => Promise<T>): Promise<T> { + if (this.closed) throw Error("wallet database is closed"); + this.waitingExclusive++; + try { + while (this.exclusive || this.shared !== 0) await this.changed(); + if (this.closed) throw Error("wallet database is closed"); + this.exclusive = true; + this.closed = true; + } finally { + this.waitingExclusive--; + this.wake(); + } + try { + return await f(); + } finally { + this.exclusive = false; + this.wake(); + } + } } /** Stable facade that resolves the current handle only after shared admission. */ export class AdmittedWalletDbHandle implements WalletDbHandle { + private closePromise: Promise<void> | undefined; + get name(): string { return this.current().name; } @@ -302,7 +330,12 @@ export class AdmittedWalletDbHandle implements WalletDbHandle { this.current().emitNotification(notification); } close(): Promise<void> { - return this.gate.runExclusive(() => this.current().close()); + if (!this.closePromise) { + this.closePromise = this.gate.runTerminalExclusive(() => + this.current().close(), + ); + } + return this.closePromise; } } @@ -919,6 +952,17 @@ async function dispatchWalletCoreApiRequest( id: string, payload: unknown, ): Promise<CoreApiResponse> { + if (ws.stopped) { + if (operation === WalletApiOperation.Shutdown) { + await ws.shutdown(); + return { type: "response", operation, id, result: {} }; + } + throw TalerError.fromDetail( + TalerErrorCode.WALLET_CORE_NOT_AVAILABLE, + {}, + "wallet core has been shut down", + ); + } const isInitOperation = isWalletInitOperation(operation); if (!isInitOperation) { if (!ws.initCalled) { @@ -1310,7 +1354,8 @@ export class InternalWalletState { private admittedDb: WalletDbHandle; - private suspendedDbRelease: (() => void) | undefined; + private taskShutdownPromise: Promise<void> | undefined; + private shutdownPromise: Promise<void> | undefined; private maintenanceNotifications: MaintenanceNotificationThrottler; @@ -1688,46 +1733,6 @@ export class InternalWalletState { } } - /** - * Prepare database for import by closing it. - */ - async suspendDatabase(): Promise<void> { - if (this.loadingDb) { - while (this.loadingDb) { - await this.loadingDbCond.wait(); - } - } - this.loadingDb = true; - this.suspendedDbRelease = await this.dbOperationGate.acquireExclusive(); - try { - await this.dbHandle.close(); - } catch (e) { - this.suspendedDbRelease(); - this.suspendedDbRelease = undefined; - this.loadingDb = false; - this.loadingDbCond.trigger(); - throw e; - } - } - - /** - * Resume database by re-opening it. - */ - async resumeDatabase(): Promise<void> { - const release = this.suspendedDbRelease; - this.suspendedDbRelease = undefined; - try { - const idb = this.idbOnly; - if (idb) { - await this.finalizeIndexedDbOpen(idb, await idb.ensureOpen()); - } - } finally { - this.loadingDb = false; - this.loadingDbCond.trigger(); - release?.(); - } - } - notify(n: WalletNotification): void { logger.trace(`Notification: ${j2s(n)}`); if (n.type === NotificationType.DatabaseMaintenanceProgress) { @@ -1769,22 +1774,30 @@ export class InternalWalletState { * Stop ongoing processing. */ stop(): void { + if (this.stopped) return; logger.trace("stopping (at internal wallet state)"); this.stopped = true; this.maintenanceNotifications.stop(); this.timerGroup.stopCurrentAndFutureTimers(); this.cryptoDispatcher.stop(); - this.taskScheduler.shutdown().catch((e) => { + this.taskShutdownPromise = this.taskScheduler.shutdown(); + this.taskShutdownPromise.catch((e) => { logger.warn(`shutdown failed: ${safeStringifyException(e)}`); }); } /** Stop processing and release the database before shutdown completes. */ async shutdown(): Promise<void> { - this.stop(); - await this.dbOperationGate.runExclusive(async () => { - await this.dbHandle.close(); - }); + if (!this.shutdownPromise) { + this.shutdownPromise = (async () => { + this.stop(); + await this.taskShutdownPromise; + await this.dbOperationGate.runTerminalExclusive(async () => { + await this.dbHandle.close(); + }); + })(); + } + await this.shutdownPromise; } /**