Files
paperclip/patches/postgres@3.4.9.patch
Devin FoleyandPaperclip 685d4faba3 Fix PostgreSQL recovery after a transaction connection closes (#13643)
Reject queued and late work from disconnected transaction and reservation
scopes. Keep closed reservations out of the open pool, and clear old
connection buffers and responses so new requests can reconnect safely.

Twelve real-PostgreSQL regression cases cover crash prevention, recovery,
and transaction isolation in both ESM and CommonJS. Database checks and
all PR CI checks pass. Greptile: 5/5, no unresolved comments.

Co-Authored-By: Paperclip <noreply@paperclip.ing>
2026-09-18 16:27:28 -07:00

171 lines
5.7 KiB
Diff

diff --git a/cjs/src/connection.js b/cjs/src/connection.js
index 07f6716702ac888215c86ad6e0071aaf55ae519c..dc24963ce2dd0a219f366f80ba9847548ba23a52 100644
--- a/cjs/src/connection.js
+++ b/cjs/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/cjs/src/index.js b/cjs/src/index.js
index f09c61c7e998f411560b95c9a8e1ed6a97972c11..04dded5c7e2cdc34ac01b29a87cc93248ec8c9c6 100644
--- a/cjs/src/index.js
+++ b/cjs/src/index.js
@@ -216,8 +216,19 @@ function Postgres(a, b) {
: move(c, reserved)
c.reserved.release = true
+ let closedError
+ c.onclose = error => {
+ closedError = error
+ while (queue.length)
+ queue.shift().reject(error)
+ }
+
const sql = Sql(handler)
sql.release = () => {
+ if (closedError)
+ return
+ closedError = Errors.connection('CONNECTION_RELEASED', options)
+ c.onclose = null
c.reserved = null
onopen(c)
}
@@ -225,6 +236,8 @@ function Postgres(a, b) {
return sql
function handler(q) {
+ if (closedError)
+ return q.reject(closedError)
c.queue === full
? queue.push(q)
: c.execute(q) || move(c, full)
@@ -236,13 +249,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 +309,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)
diff --git a/src/connection.js b/src/connection.js
index 1b1cccde43b4570d5d071d6ffaab2669b3c2065a..c95c326a2747c680c2c22b04b1c5f2e27a4e8c90 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 c7fba3dace20173f4cb13f3f138bb6d213b43683..f0888986add8ba02c8853f792ae42dd1cec88c82 100644
--- a/src/index.js
+++ b/src/index.js
@@ -216,8 +216,19 @@ function Postgres(a, b) {
: move(c, reserved)
c.reserved.release = true
+ let closedError
+ c.onclose = error => {
+ closedError = error
+ while (queue.length)
+ queue.shift().reject(error)
+ }
+
const sql = Sql(handler)
sql.release = () => {
+ if (closedError)
+ return
+ closedError = Errors.connection('CONNECTION_RELEASED', options)
+ c.onclose = null
c.reserved = null
onopen(c)
}
@@ -225,6 +236,8 @@ function Postgres(a, b) {
return sql
function handler(q) {
+ if (closedError)
+ return q.reject(closedError)
c.queue === full
? queue.push(q)
: c.execute(q) || move(c, full)
@@ -236,13 +249,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 +309,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)