From 98d8a6ccac82d1605ac10dcea17f351c385814e5 Mon Sep 17 00:00:00 2001 From: Devin Foley Date: Wed, 30 Sep 2026 18:34:00 -0700 Subject: [PATCH] Stop replaying ambiguous database disconnects (#14773) ## Thinking Path > - Paperclip stores agent work and control state in PostgreSQL. > - Its database client must not repeat a mutation after an uncertain result. > - The global retry wrapper treated `write CONNECTION_CLOSED` as proof that PostgreSQL never received a statement. > - postgres.js also uses that message when the connection closes after statement delivery. > - This pull request removes that global replay and tests the actual driver over a local wire connection. > - Callers retain control of retries when they can prove the complete operation is idempotent. ## Linked Issues or Issue Description Follow-up to #13417. Preserve the transaction disconnect handling from #13643 and the explicit actor synchronization retries introduced in #12773. Searched open and closed issues and PRs for database retries, disconnects, and `CONNECTION_CLOSED`. The open circuit-breaker proposal #11142 addresses outage queue growth; it does not establish whether an already-sent statement can be replayed. **What happened?** The database wrapper replayed an arbitrary statement up to three times after `write CONNECTION_CLOSED`. The driver adds `write ` to connection-close errors even after the peer receives the statement. A local protocol peer receives the same submitted INSERT three times when it drops each response. A committed write could therefore execute more than once. **Expected behavior** An ambiguous statement result must fail without automatic replay. A subsequent operation must be able to reconnect. **Steps to reproduce** Run the new wire regression against the parent commit. The six Simple Query cases and the parameterized Drizzle case receive three executions instead of one. The named prepared-client case was already safe and stays covered. The peer reads the entire statement and then closes the connection. This demonstrates repeated delivery with the real driver; it does not claim that a historical incident duplicated a committed write. **Paperclip version or commit** Reproduced on source commit `018993140f` with the patched postgres.js 3.4.9 dependency. **Deployment mode** Built from source with a local PostgreSQL protocol peer. No live provider or customer database is used. ## What Changed - Pass the original postgres.js client to Drizzle and remove the global statement replay wrapper. - Add eight wire regressions: six Simple Query cases for INSERT, side-effect-capable SELECT, and a data-changing CTE, plus parameterized Drizzle and named prepared-client cases. The extended peer processes Parse, Describe, Bind, and Execute, verifies bound parameters, and drops the response only after Execute. Each case checks one delivery and recovery on a fresh query. - Document ambiguous outcomes and the retry compatibility tradeoff. Keep explicit idempotent actor-sync retries and disconnected-transaction handling unchanged. ## Verification - Before the fix: the six Simple Query cases and the parameterized Drizzle case failed with three executions instead of one. The named prepared-client case was already safe. All eight wire cases pass on this branch. - Final focused client, pool teardown, configuration, and actor-sync retry checks: 28 tests passed. `pnpm --filter @paperclipai/db typecheck` also passed after the test-only follow-up. - First implementation head, `pnpm exec vitest run --project @paperclipai/db`: all 158 tests passed across 45 files, including real PostgreSQL transaction/reserved-connection recovery. The local embedded dependency's symlinks were hydrated before this run. - `pnpm -r typecheck`: passed. - `pnpm build`: passed. - The complete local `pnpm test:run` did not finish; no complete local suite pass is claimed. All CI test, typecheck, and build gates passed on the first implementation head `11f8b22d90`. Final-head CI is pending after the test-only follow-up. - `git diff --check`: passed. Reviewed the diff for secrets, personal data, generated output, and run artifacts. ## Risks Some transient statement failures that the global wrapper previously replayed now reach the caller. Operation owners must retry only when they have an idempotency guarantee or a durable receipt that prevents duplicate effects. A connection error is not proof that a write failed to commit. There is no SQL-text retry heuristic, new suppression, schema change, or migration. This change prevents unsafe replay; it does not prevent network disconnects. ## Model Used OpenAI Codex / GPT-6, with reasoning, repository inspection, code execution, and local protocol tests. The exact backend model ID and context-window size are not exposed in this session. ## Checklist - [x] I have included a thinking path that traces from project context to this change - [x] I have specified the model used (with version and capability details) - [x] I have checked ROADMAP.md and confirmed this PR does not duplicate planned core work - [x] I have searched GitHub for duplicate or related PRs and linked them above - [x] I have either (a) linked existing issues with `Fixes: #` / `Closes #` / `Refs #` OR (b) described the issue in-PR following the relevant issue template - [x] I have not referenced internal/instance-local Paperclip issues or links (only public GitHub `#NNN` / `github.com/paperclipai/paperclip` URLs) - [x] My branch name describes the change (e.g. `docs/...`, `fix/...`) and contains no internal Paperclip ticket id or instance-derived details - [x] I have run tests locally and they pass - [x] I have added or updated tests where applicable - [x] I have updated relevant documentation to reflect my changes - [x] I have considered and documented any risks above - [ ] All Paperclip CI gates are green - [ ] Greptile is 5/5 with no open P2s, recommendations, or follow-ups - [x] I will address all Greptile and reviewer comments before requesting merge Final verification (September30): every final-head CI check passed at `3004c5bda39c985c3557547ec45e33870ce5d010`. Greptile scored5/5 on this head, all review threads are resolved, and the branch is mergeable. Eight real-wire regressions cover simple, parameterized Drizzle, and named prepared queries. Full local suite did not produce a completed result; the complete CI matrix passed. This public PR remains open for maintainer merge. --------- Co-authored-by: Paperclip --- doc/DATABASE.md | 14 +- .../db/src/client-disconnect-replay.test.ts | 177 ++++++++++++++++ packages/db/src/client.ts | 10 +- packages/db/src/transient-write-retry.test.ts | 191 ------------------ packages/db/src/transient-write-retry.ts | 90 --------- 5 files changed, 195 insertions(+), 287 deletions(-) create mode 100644 packages/db/src/client-disconnect-replay.test.ts delete mode 100644 packages/db/src/transient-write-retry.test.ts delete mode 100644 packages/db/src/transient-write-retry.ts diff --git a/doc/DATABASE.md b/doc/DATABASE.md index 3a348cbe52..0f7fd97f60 100644 --- a/doc/DATABASE.md +++ b/doc/DATABASE.md @@ -144,7 +144,19 @@ DATABASE_URL=postgres://postgres.[PROJECT-REF]:[PASSWORD]@...5432/postgres \ See [Supabase pricing](https://supabase.com/pricing) for current details. -## Connection loss during a transaction +## Connection loss and retries + +The database client does not replay arbitrary statements after a disconnect. +PostgreSQL may have committed a statement before the connection loses its +response. The postgres.js message `write CONNECTION_CLOSED` does not prove +that the statement was never sent: the driver also uses it when an in-flight +query loses its connection. SQL text cannot establish replay safety either; +a `SELECT` can call a function with side effects. + +The affected operation fails and a new operation can reconnect through the +pool. Callers may retry only when the complete operation is idempotent or has +a durable receipt that prevents duplicate effects. Some transient statement +failures therefore reach the caller instead of being retried automatically. When a database connection closes, its transaction fails. Paperclip does not replay that transaction. New requests can use a fresh connection from the pool. diff --git a/packages/db/src/client-disconnect-replay.test.ts b/packages/db/src/client-disconnect-replay.test.ts new file mode 100644 index 0000000000..e03c22146d --- /dev/null +++ b/packages/db/src/client-disconnect-replay.test.ts @@ -0,0 +1,177 @@ +import net from "node:net"; +import { sql } from "drizzle-orm"; +import type { Sql } from "postgres"; +import { describe, expect, it } from "vitest"; +import { closeRegisteredClients, createDb } from "./client.js"; + +/** A wire peer that receives the statement, then loses its acknowledgement. */ +async function startDisconnectingServer() { + const sockets = new Set(); + let receivedEffects = 0; + const extendedExecutions: Array<{ statement: string; prepared: boolean; parameters: Array }> = []; + const ready = Buffer.from([0x5a, 0, 0, 0, 5, 0x49]); + const authOk = Buffer.from([0x52, 0, 0, 0, 8, 0, 0, 0, 0]); + const parseComplete = Buffer.from([0x31, 0, 0, 0, 4]); + const bindComplete = Buffer.from([0x32, 0, 0, 0, 4]); + const emptyRowDescription = Buffer.from([0x54, 0, 0, 0, 6, 0, 0]); + const commandComplete = Buffer.from([0x43, 0, 0, 0, 0x0d, ...Buffer.from("SELECT 0\0")]); + const server = net.createServer((socket) => { + sockets.add(socket); + socket.on("close", () => sockets.delete(socket)); + socket.on("error", () => {}); + let greeted = false; + let pending = Buffer.alloc(0); + const statements = new Map(); + const portals = new Map }>(); + socket.on("data", (chunk) => { + pending = Buffer.concat([pending, chunk]); + while (pending.length >= (greeted ? 5 : 4)) { + const size = greeted ? pending.readUInt32BE(1) + 1 : pending.readUInt32BE(0); + if (pending.length < size) return; + const message = pending.subarray(0, size); + pending = pending.subarray(size); + if (!greeted) { + greeted = true; + socket.write(Buffer.concat([authOk, ready])); + } else if (message[0] === 0x51) { // Simple Query + const statement = message.subarray(5, -1).toString(); + if (statement.includes("disconnect_effect")) { + receivedEffects += 1; + // The whole statement reached the peer. Its outcome is now + // ambiguous to the client, even though the error says "write". + socket.destroy(); + return; + } + socket.write(Buffer.concat([commandComplete, ready])); + } else if (message[0] === 0x50) { // Parse + const nameEnd = message.indexOf(0, 5); + const textEnd = message.indexOf(0, nameEnd + 1); + const name = message.toString("utf8", 5, nameEnd); + const text = message.toString("utf8", nameEnd + 1, textEnd); + const count = message.readUInt16BE(textEnd + 1); + const parameterTypes = Array.from({ length: count }, (_, index) => + message.readUInt32BE(textEnd + 3 + index * 4) || 23); // int4 for this fixture's integer parameters. + statements.set(name, { text, parameterTypes }); + socket.write(parseComplete); + } else if (message[0] === 0x44) { // Describe statement + const name = message.toString("utf8", 6, message.length - 1); + const statement = statements.get(name)!; + const description = Buffer.alloc(7 + statement.parameterTypes.length * 4); + description[0] = 0x74; // ParameterDescription + description.writeUInt32BE(description.length - 1, 1); + description.writeUInt16BE(statement.parameterTypes.length, 5); + statement.parameterTypes.forEach((type, index) => description.writeUInt32BE(type, 7 + index * 4)); + socket.write(Buffer.concat([description, emptyRowDescription])); + } else if (message[0] === 0x42) { // Bind + const portalEnd = message.indexOf(0, 5); + const nameEnd = message.indexOf(0, portalEnd + 1); + const portal = message.toString("utf8", 5, portalEnd); + const name = message.toString("utf8", portalEnd + 1, nameEnd); + let offset = nameEnd + 3 + message.readUInt16BE(nameEnd + 1) * 2; + const count = message.readUInt16BE(offset); + offset += 2; + const parameters: Array = []; + for (let index = 0; index < count; index++) { + const length = message.readInt32BE(offset); + offset += 4; + parameters.push(length === -1 ? null : message.toString("utf8", offset, offset + length)); + if (length !== -1) offset += length; + } + portals.set(portal, { statement: statements.get(name)!.text, prepared: name.length > 0, parameters }); + socket.write(bindComplete); + } else if (message[0] === 0x45) { // Execute: Parse/Describe alone cannot execute an effect. + const execution = portals.get(message.toString("utf8", 5, message.indexOf(0, 5)))!; + extendedExecutions.push(execution); + if (execution.statement.includes("disconnect_effect")) { + receivedEffects += 1; + socket.destroy(); + return; + } + socket.write(commandComplete); + } else if (message[0] === 0x53) { // Sync + socket.write(ready); + } + } + }); + }); + await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); + const port = (server.address() as net.AddressInfo).port; + return { + url: `postgres://test:test@127.0.0.1:${port}/test`, + receivedEffects: () => receivedEffects, + extendedExecutions: () => extendedExecutions, + close: async () => { + for (const socket of sockets) socket.destroy(); + await new Promise((resolve) => server.close(() => resolve())); + }, + }; +} + +describe("database connection loss after statement delivery", () => { + it.each(["drizzle", "prepared"] as const)("does not replay a parameterized mutation (%s)", async (surface) => { + const peer = await startDisconnectingServer(); + try { + // Keep the default client settings. Drizzle uses unsafe(), which selects + // unnamed statements; tagged client queries use named prepared statements. + const db = createDb(peer.url, { maxConnections: 1, connectTimeoutSeconds: 2 }); + const client = db.$client; + const pending = surface === "drizzle" + ? db.execute(sql`insert into disconnect_effect(value) values (${42})`) + : client`insert into disconnect_effect(value) values (${42})`; + const error = await pending.then(() => null, (failure: unknown) => failure); + expect(error).toBeInstanceOf(Error); + const driverError = error instanceof Error && error.cause ? error.cause : error; + expect(driverError).toMatchObject({ + code: "CONNECTION_CLOSED", + message: expect.stringContaining("write CONNECTION_CLOSED"), + }); + expect(peer.receivedEffects()).toBe(1); + expect(peer.extendedExecutions().filter((execution) => execution.statement.includes("disconnect_effect"))) + .toEqual([{ statement: "insert into disconnect_effect(value) values ($1)", prepared: surface === "prepared", parameters: ["42"] }]); + + const freshQuery = surface === "drizzle" ? db.execute(sql`select ${7}`) : client`select ${7}`; + await expect(freshQuery).resolves.toBeDefined(); + expect(peer.extendedExecutions().at(-1)) + .toEqual({ statement: "select $1", prepared: surface === "prepared", parameters: ["7"] }); + expect(peer.receivedEffects()).toBe(1); + } finally { + await peer.close(); + await closeRegisteredClients(peer.url); + } + }); + + for (const statement of [ + "insert into disconnect_effect default values", + "select disconnect_effect()", + "with effect as (insert into disconnect_effect default values returning *) select * from effect", + ]) { + it.each(["rows", "values"] as const)(`does not replay ${statement} (%s)`, async (shape) => { + const peer = await startDisconnectingServer(); + try { + const db = createDb(peer.url, { maxConnections: 1, connectTimeoutSeconds: 2, prepare: false }); + const client = (db as unknown as { $client: Sql }).$client; + const pending = shape === "rows" + ? db.execute(sql.raw(statement)) + : client.unsafe(statement).values(); + // This is the real postgres.js error, not a mock that invents a + // distinct "read CONNECTION_CLOSED" shape for an in-flight query. + const error = await pending.then(() => null, (failure: unknown) => failure); + expect(error).toBeInstanceOf(Error); + const driverError = error instanceof Error && error.cause ? error.cause : error; + expect(driverError).toMatchObject({ + code: "CONNECTION_CLOSED", + message: expect.stringContaining("write CONNECTION_CLOSED"), + }); + expect(peer.receivedEffects()).toBe(1); + + // Failing the ambiguous statement must not poison the pool. A new + // operation can reconnect, without resubmitting the failed effect. + await expect(db.execute(sql`select 0`)).resolves.toBeDefined(); + expect(peer.receivedEffects()).toBe(1); + } finally { + await peer.close(); + await closeRegisteredClients(peer.url); + } + }); + } +}); diff --git a/packages/db/src/client.ts b/packages/db/src/client.ts index ecfd795442..d71aeff6c9 100644 --- a/packages/db/src/client.ts +++ b/packages/db/src/client.ts @@ -5,7 +5,6 @@ import { readFile, readdir } from "node:fs/promises"; import { fileURLToPath } from "node:url"; import postgres from "postgres"; import * as schema from "./schema/index.js"; -import { withTransientWriteRetry } from "./transient-write-retry.js"; const MIGRATIONS_FOLDER = fileURLToPath(new URL("./migrations", import.meta.url)); const DRIZZLE_MIGRATIONS_TABLE = "__drizzle_migrations"; @@ -273,10 +272,11 @@ export function createDb(url: string, options?: DatabaseClientOptions) { const sql = postgres(url, postgresJsOptions(resolved)); const key = hostPortKeyOrNull(url); if (key) registerClient(key, sql); - // The registry keeps the real client (teardown must end the actual pool); - // drizzle gets the retrying face so a pooler-recycled socket replays the - // query instead of failing the request that happened to draw it. - const db = drizzlePg(withTransientWriteRetry(sql), { schema }); + // A disconnect can lose the response after a statement has committed. + // postgres.js calls that error "write CONNECTION_CLOSED" too, so the + // message cannot establish that replay is safe. Leave retries to callers + // that know the complete operation is idempotent. + const db = drizzlePg(sql, { schema }); dedicatedDbFactories.set(db, () => createDb(url, { ...resolved, maxConnections: 1, applicationName: "paperclip-workspace-finalization-lock", })); diff --git a/packages/db/src/transient-write-retry.test.ts b/packages/db/src/transient-write-retry.test.ts deleted file mode 100644 index b144cb2c77..0000000000 --- a/packages/db/src/transient-write-retry.test.ts +++ /dev/null @@ -1,191 +0,0 @@ -import net from "node:net"; -import { sql as drizzleSql } from "drizzle-orm"; -import { afterEach, describe, expect, it } from "vitest"; -import type { Sql } from "postgres"; -import { closeRegisteredClients, createDb } from "./client.js"; -import { isTransientWritePhaseError, withTransientWriteRetry } from "./transient-write-retry.js"; - -/** The exact shape postgres.js raises when the socket write fails. */ -function writeClosedError(): Error { - const error = new Error("write CONNECTION_CLOSED db.example.internal:5432"); - (error as Error & { code: string }).code = "CONNECTION_CLOSED"; - return error; -} - -/** A stub standing in for the root postgres.js client's `unsafe` surface. */ -function stubSql(behavior: { failures: number; error?: () => Error }) { - let calls = 0; - const state = { - get calls() { - return calls; - }, - }; - const unsafe = (_query: string, _parameters?: unknown[]) => { - calls += 1; - const attempt = calls; - const shouldFail = attempt <= behavior.failures; - const outcome = () => - shouldFail - ? Promise.reject((behavior.error ?? writeClosedError)()) - : Promise.resolve([{ attempt }]); - return { - then: (onFulfilled?: ((value: unknown) => unknown) | null, onRejected?: ((reason: unknown) => unknown) | null) => - outcome().then(onFulfilled, onRejected), - values: () => outcome().then(() => [[attempt]]), - }; - }; - const sql = { unsafe } as unknown as Sql; - return { sql, state }; -} - -describe("isTransientWritePhaseError", () => { - it("matches only the write-phase connection-closed shape", () => { - expect(isTransientWritePhaseError(writeClosedError())).toBe(true); - - const midQuery = new Error("read CONNECTION_CLOSED db.example.internal:5432"); - (midQuery as Error & { code: string }).code = "CONNECTION_CLOSED"; - expect(isTransientWritePhaseError(midQuery)).toBe(false); - - const ended = new Error("write CONNECTION_ENDED db.example.internal:5432"); - (ended as Error & { code: string }).code = "CONNECTION_ENDED"; - expect(isTransientWritePhaseError(ended)).toBe(false); - - expect(isTransientWritePhaseError(new Error("write CONNECTION_CLOSED x:1"))).toBe(false); // no code - expect(isTransientWritePhaseError(null)).toBe(false); - }); -}); - -describe("withTransientWriteRetry", () => { - it("replays a query whose socket write failed, and returns the replay's rows", async () => { - const { sql, state } = stubSql({ failures: 1 }); - const rows = await withTransientWriteRetry(sql).unsafe("select 1", []); - expect(rows).toEqual([{ attempt: 2 }]); - expect(state.calls).toBe(2); - }); - - it("replays the .values() form drizzle uses", async () => { - const { sql, state } = stubSql({ failures: 2 }); - const values = await withTransientWriteRetry(sql).unsafe("select 1", []).values(); - expect(values).toEqual([[3]]); - expect(state.calls).toBe(3); - }); - - it("gives up after the attempt budget and surfaces the driver error", async () => { - const { sql, state } = stubSql({ failures: Number.POSITIVE_INFINITY }); - await expect(withTransientWriteRetry(sql).unsafe("select 1", [])).rejects.toThrow( - "write CONNECTION_CLOSED", - ); - expect(state.calls).toBe(3); - }); - - it("does not replay an error that may have reached the server", async () => { - const midQuery = () => { - const error = new Error("read CONNECTION_CLOSED db.example.internal:5432"); - (error as Error & { code: string }).code = "CONNECTION_CLOSED"; - return error; - }; - const { sql, state } = stubSql({ failures: Number.POSITIVE_INFINITY, error: midQuery }); - await expect(withTransientWriteRetry(sql).unsafe("select 1", [])).rejects.toThrow( - "read CONNECTION_CLOSED", - ); - expect(state.calls).toBe(1); - }); - - it("runs one execution no matter how many handlers attach to the pending query", async () => { - const { sql, state } = stubSql({ failures: 0 }); - const pending = withTransientWriteRetry(sql).unsafe("select 1", []); - await Promise.all([pending, pending.catch(() => undefined), pending.then((rows) => rows)]); - expect(state.calls).toBe(1); - }); - - it("never executes a query twice when a caller both awaits it and takes its values", async () => { - // `.values()` selects the row shape of the one execution, like the driver. - // Starting a second execution here would repeat a mutation's effect. - const { sql, state } = stubSql({ failures: 0 }); - const pending = withTransientWriteRetry(sql).unsafe("insert into t values (1)", []); - const rows = await pending; - const values = await pending.values(); - expect(state.calls).toBe(1); - expect(values).toBe(rows); - }); - - it("takes the values shape when it is chosen before the query runs", async () => { - const { sql, state } = stubSql({ failures: 0 }); - const pending = withTransientWriteRetry(sql).unsafe("select 1", []); - expect(await pending.values()).toEqual([[1]]); - expect(state.calls).toBe(1); - }); -}); - -describe("createDb with the retrying client", () => { - // The same minimal wire-protocol fake as client-teardown-registry.test.ts: - // enough of the startup and query flow for drizzle to run a real query - // through the proxied client, proving the retry face is transparent. - function startFakePostgresServer(): Promise<{ server: net.Server; port: number }> { - const authOk = Buffer.from([0x52, 0, 0, 0, 8, 0, 0, 0, 0]); - const readyForQuery = Buffer.from([0x5a, 0, 0, 0, 5, 0x49]); - const emptyQueryReply = Buffer.concat([ - Buffer.from([0x31, 0, 0, 0, 4]), - Buffer.from([0x32, 0, 0, 0, 4]), - Buffer.from([0x54, 0, 0, 0, 6, 0, 0]), - Buffer.from([0x43, 0, 0, 0, 0x0d, 0x53, 0x45, 0x4c, 0x45, 0x43, 0x54, 0x20, 0x30, 0]), - ]); - const server = net.createServer((socket) => { - let greeted = false; - socket.on("data", () => { - if (!greeted) { - greeted = true; - socket.write(Buffer.concat([authOk, readyForQuery])); - return; - } - socket.write(Buffer.concat([emptyQueryReply, readyForQuery])); - }); - socket.on("error", () => {}); - }); - return new Promise((resolve) => { - server.listen(0, "127.0.0.1", () => { - resolve({ server, port: (server.address() as net.AddressInfo).port }); - }); - }); - } - - let server: net.Server | null = null; - let url: string | null = null; - - afterEach(async () => { - if (url) await closeRegisteredClients(url); - if (server) await new Promise((resolve) => server!.close(resolve)); - server = null; - url = null; - }); - - it("still answers ordinary drizzle queries through the proxy", async () => { - const started = await startFakePostgresServer(); - server = started.server; - url = `postgres://test:test@127.0.0.1:${started.port}/test`; - const db = createDb(url, { connectTimeoutSeconds: 5, prepare: false }); - await expect(db.execute(drizzleSql`select 0`)).resolves.toBeDefined(); - }); - - it("leaves every other client surface reachable through the proxy", async () => { - // Drizzle itself only awaits `unsafe()` or takes its `.values()`, but the - // client is reachable as `db.$client`, and callers use it as a tagged - // template, open transactions on it, and end it. A wrapper that broke any - // of those would fail far from here, so pin them against the real driver. - const started = await startFakePostgresServer(); - server = started.server; - url = `postgres://test:test@127.0.0.1:${started.port}/test`; - const db = createDb(url, { connectTimeoutSeconds: 5, prepare: false }); - const client = (db as unknown as { $client: Sql }).$client; - - await expect(client`select 1`).resolves.toBeDefined(); - await expect(client.unsafe("select 1", [])).resolves.toBeDefined(); - await expect( - client.begin(async (tx) => { - await tx.unsafe("select 1", []); - return "transaction result"; - }), - ).resolves.toBe("transaction result"); - await expect(client.end({ timeout: 1 })).resolves.toBeUndefined(); - }); -}); diff --git a/packages/db/src/transient-write-retry.ts b/packages/db/src/transient-write-retry.ts deleted file mode 100644 index 6d65751345..0000000000 --- a/packages/db/src/transient-write-retry.ts +++ /dev/null @@ -1,90 +0,0 @@ -import type { Sql } from "postgres"; - -/** - * Replays a query whose bytes never reached the server. - * - * Behind a connection pooler (Neon's PgBouncer), the server side of an idle - * connection can be recycled while the client still holds the socket. The - * next query then fails at the socket write — postgres.js reports it as - * `write CONNECTION_CLOSED :` with `code: "CONNECTION_CLOSED"`. - * Because the write itself failed, the server never saw the query, so - * replaying it on a fresh connection cannot double-execute anything — the - * replay is safe for reads and writes alike. Every other error, including a - * connection lost mid-query (where the server may have acted), propagates - * untouched. - * - * Drizzle routes every non-transactional query through the root client's - * `unsafe`; queries inside `db.transaction()` run on the scoped client - * `sql.begin()` hands out, which this wrapper deliberately does not touch — - * a transaction that loses its connection must abort, not replay. - */ -const TRANSIENT_WRITE_ATTEMPTS = 3; -const TRANSIENT_WRITE_BACKOFF_MS = 50; - -export function isTransientWritePhaseError(error: unknown): boolean { - return ( - error instanceof Error && - (error as { code?: unknown }).code === "CONNECTION_CLOSED" && - error.message.startsWith("write CONNECTION_CLOSED") - ); -} - -type UnsafeFn = (query: string, parameters?: unknown[]) => { - values: () => Promise; -} & PromiseLike; - -export function withTransientWriteRetry(sql: T): T { - const runWithRetry = async (execute: () => Promise): Promise => { - for (let attempt = 0; ; attempt++) { - try { - return await execute(); - } catch (error) { - if (attempt >= TRANSIENT_WRITE_ATTEMPTS - 1 || !isTransientWritePhaseError(error)) { - throw error; - } - await new Promise((resolve) => setTimeout(resolve, TRANSIENT_WRITE_BACKOFF_MS * (attempt + 1))); - } - } - }; - - return new Proxy(sql, { - get(target, property, receiver) { - if (property !== "unsafe") return Reflect.get(target, property, receiver); - const unsafe = target.unsafe.bind(target) as UnsafeFn; - const retryingUnsafe = (query: string, parameters?: unknown[], ...rest: unknown[]) => { - if (rest.length > 0) { - // An options argument selects driver surfaces this wrapper does not - // model; leave those calls exactly as they were. - return (target.unsafe as (...args: unknown[]) => unknown)(query, parameters, ...rest); - } - // postgres.js queries execute lazily, exactly once, and `.values()` - // selects the row shape of that one execution rather than starting a - // second one. Mirror both halves: the execution starts when a consumer - // first settles the query, `.values()` marks the shape and returns the - // same pending query, and every later handler joins the same run. A - // shape chosen after the run started cannot change it — the same as - // the driver, and the reason a mutation can never execute twice here. - let started: Promise | undefined; - let wantsValues = false; - const runOnce = () => - (started ??= runWithRetry(() => - wantsValues - ? unsafe(query, parameters).values() - : Promise.resolve(unsafe(query, parameters)), - )); - const pending = { - then: (onFulfilled?: ((value: unknown) => unknown) | null, onRejected?: ((reason: unknown) => unknown) | null) => - runOnce().then(onFulfilled, onRejected), - catch: (onRejected?: ((reason: unknown) => unknown) | null) => runOnce().catch(onRejected), - finally: (onFinally?: (() => void) | null) => runOnce().finally(onFinally), - values: () => { - if (!started) wantsValues = true; - return pending; - }, - }; - return pending; - }; - return retryingUnsafe; - }, - }) as T; -}