spb/datacenterindex
Public
HTML 53.9%
TypeScript 44.5%
JavaScript 0.6%
SQL 0.5%
1/**2 * Typed environment for the worker. No dotenv dependency: values come from process.env; in development3 * `loadEnvFile()` reads the repo-root `.env` through node:process.loadEnvFile (Node ≥ 20.12) when present.4 */5import { existsSync } from "node:fs";6import { dirname, join, resolve } from "node:path";7import { fileURLToPath } from "node:url";8import process from "node:process";9import { hostname as osHostname } from "node:os";1011export type QueueName = "crawl" | "maintenance";12export const ALL_QUEUES: readonly QueueName[] = ["crawl", "maintenance"];1314export interface WorkerEnv {15 databaseUrl: string;16 redisUrl: string;17 clickhouseUrl: string;18 clickhouseDb: string;19 s3Endpoint: string;20 s3Bucket: string;21 s3AccessKey: string;22 s3SecretKey: string;23 s3Region: string;24 /** directory holding connector YAML files */25 configDir: string;26 /** BullMQ crawl worker concurrency (jobs in parallel) */27 crawlConcurrency: number;28 maintenanceConcurrency: number;29 workerPort: number;30 /** premium credit budgets (daily) — also read by the fetchers themselves */31 scrapflyDailyBudget: number;32 firecrawlDailyBudget: number;33 /** minutes between scheduler ticks */34 schedulerIntervalMs: number;35 /** queues this process consumes (DCI_QUEUES=crawl,maintenance) */36 queues: QueueName[];37 /** run the scheduler loop inside the worker (DCI_SCHEDULER=0 to rely on the dedicated scheduler service) */38 embeddedScheduler: boolean;39 /** hard deadline for a graceful shutdown before the process exits on its own (must stay < compose stop_grace_period) */40 shutdownTimeoutMs: number;41 /** how often the shared Redis budget counters are re-read */42 budgetRefreshMs: number;43 /** runs still `running` after this many hours are considered orphaned */44 orphanRunHours: number;45 logLevel: "debug" | "info" | "warn" | "error";46 hostname: string;47}4849/** Walk up from `start` until a pnpm-workspace.yaml is found (repo root). */50export function findRepoRoot(start: string = process.cwd()): string {51 let dir = resolve(start);52 for (let i = 0; i < 8; i++) {53 if (existsSync(join(dir, "pnpm-workspace.yaml"))) return dir;54 const parent = dirname(dir);55 if (parent === dir) break;56 dir = parent;57 }58 // fall back to the location of this file (apps/worker/src → repo root)59 const here = dirname(fileURLToPath(import.meta.url));60 return resolve(here, "..", "..", "..");61}6263let envFileLoaded = false;64/**65 * Load `<repo>/.env` (or the given path) into process.env without overriding already-set variables.66 * Returns true when a file was loaded. Safe to call several times.67 */68export function loadEnvFile(file = ".env"): boolean {69 if (envFileLoaded) return true;70 const path = file.startsWith("/") ? file : join(findRepoRoot(), file);71 if (!existsSync(path)) return false;72 const loader = (process as unknown as { loadEnvFile?: (p: string) => void }).loadEnvFile;73 if (typeof loader !== "function") return false;74 // node's loadEnvFile does not override existing variables — exactly what we want.75 loader.call(process, path);76 envFileLoaded = true;77 return true;78}7980function num(v: string | undefined, fallback: number): number {81 const n = Number(v);82 return Number.isFinite(n) && v !== undefined && v !== "" ? n : fallback;83}8485/** "crawl,maintenance" → ["crawl","maintenance"]; unknown names are ignored; empty/absent → both. */86export function parseQueues(v: string | undefined): QueueName[] {87 const names = (v ?? "").split(",").map((s) => s.trim().toLowerCase()).filter(Boolean);88 const out = ALL_QUEUES.filter((q) => names.includes(q));89 return out.length ? [...out] : [...ALL_QUEUES];90}9192let cached: WorkerEnv | null = null;93export function getEnv(): WorkerEnv {94 if (cached) return cached;95 const e = process.env;96 const root = findRepoRoot();97 const level = (e.DCI_LOG_LEVEL ?? "info").toLowerCase();98 const queues = parseQueues(e.DCI_QUEUES);99 cached = {100 databaseUrl: e.DATABASE_URL ?? "postgres://dci:dci@127.0.0.1:5432/dci",101 redisUrl: e.REDIS_URL ?? "redis://127.0.0.1:6379/0",102 clickhouseUrl: e.CLICKHOUSE_URL ?? "http://127.0.0.1:8123",103 clickhouseDb: e.CLICKHOUSE_DB ?? "dci",104 s3Endpoint: e.S3_ENDPOINT ?? "http://127.0.0.1:9000",105 s3Bucket: e.S3_BUCKET ?? "dci-raw",106 s3AccessKey: e.S3_ACCESS_KEY ?? "dci",107 s3SecretKey: e.S3_SECRET_KEY ?? "",108 s3Region: e.S3_REGION ?? "us-east-1",109 configDir: e.DCI_CONFIG_DIR ?? join(root, "config", "connectors"),110 crawlConcurrency: Math.max(1, num(e.DCI_CRAWL_CONCURRENCY, 3)),111 maintenanceConcurrency: Math.max(1, num(e.DCI_MAINTENANCE_CONCURRENCY, 1)),112 workerPort: num(e.WORKER_PORT, 8320),113 scrapflyDailyBudget: num(e.DCI_SCRAPFLY_DAILY_BUDGET, 400),114 firecrawlDailyBudget: num(e.DCI_FIRECRAWL_DAILY_BUDGET, 200),115 schedulerIntervalMs: Math.max(5_000, num(e.DCI_SCHEDULER_INTERVAL_MS, 60_000)),116 queues,117 // a maintenance-only worker (data node) leaves scheduling to the dedicated scheduler service118 embeddedScheduler: e.DCI_SCHEDULER !== undefined ? e.DCI_SCHEDULER !== "0" && e.DCI_SCHEDULER.toLowerCase() !== "false" : queues.includes("crawl"),119 shutdownTimeoutMs: Math.max(5_000, num(e.DCI_SHUTDOWN_TIMEOUT_MS, 50_000)),120 budgetRefreshMs: Math.max(1_000, num(e.DCI_BUDGET_REFRESH_MS, 15_000)),121 orphanRunHours: Math.max(0.25, num(e.DCI_ORPHAN_RUN_HOURS, 6)),122 logLevel: level === "debug" || level === "warn" || level === "error" ? level : "info",123 hostname: e.HOSTNAME ?? osHostname(),124 };125 return cached;126}127128/** Human-readable list of the variables we read (for `dci doctor` / README). */129export const ENV_DOC: Array<[name: string, def: string, doc: string]> = [130 ["DATABASE_URL", "postgres://dci:dci@127.0.0.1:5432/dci", "Postgres (source of truth)"],131 ["REDIS_URL", "redis://127.0.0.1:6379/0", "Redis for BullMQ queues, locks, heartbeat, shared budgets"],132 ["CLICKHOUSE_URL", "http://127.0.0.1:8123", "ClickHouse HTTP endpoint (analytics; optional)"],133 ["CLICKHOUSE_DB", "dci", "ClickHouse database"],134 ["S3_ENDPOINT", "http://127.0.0.1:9000", "S3/MinIO endpoint for raw bodies"],135 ["S3_BUCKET", "dci-raw", "Bucket for raw/<connector>/<doc>/<hash>.<ext>.zst"],136 ["S3_ACCESS_KEY / S3_SECRET_KEY", "dci / —", "S3 credentials"],137 ["S3_REGION", "us-east-1", "S3 region (any value for MinIO)"],138 ["SCRAPFLY_API_KEY / FIRECRAWL_API_KEY", "—", "premium fetchers (L4 / L3); absent = never escalate"],139 ["DCI_SCRAPFLY_DAILY_BUDGET / DCI_FIRECRAWL_DAILY_BUDGET", "400 / 200", "daily premium credit caps (shared across workers via Redis dci:budget:<provider>:<day>)"],140 ["DCI_CONFIG_DIR", "<repo>/config/connectors", "connector YAML directory"],141 ["DCI_QUEUES", "crawl,maintenance", "queues consumed by this process"],142 ["DCI_SCHEDULER", "1 when DCI_QUEUES includes crawl", "0 = do not run the scheduler loop in this worker"],143 ["DCI_CRAWL_CONCURRENCY", "3", "BullMQ crawl jobs processed in parallel by this worker"],144 ["DCI_MAINTENANCE_CONCURRENCY", "1", "maintenance jobs in parallel"],145 ["DCI_SCHEDULER_INTERVAL_MS", "60000", "scheduler tick"],146 ["DCI_SHUTDOWN_TIMEOUT_MS", "50000", "graceful shutdown deadline (keep below compose stop_grace_period)"],147 ["DCI_BUDGET_REFRESH_MS", "15000", "refresh interval of the shared Redis budget counters"],148 ["DCI_ORPHAN_RUN_HOURS", "6", "runs still `running` after this are marked aborted at start / cleanup"],149 ["WORKER_PORT", "8320", "/healthz and /metrics HTTP port"],150 ["DCI_LOG_LEVEL", "info", "debug | info | warn | error"],151 ["DCI_USER_AGENT", "DataCenterIndexBot/0.1 (…)", "bot identity for L1 fetches"],152 ["DCI_CLICKHOUSE_OPTIONAL", "1", "set to 0 to make ClickHouse insert failures fatal"],153];154