@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_checkFiles (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:<provider>:<day>) and per-connector daily caps (dci:budget:connector:<id>:<day>), 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/<connector>/<docId>/<hash>.<ext>.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 <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) |
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:<id>), heartbeat dci:worker:status, HTTP /healthz + /metrics, graceful shutdown with deadline |
cli.ts |
pnpm dci … |
connectors/index.ts |
registers every connectors/<group>/index.ts (register()), tolerant of missing/broken groups |
ingest/, rankings.ts, metrics.ts, parsers/, connectors/<group>/ |
owned by the ingest / connector agents |
Scheduling and change detection
- Every discovered URL becomes a
documentsrow (id = stableId("document", canonicalUrl)), due immediately. The schedule group travels indiscovered_fromas"<group>|<origin url>"; priority / page type / minimum fetch level only ratchet up on re-discovery. A sitemaplastmodnewer thanlast_fetchedmakes the page due now. - Selection: not quarantined,
next_check <= now(), ordered by priority desc then oldest check first. - After each fetch (
recordFetch):change_frequency_scoreis 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.<group>: neverstops scheduling. - Conditional GET: stored
etag/last_modifiedare 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_versionsrow withdiffSummary(old text from the archive vs new) anddetectedChangesfrom ingest; significance = max field-level significance, else 10/20/30 by diff ratio; ClickHousepage_changesrow. - 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 whilefetch.maxCreditsPerRun, the connector'sfetch.maxCreditsPerDayand 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. Seedocs/CRAWL-OPERATIONS.md§4. - One run per connector at a time across all workers (Redis
dci:run-lock:<connector>): 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 <id> [--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 <id> [--dry-run]
pnpm dci docs <id> [--due] [--limit n] [--group g] [--all]
pnpm dci inspect <url-or-docId> [--body] document + latest version + provenance rows
pnpm dci reprocess <id> [--limit n] re-extract from archived bodies, no network
pnpm dci stats | doctor [--dry-run] | rank | metrics | refresh-stats
pnpm dci enqueue <id> [--task t] [--group g] | pause <id> | resume <id>
pnpm dci scheduler [--once] scheduler loop only
pnpm dci worker = apps/worker/src/main.ts
pnpm dci validate <file.yaml>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 |
<repo>/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 <repo>/.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).