diff --git a/app/store.ts b/app/store.ts index c8592b8..efa8261 100644 --- a/app/store.ts +++ b/app/store.ts @@ -138,6 +138,15 @@ async function withAppLock( work: (client: pg.PoolClient) => Promise, ): Promise { const client = await db().connect(); + // The pool listens for "error" only on an idle client. If Postgres closes + // the connection during the transaction, for example in a restart or a + // failover, the client emits "error". If no listener gets the event, Node + // stops the process. The query fails too, and its error goes to the + // caller, so this listener only logs the error. + const logError = (error: Error) => { + console.error("Lost a Postgres connection in a transaction:", error); + }; + client.on("error", logError); try { await client.query("begin"); // The lock must come before the reads: a statement sees only the rows @@ -153,6 +162,9 @@ async function withAppLock( await client.query("rollback").catch(() => {}); throw error; } finally { + // release() attaches the listener of the pool again. If this listener + // remains, each transaction on the client adds one more. + client.off("error", logError); client.release(); } } diff --git a/tests/store.test.ts b/tests/store.test.ts index abf5ed7..2b3684c 100644 --- a/tests/store.test.ts +++ b/tests/store.test.ts @@ -1,11 +1,15 @@ /** - * Postgres can close an idle connection of the pool, for example in a - * restart. The pool then emits "error", and an "error" event with no listener - * stops the gateway and the workflows host. No test connects to Postgres: the - * pool opens a connection only for a query. + * Postgres can close a connection of the pool, for example in a restart. For + * an idle connection, the pool emits "error". For a connection that a + * transaction holds, the client emits "error". An "error" event with no + * listener stops the gateway and the workflows host. No test connects to + * Postgres: the pool opens a connection only for a query, and the transaction + * gets a fake client. */ +import { EventEmitter } from "node:events"; +import type { PoolClient } from "pg"; import { afterEach, describe, expect, it, vi } from "vitest"; -import { db } from "../app/store.js"; +import { claimDelete, db } from "../app/store.js"; describe("db", () => { afterEach(() => { @@ -27,3 +31,53 @@ describe("db", () => { ); }); }); + +describe("the transaction of an app lock", () => { + afterEach(() => { + vi.restoreAllMocks(); + vi.unstubAllEnvs(); + }); + + it("logs the error of a closed connection and gives the caller the error of the query", async () => { + vi.stubEnv("DATABASE_URL", "postgres://factory@127.0.0.1:5432/factory"); + const logged = vi.spyOn(console, "error").mockImplementation(() => {}); + const terminated = new Error( + "terminating connection due to administrator command", + ); + const closed = new Error("Connection terminated unexpectedly"); + // As in pg after pg_terminate_backend(): the lock query fails, and the + // rollback waits. Then the socket closes, and the client emits "error". + let failRollback: (error: Error) => void = () => {}; + const rollback = new Promise((_, reject) => { + failRollback = reject; + }); + const client = Object.assign(new EventEmitter(), { + release: vi.fn(), + query: vi.fn(async (text: string) => { + if (text === "begin") return { rows: [] }; + if (text === "rollback") return rollback; + throw terminated; + }), + }); + vi.spyOn(db(), "connect").mockImplementation( + async () => client as unknown as PoolClient, + ); + + const claim = claimDelete("demo", "furniture-catalog"); + await vi.waitFor(() => + expect(client.query).toHaveBeenLastCalledWith("rollback"), + ); + expect(() => client.emit("error", closed)).not.toThrow(); + failRollback(closed); + + await expect(claim).rejects.toBe(terminated); + expect(logged).toHaveBeenCalledWith( + "Lost a Postgres connection in a transaction:", + closed, + ); + // The pool gives this client to the next transaction, so no listener + // can remain. + expect(client.listenerCount("error")).toBe(0); + expect(client.release).toHaveBeenCalledOnce(); + }); +});