import { Queue } from "bullmq"; import type { Redis } from "ioredis"; import { computeOpsSummary, type OpsQueueSummary } from "@app/database/ops-summary"; import { JOB_QUEUE_NAME } from "@app/shared-types"; import type { WorkerContext } from "../context.js"; const CACHE_KEY = "ops:summary:cache"; const CACHE_TS_KEY = "ops:summary:computed_at"; const CACHE_TTL_SECONDS = 120; const BUDGET_KEY_PREFIX = "gemini:daily:budget:"; function budgetKey(): string { return `${BUDGET_KEY_PREFIX}${new Date().toISOString().slice(0, 10)}`; } async function fetchQueueSummary(redis: Redis): Promise { const queue = new Queue(JOB_QUEUE_NAME, { connection: redis }); try { const [counts, waiting, workers, failed] = await Promise.all([ queue.getJobCounts("wait", "active", "delayed", "completed", "failed"), queue.getJobs(["wait"], 0, 0, true), queue.getWorkers(), queue.getJobs(["failed"], 0, 1000, true), ]); const oldestWaiting = waiting[0]; const dayAgo = Date.now() - 24 * 60 * 60 * 1000; const failed24h = failed.filter( (job) => (typeof job.finishedOn === "number" && job.finishedOn >= dayAgo) || (typeof job.timestamp === "number" && job.timestamp >= dayAgo), ).length; return { vantande: counts.wait ?? 0, aktiva: counts.active ?? 0, misslyckade_24h: failed24h, aldsta_vantande_sek: oldestWaiting ? Math.max(0, Math.floor((Date.now() - oldestWaiting.timestamp) / 1000)) : null, workers_ok: workers.length > 0, }; } finally { await queue.close(); } } /** * Compute the full ops summary and write it to Redis. Called by the * REFRESH_OPS_SUMMARY_CACHE worker job every 60 s. */ export async function refreshOpsSummaryCache(ctx: WorkerContext): Promise { if (!ctx.redis) { throw new Error("[refreshOpsSummaryCache] Redis krävs för cache- och kö-aggregering."); } const redis = ctx.redis; const [queueSummary, dailySpendRaw] = await Promise.all([ fetchQueueSummary(redis), redis.get(budgetKey()), ]); const dailySpendUsd = dailySpendRaw ? Number(dailySpendRaw) : null; const budgetUsd = Number(process.env.GEMINI_DAILY_BUDGET_USD ?? 0); const summary = await computeOpsSummary({ db: ctx.db, budgetUsd, dailySpendUsd, queueSummary, }); const json = JSON.stringify(summary); await redis.set(CACHE_KEY, json, "EX", CACHE_TTL_SECONDS); await redis.set(CACHE_TS_KEY, summary.as_of, "EX", CACHE_TTL_SECONDS); } export { CACHE_KEY, CACHE_TS_KEY };