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%
14.7 KB · 247 lines typescript
Raw Blame History
1/**2 * Worker process: registers code-backed connectors, syncs YAML configs into Postgres, ensures the raw bucket and3 * ClickHouse tables, installs the shared Redis budget store, then runs BullMQ Workers for the queues listed in4 * DCI_QUEUES (`crawl`, `maintenance`; default both) plus the scheduler loop (unless DCI_SCHEDULER=0 or the process5 * consumes no crawl queue). Heartbeat in Redis (dci:worker:status) and HTTP /healthz (JSON) + /metrics (Prometheus,6 * prom.ts contract) on WORKER_PORT.7 *8 * One run per connector at a time across all workers: the crawl processor takes `dci:run-lock:<connector>`; a job9 * that finds the lock held is moved back to delayed (BullMQ DelayedError) instead of running concurrently.10 *11 * SIGTERM/SIGINT: stop taking jobs, let in-flight documents finish (runs end as `aborted`), exit. If that takes12 * longer than DCI_SHUTDOWN_TIMEOUT_MS (default 50 s, below compose's stop_grace_period), in-flight runs are marked13 * aborted in Postgres and the process exits anyway.14 */15import { DelayedError, Worker, type Job } from "bullmq";16import { ensureClickHouse } from "@dci/db/clickhouse";17import { closeDb } from "@dci/db";18import { installBudgetStore, budgetStore, uninstallBudgetStore } from "./budget.js";19import { registerAllConnectors } from "./connectors/index.js";20import { loadAllConnectors, syncConnectorsToDb } from "./configs.js";21import { budgetSnapshot, creditsToday } from "./context.js";22import { getEnv, loadEnvFile, type QueueName } from "./env.js";23import { startHealthHttp } from "./health-http.js";24import { computeDailyMetrics, refreshStats } from "./metrics.js";25import { abortActiveRunsInDb, activeRunIds, isAbortRequested, requestAbort, runConnector, type RunResult } from "./pipeline.js";26import { metrics } from "./prom.js";27import { computeRankings } from "./rankings.js";28import { 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";29import { ensureBucket } from "./storage.js";30import { cleanupMaintenance, reconcileOrphanRuns } from "./maintenance.js";31import { snapshotAndCheck, qualitySweep, dataGaps } from "./quality.js";32import { traceDocument } from "./trace.js";3334export interface WorkerStatus {35  startedAt: string;36  host: string;37  pid: number;38  version: string;39  queues: QueueName[];40  embeddedScheduler: boolean;41  runningJobs: Array<{ id: string; name: string; connectorId?: string; task?: string; group?: string; since: string }>;42  activeRuns: string[];43  creditsToday: number;44  budgets: ReturnType<typeof budgetSnapshot>;45  /** shared (Redis) usage today per provider, when the store is installed */46  sharedBudgets: Record<string, number> | null;47  counters: Record<string, number>;48  lastScheduler: { at: string; enqueued: number; checked: number } | null;49  shuttingDown: boolean;50}5152const VERSION = "0.2.0";53const counters: Record<string, number> = { 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 };54const running = new Map<string, WorkerStatus["runningJobs"][number]>();55let lastScheduler: WorkerStatus["lastScheduler"] = null;56let shuttingDown = false;57const startedAt = new Date().toISOString();5859function accountRun(r: RunResult): void {60  counters[`runs_${r.status}`] = (counters[`runs_${r.status}`] ?? 0) + 1;61  counters.documents_fetched! += r.stats.fetched;62  counters.documents_changed! += r.stats.changed;63  counters.documents_failed! += r.stats.failed;64  counters.entities_valid! += r.stats.valid;65  counters.credits_spent! += r.stats.credits;66  counters.events! += r.stats.events;67}6869async function status(): Promise<WorkerStatus> {70  const env = getEnv();71  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 };72}7374/** Refresh the snapshot gauges right before a Prometheus scrape. */75async function refreshGauges(): Promise<void> {76  const env = getEnv();77  metrics.up.set(shuttingDown ? 0 : 1);78  metrics.runningJobs.set(running.size);79  metrics.uptime.set((Date.now() - Date.parse(startedAt)) / 1000);80  const b = budgetSnapshot();81  metrics.dailyBudget.set(env.scrapflyDailyBudget, { provider: "scrapfly" });82  metrics.dailyBudget.set(env.firecrawlDailyBudget, { provider: "firecrawl" });83  metrics.dailyUsed.set(b.scrapfly.used, { provider: "scrapfly" });84  metrics.dailyUsed.set(b.firecrawl.used, { provider: "firecrawl" });85  await publishQueueMetrics();86}8788async function heartbeat(): Promise<void> {89  try {90    const s = await status();91    await getRedis().set(WORKER_STATUS_KEY, JSON.stringify({ ...s, at: new Date().toISOString() }), "EX", 90);92  } catch (e) {93    console.error(`[worker] heartbeat failed: ${(e as Error).message}`);94  }95}9697export async function startWorker(opts: { withScheduler?: boolean } = {}): Promise<{ stop: () => Promise<void> }> {98  loadEnvFile();99  const env = getEnv();100  const reg = await registerAllConnectors();101  if (reg.failed.length) for (const f of reg.failed) console.error(`[worker] connector group ${f.group} failed to register: ${f.error}`);102  console.error(`[worker] connector groups: ${reg.registered.join(", ") || "none"}`);103  const loaded = loadAllConnectors();104  const sync = await syncConnectorsToDb(loaded);105  console.error(`[worker] synced ${sync.connectors} connectors / ${sync.sources} sources${sync.disabledInDb.length ? ` (disabled without YAML: ${sync.disabledInDb.join(", ")})` : ""}`);106  const orphans = await reconcileOrphanRuns(env.orphanRunHours).catch((e: Error) => { console.error(`[worker] orphan reconciliation: ${e.message}`); return [] as string[]; });107  if (orphans.length) console.error(`[worker] aborted ${orphans.length} orphaned run(s): ${orphans.join(", ")}`);108  await ensureBucket().catch((e: Error) => console.error(`[worker] storage: ${e.message}`));109  await ensureClickHouse().catch((e: Error) => console.error(`[worker] clickhouse (optional): ${e.message}`));110111  const connection = getRedis();112  await installBudgetStore(connection, env.budgetRefreshMs).catch((e: Error) => console.error(`[worker] budget store: ${e.message} (in-process budgets only)`));113114  const workers: Worker[] = [];115  if (env.queues.includes("crawl")) {116    const crawlWorker = new Worker<CrawlJobData>(117      CRAWL_QUEUE,118      async (job: Job<CrawlJobData>, token?: string) => {119        if (shuttingDown) throw new Error("worker shutting down");120        const d = job.data;121        // one run per connector at a time (across workers): otherwise two runs fetch the same documents and both122        // record a "change" against a stale content hash123        const lockKey = runLockKey(d.connectorId);124        const lockToken = `${env.hostname}:${process.pid}:${job.id}`;125        const got = await connection.set(lockKey, lockToken, "EX", RUN_LOCK_TTL_SECONDS, "NX");126        if (!got) {127          counters.jobs_deferred_lock!++;128          metrics.jobs.inc({ queue: CRAWL_QUEUE, result: "deferred" });129          const holder = await connection.get(lockKey).catch(() => null);130          console.error(`[worker] ${job.name} #${job.id}: ${d.connectorId} already running (${holder ?? "?"}) — retrying in 60 s`);131          await job.moveToDelayed(Date.now() + 60_000, token);132          throw new DelayedError();133        }134        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() });135        try {136          const r = await runConnector(d.connectorId, { task: d.task, group: d.group, limit: d.limit, force: d.force, urls: d.urls });137          accountRun(r);138          if (r.status === "failed") throw new Error(r.error ?? "run failed");139          return { runId: r.runId, status: r.status, stats: r.stats };140        } finally {141          running.delete(String(job.id));142          // release only our own lock (a crashed predecessor's lock expires via TTL)143          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);144        }145      },146      { connection, prefix: QUEUE_PREFIX, concurrency: env.crawlConcurrency, lockDuration: 120_000, stalledInterval: 60_000, maxStalledCount: 1 },147    );148    workers.push(crawlWorker as Worker);149  }150  if (env.queues.includes("maintenance")) {151    const maintWorker = new Worker<MaintenanceJobData>(152      MAINTENANCE_QUEUE,153      async (job: Job<MaintenanceJobData>) => {154        running.set(String(job.id), { id: String(job.id), name: job.name, task: job.data.kind, since: new Date().toISOString() });155        try {156          switch (job.data.kind) {157            case "rankings": return await computeRankings();158            case "metrics": return await computeDailyMetrics();159            case "refresh-stats": return await refreshStats();160            case "cleanup": return await cleanupMaintenance();161            case "snapshot": return await snapshotAndCheck();162            case "quality": return await qualitySweep();163            default: throw new Error(`unknown maintenance kind ${String((job.data as { kind?: string }).kind)}`);164          }165        } finally { running.delete(String(job.id)); }166      },167      { connection, prefix: QUEUE_PREFIX, concurrency: env.maintenanceConcurrency, lockDuration: 600_000, stalledInterval: 60_000, maxStalledCount: 1 },168    );169    workers.push(maintWorker as Worker);170  }171  for (const w of workers) {172    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`); });173    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}`); });174    w.on("stalled", (jobId) => { metrics.jobs.inc({ queue: w.name, result: "stalled" }); console.error(`[worker] ${w.name} #${jobId} stalled — BullMQ will retry it once`); });175    w.on("error", (err) => console.error(`[worker] ${w.name} error: ${err.message}`));176  }177178  const withScheduler = opts.withScheduler ?? env.embeddedScheduler;179  const scheduler = !withScheduler ? null : startSchedulerLoop({180    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(", ")}`); },181    onError: (e) => console.error(`[scheduler] tick failed: ${e.message}`),182  });183  const http = startHealthHttp({184    port: env.workerPort,185    role: "worker",186    health: async () => ({ ...(await status()), queueDepth: await queueSnapshot().catch(() => []) }),187    ready: () => !shuttingDown,188    beforeScrape: refreshGauges,189    log: (m) => console.error(`[worker] ${m}`),190    routes: {191      // extraction debugger (internal network only; the API proxies it behind the admin token)192      "/trace": async (url) => { const id = url.pathname.split("/")[2]; if (!id) return { status: 400, body: { error: "usage: /trace/<document-id>" } }; return { status: 200, body: await traceDocument(decodeURIComponent(id), { live: url.searchParams.get("live") === "1" }) }; },193      "/data-gaps": async () => ({ status: 200, body: await dataGaps() }),194    },195  });196  await heartbeat();197  const hb = setInterval(() => void heartbeat(), 15_000);198  console.error(`[worker] ready · queues ${env.queues.join(",")} · crawl concurrency ${env.crawlConcurrency} · scheduler ${withScheduler ? "embedded" : "external"} · ${loaded.length} connectors · host ${env.hostname}`);199200  const stop = async () => {201    if (shuttingDown) return;202    shuttingDown = true;203    requestAbort("SIGTERM");204    console.error(`[worker] shutting down: finishing in-flight documents (${running.size} job(s), deadline ${env.shutdownTimeoutMs} ms)…`);205    scheduler?.stop();206    clearInterval(hb);207    const deadline = new Promise<"timeout">((resolve) => setTimeout(() => resolve("timeout"), env.shutdownTimeoutMs).unref());208    const closed = Promise.allSettled(workers.map((w) => w.close())).then(() => "closed" as const);209    const outcome = await Promise.race([closed, deadline]);210    if (outcome === "timeout") {211      const n = await Promise.race([abortActiveRunsInDb("shutdown deadline reached").catch(() => 0), new Promise<number>((r) => setTimeout(() => r(-1), 5_000).unref())]);212      console.error(`[worker] shutdown deadline reached — ${n < 0 ? "could not mark" : n} in-flight run(s) marked aborted, forcing exit`);213      void Promise.allSettled(workers.map((w) => w.close(true)));214    }215    // cleanup must never keep a stopping process alive: a query or a blocked Redis command still in flight216    // (postgres.js `end()` waits for active queries) is cut by the 5 s race — the caller exits right after.217    const cleanup = async () => {218      await heartbeat();219      http.close();220      uninstallBudgetStore();221      await closeScheduler();222      await closeDb().catch(() => undefined);223    };224    const done = await Promise.race([cleanup().then(() => "clean" as const), new Promise<"cut">((r) => setTimeout(() => r("cut"), 5_000).unref())]);225    console.error(`[worker] stopped (${done === "clean" ? "clean" : "cleanup cut short after 5 s"})`);226  };227  return { stop };228}229230const isMain = process.argv[1] && /main\.(ts|js)$/.test(process.argv[1]);231if (isMain) {232  process.on("unhandledRejection", (e) => console.error(`[worker] unhandled rejection: ${(e as Error)?.stack ?? String(e)}`));233  process.on("uncaughtException", (e) => { console.error(`[worker] uncaught exception: ${e.stack ?? e.message}`); });234  startWorker()235    .then(({ stop }) => {236      const onSignal = (sig: string) => {237        console.error(`[worker] ${sig}`);238        // absolute ceiling: graceful deadline + 10 s of cleanup, still below compose's stop_grace_period (90 s)239        setTimeout(() => { console.error("[worker] hard exit (shutdown did not complete in time)"); process.exit(1); }, getEnv().shutdownTimeoutMs + 10_000).unref();240        void stop().then(() => process.exit(isAbortRequested() ? 0 : 0), (e: Error) => { console.error(`[worker] stop failed: ${e.message}`); process.exit(1); });241      };242      process.once("SIGTERM", () => onSignal("SIGTERM"));243      process.once("SIGINT", () => onSignal("SIGINT"));244    })245    .catch((e) => { console.error(`[worker] fatal: ${(e as Error).stack ?? (e as Error).message}`); process.exit(1); });246}247