Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 12 additions & 0 deletions app/store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -138,6 +138,15 @@ async function withAppLock<T>(
work: (client: pg.PoolClient) => Promise<T>,
): Promise<T> {
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
Expand All @@ -153,6 +162,9 @@ async function withAppLock<T>(
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();
}
}
Expand Down
64 changes: 59 additions & 5 deletions tests/store.test.ts
Original file line number Diff line number Diff line change
@@ -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(() => {
Expand All @@ -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();
});
});
Loading