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 { claimCatalogSyncJob, processCatalogSyncJob } from "@/lib/catalog/sync"; let stopping = false; const stop = () => { stopping = true; }; process.once("SIGINT", stop); process.once("SIGTERM", stop); function wait(milliseconds: number): Promise { 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 { 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); while (!stopping) { const catalogJob = await claimCatalogSyncJob(); if (catalogJob) { try { await processCatalogSyncJob(catalogJob.id); } catch (cause) { console.error("Catalog sync failed", cause); } continue; } const events = await claimOutboxEvents(batchSize, leaseMs); if (events.length === 0) { await wait(pollMs); continue; } await Promise.all( events.map(async (event) => { try { await processOutboxEvent(event, transport); await markOutboxProcessed(event.id); } catch (cause) { await releaseOutboxEvent(event.id, event.attempts, cause); } }), ); } } try { await run(); } finally { await Promise.all([closeDb(), closeRedisClient()]); }