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%
16.4 KB · 277 lines typescript
Raw Blame History
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