spb/datacenterindex
Public
HTML 53.9%
TypeScript 44.5%
JavaScript 0.6%
SQL 0.5%
1/**2 * BullMQ scheduling on Redis.3 * queue `crawl` (keys dci:crawl:*) jobs { connectorId, task, group?, limit?, force? }4 * queue `maintenance` (keys dci:maintenance:*) jobs { kind: rankings | metrics | refresh-stats | cleanup }5 *6 * The scheduler loop (every DCI_SCHEDULER_INTERVAL_MS, default 60 s) enqueues one crawl job per connector group7 * with due documents (jobId `<connector>__<group>` dedupes), a `discover` job when the connector's discovery8 * cadence is due (a `full` job for a connector that was never discovered), and keeps the repeatable maintenance9 * jobs registered (daily metrics 00:10 UTC, rankings 00:30 UTC, refresh-stats hourly, cleanup 01:00 UTC).10 * connectors.enabled / paused are honored. While a `full` / `discover` job of a connector is pending, no per-group11 * job is enqueued for it (the full job crawls everything it registers; two runs on one connector would double12 * fetches and record phantom "changes").13 *14 * BullMQ forbids ":" in queue names and custom job ids, hence the `dci` prefix + `__` separators.15 *16 * Run as a process (`dci-entrypoint scheduler`, `pnpm --filter @dci/worker scheduler`): the loop plus HTTP17 * /healthz + /metrics on WORKER_PORT. The crawl worker also runs the loop unless DCI_SCHEDULER=0; the Redis lock18 * `dci:scheduler:lock` makes concurrent loops harmless.19 */20import { Queue, type JobsOptions } from "bullmq";21import { Redis } from "ioredis";22import { intervalMs } from "@dci/connectors";23import { getDb, closeDb, connectors as connectorsTable, sql } from "@dci/db";24import { getEnv, loadEnvFile } from "./env.js";25import { nextCheckByGroup } from "./documents.js";26import { startHealthHttp } from "./health-http.js";27import type { RunTask } from "./pipeline.js";28import { metrics } from "./prom.js";29import { minIntervalMs } from "./scheduling.js";3031export const QUEUE_PREFIX = "dci";32export const CRAWL_QUEUE = "crawl";33export const MAINTENANCE_QUEUE = "maintenance";34export const WORKER_STATUS_KEY = "dci:worker:status";35export const SCHEDULER_STATUS_KEY = "dci:scheduler:status";36const SCHEDULER_LOCK_KEY = "dci:scheduler:lock";37/** Redis lock held by a worker while it runs a connector — one run per connector at a time across all workers. */38export const RUN_LOCK_PREFIX = "dci:run-lock";39export const RUN_LOCK_TTL_SECONDS = 4 * 3600;4041export interface CrawlJobData { connectorId: string; task: RunTask; group?: string; limit?: number; force?: boolean; urls?: string[]; requestedBy?: string }42export type MaintenanceKind = "rankings" | "metrics" | "refresh-stats" | "cleanup" | "snapshot" | "quality";43export interface MaintenanceJobData { kind: MaintenanceKind; requestedBy?: string }4445let _redis: Redis | null = null;46/** Shared ioredis connection for BullMQ (maxRetriesPerRequest must be null for BullMQ). Reconnects on its own. */47export function getRedis(): Redis {48 if (!_redis) {49 _redis = new Redis(getEnv().redisUrl, { maxRetriesPerRequest: null, enableReadyCheck: false, lazyConnect: false, retryStrategy: (times) => Math.min(30_000, 200 * 2 ** Math.min(times, 8)) });50 let down = false;51 // an 'error' event without listener would be printed as an unhandled error by ioredis on every retry52 _redis.on("error", (e) => { if (!down) console.error(`[redis] ${e.message} — reconnecting`); down = true; });53 _redis.on("ready", () => { if (down) console.error("[redis] connection restored"); down = false; });54 }55 return _redis;56}5758let _crawlQueue: Queue<CrawlJobData> | null = null;59let _maintQueue: Queue<MaintenanceJobData> | null = null;60export function crawlQueue(): Queue<CrawlJobData> {61 if (!_crawlQueue) _crawlQueue = new Queue<CrawlJobData>(CRAWL_QUEUE, { connection: getRedis(), prefix: QUEUE_PREFIX, defaultJobOptions: { removeOnComplete: { count: 500, age: 86_400 }, removeOnFail: { count: 500, age: 7 * 86_400 }, attempts: 1 } });62 return _crawlQueue;63}64export function maintenanceQueue(): Queue<MaintenanceJobData> {65 if (!_maintQueue) _maintQueue = new Queue<MaintenanceJobData>(MAINTENANCE_QUEUE, { connection: getRedis(), prefix: QUEUE_PREFIX, defaultJobOptions: { removeOnComplete: { count: 200 }, removeOnFail: { count: 200 }, attempts: 1 } });66 return _maintQueue;67}6869export function crawlJobId(connectorId: string, task: RunTask, group?: string): string {70 const clean = (s: string) => s.replace(/[^a-zA-Z0-9_-]/g, "-");71 return task === "discover" ? `${clean(connectorId)}__discover` : `${clean(connectorId)}__${clean(group ?? "all")}`;72}73export function runLockKey(connectorId: string): string { return `${RUN_LOCK_PREFIX}:${connectorId}`; }7475const PENDING_STATES = new Set(["waiting", "active", "delayed", "prioritized", "waiting-children"]);7677/** Enqueue a crawl job (used by scheduler loop, CLI and API). Returns the job id; duplicates are coalesced. */78export async function enqueueRun(data: CrawlJobData, opts: { priority?: number; jobId?: string; delayMs?: number } = {}): Promise<{ jobId: string; queued: boolean }> {79 const jobId = opts.jobId ?? crawlJobId(data.connectorId, data.task, data.group);80 const q = crawlQueue();81 const existing = await q.getJob(jobId);82 if (existing) {83 const state = await existing.getState();84 if (PENDING_STATES.has(state)) return { jobId, queued: false };85 await existing.remove().catch(() => undefined);86 }87 const jobOpts: JobsOptions = { jobId, priority: opts.priority ?? 10 };88 if (opts.delayMs) jobOpts.delay = opts.delayMs;89 await q.add(`${data.task}:${data.connectorId}`, data, jobOpts);90 return { jobId, queued: true };91}9293/** True when a whole-connector job (`full` = `<id>__all`, or `<id>__discover`) is pending for this connector. */94async function wholeConnectorJobPending(connectorId: string): Promise<boolean> {95 const q = crawlQueue();96 for (const id of [crawlJobId(connectorId, "full"), crawlJobId(connectorId, "discover")]) {97 const job = await q.getJob(id);98 if (job && PENDING_STATES.has(await job.getState())) return true;99 }100 return false;101}102103export async function enqueueMaintenance(kind: MaintenanceKind, requestedBy = "cli"): Promise<string> {104 const jobId = `manual__${kind}__${Date.now().toString(36)}`;105 await maintenanceQueue().add(kind, { kind, requestedBy }, { jobId });106 return jobId;107}108109/** Repeatable maintenance schedulers (idempotent upserts). */110export async function ensureMaintenanceSchedulers(): Promise<void> {111 const q = maintenanceQueue();112 const tz = "UTC";113 await q.upsertJobScheduler("daily-metrics", { pattern: "10 0 * * *", tz }, { name: "metrics", data: { kind: "metrics", requestedBy: "scheduler" } });114 await q.upsertJobScheduler("daily-rankings", { pattern: "30 0 * * *", tz }, { name: "rankings", data: { kind: "rankings", requestedBy: "scheduler" } });115 await q.upsertJobScheduler("hourly-refresh-stats", { pattern: "0 * * * *", tz }, { name: "refresh-stats", data: { kind: "refresh-stats", requestedBy: "scheduler" } });116 await q.upsertJobScheduler("daily-cleanup", { pattern: "0 1 * * *", tz }, { name: "cleanup", data: { kind: "cleanup", requestedBy: "scheduler" } });117 // daily snapshots + regression checks after rankings; quality sweep every 6 hours118 await q.upsertJobScheduler("daily-snapshot", { pattern: "45 0 * * *", tz }, { name: "snapshot", data: { kind: "snapshot", requestedBy: "scheduler" } });119 await q.upsertJobScheduler("quality-sweep", { pattern: "20 */6 * * *", tz }, { name: "quality", data: { kind: "quality", requestedBy: "scheduler" } });120}121122export interface TickResult { checked: number; enqueued: string[]; skipped: string[]; lockHeldElsewhere: boolean }123124interface ConnectorSchedRow { id: string; enabled: boolean; paused: boolean; schedule: Record<string, string>; config: Record<string, unknown>; lastDiscoverAt: string | null; docCount: number }125126async function schedulableConnectors(): Promise<ConnectorSchedRow[]> {127 const rows = await getDb().execute<{ id: string; enabled: boolean; paused: boolean; schedule: Record<string, string>; config: Record<string, unknown>; last_discover_at: string | null; doc_count: number }>(sql`128 select c.id, c.enabled, c.paused, c.schedule, c.config,129 (select value #>> '{}' from connector_state s where s.connector_id = c.id and s.key = 'lastDiscoverAt') as last_discover_at,130 (select count(*)::int from documents d where d.connector_id = c.id) as doc_count131 from connectors c order by c.id`);132 return rows.map((r) => ({ id: r.id, enabled: r.enabled, paused: r.paused, schedule: r.schedule ?? {}, config: r.config ?? {}, lastDiscoverAt: r.last_discover_at, docCount: Number(r.doc_count) }));133}134135/** Discovery cadence: schedule.discovery → schedule.sitemap → shortest group interval (≥ 6 h). */136export function discoveryIntervalMs(schedule: Record<string, string>): number {137 for (const k of ["discovery", "sitemap", "seed"]) { const spec = schedule[k]; if (spec) { try { return Math.max(6 * 3_600_000, intervalMs(spec)); } catch { /* fallthrough */ } } }138 return Math.max(6 * 3_600_000, minIntervalMs(schedule));139}140141/**142 * Is a discovery pass due? `lastDiscoverAt` is authoritative: a connector whose discovery legitimately yields zero143 * URLs (blocked feed, filter matching nothing) must wait for the next cadence like any other, not run every tick.144 */145export function discoveryDue(c: { docCount: number; lastDiscoverAt: string | null }, discEveryMs: number, now: Date): { due: boolean; task: RunTask } {146 const task: RunTask = c.docCount === 0 ? "full" : "discover";147 const last = c.lastDiscoverAt ? Date.parse(c.lastDiscoverAt) : NaN;148 if (!Number.isFinite(last)) return { due: true, task };149 return { due: now.getTime() - last >= discEveryMs, task };150}151152/** One scheduler pass. Uses a short Redis lock so several processes do not double-enqueue. */153export async function schedulerTick(opts: { force?: boolean; now?: Date } = {}): Promise<TickResult> {154 const res: TickResult = { checked: 0, enqueued: [], skipped: [], lockHeldElsewhere: false };155 const redis = getRedis();156 const now = opts.now ?? new Date();157 const token = `${process.pid}:${now.getTime()}`;158 const got = await redis.set(SCHEDULER_LOCK_KEY, token, "EX", 55, "NX");159 if (!got && !opts.force) { res.lockHeldElsewhere = true; metrics.schedulerTicks.inc({ result: "lock_held" }); return res; }160 try {161 const rows = await schedulableConnectors();162 for (const c of rows) {163 res.checked++;164 if (!c.enabled || c.paused) { res.skipped.push(`${c.id}:${!c.enabled ? "disabled" : "paused"}`); continue; }165 // discovery (or first ever run) when the discovery cadence elapsed166 const disc = discoveryDue(c, discoveryIntervalMs(c.schedule), now);167 if (disc.due) {168 const r = await enqueueRun({ connectorId: c.id, task: disc.task, requestedBy: "scheduler" }, { priority: disc.task === "full" ? 5 : 20 });169 if (r.queued) res.enqueued.push(r.jobId);170 }171 // per-group crawl jobs for groups with due documents — unless a whole-connector job is pending172 if (c.docCount === 0) continue;173 if (await wholeConnectorJobPending(c.id)) { res.skipped.push(`${c.id}:full-pending`); continue; }174 const groups = await nextCheckByGroup(c.id);175 for (const g of groups) {176 if (g.due <= 0) continue;177 const r = await enqueueRun({ connectorId: c.id, task: "crawl", group: g.group, requestedBy: "scheduler" }, { priority: 10 });178 if (r.queued) res.enqueued.push(r.jobId);179 }180 }181 await ensureMaintenanceSchedulers();182 metrics.schedulerTicks.inc({ result: "ok" });183 metrics.schedulerEnqueued.inc(undefined, res.enqueued.length);184 metrics.schedulerLastTick.set(Date.now() / 1000);185 } catch (e) {186 metrics.schedulerTicks.inc({ result: "error" });187 throw e;188 } finally {189 const cur = await redis.get(SCHEDULER_LOCK_KEY).catch(() => null);190 if (cur === token) await redis.del(SCHEDULER_LOCK_KEY).catch(() => undefined);191 }192 return res;193}194195/** Run the scheduler loop until the returned stop() is called. */196export function startSchedulerLoop(opts: { intervalMs?: number; onTick?: (r: TickResult) => void; onError?: (e: Error) => void } = {}): { stop: () => void } {197 const every = opts.intervalMs ?? getEnv().schedulerIntervalMs;198 let stopped = false;199 let timer: NodeJS.Timeout | null = null;200 const loop = async () => {201 if (stopped) return;202 try { const r = await schedulerTick(); opts.onTick?.(r); } catch (e) { opts.onError?.(e as Error); }203 if (!stopped) timer = setTimeout(loop, every);204 };205 void loop();206 return { stop: () => { stopped = true; if (timer) clearTimeout(timer); } };207}208209export interface QueueSnapshot { name: string; waiting: number; active: number; delayed: number; failed: number; completed: number; prioritized: number }210export const QUEUE_STATES = ["waiting", "active", "delayed", "failed", "completed", "prioritized"] as const;211export async function queueSnapshot(): Promise<QueueSnapshot[]> {212 const out: QueueSnapshot[] = [];213 for (const q of [crawlQueue(), maintenanceQueue()]) {214 const c = await q.getJobCounts(...QUEUE_STATES);215 out.push({ name: q.name, waiting: c.waiting ?? 0, active: c.active ?? 0, delayed: c.delayed ?? 0, failed: c.failed ?? 0, completed: c.completed ?? 0, prioritized: c.prioritized ?? 0 });216 }217 return out;218}219220/** Publish queue depths into the `dci_queue_jobs{queue,state}` gauge (called before every /metrics scrape). */221export async function publishQueueMetrics(): Promise<QueueSnapshot[]> {222 const snap = await queueSnapshot();223 metrics.queueJobs.reset();224 for (const q of snap) for (const s of QUEUE_STATES) metrics.queueJobs.set(q[s], { queue: q.name, state: s });225 return snap;226}227228export async function pauseConnector(id: string, paused: boolean): Promise<void> {229 await getDb().update(connectorsTable).set({ paused, updatedAt: new Date().toISOString() }).where(sql`${connectorsTable.id} = ${id}`);230}231232export async function closeScheduler(): Promise<void> {233 await _crawlQueue?.close().catch(() => undefined);234 await _maintQueue?.close().catch(() => undefined);235 _crawlQueue = _maintQueue = null;236 if (_redis) { await _redis.quit().catch(() => undefined); _redis = null; }237}238239/* ---------- standalone scheduler process ---------- */240export async function startSchedulerProcess(): Promise<{ stop: () => Promise<void> }> {241 loadEnvFile();242 const env = getEnv();243 const startedAt = new Date().toISOString();244 let last: { at: string; checked: number; enqueued: string[]; skipped: number; lockHeldElsewhere: boolean; error: string | null } | null = null;245 let ticks = 0;246 let stopping = false;247 const loop = startSchedulerLoop({248 onTick: (r) => { ticks++; last = { at: new Date().toISOString(), checked: r.checked, enqueued: r.enqueued, skipped: r.skipped.length, lockHeldElsewhere: r.lockHeldElsewhere, error: null }; if (r.enqueued.length) console.error(`[scheduler] enqueued ${r.enqueued.join(", ")}`); },249 onError: (e) => { last = { at: new Date().toISOString(), checked: 0, enqueued: [], skipped: 0, lockHeldElsewhere: false, error: e.message }; console.error(`[scheduler] tick failed: ${e.message}`); },250 });251 const health = async () => ({ startedAt, pid: process.pid, host: env.hostname, intervalMs: env.schedulerIntervalMs, ticks, lastTick: last, queues: await queueSnapshot().catch((e: Error) => ({ error: e.message })) });252 const http = startHealthHttp({ port: env.workerPort, role: "scheduler", health, ready: () => !stopping && (last === null || last.error === null || Date.now() - Date.parse(last.at) < 3 * env.schedulerIntervalMs), beforeScrape: async () => { metrics.up.set(stopping ? 0 : 1); metrics.uptime.set((Date.now() - Date.parse(startedAt)) / 1000); await publishQueueMetrics(); } });253 const hb = setInterval(() => void getRedis().set(SCHEDULER_STATUS_KEY, JSON.stringify({ at: new Date().toISOString(), pid: process.pid, host: env.hostname, ticks, last }), "EX", 180).catch(() => undefined), 15_000);254 console.error(`[scheduler] loop every ${env.schedulerIntervalMs} ms · host ${env.hostname}`);255 return {256 stop: async () => {257 stopping = true;258 loop.stop();259 clearInterval(hb);260 http.close();261 await closeScheduler();262 await closeDb().catch(() => undefined);263 },264 };265}266267const isMain = process.argv[1] && /scheduler\.(ts|js)$/.test(process.argv[1]);268if (isMain) {269 startSchedulerProcess()270 .then(({ stop }) => {271 const onSignal = (sig: string) => { console.error(`[scheduler] ${sig}`); void stop().then(() => process.exit(0)); setTimeout(() => process.exit(0), 10_000).unref(); };272 process.on("SIGTERM", () => onSignal("SIGTERM"));273 process.on("SIGINT", () => onSignal("SIGINT"));274 })275 .catch((e) => { console.error(`[scheduler] fatal: ${(e as Error).stack ?? (e as Error).message}`); process.exit(1); });276}277