107 lines
3.5 KiB
TypeScript
107 lines
3.5 KiB
TypeScript
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);
|
|
});
|
|
});
|