diff --git a/scripts/outbox-worker.ts b/scripts/outbox-worker.ts index 6d8ec41..0c90446 100644 --- a/scripts/outbox-worker.ts +++ b/scripts/outbox-worker.ts @@ -74,7 +74,15 @@ async function run(): Promise { } try { - await run(); + while (!stopping) { + try { + await run(); + } catch (cause) { + console.error("Outbox worker failed; retrying in 5 seconds", cause); + await Promise.all([closeDb(), closeRedisClient()]); + if (!stopping) await wait(5_000); + } + } } finally { await Promise.all([closeDb(), closeRedisClient()]); } diff --git a/tests/outbox-worker.test.ts b/tests/outbox-worker.test.ts new file mode 100644 index 0000000..b635021 --- /dev/null +++ b/tests/outbox-worker.test.ts @@ -0,0 +1,106 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; + +const mocks = vi.hoisted(() => ({ + closeDb: vi.fn(), + closeRedisClient: vi.fn(), + getRedisClient: vi.fn(), + claimOutboxEvents: vi.fn(), + markOutboxProcessed: vi.fn(), + releaseOutboxEvent: vi.fn(), + processOutboxEvent: vi.fn(), +})); + +vi.mock("@/db/client", () => ({ closeDb: mocks.closeDb })); +vi.mock("@/lib/redis/client", () => ({ + getRedisClient: mocks.getRedisClient, + closeRedisClient: mocks.closeRedisClient, +})); +vi.mock("@/lib/outbox/repository", () => ({ + claimOutboxEvents: mocks.claimOutboxEvents, + markOutboxProcessed: mocks.markOutboxProcessed, + releaseOutboxEvent: mocks.releaseOutboxEvent, +})); +vi.mock("@/lib/outbox/processor", () => ({ + createRedisOutboxTransport: () => ({}), + processOutboxEvent: mocks.processOutboxEvent, +})); +vi.mock("@/lib/commission/cleanup", () => ({ + removeExpiredUnpaidCheckouts: vi.fn(), +})); +vi.mock("@/lib/guides/notifications", () => ({ + sendGuidePublicationNotifications: vi.fn(), +})); + +describe("outbox worker recovery", () => { + let stop: () => void; + + beforeEach(() => { + vi.resetModules(); + vi.resetAllMocks(); + vi.useFakeTimers(); + vi.spyOn(console, "error").mockImplementation(() => undefined); + const once = process.once.bind(process); + vi.spyOn(process, "once").mockImplementation(((event: string | symbol, listener: () => void) => { + if (event === "SIGINT" || event === "SIGTERM") { + stop = listener; + return process; + } + return once(event, listener); + }) as typeof process.once); + mocks.getRedisClient.mockResolvedValue({}); + mocks.claimOutboxEvents.mockResolvedValue([ + { id: 1, attempts: 1, eventType: "guide.updated" }, + ]); + mocks.markOutboxProcessed.mockImplementation(async () => stop()); + }); + + afterEach(() => { + vi.useRealTimers(); + vi.restoreAllMocks(); + }); + + it.each(["connection", "polling"])( + "recovers from a temporary %s failure and processes the next event", + async (failure) => { + const cause = new Error("Connection is closed."); + if (failure === "connection") { + mocks.getRedisClient.mockRejectedValueOnce(cause); + } else { + mocks.claimOutboxEvents.mockRejectedValueOnce(cause); + } + const worker = import("../scripts/outbox-worker"); + + await vi.waitFor(() => { + expect(mocks.closeRedisClient).toHaveBeenCalledTimes(1); + }); + expect(mocks.markOutboxProcessed).not.toHaveBeenCalled(); + await vi.advanceTimersByTimeAsync(5_000); + await worker; + + expect(mocks.getRedisClient).toHaveBeenCalledTimes(2); + expect(mocks.markOutboxProcessed).toHaveBeenCalledWith(1); + expect(mocks.closeDb).toHaveBeenCalledTimes(2); + expect(mocks.closeRedisClient).toHaveBeenCalledTimes(2); + expect(console.error).toHaveBeenCalledWith( + "Outbox worker failed; retrying in 5 seconds", + cause, + ); + }, + ); + + it("stops during the retry delay without opening another connection", async () => { + mocks.getRedisClient.mockRejectedValueOnce(new Error("Connection is closed.")); + const worker = import("../scripts/outbox-worker"); + await vi.waitFor(() => { + expect(mocks.closeRedisClient).toHaveBeenCalledTimes(1); + }); + + stop(); + await vi.advanceTimersByTimeAsync(5_000); + await worker; + + expect(mocks.getRedisClient).toHaveBeenCalledTimes(1); + expect(mocks.claimOutboxEvents).not.toHaveBeenCalled(); + expect(mocks.closeRedisClient).toHaveBeenCalledTimes(2); + }); +});