From 685d4faba3715197fef85e7634f415b589f67812 Mon Sep 17 00:00:00 2001 From: Devin Foley Date: Fri, 18 Sep 2026 16:27:28 -0700 Subject: [PATCH] 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 --- doc/DATABASE.md | 15 ++ package.json | 3 +- .../postgres-connection-recovery.mjs | 122 +++++++++++++ .../src/postgres-connection-recovery.test.ts | 42 +++++ patches/postgres@3.4.9.patch | 170 ++++++++++++++++++ pnpm-workspace.yaml | 1 + 6 files changed, 352 insertions(+), 1 deletion(-) create mode 100644 packages/db/src/__fixtures__/postgres-connection-recovery.mjs create mode 100644 packages/db/src/postgres-connection-recovery.test.ts create mode 100644 patches/postgres@3.4.9.patch diff --git a/doc/DATABASE.md b/doc/DATABASE.md index 53a9401f19..1bd7c6bacd 100644 --- a/doc/DATABASE.md +++ b/doc/DATABASE.md @@ -144,6 +144,21 @@ DATABASE_URL=postgres://postgres.[PROJECT-REF]:[PASSWORD]@...5432/postgres \ See [Supabase pricing](https://supabase.com/pricing) for current details. +## Connection loss during a transaction + +When a database connection closes, its transaction fails. Paperclip does not +replay that transaction. New requests can use a fresh connection from the pool. +Queries from the failed transaction must keep failing, even after the pool +reconnects. + +Source builds carry `patches/postgres@3.4.9.patch` for this behavior. It rejects +queued and later queries from a disconnected transaction or reserved connection, +and prevents a released, closed connection from returning to the open pool. +The patch covers both ESM and CommonJS. The regression suite terminates real +PostgreSQL backends and checks rejection, pool recovery, and transaction isolation. +Remove the patch when an upstream release passes these tests. Installs of the +unmodified `postgres` package outside this workspace do not include the patch. + ## Switching between modes The database mode is controlled by `DATABASE_URL`: diff --git a/package.json b/package.json index cb34da76f1..92b33a0d25 100644 --- a/package.json +++ b/package.json @@ -112,7 +112,8 @@ "@chat-adapter/slack@4.39.0": "patches/@chat-adapter__slack@4.39.0.patch", "@chat-adapter/discord@4.39.0": "patches/@chat-adapter__discord@4.39.0.patch", "@chat-adapter/github@4.39.0": "patches/@chat-adapter__github@4.39.0.patch", - "@discordjs/ws@1.2.3": "patches/@discordjs__ws@1.2.3.patch" + "@discordjs/ws@1.2.3": "patches/@discordjs__ws@1.2.3.patch", + "postgres@3.4.9": "patches/postgres@3.4.9.patch" }, "overrides": { "@agentclientprotocol/codex-acp@1.6.2>@openai/codex": "0.153.4", diff --git a/packages/db/src/__fixtures__/postgres-connection-recovery.mjs b/packages/db/src/__fixtures__/postgres-connection-recovery.mjs new file mode 100644 index 0000000000..64ee9db814 --- /dev/null +++ b/packages/db/src/__fixtures__/postgres-connection-recovery.mjs @@ -0,0 +1,122 @@ +import assert from "node:assert/strict"; +import { createRequire } from "node:module"; +import { setImmediate } from "node:timers/promises"; + +const postgres = process.env.DRIVER_FORMAT === "cjs" + ? createRequire(import.meta.url)("postgres") + : (await import("postgres")).default; +const mode = process.env.RECOVERY_CASE; +const closed = Promise.withResolvers(); +const sql = postgres(process.env.TEST_DATABASE_URL, { + max: 1, + max_pipeline: 1, + backoff: 0, + onnotice() {}, + onclose: () => closed.resolve(), +}); +const control = postgres(process.env.TEST_DATABASE_URL, { max: 1, onnotice() {} }); +const capture = (query) => Promise.resolve(query).then( + () => { throw new Error("Expected the closed connection to reject"); }, + (error) => error, +); +const terminate = async (pid) => { + await control.unsafe("select pg_terminate_backend($1)", [pid]); + await closed.promise; +}; + +try { + await control.unsafe("drop table if exists connection_recovery_probe"); + await control.unsafe("create table connection_recovery_probe (id integer)"); + + if (mode === "transaction" || mode === "transaction-closed" || mode === "savepoint") { + const started = Promise.withResolvers(); + const resume = Promise.withResolvers(); + const callbackDone = Promise.withResolvers(); + let lateError; + const work = async (tx) => { + const [row] = await tx.unsafe("select pg_backend_pid() as pid"); + await tx.unsafe("insert into connection_recovery_probe values (1)"); + started.resolve(row.pid); + await resume.promise; + try { + await tx.unsafe("insert into connection_recovery_probe values (2)"); + } catch (error) { + lateError = error; + throw error; + } finally { + callbackDone.resolve(); + } + }; + const failed = capture(sql.begin((tx) => mode === "savepoint" ? tx.savepoint(work) : work(tx))); + await terminate(await started.promise); + assert.equal((await failed).code, "CONNECTION_CLOSED"); + + const finishOldCallback = async () => { + resume.resolve(); + await callbackDone.promise; + await setImmediate(); + await setImmediate(); + assert.equal(lateError?.code, "CONNECTION_CLOSED"); + }; + if (mode === "transaction-closed") await finishOldCallback(); + + // Reuse the physical pool slot before the old callback resumes. The old + // query and its automatic rollback must not run in this new transaction. + await sql.begin(async (tx) => { + await tx.unsafe("insert into connection_recovery_probe values (3)"); + await finishOldCallback(); + await tx.unsafe("insert into connection_recovery_probe values (4)"); + }); + const rows = await sql.unsafe("select id from connection_recovery_probe order by id"); + assert.deepEqual(rows.map((row) => row.id), [3, 4]); + } else if (mode === "reserve") { + const reserved = await sql.reserve(); + const [row] = await reserved.unsafe("select pg_backend_pid() as pid"); + await terminate(row.pid); + assert.equal((await capture(reserved.unsafe("select 1"))).code, "CONNECTION_CLOSED"); + reserved.release(); + const fresh = await sql.reserve(); + try { + const [freshRow] = await fresh.unsafe("select pg_backend_pid() as pid"); + assert.notEqual(freshRow.pid, row.pid); + reserved.release(); + assert.equal((await capture(reserved.unsafe("select 1"))).code, "CONNECTION_CLOSED"); + assert.equal((await fresh.unsafe("select 42 as value"))[0].value, 42); + } finally { + fresh.release(); + } + } else if (mode === "transaction-queue" || mode === "reserve-queue") { + const started = Promise.withResolvers(); + let pending; + const work = async (tx) => { + const [row] = await tx.unsafe("select pg_backend_pid() as pid"); + // With max_pipeline=1, later queries stay in the scope's local queue. + // Disconnect after PostgreSQL starts the first query. + pending = [0, 1, 2].map(() => capture(tx.unsafe("select pg_sleep(30)"))); + started.resolve(row.pid); + const errors = await Promise.all(pending); + assert.ok(errors.every((error) => ["57P01", "CONNECTION_CLOSED"].includes(error.code))); + }; + const reserved = mode === "reserve-queue" ? await sql.reserve() : null; + const task = reserved ? work(reserved) : capture(sql.begin(work)); + const pid = await started.promise; + while (true) { + const [row] = await control.unsafe("select wait_event from pg_stat_activity where pid = $1", [pid]); + if (row?.wait_event === "PgSleep") break; + await setImmediate(); + } + await terminate(pid); + await task; + await Promise.all(pending); + reserved?.release(); + } else { + throw new Error("Unknown recovery case"); + } + + assert.equal((await sql.unsafe("select 42 as value"))[0].value, 42); + await setImmediate(); + console.log("recovered"); +} finally { + await sql.end({ timeout: 1 }); + await control.end({ timeout: 1 }); +} diff --git a/packages/db/src/postgres-connection-recovery.test.ts b/packages/db/src/postgres-connection-recovery.test.ts new file mode 100644 index 0000000000..ddb6bb325b --- /dev/null +++ b/packages/db/src/postgres-connection-recovery.test.ts @@ -0,0 +1,42 @@ +import { execFile } from "node:child_process"; +import { fileURLToPath } from "node:url"; +import { promisify } from "node:util"; +import { afterAll, beforeAll, describe, expect, it } from "vitest"; +import { + EMBEDDED_POSTGRES_TEST_TIMEOUT_MS, + getEmbeddedPostgresTestSupport, + startEmbeddedPostgresTestDatabase, + type EmbeddedPostgresTestDatabase, +} from "./test-embedded-postgres.js"; + +const support = await getEmbeddedPostgresTestSupport(); +const run = promisify(execFile); +const fixture = fileURLToPath(new URL("./__fixtures__/postgres-connection-recovery.mjs", import.meta.url)); + +describe.skipIf(!support.supported)("postgres connection recovery", () => { + let database: EmbeddedPostgresTestDatabase; + beforeAll(async () => { + database = await startEmbeddedPostgresTestDatabase("paperclip-driver-recovery-"); + }, EMBEDDED_POSTGRES_TEST_TIMEOUT_MS); + afterAll(async () => { await database?.cleanup(); }); + + for (const format of ["esm", "cjs"]) { + for (const scenario of ["transaction", "transaction-closed", "savepoint", "reserve", "transaction-queue", "reserve-queue"]) { + it(`${format}: rejects disconnected ${scenario} work and recovers the pool`, async () => { + // A driver timer used to crash the whole process. A child also detects + // hangs without leaving broken connections in the Vitest worker. + const { stdout, stderr } = await run(process.execPath, [fixture], { + env: { + ...process.env, + TEST_DATABASE_URL: database.connectionString, + DRIVER_FORMAT: format, + RECOVERY_CASE: scenario, + }, + timeout: 15_000, + }); + expect(stdout.trim()).toBe("recovered"); + expect(stderr).toBe(""); + }, 20_000); + } + } +}); diff --git a/patches/postgres@3.4.9.patch b/patches/postgres@3.4.9.patch new file mode 100644 index 0000000000..ab27a64ba3 --- /dev/null +++ b/patches/postgres@3.4.9.patch @@ -0,0 +1,170 @@ +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) diff --git a/pnpm-workspace.yaml b/pnpm-workspace.yaml index 026c899c97..3518e9086f 100644 --- a/pnpm-workspace.yaml +++ b/pnpm-workspace.yaml @@ -16,6 +16,7 @@ packages: # Keep in sync with package.json#pnpm.patchedDependencies. Newer pnpm # versions read patch configuration only from the workspace manifest. patchedDependencies: + postgres@3.4.9: patches/postgres@3.4.9.patch embedded-postgres@18.1.0-beta.16: patches/embedded-postgres@18.1.0-beta.16.patch acpx@0.12.0: patches/acpx@0.12.0.patch acpx@0.13.1: patches/acpx@0.13.1.patch