/** * BullMQ producers for the worker queues. Mirrors apps/worker/src/scheduler.ts: prefix `dci`, queues `crawl` and * `maintenance` (Redis keys `dci:crawl:*`, `dci:maintenance:*`). BullMQ forbids ":" in queue names and custom job * ids, hence `manual____` ids. Lazy: Redis is touched on first use. */ import { Queue } from "bullmq"; import { getQueueRedis } from "./redis.js"; export const QUEUE_PREFIX = "dci"; export const CRAWL_QUEUE = "crawl"; export const MAINTENANCE_QUEUE = "maintenance"; export const INGEST_QUEUE = "ingest"; export type CrawlTask = "discover" | "crawl" | "full" | "reprocess"; export interface CrawlJobData { connectorId: string; task: CrawlTask; group?: string; limit?: number; /** explicit URLs (document reprocess / refetch); honoured by the worker when supported */ urls?: string[]; force?: boolean; requestedBy?: string; } export type MaintenanceKind = "rankings" | "metrics" | "refresh-stats" | "cleanup" | "quality" | "snapshot"; export interface MaintenanceJobData { kind: MaintenanceKind; task: MaintenanceKind; requestedBy?: string; [k: string]: unknown } const queues = new Map(); function queue(name: string): Queue { let q = queues.get(name); if (!q) { q = new Queue(name, { connection: getQueueRedis(), prefix: QUEUE_PREFIX, defaultJobOptions: { removeOnComplete: { count: 500, age: 86_400 }, removeOnFail: { count: 500, age: 7 * 86_400 }, attempts: 1 } }); queues.set(name, q); } return q; } export function manualJobId(connectorId: string, suffix?: string): string { const clean = (s: string) => s.replace(/[^a-zA-Z0-9_-]/g, "-"); return `manual__${clean(connectorId)}${suffix ? `__${clean(suffix)}` : ""}__${Date.now()}`; } export async function enqueueCrawl(data: CrawlJobData, opts: { jobId?: string; priority?: number } = {}): Promise<{ id: string; name: string; queue: string }> { const jobId = opts.jobId ?? manualJobId(data.connectorId); const job = await queue(CRAWL_QUEUE).add(`${data.task}:${data.connectorId}`, data, { jobId, priority: opts.priority ?? 1 }); return { id: job.id ?? jobId, name: job.name, queue: `${QUEUE_PREFIX}:${CRAWL_QUEUE}` }; } export async function enqueueMaintenance(kind: MaintenanceKind, extra: Record = {}): Promise<{ id: string; name: string; queue: string }> { const jobId = manualJobId(kind); const data: MaintenanceJobData = { ...extra, kind, task: kind, requestedBy: "admin" }; const job = await queue(MAINTENANCE_QUEUE).add(kind, data, { jobId }); return { id: job.id ?? jobId, name: job.name, queue: `${QUEUE_PREFIX}:${MAINTENANCE_QUEUE}` }; } export async function queueCounts(): Promise | { error: string }>> { const out: Record | { error: string }> = {}; for (const name of [CRAWL_QUEUE, MAINTENANCE_QUEUE, INGEST_QUEUE]) { try { const counts = await Promise.race([queue(name).getJobCounts("waiting", "active", "completed", "failed", "delayed", "paused", "prioritized"), new Promise((_, rej) => setTimeout(() => rej(new Error("timeout")), 2_000))]); out[`${QUEUE_PREFIX}:${name}`] = counts; } catch (e) { out[`${QUEUE_PREFIX}:${name}`] = { error: (e as Error).message }; } } return out; } export async function closeQueues(): Promise { await Promise.allSettled([...queues.values()].map((q) => q.close())); queues.clear(); }