import "dotenv/config"; import { Queue, Worker, type Job } from "bullmq"; import IORedis from "ioredis"; 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 { processExpiryNotifications, processMealBoxReminders, processMemorySync, processOutbox, processRetention, processStoreNotification, processSubscriptionSweep, processTrainingExport, } from "./processors/maintenance.js"; /** * workern. Konsumerar kön "-jobs" (spec §50): * App → signed S3 upload → API job → worker → AAMOS → result → app * * Jobbnamn = jobbtyp (spec §54). Deterministiska jobb rör aldrig AAMOS; * AI-jobb går alltid via de typade kontrakten i @app/ai-contracts. */ import { JOB_QUEUE_NAME as QUEUE_NAME } from "@app/shared-types"; const connection = new IORedis(process.env.REDIS_URL ?? "redis://localhost:6379", { maxRetriesPerRequest: null, }); const ctx = createContext(); 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 "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; } // 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), }, ); worker.on("completed", (job) => log(`✓ ${job.name} (${job.id})`)); worker.on("failed", (job, err) => { console.error(`✗ ${job?.name} (${job?.id}): ${err.message}`); void alertWebhook(`Worker-jobb misslyckades: ${job?.name} (${job?.id}): ${err.message}`); }); /** 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) --- const queue = new Queue(QUEUE_NAME, { connection }); async function registerRepeatableJobs() { await queue.upsertJobScheduler( "scheduler-outbox", { every: 30_000 }, { name: "PUBLISH_OUTBOX", data: { jobType: "PUBLISH_OUTBOX" } }, ); await queue.upsertJobScheduler( "scheduler-expiry", { pattern: "0 7 * * *", tz: "Europe/Stockholm" }, { name: "SEND_EXPIRY_NOTIFICATION", data: { jobType: "SEND_EXPIRY_NOTIFICATION" } }, ); await queue.upsertJobScheduler( "scheduler-memory", { pattern: "30 3 * * *", tz: "Europe/Stockholm" }, { name: "UPDATE_USER_MEMORY", data: { jobType: "UPDATE_USER_MEMORY" } }, ); await queue.upsertJobScheduler( "scheduler-subs", { pattern: "15 4 * * *", tz: "Europe/Stockholm" }, { name: "VERIFY_SUBSCRIPTION", data: { jobType: "VERIFY_SUBSCRIPTION" } }, ); await queue.upsertJobScheduler( "scheduler-retention", { pattern: "45 4 * * *", tz: "Europe/Stockholm" }, { name: "RUN_RETENTION", data: { jobType: "RUN_RETENTION" } }, ); await queue.upsertJobScheduler( "scheduler-training", { pattern: "45 2 * * 0", tz: "Europe/Stockholm" }, { name: "BUILD_TRAINING_SAMPLE", data: { jobType: "BUILD_TRAINING_SAMPLE" } }, ); } registerRepeatableJobs() .then(() => log("Worker igång. Väntar på jobb …")) .catch((err) => { console.error("Kunde inte registrera återkommande jobb:", err); }); const shutdown = async () => { log("Stänger ner …"); await worker.close(); await queue.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 ${new Date().toISOString()}] ${msg}`); }