fix(outbox) : retry temporary worker connection failures
This commit is contained in:
@@ -74,7 +74,15 @@ async function run(): Promise<void> {
|
||||
}
|
||||
|
||||
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()]);
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user