From 5a0ce0fee2a89126acf34d17e1cd280fb7f60e4f Mon Sep 17 00:00:00 2001 From: Anurag Goel Date: Thu, 24 Sep 2026 10:04:04 -0700 Subject: [PATCH] Log the error when Postgres closes a connection in a transaction withAppLock() checks out a client from the pool for one transaction. claimRunApp (workflows host) and claimDelete (gateway) use it. While a client is checked out, pg-pool does not listen for "error" on it, and the pool listener of #43 gets only the errors of idle clients. If Postgres closes the connection during the transaction, for example in a restart, a failover, or with pg_terminate_backend, the query fails. Then the socket closes, and the client emits "error" ("Connection terminated unexpectedly"). No listener got the event, so Node stopped the process with exit code 1. withAppLock() now attaches an "error" listener after connect() and removes it before release(). The listener logs the error and does not throw it. The error of the query still goes to the caller, so claimRunApp and claimDelete fail as before. release() attaches the listener of the pool again, and the pool removes the closed client. tests/store.test.ts gives claimDelete a fake client. The lock query fails, and the client emits "error" while the rollback waits. Without the listener, the emit throws. The test connects to no Postgres. Co-Authored-By: Claude Opus 5.5 --- app/store.ts | 12 +++++++++ tests/store.test.ts | 64 +++++++++++++++++++++++++++++++++++++++++---- 2 files changed, 71 insertions(+), 5 deletions(-) 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(); + }); +});