From d3de55007cbcc0b8d1a9d6c7c9073d163b303b36 Mon Sep 17 00:00:00 2001 From: realies <5107843+realies@users.noreply.github.com> Date: Sat, 26 Sep 2026 12:25:16 +0300 Subject: [PATCH] Fix reserves stranded after connection termination Keep reserve waiters queued when a connection attempt starts, including reconnects after a server-side close. Clear the closed session query state so a stale fatal error cannot reject a new reserve. Add real-server regressions for queued reserves and terminated in-flight queries. --- src/connection.js | 1 + src/index.js | 8 +++++--- tests/index.js | 36 ++++++++++++++++++++++++++++++++++++ 3 files changed, 42 insertions(+), 3 deletions(-) diff --git a/src/connection.js b/src/connection.js index 10ab1bb..2bc797e 100644 --- a/src/connection.js +++ b/src/connection.js @@ -451,6 +451,7 @@ 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 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..8b62a7f 100644 --- a/src/index.js +++ b/src/index.js @@ -205,9 +205,10 @@ function Postgres(a, b) { const c = open.length ? open.shift() : await new Promise((resolve, reject) => { - const query = { reserve: resolve, reject } - queries.push(query) - closed.length && connect(closed.shift(), query) + const query = { reserve: resolve, reject: e => (queries.remove(query), reject(e)) } + closed.length + ? connect(closed.shift(), query) + : queries.push(query) }) move(c, reserved) @@ -389,6 +390,7 @@ function Postgres(a, b) { } function connect(c, query) { + query.reserve && queries.push(query) move(c, connecting) c.connect(query) return c diff --git a/tests/index.js b/tests/index.js index 845c6ce..194f397 100644 --- a/tests/index.js +++ b/tests/index.js @@ -2666,6 +2666,42 @@ t('concurrent cursors multiple connections', async() => { return ['12233445566778', xs.sort().join('')] }) +t('reserve queued behind a terminated connection', async() => { + const pool = postgres(options) + const reserved = await pool.reserve() + const [{ pid }] = await reserved`select pg_backend_pid() as pid` + const pending = pool.reserve() + + await sql`select pg_terminate_backend(${ pid })` + + const next = await pending + const [{ x }] = await next`select 1 as x` + next.release() + await pool.end() + + return [1, x] +}) + +t('reserve after a terminated query', async() => { + const pool = postgres(options) + const reserved = await pool.reserve() + const [{ pid }] = await reserved`select pg_backend_pid() as pid` + const query = reserved`select pg_sleep(30)`.catch(e => e) + + while (!(await sql`select wait_event = 'PgSleep' as sleeping from pg_stat_activity where pid = ${ pid }`)[0].sleeping) + await delay(1) + + await sql`select pg_terminate_backend(${ pid })` + await query + + const next = await pool.reserve() + const [{ x }] = await next`select 1 as x` + next.release() + await pool.end() + + return [1, x] +}) + t('reserve connection', async() => { const reserved = await sql.reserve()