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