spb/datacenterindex
Public
HTML 53.9%
TypeScript 44.5%
JavaScript 0.6%
SQL 0.5%
1# @dci/worker — crawl runtime23The worker turns connector configs (`config/connectors/*.yaml`) into a legal, incremental crawl and feeds the4ingest layer. It owns: the production `ConnectorContext`, the URL registry (`documents`), raw archiving5(MinIO), change detection (`document_versions`, ClickHouse `page_changes`), adaptive scheduling (BullMQ),6run/health bookkeeping (`connector_runs`, `connectors`) and the `pnpm dci` CLI.78```9discover ─► registerDiscovered ─► dueDocuments ─► ctx.fetch ─► putRaw ─► fingerprint ─► skip? ─► extract10 (sitemap/rss/seeds) (documents) (priority, group) (robots, rate (zstd) (noise- (unchanged) normalize11 limit, cond. stripped) validate12 GET, L1→L4, ingestEntities13 budgets) recordVersion14 next_check15```1617## Files (apps/worker/src)1819| file | role |20|---|---|21| `env.ts` | typed `process.env` (`getEnv()`), `DCI_QUEUES` parsing, `loadEnvFile()` via `node:process.loadEnvFile` for dev |22| `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_*`) |23| `budget.ts` | `RedisBudgetStore`: shared daily premium budgets (`dci:budget:<provider>:<day>`) and per-connector daily caps (`dci:budget:connector:<id>:<day>`), installed into the fetchers with `setBudgetStore` |24| `health-http.ts` | the `/healthz` + `/metrics` HTTP server shared by the worker and the standalone scheduler |25| `storage.ts` | S3/MinIO: `ensureBucket`, `putRaw` (key `raw/<connector>/<docId>/<hash>.<ext>.zst`, zstd → gzip fallback, never uploads the same hash twice), `getRaw`, `headRaw` |26| `configs.ts` | load YAML → `Connector` (registered `implementation` or `GenericConnector`), `syncConnectorsToDb()` upserts `connectors` + `sources` (id = `stableId("source", connectorId)`) |27| `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 |28| `documents.ts` | URL registry: `registerDiscovered`, `dueDocuments`, `recordFetch` (adaptive `next_check`), `recordVersion` (line diff of text projections, significance), `updateExtraction`, stats |29| `scheduling.ts` | pure math: intervals, change-score EMA, backoff/quarantine, skip decision, health thresholds, budget level cap — unit-tested |30| `pipeline.ts` | `runConnector(id, opts)`: one run end to end, per-document isolation, inline `pLimit`, graceful abort, run status + connector health |31| `scheduler.ts` | BullMQ queues `crawl` / `maintenance` (Redis prefix `dci` → keys `dci:crawl:*`), `enqueueRun()`, 60 s scheduler tick (per-group jobs, dedupe by jobId `<connector>__<group>`, 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`) |32| `maintenance.ts` | `doctor()` diagnostics → `system_alerts`, `cleanupMaintenance()` |33| `main.ts` | worker process: register connector groups, sync configs, BullMQ Workers for `DCI_QUEUES`, per-connector run lock (`dci:run-lock:<id>`), heartbeat `dci:worker:status`, HTTP `/healthz` + `/metrics`, graceful shutdown with deadline |34| `cli.ts` | `pnpm dci …` |35| `connectors/index.ts` | registers every `connectors/<group>/index.ts` (`register()`), tolerant of missing/broken groups |36| `ingest/`, `rankings.ts`, `metrics.ts`, `parsers/`, `connectors/<group>/` | owned by the ingest / connector agents |3738## Scheduling and change detection3940* Every discovered URL becomes a `documents` row (`id = stableId("document", canonicalUrl)`), due immediately.41 The schedule group travels in `discovered_from` as `"<group>|<origin url>"`; priority / page type / minimum42 fetch level only ratchet up on re-discovery. A sitemap `lastmod` newer than `last_fetched` makes the page due now.43* Selection: not quarantined, `next_check <= now()`, ordered by priority desc then oldest check first.44* After each fetch (`recordFetch`): `change_frequency_score` is an EMA (α = 0.3) of "changed on fetch";45 `next_check = now + base(group) × scale(score)` with scale 0.5 (score ≥ 0.5), 1, 2 (< 0.25), 4 (< 0.1);46 floor 15 min. Transport/HTTP errors back off ×2^(n−1) up to ×8; 404/410 back off ×4 and quarantine after 347 consecutive ones. `schedule.<group>: never` stops scheduling.48* Conditional GET: stored `etag` / `last_modified` are injected per URL (`ctx.validators`); a 304 counts as a49 successful unchanged fetch.50* Fingerprint: `contentFingerprint` (core) over the body with runtime noise stripped (epoch cache-busters in51 URLs). Same hash + same extractor version → no extraction. Hash moved but text projection identical52 (markup-only) → also skipped, hash adopted, archive kept.53* On real change: `document_versions` row with `diffSummary` (old text from the archive vs new) and54 `detectedChanges` from ingest; significance = max field-level significance, else 10/20/30 by diff ratio;55 ClickHouse `page_changes` row.56* Premium fetchers (L3 Firecrawl, L4 Scrapfly) are used only when a page is due **and** direct fetch was57 blocked / JS-shell, only up to `fetch.maxLevel`, only while `fetch.maxCreditsPerRun`, the connector's58 `fetch.maxCreditsPerDay` and the provider daily budgets (`DCI_*_DAILY_BUDGET`, shared across workers through59 Redis) allow — never for sitemap / RSS discovery fetches, and only every 4th attempt for a document that already60 failed twice in a row. Robots disallow is a hard stop. See `docs/CRAWL-OPERATIONS.md` §4.61* One run per connector at a time across all workers (Redis `dci:run-lock:<connector>`): a job that finds the62 lock held is delayed 60 s. Without it two overlapping runs fetch the same documents and both record a phantom63 "change" against a stale hash.6465## Runs and health6667`connector_runs` gets a row per run (`task` discover | crawl | full | reprocess; `status` ok | partial | failed |68aborted; `stats` discovered, registered, selected, fetched, notModified, changed, unchanged, markupOnly,69failed, quarantined, extracted, entities, valid, rejected, created, updated, events, provenanceRows, credits,70avgMs; `log` capped at 500 lines). `connectors.health` = ok (< 20 % failed fetches) | degraded (< 50 %) |71failing; `next_run_at` = earliest `next_check` over the connector's groups.7273## CLI7475```76pnpm dci connectors table: id, kind, mode, enabled, health, last/next run, docs, yaml77pnpm dci sync YAML → connectors + sources (also disables DB rows without YAML)78pnpm dci run <id> [--task full|discover|crawl|reprocess] [--group g] [--limit n] [--dry-run] [--url u]… [--force] [--verbose] [--json]79pnpm dci run-all [--dry-run] [--task t] [--limit n] [ids…]80pnpm dci discover <id> [--dry-run]81pnpm dci docs <id> [--due] [--limit n] [--group g] [--all]82pnpm dci inspect <url-or-docId> [--body] document + latest version + provenance rows83pnpm dci reprocess <id> [--limit n] re-extract from archived bodies, no network84pnpm dci stats | doctor [--dry-run] | rank | metrics | refresh-stats85pnpm dci enqueue <id> [--task t] [--group g] | pause <id> | resume <id>86pnpm dci scheduler [--once] scheduler loop only87pnpm dci worker = apps/worker/src/main.ts88pnpm dci validate <file.yaml>89```9091Exit codes: 0 ok, 1 failure (run failed, doctor critical), 2 usage. `--dry-run` persists nothing (ingest is92called with `dryRun: true`; discovered URLs are crawled virtually in priority order).9394## Worker process9596`pnpm --filter @dci/worker run dev` (or `pnpm dci worker`): registers connector groups, syncs configs,97ensures the bucket and ClickHouse tables, installs the Redis budget store, starts BullMQ Workers for the queues in98`DCI_QUEUES` (`crawl,maintenance`; `DCI_CRAWL_CONCURRENCY`, default 3, and `DCI_MAINTENANCE_CONCURRENCY`), the99scheduler loop when the process consumes `crawl` (Redis lock `dci:scheduler:lock`; `DCI_SCHEDULER=0` to disable),100a 15 s heartbeat in `dci:worker:status` and HTTP on `WORKER_PORT` (default 8320): `GET /healthz` JSON,101`GET /metrics` Prometheus text (contract in `prom.ts`). SIGTERM/SIGINT: stop taking jobs, finish in-flight102documents, mark runs `aborted`, exit — after `DCI_SHUTDOWN_TIMEOUT_MS` (50 s) in-flight runs are marked aborted in103Postgres and the process exits anyway (keep it below compose's `stop_grace_period`).104105`pnpm --filter @dci/worker scheduler` (container role `scheduler`): only the scheduler loop plus `/healthz` +106`/metrics` on `WORKER_PORT` (8321 in `compose.data.yml`), heartbeat `dci:scheduler:status`.107108Maintenance schedulers (UTC): daily metrics 00:10, rankings 00:30, refresh-stats hourly, cleanup 01:00.109110## Environment111112| variable | default | purpose |113|---|---|---|114| `DATABASE_URL` | `postgres://dci:dci@127.0.0.1:5432/dci` | Postgres |115| `REDIS_URL` | `redis://127.0.0.1:6379/0` | BullMQ, locks, heartbeat |116| `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) |117| `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 |118| `SCRAPFLY_API_KEY` / `FIRECRAWL_API_KEY` | — | premium fetchers; absent = never escalate past L2 |119| `DCI_SCRAPFLY_DAILY_BUDGET` / `DCI_FIRECRAWL_DAILY_BUDGET` | 400 / 200 | daily credit caps |120| `DCI_CONFIG_DIR` | `<repo>/config/connectors` | YAML directory |121| `DCI_QUEUES` | `crawl,maintenance` | queues consumed by this process (`maintenance` alone = the data node's `worker-maint`) |122| `DCI_SCHEDULER` | on when `DCI_QUEUES` has `crawl` | `0` disables the embedded scheduler loop |123| `DCI_CRAWL_CONCURRENCY` / `DCI_MAINTENANCE_CONCURRENCY` | 3 / 1 | BullMQ concurrency |124| `DCI_SHUTDOWN_TIMEOUT_MS` | 50000 | graceful shutdown deadline |125| `DCI_BUDGET_REFRESH_MS` | 15000 | refresh of the shared Redis budget counters |126| `DCI_ORPHAN_RUN_HOURS` | 6 | `running` runs older than this are marked `aborted` at start / cleanup |127| `DCI_SCHEDULER_INTERVAL_MS` | 60000 | scheduler tick |128| `WORKER_PORT` | 8320 | `/healthz`, `/metrics` |129| `DCI_LOG_LEVEL` | info | debug \| info \| warn \| error |130| `DCI_USER_AGENT` | `DataCenterIndexBot/0.1 (…)` | L1 identity |131132Local dev: `set -a; source .env; set +a` (or rely on `loadEnvFile()`, which reads `<repo>/.env` without133overriding existing variables).134135## Tests136137`pnpm --filter @dci/worker test` — `scheduling.test.ts` covers interval resolution, EMA/scale, backoff and138quarantine, fetch bookkeeping, fingerprint noise handling, skip logic, significance, health, budget cap and139`pLimit`; `prom.test.ts` covers the Prometheus exposition format, `DCI_QUEUES` parsing, the budget key format,140the discovery/error-gate cost rules, the scheduler's discovery cadence and the priority-aware discovery cap.141142## Smoke test143144`config/connectors/_runtime-smoke.yaml.disabled` is a sitemap-only Digital Realty connector (L1–L2, no145credits). Copy it to `digitalrealty-smoke.yaml`, then `pnpm dci sync && pnpm dci run digitalrealty-smoke146--dry-run --limit 2`. Remove the copy afterwards (`pnpm dci sync` disables the DB row).147