Skip to content

Commit 8ae3e5e

Browse files
committed
fix(billing): make Autumn ingestion replay-safe
1 parent 11f240b commit 8ae3e5e

35 files changed

Lines changed: 2159 additions & 142 deletions
Lines changed: 136 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,136 @@
1+
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
2+
3+
const state = vi.hoisted(() => ({ error: vi.fn(), warn: vi.fn() }));
4+
5+
vi.mock("evlog", () => ({ log: { error: state.error, warn: state.warn } }));
6+
vi.mock("@/routes/webhooks/autumn", () => ({
7+
replayDeferredAutumnWebhooks: vi.fn(async () => ({
8+
completed: 0,
9+
deadLettered: 0,
10+
deferred: 0,
11+
failed: [],
12+
})),
13+
}));
14+
vi.mock("@/routes/webhooks/autumn-inbox", () => ({
15+
deleteCompletedAutumnWebhooks: vi.fn(async () => 0),
16+
deleteDeadLetterAutumnWebhooks: vi.fn(async () => 0),
17+
listUnalertedAutumnWebhookDeadLetters: vi.fn(async () => []),
18+
markAutumnWebhookDeadLettersAlerted: vi.fn(async () => 0),
19+
}));
20+
21+
import { startAutumnWebhookReplayLoop } from "./autumn-webhook-replay";
22+
import { runAutumnWebhookMaintenance } from "./autumn-webhook-replay";
23+
import { replayDeferredAutumnWebhooks } from "@/routes/webhooks/autumn";
24+
import {
25+
deleteCompletedAutumnWebhooks,
26+
deleteDeadLetterAutumnWebhooks,
27+
listUnalertedAutumnWebhookDeadLetters,
28+
markAutumnWebhookDeadLettersAlerted,
29+
} from "@/routes/webhooks/autumn-inbox";
30+
31+
beforeEach(() => {
32+
vi.useFakeTimers();
33+
state.error.mockClear();
34+
state.warn.mockClear();
35+
vi.mocked(deleteCompletedAutumnWebhooks).mockClear();
36+
vi.mocked(deleteDeadLetterAutumnWebhooks).mockClear();
37+
vi.mocked(listUnalertedAutumnWebhookDeadLetters).mockClear();
38+
vi.mocked(markAutumnWebhookDeadLettersAlerted).mockClear();
39+
vi.mocked(replayDeferredAutumnWebhooks).mockClear();
40+
});
41+
42+
afterEach(() => {
43+
vi.useRealTimers();
44+
});
45+
46+
describe("Autumn webhook replay loop", () => {
47+
it("runs bounded replay and retention maintenance in one shot", async () => {
48+
await expect(runAutumnWebhookMaintenance()).resolves.toEqual({
49+
completed: 0,
50+
deadLettered: 0,
51+
deadLetters: 0,
52+
deferred: 0,
53+
deleted: 0,
54+
failed: [],
55+
});
56+
expect(replayDeferredAutumnWebhooks).toHaveBeenCalledWith(100);
57+
expect(deleteCompletedAutumnWebhooks).toHaveBeenCalledWith({ limit: 100 });
58+
expect(deleteDeadLetterAutumnWebhooks).toHaveBeenCalledWith({ limit: 100 });
59+
});
60+
61+
it("reports item-level replay failures that remain queued", async () => {
62+
vi.mocked(replayDeferredAutumnWebhooks).mockResolvedValueOnce({
63+
completed: 0,
64+
deadLettered: 0,
65+
deferred: 0,
66+
failed: ["msg-failed"],
67+
});
68+
69+
await runAutumnWebhookMaintenance();
70+
71+
expect(state.warn).toHaveBeenCalledWith(
72+
expect.objectContaining({
73+
component: "autumn_webhook_replay",
74+
failed_count: 1,
75+
})
76+
);
77+
});
78+
79+
it("alerts once for newly dead-lettered webhooks before retention", async () => {
80+
vi.mocked(listUnalertedAutumnWebhookDeadLetters).mockResolvedValueOnce([
81+
{
82+
attempts: 12,
83+
deadLetteredAt: new Date("2026-07-01T00:00:00.000Z"),
84+
errorMessage: "provider unavailable",
85+
id: "msg-dead",
86+
type: "balances.limit_reached",
87+
},
88+
]);
89+
90+
await runAutumnWebhookMaintenance();
91+
92+
expect(state.error).toHaveBeenCalledWith(
93+
expect.objectContaining({
94+
component: "autumn_webhook_replay",
95+
dead_letter_count: 1,
96+
dead_letter_ids: ["msg-dead"],
97+
})
98+
);
99+
expect(markAutumnWebhookDeadLettersAlerted).toHaveBeenCalledWith([
100+
"msg-dead",
101+
]);
102+
});
103+
104+
it("runs immediately, repeats every minute, and stops deterministically", async () => {
105+
const maintenance = vi.fn(async () => undefined);
106+
const loop = startAutumnWebhookReplayLoop(maintenance);
107+
108+
await loop.run();
109+
expect(maintenance).toHaveBeenCalledTimes(1);
110+
await vi.advanceTimersByTimeAsync(60_000);
111+
expect(maintenance).toHaveBeenCalledTimes(2);
112+
113+
loop.stop();
114+
await vi.advanceTimersByTimeAsync(120_000);
115+
expect(maintenance).toHaveBeenCalledTimes(2);
116+
});
117+
118+
it("logs maintenance failures without rejecting or stopping the loop", async () => {
119+
const maintenance = vi.fn(async () => {
120+
throw new Error("database unavailable");
121+
});
122+
const loop = startAutumnWebhookReplayLoop(maintenance);
123+
124+
await expect(loop.run()).resolves.toBeUndefined();
125+
expect(state.error).toHaveBeenCalledWith(
126+
expect.objectContaining({
127+
component: "autumn_webhook_replay",
128+
error_message: "database unavailable",
129+
})
130+
);
131+
132+
await vi.advanceTimersByTimeAsync(60_000);
133+
expect(maintenance).toHaveBeenCalledTimes(2);
134+
loop.stop();
135+
});
136+
});
Lines changed: 103 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,103 @@
1+
import { log } from "evlog";
2+
import {
3+
deleteCompletedAutumnWebhooks,
4+
deleteDeadLetterAutumnWebhooks,
5+
listUnalertedAutumnWebhookDeadLetters,
6+
markAutumnWebhookDeadLettersAlerted,
7+
} from "@/routes/webhooks/autumn-inbox";
8+
import { replayDeferredAutumnWebhooks } from "@/routes/webhooks/autumn";
9+
10+
const REPLAY_INTERVAL_MS = 60_000;
11+
const REPLAY_BATCH_SIZE = 100;
12+
13+
export interface AutumnWebhookReplayLoop {
14+
run(): Promise<void>;
15+
stop(): void;
16+
}
17+
18+
export async function runAutumnWebhookMaintenance(): Promise<{
19+
completed: number;
20+
deadLettered: number;
21+
deadLetters: number;
22+
deferred: number;
23+
deleted: number;
24+
failed: string[];
25+
}> {
26+
const replay = await replayDeferredAutumnWebhooks(REPLAY_BATCH_SIZE);
27+
if (replay.failed.length > 0) {
28+
log.warn({
29+
service: "api",
30+
component: "autumn_webhook_replay",
31+
message: "Autumn webhook replay batch had failures",
32+
failed_count: replay.failed.length,
33+
});
34+
}
35+
36+
const deadLetters = await listUnalertedAutumnWebhookDeadLetters({
37+
limit: REPLAY_BATCH_SIZE,
38+
});
39+
if (deadLetters.length > 0) {
40+
log.error({
41+
service: "api",
42+
component: "autumn_webhook_replay",
43+
dead_letter_count: deadLetters.length,
44+
dead_letter_ids: deadLetters.map((row) => row.id),
45+
error_message: "Autumn webhooks exhausted replay attempts",
46+
});
47+
await markAutumnWebhookDeadLettersAlerted(deadLetters.map((row) => row.id));
48+
}
49+
50+
const deletedRows = await Promise.all([
51+
deleteCompletedAutumnWebhooks({ limit: REPLAY_BATCH_SIZE }),
52+
deleteDeadLetterAutumnWebhooks({ limit: REPLAY_BATCH_SIZE }),
53+
]);
54+
return {
55+
...replay,
56+
deadLetters: deadLetters.length,
57+
deleted: deletedRows.reduce((total, count) => total + count, 0),
58+
};
59+
}
60+
61+
export function startAutumnWebhookReplayLoop(
62+
maintenance: () => Promise<unknown> = runAutumnWebhookMaintenance
63+
): AutumnWebhookReplayLoop {
64+
let active: Promise<void> | null = null;
65+
let stopped = false;
66+
67+
const run = (): Promise<void> => {
68+
if (stopped) {
69+
return Promise.resolve();
70+
}
71+
if (active) {
72+
return active;
73+
}
74+
active = Promise.resolve()
75+
.then(maintenance)
76+
.then(() => undefined)
77+
.catch((error) => {
78+
log.error({
79+
service: "api",
80+
component: "autumn_webhook_replay",
81+
error_message: error instanceof Error ? error.message : String(error),
82+
});
83+
})
84+
.finally(() => {
85+
active = null;
86+
});
87+
return active;
88+
};
89+
90+
run().catch(() => undefined);
91+
const timer = setInterval(() => {
92+
run().catch(() => undefined);
93+
}, REPLAY_INTERVAL_MS);
94+
timer.unref?.();
95+
96+
return {
97+
run,
98+
stop: () => {
99+
stopped = true;
100+
clearInterval(timer);
101+
},
102+
};
103+
}

