SPB Git forge
38commits 1branches 0releases
338.7 MBsize
maindefault branch
3 h agolast push
HTML 53.9% TypeScript 44.5% JavaScript 0.6% SQL 0.5%
3.4 KB · 75 lines typescript
Raw Blame History
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