spb/datacenterindex
Public
HTML 53.9%
TypeScript 44.5%
JavaScript 0.6%
SQL 0.5%
1/**2 * BullMQ producers for the worker queues. Mirrors apps/worker/src/scheduler.ts: prefix `dci`, queues `crawl` and3 * `maintenance` (Redis keys `dci:crawl:*`, `dci:maintenance:*`). BullMQ forbids ":" in queue names and custom job4 * ids, hence `manual__<connector>__<ts>` ids. Lazy: Redis is touched on first use.5 */6import { Queue } from "bullmq";7import { getQueueRedis } from "./redis.js";89export const QUEUE_PREFIX = "dci";10export const CRAWL_QUEUE = "crawl";11export const MAINTENANCE_QUEUE = "maintenance";12export const INGEST_QUEUE = "ingest";1314export type CrawlTask = "discover" | "crawl" | "full" | "reprocess";1516export interface CrawlJobData {17 connectorId: string;18 task: CrawlTask;19 group?: string;20 limit?: number;21 /** explicit URLs (document reprocess / refetch); honoured by the worker when supported */22 urls?: string[];23 force?: boolean;24 requestedBy?: string;25}2627export type MaintenanceKind = "rankings" | "metrics" | "refresh-stats" | "cleanup" | "quality" | "snapshot";28export interface MaintenanceJobData { kind: MaintenanceKind; task: MaintenanceKind; requestedBy?: string; [k: string]: unknown }2930const queues = new Map<string, Queue>();31function queue(name: string): Queue {32 let q = queues.get(name);33 if (!q) {34 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 } });35 queues.set(name, q);36 }37 return q;38}3940export function manualJobId(connectorId: string, suffix?: string): string {41 const clean = (s: string) => s.replace(/[^a-zA-Z0-9_-]/g, "-");42 return `manual__${clean(connectorId)}${suffix ? `__${clean(suffix)}` : ""}__${Date.now()}`;43}4445export async function enqueueCrawl(data: CrawlJobData, opts: { jobId?: string; priority?: number } = {}): Promise<{ id: string; name: string; queue: string }> {46 const jobId = opts.jobId ?? manualJobId(data.connectorId);47 const job = await queue(CRAWL_QUEUE).add(`${data.task}:${data.connectorId}`, data, { jobId, priority: opts.priority ?? 1 });48 return { id: job.id ?? jobId, name: job.name, queue: `${QUEUE_PREFIX}:${CRAWL_QUEUE}` };49}5051export async function enqueueMaintenance(kind: MaintenanceKind, extra: Record<string, unknown> = {}): Promise<{ id: string; name: string; queue: string }> {52 const jobId = manualJobId(kind);53 const data: MaintenanceJobData = { ...extra, kind, task: kind, requestedBy: "admin" };54 const job = await queue(MAINTENANCE_QUEUE).add(kind, data, { jobId });55 return { id: job.id ?? jobId, name: job.name, queue: `${QUEUE_PREFIX}:${MAINTENANCE_QUEUE}` };56}5758export async function queueCounts(): Promise<Record<string, Record<string, number> | { error: string }>> {59 const out: Record<string, Record<string, number> | { error: string }> = {};60 for (const name of [CRAWL_QUEUE, MAINTENANCE_QUEUE, INGEST_QUEUE]) {61 try {62 const counts = await Promise.race([queue(name).getJobCounts("waiting", "active", "completed", "failed", "delayed", "paused", "prioritized"), new Promise<never>((_, rej) => setTimeout(() => rej(new Error("timeout")), 2_000))]);63 out[`${QUEUE_PREFIX}:${name}`] = counts;64 } catch (e) {65 out[`${QUEUE_PREFIX}:${name}`] = { error: (e as Error).message };66 }67 }68 return out;69}7071export async function closeQueues(): Promise<void> {72 await Promise.allSettled([...queues.values()].map((q) => q.close()));73 queues.clear();74}75