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); }); });