apps/api/src/bootstrap/shutdown.ts

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -15,12 +15,12 @@ export function warmPostgresPool() {
1515
);
1616
}
1717

18-
export function registerShutdownHooks() {
19-
process.on("SIGINT", () => shutdownApi("SIGINT"));
20-
process.on("SIGTERM", () => shutdownApi("SIGTERM"));
18+
export function registerShutdownHooks(beforeShutdown?: () => void) {
19+
process.on("SIGINT", () => shutdownApi("SIGINT", beforeShutdown));
20+
process.on("SIGTERM", () => shutdownApi("SIGTERM", beforeShutdown));
2121
}
2222

23-
async function shutdownApi(signal: string) {
23+
async function shutdownApi(signal: string, beforeShutdown?: () => void) {
2424
if (shuttingDown) {
2525
log.info({
2626
lifecycle: "shutdown",
@@ -31,6 +31,7 @@ async function shutdownApi(signal: string) {
3131
}
3232

3333
shuttingDown = true;
34+
beforeShutdown?.();
3435
const timeout = setTimeout(() => {
3536
log.error({
3637
lifecycle: "shutdown",

apps/api/src/index.ts

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@ import cors from "@elysiajs/cors";
44
import { Elysia } from "elysia";
55
import { evlog } from "evlog/elysia";
66
import { handleAutumnRequest } from "@/billing/autumn";
7+
import { startAutumnWebhookReplayLoop } from "@/billing/autumn-webhook-replay";
78
import { configureApiInstrumentation } from "@/bootstrap/instrumentation";
89
import { configureApiLogger } from "@/bootstrap/logger";
910
import { registerProcessErrorHandlers } from "@/bootstrap/process-errors";
@@ -115,8 +116,9 @@ const app = new Elysia({ precompile: true })
115116
.all("/*", handleOpenApiEndpoint, { parse: "none" })
116117
.onError(handleAppError);
117118

119+
const autumnWebhookReplay = startAutumnWebhookReplayLoop();
118120
warmPostgresPool();
119-
registerShutdownHooks();
121+
registerShutdownHooks(autumnWebhookReplay.stop);
120122

121123
export default {
122124
fetch: app.fetch,

0 commit comments

Comments
 (0)