/** * Worker process: registers code-backed connectors, syncs YAML configs into Postgres, ensures the raw bucket and * ClickHouse tables, installs the shared Redis budget store, then runs BullMQ Workers for the queues listed in * DCI_QUEUES (`crawl`, `maintenance`; default both) plus the scheduler loop (unless DCI_SCHEDULER=0 or the process * consumes no crawl queue). Heartbeat in Redis (dci:worker:status) and HTTP /healthz (JSON) + /metrics (Prometheus, * prom.ts contract) on WORKER_PORT. * * One run per connector at a time across all workers: the crawl processor takes `dci:run-lock:`; a job * that finds the lock held is moved back to delayed (BullMQ DelayedError) instead of running concurrently. * * SIGTERM/SIGINT: stop taking jobs, let in-flight documents finish (runs end as `aborted`), exit. If that takes * longer than DCI_SHUTDOWN_TIMEOUT_MS (default 50 s, below compose's stop_grace_period), in-flight runs are marked * aborted in Postgres and the process exits anyway. */ import { DelayedError, Worker, type Job } from "bullmq"; import { ensureClickHouse } from "@dci/db/clickhouse"; import { closeDb } from "@dci/db"; import { installBudgetStore, budgetStore, uninstallBudgetStore } from "./budget.js"; import { registerAllConnectors } from "./connectors/index.js"; import { loadAllConnectors, syncConnectorsToDb } from "./configs.js"; import { budgetSnapshot, creditsToday } from "./context.js"; import { getEnv, loadEnvFile, type QueueName } from "./env.js"; import { startHealthHttp } from "./health-http.js"; import { computeDailyMetrics, refreshStats } from "./metrics.js"; import { abortActiveRunsInDb, activeRunIds, isAbortRequested, requestAbort, runConnector, type RunResult } from "./pipeline.js"; import { metrics } from "./prom.js"; import { computeRankings } from "./rankings.js"; import { CRAWL_QUEUE, MAINTENANCE_QUEUE, QUEUE_PREFIX, RUN_LOCK_TTL_SECONDS, WORKER_STATUS_KEY, closeScheduler, getRedis, publishQueueMetrics, queueSnapshot, runLockKey, startSchedulerLoop, type CrawlJobData, type MaintenanceJobData } from "./scheduler.js"; import { ensureBucket } from "./storage.js"; import { cleanupMaintenance, reconcileOrphanRuns } from "./maintenance.js"; import { snapshotAndCheck, qualitySweep, dataGaps } from "./quality.js"; import { traceDocument } from "./trace.js"; export interface WorkerStatus { startedAt: string; host: string; pid: number; version: string; queues: QueueName[]; embeddedScheduler: boolean; runningJobs: Array<{ id: string; name: string; connectorId?: string; task?: string; group?: string; since: string }>; activeRuns: string[]; creditsToday: number; budgets: ReturnType; /** shared (Redis) usage today per provider, when the store is installed */ sharedBudgets: Record | null; counters: Record; lastScheduler: { at: string; enqueued: number; checked: number } | null; shuttingDown: boolean; } const VERSION = "0.2.0"; const counters: Record = { jobs_completed: 0, jobs_failed: 0, jobs_deferred_lock: 0, runs_ok: 0, runs_partial: 0, runs_failed: 0, runs_aborted: 0, documents_fetched: 0, documents_changed: 0, documents_failed: 0, entities_valid: 0, credits_spent: 0, events: 0 }; const running = new Map(); let lastScheduler: WorkerStatus["lastScheduler"] = null; let shuttingDown = false; const startedAt = new Date().toISOString(); function accountRun(r: RunResult): void { counters[`runs_${r.status}`] = (counters[`runs_${r.status}`] ?? 0) + 1; counters.documents_fetched! += r.stats.fetched; counters.documents_changed! += r.stats.changed; counters.documents_failed! += r.stats.failed; counters.entities_valid! += r.stats.valid; counters.credits_spent! += r.stats.credits; counters.events! += r.stats.events; } async function status(): Promise { const env = getEnv(); return { startedAt, host: env.hostname, pid: process.pid, version: VERSION, queues: env.queues, embeddedScheduler: env.embeddedScheduler, runningJobs: [...running.values()], activeRuns: activeRunIds(), creditsToday: await creditsToday().catch(() => -1), budgets: budgetSnapshot(), sharedBudgets: budgetStore()?.snapshot() ?? null, counters, lastScheduler, shuttingDown }; } /** Refresh the snapshot gauges right before a Prometheus scrape. */ async function refreshGauges(): Promise { const env = getEnv(); metrics.up.set(shuttingDown ? 0 : 1); metrics.runningJobs.set(running.size); metrics.uptime.set((Date.now() - Date.parse(startedAt)) / 1000); const b = budgetSnapshot(); metrics.dailyBudget.set(env.scrapflyDailyBudget, { provider: "scrapfly" }); metrics.dailyBudget.set(env.firecrawlDailyBudget, { provider: "firecrawl" }); metrics.dailyUsed.set(b.scrapfly.used, { provider: "scrapfly" }); metrics.dailyUsed.set(b.firecrawl.used, { provider: "firecrawl" }); await publishQueueMetrics(); } async function heartbeat(): Promise { try { const s = await status(); await getRedis().set(WORKER_STATUS_KEY, JSON.stringify({ ...s, at: new Date().toISOString() }), "EX", 90); } catch (e) { console.error(`[worker] heartbeat failed: ${(e as Error).message}`); } } export async function startWorker(opts: { withScheduler?: boolean } = {}): Promise<{ stop: () => Promise }> { loadEnvFile(); const env = getEnv(); const reg = await registerAllConnectors(); if (reg.failed.length) for (const f of reg.failed) console.error(`[worker] connector group ${f.group} failed to register: ${f.error}`); console.error(`[worker] connector groups: ${reg.registered.join(", ") || "none"}`); const loaded = loadAllConnectors(); const sync = await syncConnectorsToDb(loaded); console.error(`[worker] synced ${sync.connectors} connectors / ${sync.sources} sources${sync.disabledInDb.length ? ` (disabled without YAML: ${sync.disabledInDb.join(", ")})` : ""}`); const orphans = await reconcileOrphanRuns(env.orphanRunHours).catch((e: Error) => { console.error(`[worker] orphan reconciliation: ${e.message}`); return [] as string[]; }); if (orphans.length) console.error(`[worker] aborted ${orphans.length} orphaned run(s): ${orphans.join(", ")}`); await ensureBucket().catch((e: Error) => console.error(`[worker] storage: ${e.message}`)); await ensureClickHouse().catch((e: Error) => console.error(`[worker] clickhouse (optional): ${e.message}`)); const connection = getRedis(); await installBudgetStore(connection, env.budgetRefreshMs).catch((e: Error) => console.error(`[worker] budget store: ${e.message} (in-process budgets only)`)); const workers: Worker[] = []; if (env.queues.includes("crawl")) { const crawlWorker = new Worker( CRAWL_QUEUE, async (job: Job, token?: string) => { if (shuttingDown) throw new Error("worker shutting down"); const d = job.data; // one run per connector at a time (across workers): otherwise two runs fetch the same documents and both // record a "change" against a stale content hash const lockKey = runLockKey(d.connectorId); const lockToken = `${env.hostname}:${process.pid}:${job.id}`; const got = await connection.set(lockKey, lockToken, "EX", RUN_LOCK_TTL_SECONDS, "NX"); if (!got) { counters.jobs_deferred_lock!++; metrics.jobs.inc({ queue: CRAWL_QUEUE, result: "deferred" }); const holder = await connection.get(lockKey).catch(() => null); console.error(`[worker] ${job.name} #${job.id}: ${d.connectorId} already running (${holder ?? "?"}) — retrying in 60 s`); await job.moveToDelayed(Date.now() + 60_000, token); throw new DelayedError(); } running.set(String(job.id), { id: String(job.id), name: job.name, connectorId: d.connectorId, task: d.task, group: d.group, since: new Date().toISOString() }); try { const r = await runConnector(d.connectorId, { task: d.task, group: d.group, limit: d.limit, force: d.force, urls: d.urls }); accountRun(r); if (r.status === "failed") throw new Error(r.error ?? "run failed"); return { runId: r.runId, status: r.status, stats: r.stats }; } finally { running.delete(String(job.id)); // release only our own lock (a crashed predecessor's lock expires via TTL) await connection.eval("if redis.call('get', KEYS[1]) == ARGV[1] then return redis.call('del', KEYS[1]) else return 0 end", 1, lockKey, lockToken).catch(() => undefined); } }, { connection, prefix: QUEUE_PREFIX, concurrency: env.crawlConcurrency, lockDuration: 120_000, stalledInterval: 60_000, maxStalledCount: 1 }, ); workers.push(crawlWorker as Worker); } if (env.queues.includes("maintenance")) { const maintWorker = new Worker( MAINTENANCE_QUEUE, async (job: Job) => { running.set(String(job.id), { id: String(job.id), name: job.name, task: job.data.kind, since: new Date().toISOString() }); try { switch (job.data.kind) { case "rankings": return await computeRankings(); case "metrics": return await computeDailyMetrics(); case "refresh-stats": return await refreshStats(); case "cleanup": return await cleanupMaintenance(); case "snapshot": return await snapshotAndCheck(); case "quality": return await qualitySweep(); default: throw new Error(`unknown maintenance kind ${String((job.data as { kind?: string }).kind)}`); } } finally { running.delete(String(job.id)); } }, { connection, prefix: QUEUE_PREFIX, concurrency: env.maintenanceConcurrency, lockDuration: 600_000, stalledInterval: 60_000, maxStalledCount: 1 }, ); workers.push(maintWorker as Worker); } for (const w of workers) { w.on("completed", (job) => { counters.jobs_completed!++; metrics.jobs.inc({ queue: w.name, result: "completed" }); console.error(`[worker] ${w.name} ${job.name} #${job.id} completed`); }); w.on("failed", (job, err) => { if (err instanceof DelayedError) return; counters.jobs_failed!++; metrics.jobs.inc({ queue: w.name, result: "failed" }); console.error(`[worker] ${w.name} ${job?.name ?? "?"} #${job?.id ?? "?"} failed: ${err.message}`); }); w.on("stalled", (jobId) => { metrics.jobs.inc({ queue: w.name, result: "stalled" }); console.error(`[worker] ${w.name} #${jobId} stalled — BullMQ will retry it once`); }); w.on("error", (err) => console.error(`[worker] ${w.name} error: ${err.message}`)); } const withScheduler = opts.withScheduler ?? env.embeddedScheduler; const scheduler = !withScheduler ? null : startSchedulerLoop({ onTick: (r) => { lastScheduler = { at: new Date().toISOString(), enqueued: r.enqueued.length, checked: r.checked }; if (r.enqueued.length) console.error(`[scheduler] enqueued ${r.enqueued.join(", ")}`); }, onError: (e) => console.error(`[scheduler] tick failed: ${e.message}`), }); const http = startHealthHttp({ port: env.workerPort, role: "worker", health: async () => ({ ...(await status()), queueDepth: await queueSnapshot().catch(() => []) }), ready: () => !shuttingDown, beforeScrape: refreshGauges, log: (m) => console.error(`[worker] ${m}`), routes: { // extraction debugger (internal network only; the API proxies it behind the admin token) "/trace": async (url) => { const id = url.pathname.split("/")[2]; if (!id) return { status: 400, body: { error: "usage: /trace/" } }; return { status: 200, body: await traceDocument(decodeURIComponent(id), { live: url.searchParams.get("live") === "1" }) }; }, "/data-gaps": async () => ({ status: 200, body: await dataGaps() }), }, }); await heartbeat(); const hb = setInterval(() => void heartbeat(), 15_000); console.error(`[worker] ready · queues ${env.queues.join(",")} · crawl concurrency ${env.crawlConcurrency} · scheduler ${withScheduler ? "embedded" : "external"} · ${loaded.length} connectors · host ${env.hostname}`); const stop = async () => { if (shuttingDown) return; shuttingDown = true; requestAbort("SIGTERM"); console.error(`[worker] shutting down: finishing in-flight documents (${running.size} job(s), deadline ${env.shutdownTimeoutMs} ms)…`); scheduler?.stop(); clearInterval(hb); const deadline = new Promise<"timeout">((resolve) => setTimeout(() => resolve("timeout"), env.shutdownTimeoutMs).unref()); const closed = Promise.allSettled(workers.map((w) => w.close())).then(() => "closed" as const); const outcome = await Promise.race([closed, deadline]); if (outcome === "timeout") { const n = await Promise.race([abortActiveRunsInDb("shutdown deadline reached").catch(() => 0), new Promise((r) => setTimeout(() => r(-1), 5_000).unref())]); console.error(`[worker] shutdown deadline reached — ${n < 0 ? "could not mark" : n} in-flight run(s) marked aborted, forcing exit`); void Promise.allSettled(workers.map((w) => w.close(true))); } // cleanup must never keep a stopping process alive: a query or a blocked Redis command still in flight // (postgres.js `end()` waits for active queries) is cut by the 5 s race — the caller exits right after. const cleanup = async () => { await heartbeat(); http.close(); uninstallBudgetStore(); await closeScheduler(); await closeDb().catch(() => undefined); }; const done = await Promise.race([cleanup().then(() => "clean" as const), new Promise<"cut">((r) => setTimeout(() => r("cut"), 5_000).unref())]); console.error(`[worker] stopped (${done === "clean" ? "clean" : "cleanup cut short after 5 s"})`); }; return { stop }; } const isMain = process.argv[1] && /main\.(ts|js)$/.test(process.argv[1]); if (isMain) { process.on("unhandledRejection", (e) => console.error(`[worker] unhandled rejection: ${(e as Error)?.stack ?? String(e)}`)); process.on("uncaughtException", (e) => { console.error(`[worker] uncaught exception: ${e.stack ?? e.message}`); }); startWorker() .then(({ stop }) => { const onSignal = (sig: string) => { console.error(`[worker] ${sig}`); // absolute ceiling: graceful deadline + 10 s of cleanup, still below compose's stop_grace_period (90 s) setTimeout(() => { console.error("[worker] hard exit (shutdown did not complete in time)"); process.exit(1); }, getEnv().shutdownTimeoutMs + 10_000).unref(); void stop().then(() => process.exit(isAbortRequested() ? 0 : 0), (e: Error) => { console.error(`[worker] stop failed: ${e.message}`); process.exit(1); }); }; process.once("SIGTERM", () => onSignal("SIGTERM")); process.once("SIGINT", () => onSignal("SIGINT")); }) .catch((e) => { console.error(`[worker] fatal: ${(e as Error).stack ?? (e as Error).message}`); process.exit(1); }); }