89 lines
2.5 KiB
TypeScript
89 lines
2.5 KiB
TypeScript
import { closeDb } from "@/db/client";
|
|
import {
|
|
claimOutboxEvents,
|
|
markOutboxProcessed,
|
|
releaseOutboxEvent,
|
|
} from "@/lib/outbox/repository";
|
|
import {
|
|
createRedisOutboxTransport,
|
|
processOutboxEvent,
|
|
} from "@/lib/outbox/processor";
|
|
import { closeRedisClient, getRedisClient } from "@/lib/redis/client";
|
|
import { removeExpiredUnpaidCheckouts } from "@/lib/commission/cleanup";
|
|
import { sendGuidePublicationNotifications } from "@/lib/guides/notifications";
|
|
|
|
let stopping = false;
|
|
const stop = () => {
|
|
stopping = true;
|
|
};
|
|
process.once("SIGINT", stop);
|
|
process.once("SIGTERM", stop);
|
|
|
|
function wait(milliseconds: number): Promise<void> {
|
|
return new Promise((resolve) => setTimeout(resolve, milliseconds));
|
|
}
|
|
|
|
function positiveInteger(
|
|
name: string,
|
|
fallback: number,
|
|
maximum: number,
|
|
): number {
|
|
const parsed = Number(process.env[name]);
|
|
return Number.isSafeInteger(parsed) && parsed > 0
|
|
? Math.min(parsed, maximum)
|
|
: fallback;
|
|
}
|
|
|
|
async function run(): Promise<void> {
|
|
const redis = await getRedisClient();
|
|
const transport = createRedisOutboxTransport(redis);
|
|
const batchSize = positiveInteger("OUTBOX_BATCH_SIZE", 20, 100);
|
|
const pollMs = positiveInteger("OUTBOX_POLL_MS", 1_000, 60_000);
|
|
const leaseMs = positiveInteger("OUTBOX_LEASE_MS", 30_000, 300_000);
|
|
let nextCheckoutCleanup = 0;
|
|
|
|
while (!stopping) {
|
|
if (Date.now() >= nextCheckoutCleanup) {
|
|
nextCheckoutCleanup = Date.now() + 60_000;
|
|
try {
|
|
await removeExpiredUnpaidCheckouts();
|
|
} catch (cause) {
|
|
console.error("Unable to remove expired unpaid checkouts", cause);
|
|
}
|
|
}
|
|
const events = await claimOutboxEvents(batchSize, leaseMs);
|
|
if (events.length === 0) {
|
|
await wait(pollMs);
|
|
continue;
|
|
}
|
|
await Promise.all(
|
|
events.map(async (event) => {
|
|
try {
|
|
if (event.eventType === "guide.published") {
|
|
await sendGuidePublicationNotifications(event);
|
|
} else {
|
|
await processOutboxEvent(event, transport);
|
|
}
|
|
await markOutboxProcessed(event.id);
|
|
} catch (cause) {
|
|
await releaseOutboxEvent(event.id, event.attempts, cause);
|
|
}
|
|
}),
|
|
);
|
|
}
|
|
}
|
|
|
|
try {
|
|
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()]);
|
|
}
|