Files
Cibello-app/apps/worker/src/index.ts
T
Sven (AAMOS AI) af05f084b8 feat(worker): EOC push-agent for ops summary
- Add EOC_URL + EOC_PUSH_TOKEN config
- Add pushOpsSummaryToEoc processor
- Schedule PUSH_OPS_SUMMARY_TO_EOC every 2 minutes
- Add unit tests
2026-08-11 17:03:53 +07:00

376 lines
12 KiB
TypeScript
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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 "<slug>-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<string, unknown>;
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<void> {
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}`);
}