Files
buzz-sheet/scripts/outbox-worker.ts
T
gunshiz dc5c667b4e
CI / Verify (push) Successful in 1m0s
CI / Build immutable images and deploy (push) Failing after 2m38s
feat(guides): add structured image-first editor
2026-08-29 13:30:52 +00:00

76 lines
2.0 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 { 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<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);
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()]);
}