From 0282301130213f481669e0f5e2750fff3bf0fe45 Mon Sep 17 00:00:00 2001 From: Adam Blumenfeld Date: Tue, 8 Sep 2026 23:56:26 -0700 Subject: [PATCH 1/2] Add regression tests for disconnected transactions --- tests/index.js | 75 ++++++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 75 insertions(+) diff --git a/tests/index.js b/tests/index.js index 845c6ce..cf517d7 100644 --- a/tests/index.js +++ b/tests/index.js @@ -2501,6 +2501,81 @@ t('Ensure transactions throw if connection is closed dwhile there is no query', return ['CONNECTION_CLOSED', x.code] }) +t('Disconnect rejects queued transaction queries and allows reconnect', async() => { + const sql = postgres({ ...options, max_pipeline: 1, fetch_types: false }) + let queries + + try { + const error = await sql.begin(sql => { + queries = [ + sql`select pg_terminate_backend(pg_backend_pid())`.execute(), + sql`select 1`.execute(), + sql`select 2`.execute() + ] + return Promise.allSettled(queries) + }).catch(x => x) + + const results = await Promise.allSettled(queries) + const [{ x }] = await sql`select 1 as x` + return [ + 'CONNECTION_CLOSED,rejected,rejected,rejected,1', + [error.code, ...results.map(x => x.status), x].join(',') + ] + } finally { + await sql.end({ timeout: 0 }) + } +}) + +t('Disconnected transaction cannot query a reused connection', async() => + withDisconnectedTransaction(({ sql, disconnected }) => sql.begin(async sql => { + await sql`select set_config('postgres_js.test', 'replacement', true)` + const result = await disconnected`select current_setting('postgres_js.test') as x`.catch(x => x) + return ['CONNECTION_CLOSED', result.code] + })) +) + +t('Disconnected transaction cannot commit a reused connection', async() => + finishDisconnectedTransaction() +) + +t('Disconnected transaction cannot roll back a reused connection', async() => + finishDisconnectedTransaction(new Error('original callback failed')) +) + +function finishDisconnectedTransaction(error) { + return withDisconnectedTransaction(({ sql, finish }) => sql.begin(async sql => { + const [{ x: before }] = await sql`select txid_current()::text as x` + finish(error) + await new Promise(resolve => setImmediate(resolve)) + const [{ x: after }] = await sql`select txid_current()::text as x` + return [before, after] + })) +} + +async function withDisconnectedTransaction(fn) { + const pool = postgres({ ...options, fetch_types: false }) + let finish + , ready + const gate = new Promise((resolve, reject) => finish = error => error ? reject(error) : resolve()) + const connected = new Promise(resolve => ready = resolve) + const failed = pool.begin(async sql => { + const [{ pid }] = await sql`select pg_backend_pid() as pid` + ready({ disconnected: sql, pid }) + await gate + }).catch(x => x) + + try { + const { disconnected, pid } = await Promise.race([connected, failed.then(error => { throw error })]) + await sql`select pg_terminate_backend(${ pid }::int)` + await failed + return await fn({ sql: pool, disconnected, finish }) + } finally { + finish() + await new Promise(resolve => setImmediate(resolve)) + await pool.end({ timeout: 0 }) + } +} + t('Custom socket', {}, async() => { let result const sql = postgres({ From 18cffe016927e01481b9e746587da789f4260a4c Mon Sep 17 00:00:00 2001 From: Adam Blumenfeld Date: Tue, 8 Sep 2026 23:56:26 -0700 Subject: [PATCH 2/2] Reject transaction work after its connection closes --- src/connection.js | 4 ++++ src/index.js | 10 +++++++++- 2 files changed, 13 insertions(+), 1 deletion(-) diff --git a/src/connection.js b/src/connection.js index 10ab1bb..a5f9ea1 100644 --- a/src/connection.js +++ b/src/connection.js @@ -438,6 +438,7 @@ function Connection(options, queues = {}, { onopen = noop, onend = noop, onclose remaining = 0 incomings = null clearImmediate(nextWriteTimer) + chunk = nextWriteTimer = null socket.removeListener('data', data) socket.removeListener('connect', connected) idleTimer.cancel() @@ -451,6 +452,9 @@ function Connection(options, queues = {}, { onopen = noop, onend = noop, onclose return reconnect() !hadError && (query || sent.length) && error(Errors.connection('CONNECTION_CLOSED', options, socket)) + query = results = errorResponse = null + result = new Result() + rows = 0 closedTime = performance.now() hadError && options.shared.retries++ delay = (typeof backoff === 'function' ? backoff(options.shared.retries) : backoff) * 1000 diff --git a/src/index.js b/src/index.js index c7fba3d..d6b5aee 100644 --- a/src/index.js +++ b/src/index.js @@ -236,13 +236,19 @@ function Postgres(a, b) { const queries = Queue() let savepoints = 0 , connection + , closedError , prepare = null try { await sql.unsafe('begin ' + options.replace(/[^a-z ]/ig, ''), [], { onexecute }).execute() return await Promise.race([ scope(connection, fn), - new Promise((_, reject) => connection.onclose = reject) + new Promise((_, reject) => connection.onclose = error => { + closedError = error + while (queries.length) + queries.shift().reject(error) + reject(error) + }) ]) } catch (error) { throw error @@ -290,6 +296,8 @@ function Postgres(a, b) { function handler(q) { q.catch(e => uncaughtError || (uncaughtError = e)) + if (closedError) + return q.reject(closedError) c.queue === full ? queries.push(q) : c.execute(q) || move(c, full)