/** * BullMQ scheduling on Redis. * queue `crawl` (keys dci:crawl:*) jobs { connectorId, task, group?, limit?, force? } * queue `maintenance` (keys dci:maintenance:*) jobs { kind: rankings | metrics | refresh-stats | cleanup } * * The scheduler loop (every DCI_SCHEDULER_INTERVAL_MS, default 60 s) enqueues one crawl job per connector group * with due documents (jobId `__` dedupes), a `discover` job when the connector's discovery * cadence is due (a `full` job for a connector that was never discovered), and keeps the repeatable maintenance * jobs registered (daily metrics 00:10 UTC, rankings 00:30 UTC, refresh-stats hourly, cleanup 01:00 UTC). * connectors.enabled / paused are honored. While a `full` / `discover` job of a connector is pending, no per-group * job is enqueued for it (the full job crawls everything it registers; two runs on one connector would double * fetches and record phantom "changes"). * * BullMQ forbids ":" in queue names and custom job ids, hence the `dci` prefix + `__` separators. * * Run as a process (`dci-entrypoint scheduler`, `pnpm --filter @dci/worker scheduler`): the loop plus HTTP * /healthz + /metrics on WORKER_PORT. The crawl worker also runs the loop unless DCI_SCHEDULER=0; the Redis lock * `dci:scheduler:lock` makes concurrent loops harmless. */ import { Queue, type JobsOptions } from "bullmq"; import { Redis } from "ioredis"; import { intervalMs } from "@dci/connectors"; import { getDb, closeDb, connectors as connectorsTable, sql } from "@dci/db"; import { getEnv, loadEnvFile } from "./env.js"; import { nextCheckByGroup } from "./documents.js"; import { startHealthHttp } from "./health-http.js"; import type { RunTask } from "./pipeline.js"; import { metrics } from "./prom.js"; import { minIntervalMs } from "./scheduling.js"; export const QUEUE_PREFIX = "dci"; export const CRAWL_QUEUE = "crawl"; export const MAINTENANCE_QUEUE = "maintenance"; export const WORKER_STATUS_KEY = "dci:worker:status"; export const SCHEDULER_STATUS_KEY = "dci:scheduler:status"; const SCHEDULER_LOCK_KEY = "dci:scheduler:lock"; /** Redis lock held by a worker while it runs a connector — one run per connector at a time across all workers. */ export const RUN_LOCK_PREFIX = "dci:run-lock"; export const RUN_LOCK_TTL_SECONDS = 4 * 3600; export interface CrawlJobData { connectorId: string; task: RunTask; group?: string; limit?: number; force?: boolean; urls?: string[]; requestedBy?: string } export type MaintenanceKind = "rankings" | "metrics" | "refresh-stats" | "cleanup" | "snapshot" | "quality"; export interface MaintenanceJobData { kind: MaintenanceKind; requestedBy?: string } let _redis: Redis | null = null; /** Shared ioredis connection for BullMQ (maxRetriesPerRequest must be null for BullMQ). Reconnects on its own. */ export function getRedis(): Redis { if (!_redis) { _redis = new Redis(getEnv().redisUrl, { maxRetriesPerRequest: null, enableReadyCheck: false, lazyConnect: false, retryStrategy: (times) => Math.min(30_000, 200 * 2 ** Math.min(times, 8)) }); let down = false; // an 'error' event without listener would be printed as an unhandled error by ioredis on every retry _redis.on("error", (e) => { if (!down) console.error(`[redis] ${e.message} — reconnecting`); down = true; }); _redis.on("ready", () => { if (down) console.error("[redis] connection restored"); down = false; }); } return _redis; } let _crawlQueue: Queue | null = null; let _maintQueue: Queue | null = null; export function crawlQueue(): Queue { if (!_crawlQueue) _crawlQueue = new Queue(CRAWL_QUEUE, { connection: getRedis(), prefix: QUEUE_PREFIX, defaultJobOptions: { removeOnComplete: { count: 500, age: 86_400 }, removeOnFail: { count: 500, age: 7 * 86_400 }, attempts: 1 } }); return _crawlQueue; } export function maintenanceQueue(): Queue { if (!_maintQueue) _maintQueue = new Queue(MAINTENANCE_QUEUE, { connection: getRedis(), prefix: QUEUE_PREFIX, defaultJobOptions: { removeOnComplete: { count: 200 }, removeOnFail: { count: 200 }, attempts: 1 } }); return _maintQueue; } export function crawlJobId(connectorId: string, task: RunTask, group?: string): string { const clean = (s: string) => s.replace(/[^a-zA-Z0-9_-]/g, "-"); return task === "discover" ? `${clean(connectorId)}__discover` : `${clean(connectorId)}__${clean(group ?? "all")}`; } export function runLockKey(connectorId: string): string { return `${RUN_LOCK_PREFIX}:${connectorId}`; } const PENDING_STATES = new Set(["waiting", "active", "delayed", "prioritized", "waiting-children"]); /** Enqueue a crawl job (used by scheduler loop, CLI and API). Returns the job id; duplicates are coalesced. */ export async function enqueueRun(data: CrawlJobData, opts: { priority?: number; jobId?: string; delayMs?: number } = {}): Promise<{ jobId: string; queued: boolean }> { const jobId = opts.jobId ?? crawlJobId(data.connectorId, data.task, data.group); const q = crawlQueue(); const existing = await q.getJob(jobId); if (existing) { const state = await existing.getState(); if (PENDING_STATES.has(state)) return { jobId, queued: false }; await existing.remove().catch(() => undefined); } const jobOpts: JobsOptions = { jobId, priority: opts.priority ?? 10 }; if (opts.delayMs) jobOpts.delay = opts.delayMs; await q.add(`${data.task}:${data.connectorId}`, data, jobOpts); return { jobId, queued: true }; } /** True when a whole-connector job (`full` = `__all`, or `__discover`) is pending for this connector. */ async function wholeConnectorJobPending(connectorId: string): Promise { const q = crawlQueue(); for (const id of [crawlJobId(connectorId, "full"), crawlJobId(connectorId, "discover")]) { const job = await q.getJob(id); if (job && PENDING_STATES.has(await job.getState())) return true; } return false; } export async function enqueueMaintenance(kind: MaintenanceKind, requestedBy = "cli"): Promise { const jobId = `manual__${kind}__${Date.now().toString(36)}`; await maintenanceQueue().add(kind, { kind, requestedBy }, { jobId }); return jobId; } /** Repeatable maintenance schedulers (idempotent upserts). */ export async function ensureMaintenanceSchedulers(): Promise { const q = maintenanceQueue(); const tz = "UTC"; await q.upsertJobScheduler("daily-metrics", { pattern: "10 0 * * *", tz }, { name: "metrics", data: { kind: "metrics", requestedBy: "scheduler" } }); await q.upsertJobScheduler("daily-rankings", { pattern: "30 0 * * *", tz }, { name: "rankings", data: { kind: "rankings", requestedBy: "scheduler" } }); await q.upsertJobScheduler("hourly-refresh-stats", { pattern: "0 * * * *", tz }, { name: "refresh-stats", data: { kind: "refresh-stats", requestedBy: "scheduler" } }); await q.upsertJobScheduler("daily-cleanup", { pattern: "0 1 * * *", tz }, { name: "cleanup", data: { kind: "cleanup", requestedBy: "scheduler" } }); // daily snapshots + regression checks after rankings; quality sweep every 6 hours await q.upsertJobScheduler("daily-snapshot", { pattern: "45 0 * * *", tz }, { name: "snapshot", data: { kind: "snapshot", requestedBy: "scheduler" } }); await q.upsertJobScheduler("quality-sweep", { pattern: "20 */6 * * *", tz }, { name: "quality", data: { kind: "quality", requestedBy: "scheduler" } }); } export interface TickResult { checked: number; enqueued: string[]; skipped: string[]; lockHeldElsewhere: boolean } interface ConnectorSchedRow { id: string; enabled: boolean; paused: boolean; schedule: Record; config: Record; lastDiscoverAt: string | null; docCount: number } async function schedulableConnectors(): Promise { const rows = await getDb().execute<{ id: string; enabled: boolean; paused: boolean; schedule: Record; config: Record; last_discover_at: string | null; doc_count: number }>(sql` select c.id, c.enabled, c.paused, c.schedule, c.config, (select value #>> '{}' from connector_state s where s.connector_id = c.id and s.key = 'lastDiscoverAt') as last_discover_at, (select count(*)::int from documents d where d.connector_id = c.id) as doc_count from connectors c order by c.id`); 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) })); } /** Discovery cadence: schedule.discovery → schedule.sitemap → shortest group interval (≥ 6 h). */ export function discoveryIntervalMs(schedule: Record): number { 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 */ } } } return Math.max(6 * 3_600_000, minIntervalMs(schedule)); } /** * Is a discovery pass due? `lastDiscoverAt` is authoritative: a connector whose discovery legitimately yields zero * URLs (blocked feed, filter matching nothing) must wait for the next cadence like any other, not run every tick. */ export function discoveryDue(c: { docCount: number; lastDiscoverAt: string | null }, discEveryMs: number, now: Date): { due: boolean; task: RunTask } { const task: RunTask = c.docCount === 0 ? "full" : "discover"; const last = c.lastDiscoverAt ? Date.parse(c.lastDiscoverAt) : NaN; if (!Number.isFinite(last)) return { due: true, task }; return { due: now.getTime() - last >= discEveryMs, task }; } /** One scheduler pass. Uses a short Redis lock so several processes do not double-enqueue. */ export async function schedulerTick(opts: { force?: boolean; now?: Date } = {}): Promise { const res: TickResult = { checked: 0, enqueued: [], skipped: [], lockHeldElsewhere: false }; const redis = getRedis(); const now = opts.now ?? new Date(); const token = `${process.pid}:${now.getTime()}`; const got = await redis.set(SCHEDULER_LOCK_KEY, token, "EX", 55, "NX"); if (!got && !opts.force) { res.lockHeldElsewhere = true; metrics.schedulerTicks.inc({ result: "lock_held" }); return res; } try { const rows = await schedulableConnectors(); for (const c of rows) { res.checked++; if (!c.enabled || c.paused) { res.skipped.push(`${c.id}:${!c.enabled ? "disabled" : "paused"}`); continue; } // discovery (or first ever run) when the discovery cadence elapsed const disc = discoveryDue(c, discoveryIntervalMs(c.schedule), now); if (disc.due) { const r = await enqueueRun({ connectorId: c.id, task: disc.task, requestedBy: "scheduler" }, { priority: disc.task === "full" ? 5 : 20 }); if (r.queued) res.enqueued.push(r.jobId); } // per-group crawl jobs for groups with due documents — unless a whole-connector job is pending if (c.docCount === 0) continue; if (await wholeConnectorJobPending(c.id)) { res.skipped.push(`${c.id}:full-pending`); continue; } const groups = await nextCheckByGroup(c.id); for (const g of groups) { if (g.due <= 0) continue; const r = await enqueueRun({ connectorId: c.id, task: "crawl", group: g.group, requestedBy: "scheduler" }, { priority: 10 }); if (r.queued) res.enqueued.push(r.jobId); } } await ensureMaintenanceSchedulers(); metrics.schedulerTicks.inc({ result: "ok" }); metrics.schedulerEnqueued.inc(undefined, res.enqueued.length); metrics.schedulerLastTick.set(Date.now() / 1000); } catch (e) { metrics.schedulerTicks.inc({ result: "error" }); throw e; } finally { const cur = await redis.get(SCHEDULER_LOCK_KEY).catch(() => null); if (cur === token) await redis.del(SCHEDULER_LOCK_KEY).catch(() => undefined); } return res; } /** Run the scheduler loop until the returned stop() is called. */ export function startSchedulerLoop(opts: { intervalMs?: number; onTick?: (r: TickResult) => void; onError?: (e: Error) => void } = {}): { stop: () => void } { const every = opts.intervalMs ?? getEnv().schedulerIntervalMs; let stopped = false; let timer: NodeJS.Timeout | null = null; const loop = async () => { if (stopped) return; try { const r = await schedulerTick(); opts.onTick?.(r); } catch (e) { opts.onError?.(e as Error); } if (!stopped) timer = setTimeout(loop, every); }; void loop(); return { stop: () => { stopped = true; if (timer) clearTimeout(timer); } }; } export interface QueueSnapshot { name: string; waiting: number; active: number; delayed: number; failed: number; completed: number; prioritized: number } export const QUEUE_STATES = ["waiting", "active", "delayed", "failed", "completed", "prioritized"] as const; export async function queueSnapshot(): Promise { const out: QueueSnapshot[] = []; for (const q of [crawlQueue(), maintenanceQueue()]) { const c = await q.getJobCounts(...QUEUE_STATES); 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 }); } return out; } /** Publish queue depths into the `dci_queue_jobs{queue,state}` gauge (called before every /metrics scrape). */ export async function publishQueueMetrics(): Promise { const snap = await queueSnapshot(); metrics.queueJobs.reset(); for (const q of snap) for (const s of QUEUE_STATES) metrics.queueJobs.set(q[s], { queue: q.name, state: s }); return snap; } export async function pauseConnector(id: string, paused: boolean): Promise { await getDb().update(connectorsTable).set({ paused, updatedAt: new Date().toISOString() }).where(sql`${connectorsTable.id} = ${id}`); } export async function closeScheduler(): Promise { await _crawlQueue?.close().catch(() => undefined); await _maintQueue?.close().catch(() => undefined); _crawlQueue = _maintQueue = null; if (_redis) { await _redis.quit().catch(() => undefined); _redis = null; } } /* ---------- standalone scheduler process ---------- */ export async function startSchedulerProcess(): Promise<{ stop: () => Promise }> { loadEnvFile(); const env = getEnv(); const startedAt = new Date().toISOString(); let last: { at: string; checked: number; enqueued: string[]; skipped: number; lockHeldElsewhere: boolean; error: string | null } | null = null; let ticks = 0; let stopping = false; const loop = startSchedulerLoop({ 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(", ")}`); }, onError: (e) => { last = { at: new Date().toISOString(), checked: 0, enqueued: [], skipped: 0, lockHeldElsewhere: false, error: e.message }; console.error(`[scheduler] tick failed: ${e.message}`); }, }); 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 })) }); 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(); } }); 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); console.error(`[scheduler] loop every ${env.schedulerIntervalMs} ms · host ${env.hostname}`); return { stop: async () => { stopping = true; loop.stop(); clearInterval(hb); http.close(); await closeScheduler(); await closeDb().catch(() => undefined); }, }; } const isMain = process.argv[1] && /scheduler\.(ts|js)$/.test(process.argv[1]); if (isMain) { startSchedulerProcess() .then(({ stop }) => { const onSignal = (sig: string) => { console.error(`[scheduler] ${sig}`); void stop().then(() => process.exit(0)); setTimeout(() => process.exit(0), 10_000).unref(); }; process.on("SIGTERM", () => onSignal("SIGTERM")); process.on("SIGINT", () => onSignal("SIGINT")); }) .catch((e) => { console.error(`[scheduler] fatal: ${(e as Error).stack ?? (e as Error).message}`); process.exit(1); }); }