From 51649fbf5f6c16e2ad7d6ccaee6c49d84919b94f Mon Sep 17 00:00:00 2001 From: ImBenji Date: Mon, 17 Aug 2026 14:02:59 +0100 Subject: [PATCH] fix: pin postgres transaction clients --- src/db/runtime.js | 30 +++++++++++++++++++++++------- 1 file changed, 23 insertions(+), 7 deletions(-) diff --git a/src/db/runtime.js b/src/db/runtime.js index 11f505f..54b025d 100644 --- a/src/db/runtime.js +++ b/src/db/runtime.js @@ -32,6 +32,16 @@ function querySync(pool, sql, params = []) { return result; } +function connectSync(pool) { + let done = false; + let client; + let error; + pool.connect().then((value) => { client = value; done = true; }).catch((err) => { error = err; done = true; }); + deasync.loopWhile(() => !done); + if (error) throw error; + return client; +} + function normalizeParams(params) { if (params.length === 1 && params[0] && typeof params[0] === 'object' && !Array.isArray(params[0]) && !Buffer.isBuffer(params[0])) { return params[0]; @@ -94,18 +104,19 @@ class PgCompatDb { prepare(sql) { const db = this; + const target = () => db.activeClient || db.pool; return { get(...rawParams) { const { sql: text, params } = rewriteSql(sql, normalizeParams(rawParams)); - return querySync(db.pool, text, params).rows[0]; + return querySync(target(), text, params).rows[0]; }, all(...rawParams) { const { sql: text, params } = rewriteSql(sql, normalizeParams(rawParams)); - return querySync(db.pool, text, params).rows; + return querySync(target(), text, params).rows; }, run(...rawParams) { const { sql: text, params } = rewriteSql(sql, normalizeParams(rawParams)); - const result = querySync(db.pool, text, params); + const result = querySync(target(), text, params); return { changes: result.rowCount || 0, lastInsertRowid: result.rows?.[0]?.id ?? null }; }, }; @@ -115,21 +126,26 @@ class PgCompatDb { const statements = String(sql).split(';').map((statement) => statement.trim()).filter(Boolean); for (const statement of statements) { const { sql: text, params } = rewriteSql(statement, []); - querySync(this.pool, text, params); + querySync(this.activeClient || this.pool, text, params); } } transaction(fn) { const db = this; const run = (...args) => { - querySync(db.pool, 'BEGIN', []); + const client = connectSync(db.pool); try { + db.activeClient = client; + querySync(client, 'BEGIN', []); const result = fn(...args); - querySync(db.pool, 'COMMIT', []); + querySync(client, 'COMMIT', []); return result; } catch (error) { - querySync(db.pool, 'ROLLBACK', []); + try { querySync(client, 'ROLLBACK', []); } catch (_) {} throw error; + } finally { + db.activeClient = null; + client.release(); } }; run.immediate = run;