import "dotenv/config"; import { Queue, Worker, type Job } from "bullmq"; import IORedis from "ioredis"; import { createServer } from "node:http"; import { createContext } from "./context.js"; import { processScanJob } from "./processors/scans.js"; import { processModerateRecipe } from "./processors/moderation.js"; import { processTranslateRecipe } from "./processors/translation.js"; import { processGenerateWeekPlan } from "./processors/weekplan.js"; import { processCookingSessionTimeout, processExpiryNotifications, processMealBoxReminders, processMemorySync, processOutbox, processProactiveTips, processRetention, processStoreNotification, processSubscriptionSweep, processTrainingExport, processTrustDecay, } from "./processors/maintenance.js"; import { processSafetyCanary } from "./processors/safety-canary.js"; import { refreshOpsSummaryCache } from "./processors/ops-summary.js"; import { pushOpsSummaryToEoc } from "./processors/eoc-push.js"; import { DEAD_LETTER_QUEUE_NAME, JOB_QUEUE_NAME as QUEUE_NAME, WORKER_DEFAULT_ATTEMPTS, WORKER_DEFAULT_BACKOFF_MS, } from "@app/shared-types"; /** * workern. Konsumerar kön "-jobs" (spec §50). * * Flera instanser kan köra samtidigt mot samma Redis- och DB-anslutning * utan kodändring (BullMQ hanterar fördelning). Varje instans exponerar * en liten HTTP-healthcheck så att lastbalansering kan se om den lever. */ const connection = new IORedis(process.env.REDIS_URL ?? "redis://localhost:6379", { maxRetriesPerRequest: null, }); const ctx = createContext(connection); const instanceId = process.env.WORKER_INSTANCE_ID ?? `worker-${process.pid}`; const deadLetterQueue = new Queue(DEAD_LETTER_QUEUE_NAME, { connection }); const queue = new Queue(QUEUE_NAME, { connection }); const worker = new Worker( QUEUE_NAME, async (job: Job) => { const data = job.data as Record; const jobType = (data.jobType as string) ?? job.name; log(`▶ ${jobType} (${job.id})`); switch (jobType) { case "ANALYZE_FRIDGE_IMAGE": case "ANALYZE_PANTRY_IMAGE": case "ANALYZE_MEAL_IMAGE": case "READ_RECEIPT": case "READ_NUTRITION_LABEL": case "READ_EXPIRY_DATE": return processScanJob(ctx, String(data.scanJobId)); case "MODERATE_RECIPE": return processModerateRecipe(ctx, String(data.recipeId)); case "TRANSLATE_RECIPE": return processTranslateRecipe(ctx, { recipeId: String(data.recipeId), targetLanguageTag: String(data.targetLanguageTag), }); case "GENERATE_WEEK_PLAN": return processGenerateWeekPlan(ctx, data as never); case "PROCESS_STORE_NOTIFICATION": return processStoreNotification(ctx, String(data.notificationId)); case "SEND_EXPIRY_NOTIFICATION": { const expiry = await processExpiryNotifications(ctx); const boxes = await processMealBoxReminders(ctx); log(`Notiser: ${expiry} bäst före, ${boxes} matlådor`); return; } case "SEND_PROACTIVE_TIPS": { const tips = await processProactiveTips(ctx); log(`Proaktiva puffar: ${tips}`); return; } case "UPDATE_USER_MEMORY": { const updates = await processMemorySync(ctx); log(`Minnesuppdateringar: ${updates}`); return; } case "VERIFY_SUBSCRIPTION": { const expired = await processSubscriptionSweep(ctx); log(`Prenumerationssvep: ${expired} markerade som utgångna`); return; } case "BUILD_TRAINING_SAMPLE": { const exported = await processTrainingExport(ctx); log(`Träningsexport: ${exported} korrigeringar (med samtycke)`); return; } case "RUN_RETENTION": { const removed = await processRetention(ctx); log(`Retention: ${JSON.stringify(removed)}`); return; } case "PUBLISH_OUTBOX": { const published = await processOutbox(ctx); if (published > 0) log(`Outbox: ${published} events publicerade`); return; } case "UPDATE_TRUST_STATES": { const updated = await processTrustDecay(ctx); if (updated > 0) log(`Trust decay: ${updated} items uppdaterade`); return; } case "COOKING_SESSION_TIMEOUT": { const cancelled = await processCookingSessionTimeout(ctx); if (cancelled > 0) log(`Cooking session timeout: ${cancelled} avbrutna`); return; } case "RUN_SAFETY_CANARY": { const result = await processSafetyCanary(ctx); log(`Safety canary: ${JSON.stringify(result)}`); return; } case "REFRESH_OPS_SUMMARY_CACHE": { await refreshOpsSummaryCache(ctx); log("Ops summary cache uppdaterad."); return; } case "PUSH_OPS_SUMMARY_TO_EOC": { const result = await pushOpsSummaryToEoc(ctx.redis!); log(`EOC push: ${result.ok ? "OK" : "FEL"}${result.error ? ` (${result.error})` : ""}`); return; } // Deterministiska/planerade jobb som inte kräver egen processor ännu case "NORMALIZE_PRODUCTS": case "DEDUPLICATE_INVENTORY": case "CALCULATE_NUTRITION": case "GENERATE_RECIPE_OPTIONS": case "RANK_RECIPES": case "RUN_AI_EVALUATION": log(`${jobType}: hanteras synkront i API:t eller aktiveras i senare fas`); return; default: throw new Error(`Okänd jobbtyp: ${jobType}`); } }, { connection, concurrency: Number(process.env.WORKER_CONCURRENCY ?? 5), maxStalledCount: 2, stalledInterval: 30_000, limiter: { max: Number(process.env.WORKER_RATE_LIMIT_MAX ?? 60), duration: Number(process.env.WORKER_RATE_LIMIT_DURATION_MS ?? 60_000), }, }, ); worker.on("completed", (job) => log(`✓ ${job.name} (${job.id})`)); worker.on("failed", async (job, err) => { const jobType = (job?.data?.jobType as string) ?? job?.name ?? "unknown"; const attempts = job?.opts?.attempts ?? WORKER_DEFAULT_ATTEMPTS; const attemptsMade = job?.attemptsMade ?? 0; console.error(`✗ ${jobType} (${job?.id}): ${err.message} (attempt ${attemptsMade}/${attempts})`); if (job && attemptsMade >= attempts) { try { await deadLetterQueue.add( job.name, { ...job.data, _deadLetteredAt: new Date().toISOString(), _instanceId: instanceId, _failureReason: err.message, _failedAt: new Date().toISOString(), }, { jobId: `${instanceId}:${job.id}`, removeOnComplete: { age: 30 * 24 * 60 * 60 }, // 30 dagar removeOnFail: { age: 90 * 24 * 60 * 60 }, // 90 dagar }, ); log(`→ dead letter: ${jobType} (${job.id})`); } catch (dlqErr) { console.error("Kunde inte skriva till dead letter queue:", dlqErr); } } await alertWebhook( `Worker-jobb misslyckades: ${jobType} (${job?.id}): ${err.message} (attempt ${attemptsMade}/${attempts})`, ); }); /** Larm till valfri webhook (Slack/Discord/Teams …) – fire-and-forget. */ async function alertWebhook(text: string): Promise { const url = process.env.ERROR_WEBHOOK_URL; if (!url) return; try { await fetch(url, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ text }), signal: AbortSignal.timeout(5000), }); } catch { // Larmet får aldrig fälla arbetsflödet. } } // --- Återkommande jobb via BullMQ Job Schedulers (spec §54) --- async function registerRepeatableJobs() { const baseOpts = { attempts: WORKER_DEFAULT_ATTEMPTS, backoff: { type: "exponential" as const, delay: WORKER_DEFAULT_BACKOFF_MS }, removeOnComplete: { count: 100 }, removeOnFail: { count: 100 }, }; await queue.upsertJobScheduler( "scheduler-outbox", { every: 30_000 }, { name: "PUBLISH_OUTBOX", data: { jobType: "PUBLISH_OUTBOX" }, opts: baseOpts }, ); await queue.upsertJobScheduler( "scheduler-trust", { every: 6 * 60 * 60 * 1000 }, // var 6:e timme; decay lever på dygnsskala { name: "UPDATE_TRUST_STATES", data: { jobType: "UPDATE_TRUST_STATES" }, opts: baseOpts }, ); await queue.upsertJobScheduler( "scheduler-cooking-timeout", { every: 60 * 60 * 1000 }, // varje timme räcker för 24 h-timeout { name: "COOKING_SESSION_TIMEOUT", data: { jobType: "COOKING_SESSION_TIMEOUT" }, opts: baseOpts, }, ); await queue.upsertJobScheduler( "scheduler-expiry", { pattern: "0 7 * * *", tz: "Europe/Stockholm" }, { name: "SEND_EXPIRY_NOTIFICATION", data: { jobType: "SEND_EXPIRY_NOTIFICATION" }, opts: baseOpts, }, ); await queue.upsertJobScheduler( "scheduler-memory", { pattern: "30 3 * * *", tz: "Europe/Stockholm" }, { name: "UPDATE_USER_MEMORY", data: { jobType: "UPDATE_USER_MEMORY" }, opts: baseOpts }, ); await queue.upsertJobScheduler( "scheduler-subs", { pattern: "15 4 * * *", tz: "Europe/Stockholm" }, { name: "VERIFY_SUBSCRIPTION", data: { jobType: "VERIFY_SUBSCRIPTION" }, opts: baseOpts }, ); await queue.upsertJobScheduler( "scheduler-retention", { pattern: "45 4 * * *", tz: "Europe/Stockholm" }, { name: "RUN_RETENTION", data: { jobType: "RUN_RETENTION" }, opts: baseOpts }, ); await queue.upsertJobScheduler( "scheduler-training", { pattern: "45 2 * * 0", tz: "Europe/Stockholm" }, { name: "BUILD_TRAINING_SAMPLE", data: { jobType: "BUILD_TRAINING_SAMPLE" }, opts: baseOpts }, ); await queue.upsertJobScheduler( "scheduler-safety-canary", { every: 60 * 60 * 1000 }, // varje timme { name: "RUN_SAFETY_CANARY", data: { jobType: "RUN_SAFETY_CANARY" }, opts: baseOpts }, ); await queue.upsertJobScheduler( "scheduler-ops-summary", { every: 60 * 1000 }, // var 60:e sekund { name: "REFRESH_OPS_SUMMARY_CACHE", data: { jobType: "REFRESH_OPS_SUMMARY_CACHE" }, opts: baseOpts, }, ); await queue.upsertJobScheduler( "scheduler-eoc-push", { every: 2 * 60 * 1000 }, // var 2:e minut { name: "PUSH_OPS_SUMMARY_TO_EOC", data: { jobType: "PUSH_OPS_SUMMARY_TO_EOC" }, opts: baseOpts, }, ); } // --- Minimal healthcheck-server så att flera instanser kan övervakas --- const healthPort = Number(process.env.WORKER_HEALTH_PORT ?? 4001); let processedCount = 0; let failedCount = 0; worker.on("completed", () => processedCount++); worker.on("failed", () => failedCount++); const healthServer = createServer(async (req, res) => { if (req.url === "/healthz") { const [queueCount, dlqCount, workers] = await Promise.all([ queue.getJobCounts("wait", "active", "delayed", "completed", "failed"), deadLetterQueue.getJobCounts("wait", "active", "completed", "failed"), queue.getWorkers(), ]); res.writeHead(200, { "content-type": "application/json" }); res.end( JSON.stringify({ ok: true, instanceId, uptimeSeconds: Math.floor(process.uptime()), queue: queueCount, deadLetter: dlqCount, activeWorkers: workers.length, processedSinceStart: processedCount, failedSinceStart: failedCount, }), ); return; } res.writeHead(404); res.end("not found"); }); async function start() { await registerRepeatableJobs(); log("Worker igång. Väntar på jobb …"); healthServer.listen(healthPort, () => { log(`Healthcheck på :${healthPort}`); }); } start().catch((err) => { console.error("Kunde inte starta worker:", err); process.exit(1); }); const shutdown = async () => { log("Stänger ner …"); healthServer.close(); await worker.close(); await queue.close(); await deadLetterQueue.close(); connection.disconnect(); await ctx.close(); process.exit(0); }; process.on("SIGTERM", () => void shutdown()); process.on("SIGINT", () => void shutdown()); function log(msg: string) { console.log(`[worker ${instanceId} ${new Date().toISOString()}] ${msg}`); }