import "dotenv/config"; import { Worker } from "bullmq"; import { JOB_QUEUES, workerConnection, closeJobQueues, scheduleRecurringJobs, type JobPayload } from "../src/server/jobs/queues"; import { PROCESSORS } from "../src/server/jobs/processors"; /** Craftvia background worker: `npm run worker:craftvia`. One BullMQ worker per registered queue. */ async function main() { const connection = workerConnection(); if (!connection) { console.error("[worker] REDIS_URL is not set — nothing to do."); process.exit(1); } const workers: Worker[] = []; for (const name of Object.values(JOB_QUEUES)) { const load = PROCESSORS[name]; if (!load) { console.warn(`[worker] no processor registered for ${name} — skipped`); continue; } const processor = await load(); const geocode = name === JOB_QUEUES.geocodeSite; // L13: OSM Nominatim policy — max. 1 request/s const concurrency = name === JOB_QUEUES.reportPdf ? 2 : geocode ? 1 : 4; const w = new Worker(name, async (job) => processor(job.data), { connection, concurrency, ...(geocode ? { limiter: { max: 1, duration: 1_000 } } : {}) }); w.on("failed", (job, err) => console.error(`[worker] ${name} job ${job?.id} failed:`, err.message)); workers.push(w); console.info(`[worker] listening on ${name}`); } // L10b: recurring jobs (AI log retention); a scheduling failure must not stop the queue workers await scheduleRecurringJobs(connection).then( () => console.info("[worker] recurring jobs scheduled"), (err) => console.error("[worker] scheduling recurring jobs failed:", (err as Error).message), ); const shutdown = async () => { await Promise.all(workers.map((w) => w.close())); await closeJobQueues(); process.exit(0); }; process.on("SIGTERM", shutdown); process.on("SIGINT", shutdown); } main().catch((err) => { console.error(err); process.exit(1); });