SPB Git forge
38commits 1branches 0releases
338.7 MBsize
maindefault branch
3 h agolast push
HTML 53.9% TypeScript 44.5% JavaScript 0.6% SQL 0.5%
11.8 KB

# @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.

text
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:<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 documents row (id = stableId("document", canonicalUrl)), due immediately. The schedule group travels in discovered_from as "<group>|<origin url>"; 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.<group>: 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:<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

text
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).