mirror of
https://github.com/paperclipai/paperclip.git
synced 2026-10-06 10:48:12 +02:00
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 <noreply@paperclip.ing>
This commit is contained in:
1 parent
c8f874311c
commit
98d8a6ccac
5 files changed
+195
-287
No files matched your search
@@ -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<net.Socket>();
|
||||
let receivedEffects = 0;
|
||||
const extendedExecutions: Array<{ statement: string; prepared: boolean; parameters: Array<string | null> }> = [];
|
||||
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<string, { text: string; parameterTypes: number[] }>();
|
||||
const portals = new Map<string, { statement: string; prepared: boolean; parameters: Array<string | null> }>();
|
||||
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<string | null> = [];
|
||||
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<void>((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<void>((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);
|
||||
}
|
||||
});
|
||||
}
|
||||
});
|
||||
@@ -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",
|
||||
}));
|
||||
|
||||
@@ -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();
|
||||
});
|
||||
});
|
||||
@@ -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 <host>:<port>` 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<unknown>;
|
||||
} & PromiseLike<unknown>;
|
||||
|
||||
export function withTransientWriteRetry<T extends Sql>(sql: T): T {
|
||||
const runWithRetry = async (execute: () => Promise<unknown>): Promise<unknown> => {
|
||||
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<unknown> | 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;
|
||||
}
|
||||
Reference in new issue
Block a user