# @dci/worker — crawl runtime The worker turns connector configs (`config/connectors/*.yaml`) into a legal, incremental crawl and feeds the ingest layer. It owns: the production `ConnectorContext`, the URL registry (`documents`), raw archiving (MinIO), change detection (`document_versions`, ClickHouse `page_changes`), adaptive scheduling (BullMQ), run/health bookkeeping (`connector_runs`, `connectors`) and the `pnpm dci` CLI. ``` discover ─► registerDiscovered ─► dueDocuments ─► ctx.fetch ─► putRaw ─► fingerprint ─► skip? ─► extract (sitemap/rss/seeds) (documents) (priority, group) (robots, rate (zstd) (noise- (unchanged) normalize limit, cond. stripped) validate GET, L1→L4, ingestEntities budgets) recordVersion next_check ``` ## Files (apps/worker/src) | file | role | |---|---| | `env.ts` | typed `process.env` (`getEnv()`), `DCI_QUEUES` parsing, `loadEnvFile()` via `node:process.loadEnvFile` for dev | | `prom.ts` | in-process Prometheus registry (counters, gauges, histogram) and the worker's metric contract (`dci_crawl_*`, `dci_ingest_*`, `dci_queue_jobs`, `dci_worker_*`, `dci_scheduler_*`) | | `budget.ts` | `RedisBudgetStore`: shared daily premium budgets (`dci:budget::`) and per-connector daily caps (`dci:budget:connector::`), installed into the fetchers with `setBudgetStore` | | `health-http.ts` | the `/healthz` + `/metrics` HTTP server shared by the worker and the standalone scheduler | | `storage.ts` | S3/MinIO: `ensureBucket`, `putRaw` (key `raw///..zst`, zstd → gzip fallback, never uploads the same hash twice), `getRaw`, `headRaw` | | `configs.ts` | load YAML → `Connector` (registered `implementation` or `GenericConnector`), `syncConnectorsToDb()` upserts `connectors` + `sources` (id = `stableId("source", connectorId)`) | | `context.ts` | production `ConnectorContext`: robots.txt (+ crawl-delay), per-host token bucket, escalating fetch capped by `fetch.maxLevel`, `maxCreditsPerRun`, `maxCreditsPerDay` (Redis), the provider daily budgets, the discovery rule (sitemap/RSS never premium) and the error gate (a document failing twice is direct-only until every 4th attempt), conditional GET validators, `connector_state`, `isKnownUnchanged`, capped run log → `connector_runs.log`, Prometheus fetch metrics, ClickHouse `crawl_log` row per fetch (never fatal), cost per fetcher | | `documents.ts` | URL registry: `registerDiscovered`, `dueDocuments`, `recordFetch` (adaptive `next_check`), `recordVersion` (line diff of text projections, significance), `updateExtraction`, stats | | `scheduling.ts` | pure math: intervals, change-score EMA, backoff/quarantine, skip decision, health thresholds, budget level cap — unit-tested | | `pipeline.ts` | `runConnector(id, opts)`: one run end to end, per-document isolation, inline `pLimit`, graceful abort, run status + connector health | | `scheduler.ts` | BullMQ queues `crawl` / `maintenance` (Redis prefix `dci` → keys `dci:crawl:*`), `enqueueRun()`, 60 s scheduler tick (per-group jobs, dedupe by jobId `__`, no group job while a `full`/`discover` job is pending, discovery strictly on cadence), repeatable maintenance jobs, run-lock keys; runnable as the standalone scheduler process (`/healthz` + `/metrics`) | | `maintenance.ts` | `doctor()` diagnostics → `system_alerts`, `cleanupMaintenance()` | | `main.ts` | worker process: register connector groups, sync configs, BullMQ Workers for `DCI_QUEUES`, per-connector run lock (`dci:run-lock:`), heartbeat `dci:worker:status`, HTTP `/healthz` + `/metrics`, graceful shutdown with deadline | | `cli.ts` | `pnpm dci …` | | `connectors/index.ts` | registers every `connectors//index.ts` (`register()`), tolerant of missing/broken groups | | `ingest/`, `rankings.ts`, `metrics.ts`, `parsers/`, `connectors//` | owned by the ingest / connector agents | ## Scheduling and change detection * Every discovered URL becomes a `documents` row (`id = stableId("document", canonicalUrl)`), due immediately. The schedule group travels in `discovered_from` as `"|"`; priority / page type / minimum fetch level only ratchet up on re-discovery. A sitemap `lastmod` newer than `last_fetched` makes the page due now. * Selection: not quarantined, `next_check <= now()`, ordered by priority desc then oldest check first. * After each fetch (`recordFetch`): `change_frequency_score` is an EMA (α = 0.3) of "changed on fetch"; `next_check = now + base(group) × scale(score)` with scale 0.5 (score ≥ 0.5), 1, 2 (< 0.25), 4 (< 0.1); floor 15 min. Transport/HTTP errors back off ×2^(n−1) up to ×8; 404/410 back off ×4 and quarantine after 3 consecutive ones. `schedule.: never` stops scheduling. * Conditional GET: stored `etag` / `last_modified` are injected per URL (`ctx.validators`); a 304 counts as a successful unchanged fetch. * Fingerprint: `contentFingerprint` (core) over the body with runtime noise stripped (epoch cache-busters in URLs). Same hash + same extractor version → no extraction. Hash moved but text projection identical (markup-only) → also skipped, hash adopted, archive kept. * On real change: `document_versions` row with `diffSummary` (old text from the archive vs new) and `detectedChanges` from ingest; significance = max field-level significance, else 10/20/30 by diff ratio; ClickHouse `page_changes` row. * Premium fetchers (L3 Firecrawl, L4 Scrapfly) are used only when a page is due **and** direct fetch was blocked / JS-shell, only up to `fetch.maxLevel`, only while `fetch.maxCreditsPerRun`, the connector's `fetch.maxCreditsPerDay` and the provider daily budgets (`DCI_*_DAILY_BUDGET`, shared across workers through Redis) allow — never for sitemap / RSS discovery fetches, and only every 4th attempt for a document that already failed twice in a row. Robots disallow is a hard stop. See `docs/CRAWL-OPERATIONS.md` §4. * One run per connector at a time across all workers (Redis `dci:run-lock:`): a job that finds the lock held is delayed 60 s. Without it two overlapping runs fetch the same documents and both record a phantom "change" against a stale hash. ## Runs and health `connector_runs` gets a row per run (`task` discover | crawl | full | reprocess; `status` ok | partial | failed | aborted; `stats` discovered, registered, selected, fetched, notModified, changed, unchanged, markupOnly, failed, quarantined, extracted, entities, valid, rejected, created, updated, events, provenanceRows, credits, avgMs; `log` capped at 500 lines). `connectors.health` = ok (< 20 % failed fetches) | degraded (< 50 %) | failing; `next_run_at` = earliest `next_check` over the connector's groups. ## CLI ``` pnpm dci connectors table: id, kind, mode, enabled, health, last/next run, docs, yaml pnpm dci sync YAML → connectors + sources (also disables DB rows without YAML) pnpm dci run [--task full|discover|crawl|reprocess] [--group g] [--limit n] [--dry-run] [--url u]… [--force] [--verbose] [--json] pnpm dci run-all [--dry-run] [--task t] [--limit n] [ids…] pnpm dci discover [--dry-run] pnpm dci docs [--due] [--limit n] [--group g] [--all] pnpm dci inspect [--body] document + latest version + provenance rows pnpm dci reprocess [--limit n] re-extract from archived bodies, no network pnpm dci stats | doctor [--dry-run] | rank | metrics | refresh-stats pnpm dci enqueue [--task t] [--group g] | pause | resume pnpm dci scheduler [--once] scheduler loop only pnpm dci worker = apps/worker/src/main.ts pnpm dci validate ``` Exit codes: 0 ok, 1 failure (run failed, doctor critical), 2 usage. `--dry-run` persists nothing (ingest is called with `dryRun: true`; discovered URLs are crawled virtually in priority order). ## Worker process `pnpm --filter @dci/worker run dev` (or `pnpm dci worker`): registers connector groups, syncs configs, ensures the bucket and ClickHouse tables, installs the Redis budget store, starts BullMQ Workers for the queues in `DCI_QUEUES` (`crawl,maintenance`; `DCI_CRAWL_CONCURRENCY`, default 3, and `DCI_MAINTENANCE_CONCURRENCY`), the scheduler loop when the process consumes `crawl` (Redis lock `dci:scheduler:lock`; `DCI_SCHEDULER=0` to disable), a 15 s heartbeat in `dci:worker:status` and HTTP on `WORKER_PORT` (default 8320): `GET /healthz` JSON, `GET /metrics` Prometheus text (contract in `prom.ts`). SIGTERM/SIGINT: stop taking jobs, finish in-flight documents, mark runs `aborted`, exit — after `DCI_SHUTDOWN_TIMEOUT_MS` (50 s) in-flight runs are marked aborted in Postgres and the process exits anyway (keep it below compose's `stop_grace_period`). `pnpm --filter @dci/worker scheduler` (container role `scheduler`): only the scheduler loop plus `/healthz` + `/metrics` on `WORKER_PORT` (8321 in `compose.data.yml`), heartbeat `dci:scheduler:status`. Maintenance schedulers (UTC): daily metrics 00:10, rankings 00:30, refresh-stats hourly, cleanup 01:00. ## Environment | variable | default | purpose | |---|---|---| | `DATABASE_URL` | `postgres://dci:dci@127.0.0.1:5432/dci` | Postgres | | `REDIS_URL` | `redis://127.0.0.1:6379/0` | BullMQ, locks, heartbeat | | `CLICKHOUSE_URL` / `CLICKHOUSE_DB` | `http://127.0.0.1:8123` / `dci` | analytics (optional; failures are logged, never fatal — `DCI_CLICKHOUSE_OPTIONAL=0` to make them fatal) | | `S3_ENDPOINT` / `S3_BUCKET` / `S3_ACCESS_KEY` / `S3_SECRET_KEY` / `S3_REGION` | `http://127.0.0.1:9000` / `dci-raw` / `dci` / — / `us-east-1` | raw archive | | `SCRAPFLY_API_KEY` / `FIRECRAWL_API_KEY` | — | premium fetchers; absent = never escalate past L2 | | `DCI_SCRAPFLY_DAILY_BUDGET` / `DCI_FIRECRAWL_DAILY_BUDGET` | 400 / 200 | daily credit caps | | `DCI_CONFIG_DIR` | `/config/connectors` | YAML directory | | `DCI_QUEUES` | `crawl,maintenance` | queues consumed by this process (`maintenance` alone = the data node's `worker-maint`) | | `DCI_SCHEDULER` | on when `DCI_QUEUES` has `crawl` | `0` disables the embedded scheduler loop | | `DCI_CRAWL_CONCURRENCY` / `DCI_MAINTENANCE_CONCURRENCY` | 3 / 1 | BullMQ concurrency | | `DCI_SHUTDOWN_TIMEOUT_MS` | 50000 | graceful shutdown deadline | | `DCI_BUDGET_REFRESH_MS` | 15000 | refresh of the shared Redis budget counters | | `DCI_ORPHAN_RUN_HOURS` | 6 | `running` runs older than this are marked `aborted` at start / cleanup | | `DCI_SCHEDULER_INTERVAL_MS` | 60000 | scheduler tick | | `WORKER_PORT` | 8320 | `/healthz`, `/metrics` | | `DCI_LOG_LEVEL` | info | debug \| info \| warn \| error | | `DCI_USER_AGENT` | `DataCenterIndexBot/0.1 (…)` | L1 identity | Local dev: `set -a; source .env; set +a` (or rely on `loadEnvFile()`, which reads `/.env` without overriding existing variables). ## Tests `pnpm --filter @dci/worker test` — `scheduling.test.ts` covers interval resolution, EMA/scale, backoff and quarantine, fetch bookkeeping, fingerprint noise handling, skip logic, significance, health, budget cap and `pLimit`; `prom.test.ts` covers the Prometheus exposition format, `DCI_QUEUES` parsing, the budget key format, the discovery/error-gate cost rules, the scheduler's discovery cadence and the priority-aware discovery cap. ## Smoke test `config/connectors/_runtime-smoke.yaml.disabled` is a sitemap-only Digital Realty connector (L1–L2, no credits). Copy it to `digitalrealty-smoke.yaml`, then `pnpm dci sync && pnpm dci run digitalrealty-smoke --dry-run --limit 2`. Remove the copy afterwards (`pnpm dci sync` disables the DB row).