/** * Typed environment for the worker. No dotenv dependency: values come from process.env; in development * `loadEnvFile()` reads the repo-root `.env` through node:process.loadEnvFile (Node ≥ 20.12) when present. */ import { existsSync } from "node:fs"; import { dirname, join, resolve } from "node:path"; import { fileURLToPath } from "node:url"; import process from "node:process"; import { hostname as osHostname } from "node:os"; export type QueueName = "crawl" | "maintenance"; export const ALL_QUEUES: readonly QueueName[] = ["crawl", "maintenance"]; export interface WorkerEnv { databaseUrl: string; redisUrl: string; clickhouseUrl: string; clickhouseDb: string; s3Endpoint: string; s3Bucket: string; s3AccessKey: string; s3SecretKey: string; s3Region: string; /** directory holding connector YAML files */ configDir: string; /** BullMQ crawl worker concurrency (jobs in parallel) */ crawlConcurrency: number; maintenanceConcurrency: number; workerPort: number; /** premium credit budgets (daily) — also read by the fetchers themselves */ scrapflyDailyBudget: number; firecrawlDailyBudget: number; /** minutes between scheduler ticks */ schedulerIntervalMs: number; /** queues this process consumes (DCI_QUEUES=crawl,maintenance) */ queues: QueueName[]; /** run the scheduler loop inside the worker (DCI_SCHEDULER=0 to rely on the dedicated scheduler service) */ embeddedScheduler: boolean; /** hard deadline for a graceful shutdown before the process exits on its own (must stay < compose stop_grace_period) */ shutdownTimeoutMs: number; /** how often the shared Redis budget counters are re-read */ budgetRefreshMs: number; /** runs still `running` after this many hours are considered orphaned */ orphanRunHours: number; logLevel: "debug" | "info" | "warn" | "error"; hostname: string; } /** Walk up from `start` until a pnpm-workspace.yaml is found (repo root). */ export function findRepoRoot(start: string = process.cwd()): string { let dir = resolve(start); for (let i = 0; i < 8; i++) { if (existsSync(join(dir, "pnpm-workspace.yaml"))) return dir; const parent = dirname(dir); if (parent === dir) break; dir = parent; } // fall back to the location of this file (apps/worker/src → repo root) const here = dirname(fileURLToPath(import.meta.url)); return resolve(here, "..", "..", ".."); } let envFileLoaded = false; /** * Load `/.env` (or the given path) into process.env without overriding already-set variables. * Returns true when a file was loaded. Safe to call several times. */ export function loadEnvFile(file = ".env"): boolean { if (envFileLoaded) return true; const path = file.startsWith("/") ? file : join(findRepoRoot(), file); if (!existsSync(path)) return false; const loader = (process as unknown as { loadEnvFile?: (p: string) => void }).loadEnvFile; if (typeof loader !== "function") return false; // node's loadEnvFile does not override existing variables — exactly what we want. loader.call(process, path); envFileLoaded = true; return true; } function num(v: string | undefined, fallback: number): number { const n = Number(v); return Number.isFinite(n) && v !== undefined && v !== "" ? n : fallback; } /** "crawl,maintenance" → ["crawl","maintenance"]; unknown names are ignored; empty/absent → both. */ export function parseQueues(v: string | undefined): QueueName[] { const names = (v ?? "").split(",").map((s) => s.trim().toLowerCase()).filter(Boolean); const out = ALL_QUEUES.filter((q) => names.includes(q)); return out.length ? [...out] : [...ALL_QUEUES]; } let cached: WorkerEnv | null = null; export function getEnv(): WorkerEnv { if (cached) return cached; const e = process.env; const root = findRepoRoot(); const level = (e.DCI_LOG_LEVEL ?? "info").toLowerCase(); const queues = parseQueues(e.DCI_QUEUES); cached = { databaseUrl: e.DATABASE_URL ?? "postgres://dci:dci@127.0.0.1:5432/dci", redisUrl: e.REDIS_URL ?? "redis://127.0.0.1:6379/0", clickhouseUrl: e.CLICKHOUSE_URL ?? "http://127.0.0.1:8123", clickhouseDb: e.CLICKHOUSE_DB ?? "dci", s3Endpoint: e.S3_ENDPOINT ?? "http://127.0.0.1:9000", s3Bucket: e.S3_BUCKET ?? "dci-raw", s3AccessKey: e.S3_ACCESS_KEY ?? "dci", s3SecretKey: e.S3_SECRET_KEY ?? "", s3Region: e.S3_REGION ?? "us-east-1", configDir: e.DCI_CONFIG_DIR ?? join(root, "config", "connectors"), crawlConcurrency: Math.max(1, num(e.DCI_CRAWL_CONCURRENCY, 3)), maintenanceConcurrency: Math.max(1, num(e.DCI_MAINTENANCE_CONCURRENCY, 1)), workerPort: num(e.WORKER_PORT, 8320), scrapflyDailyBudget: num(e.DCI_SCRAPFLY_DAILY_BUDGET, 400), firecrawlDailyBudget: num(e.DCI_FIRECRAWL_DAILY_BUDGET, 200), schedulerIntervalMs: Math.max(5_000, num(e.DCI_SCHEDULER_INTERVAL_MS, 60_000)), queues, // a maintenance-only worker (data node) leaves scheduling to the dedicated scheduler service embeddedScheduler: e.DCI_SCHEDULER !== undefined ? e.DCI_SCHEDULER !== "0" && e.DCI_SCHEDULER.toLowerCase() !== "false" : queues.includes("crawl"), shutdownTimeoutMs: Math.max(5_000, num(e.DCI_SHUTDOWN_TIMEOUT_MS, 50_000)), budgetRefreshMs: Math.max(1_000, num(e.DCI_BUDGET_REFRESH_MS, 15_000)), orphanRunHours: Math.max(0.25, num(e.DCI_ORPHAN_RUN_HOURS, 6)), logLevel: level === "debug" || level === "warn" || level === "error" ? level : "info", hostname: e.HOSTNAME ?? osHostname(), }; return cached; } /** Human-readable list of the variables we read (for `dci doctor` / README). */ export const ENV_DOC: Array<[name: string, def: string, doc: string]> = [ ["DATABASE_URL", "postgres://dci:dci@127.0.0.1:5432/dci", "Postgres (source of truth)"], ["REDIS_URL", "redis://127.0.0.1:6379/0", "Redis for BullMQ queues, locks, heartbeat, shared budgets"], ["CLICKHOUSE_URL", "http://127.0.0.1:8123", "ClickHouse HTTP endpoint (analytics; optional)"], ["CLICKHOUSE_DB", "dci", "ClickHouse database"], ["S3_ENDPOINT", "http://127.0.0.1:9000", "S3/MinIO endpoint for raw bodies"], ["S3_BUCKET", "dci-raw", "Bucket for raw///..zst"], ["S3_ACCESS_KEY / S3_SECRET_KEY", "dci / —", "S3 credentials"], ["S3_REGION", "us-east-1", "S3 region (any value for MinIO)"], ["SCRAPFLY_API_KEY / FIRECRAWL_API_KEY", "—", "premium fetchers (L4 / L3); absent = never escalate"], ["DCI_SCRAPFLY_DAILY_BUDGET / DCI_FIRECRAWL_DAILY_BUDGET", "400 / 200", "daily premium credit caps (shared across workers via Redis dci:budget::)"], ["DCI_CONFIG_DIR", "/config/connectors", "connector YAML directory"], ["DCI_QUEUES", "crawl,maintenance", "queues consumed by this process"], ["DCI_SCHEDULER", "1 when DCI_QUEUES includes crawl", "0 = do not run the scheduler loop in this worker"], ["DCI_CRAWL_CONCURRENCY", "3", "BullMQ crawl jobs processed in parallel by this worker"], ["DCI_MAINTENANCE_CONCURRENCY", "1", "maintenance jobs in parallel"], ["DCI_SCHEDULER_INTERVAL_MS", "60000", "scheduler tick"], ["DCI_SHUTDOWN_TIMEOUT_MS", "50000", "graceful shutdown deadline (keep below compose stop_grace_period)"], ["DCI_BUDGET_REFRESH_MS", "15000", "refresh interval of the shared Redis budget counters"], ["DCI_ORPHAN_RUN_HOURS", "6", "runs still `running` after this are marked aborted at start / cleanup"], ["WORKER_PORT", "8320", "/healthz and /metrics HTTP port"], ["DCI_LOG_LEVEL", "info", "debug | info | warn | error"], ["DCI_USER_AGENT", "DataCenterIndexBot/0.1 (…)", "bot identity for L1 fetches"], ["DCI_CLICKHOUSE_OPTIONAL", "1", "set to 0 to make ClickHouse insert failures fatal"], ];