spb/datacenterindex
Public
HTML 53.9%
TypeScript 44.5%
JavaScript 0.6%
SQL 0.5%
1/**2 * Production ConnectorContext: robots.txt, per-host rate limit + crawl-delay, escalating fetch bounded by the3 * connector's levels and premium credit budgets (per run, per connector per day, per provider per day),4 * connector_state persistence, content-hash lookups, structured logging mirrored into connector_runs.log,5 * Prometheus fetch metrics and a ClickHouse crawl_log row per fetch (never fatal).6 *7 * Cost-control rules enforced here (see docs/CRAWL-OPERATIONS.md):8 * - discovery fetches (group `sitemap` / `rss`) are direct-only: L1 → L2, never Firecrawl / Scrapfly;9 * - premium levels stop once `fetch.maxCreditsPerRun` is spent in this run, once `fetch.maxCreditsPerDay` is10 * spent by this connector today (Redis `dci:budget:connector:<id>:<day>`), or once a provider's daily budget11 * (`DCI_*_DAILY_BUDGET`, shared via Redis `dci:budget:<provider>:<day>`) is exhausted;12 * - a document that failed twice in a row is fetched direct-only until it either succeeds or reaches every 4th13 * attempt (`premiumAllowedAfterErrors`), so a hard 403 does not burn credits on every backoff cycle.14 */15import type { ConnectorContext, FetchOptions, RawDocument } from "@dci/connectors";16import { fetchWithEscalation, isAllowedByRobots, acquire, configureHost, failDoc, BROWSER_UA, scrapflyBudget, firecrawlBudget } from "@dci/connectors";17import type { FetchLevel, Provenance } from "@dci/core";18import { SOURCE_KIND_BASE } from "@dci/core";19import { getDb, connectorState, connectorRuns, sql, eq, and } from "@dci/db";20import { chInsert } from "@dci/db/clickhouse";21import { budgetStore } from "./budget.js";22import type { LoadedConnector } from "./configs.js";23import { getEnv } from "./env.js";24import { documentIdFor, storedHashFor } from "./documents.js";25import { fetchOutcome, metrics } from "./prom.js";26import { isDiscoveryGroup, maxLevelForBudget, premiumAllowedAfterErrors } from "./scheduling.js";2728export type LogLevel = "debug" | "info" | "warn" | "error";29export interface LogLine { t: string; level: LogLevel; msg: string }30const LEVEL_RANK: Record<LogLevel, number> = { debug: 0, info: 1, warn: 2, error: 3 };31export const RUN_LOG_CAP = 500;3233export interface FetcherCost { count: number; credits: number; ms: number; errors: number }3435/** Per-URL hints the pipeline injects before calling connector.fetch (connectors call ctx.fetch without them). */36export interface FetchHint { errorCount: number }3738export interface RunContext extends ConnectorContext {39 readonly loaded: LoadedConnector;40 /** premium credits spent in this run */41 credits: number;42 fetchCount: number;43 notModifiedCount: number;44 fetchErrorCount: number;45 robotsBlocked: number;46 /** fetches whose level was capped by a budget rule (run / connector-day / provider-day / error gate) */47 budgetCapped: number;48 costByFetcher: Record<string, FetcherCost>;49 logs: LogLine[];50 /** persist the (capped) log into connector_runs.log — no-op in dry run */51 flushLog(): Promise<void>;52 /** connector_state snapshot written during a dry run (kept in memory only) */53 readonly dryState: Map<string, unknown>;54 /** conditional-GET validators injected by the pipeline per URL */55 readonly validators: Map<string, { etag: string | null; lastModified: string | null }>;56 /** document scheduling hints injected by the pipeline per URL (error gate) */57 readonly hints: Map<string, FetchHint>;58}5960export interface CreateContextOptions {61 runId: string;62 dryRun: boolean;63 /** forward log lines (CLI pretty printing) */64 onLog?: (line: LogLine, extra?: Record<string, unknown>) => void;65 /** minimum level printed to stderr */66 logLevel?: LogLevel;67}6869function toChTs(iso: string): string { return iso.replace("T", " ").replace("Z", ""); }7071export function createRunContext(loaded: LoadedConnector, opts: CreateContextOptions): RunContext {72 const { cfg } = loaded;73 const env = getEnv();74 const minLevel = opts.logLevel ?? env.logLevel;75 const host = cfg.domain.replace(/^https?:\/\//, "").split("/")[0]!;76 configureHost(host, cfg.fetch.rpm, cfg.fetch.concurrency);77 const dryState = new Map<string, unknown>();78 const validators = new Map<string, { etag: string | null; lastModified: string | null }>();79 const hints = new Map<string, FetchHint>();80 const logs: LogLine[] = [];81 let dirtyLog = false;82 let flushing: Promise<void> | null = null;83 // credits this run added to the connector's daily counter (the Redis read may lag a few seconds)84 let dayCreditsThisRun = 0;85 let dayCapLogged = false;8687 const ctx: RunContext = {88 loaded,89 connectorId: cfg.id,90 sourceId: loaded.sourceId,91 runId: opts.runId,92 dryRun: opts.dryRun,93 env: process.env,94 credits: 0,95 fetchCount: 0,96 notModifiedCount: 0,97 fetchErrorCount: 0,98 robotsBlocked: 0,99 budgetCapped: 0,100 costByFetcher: {},101 logs,102 dryState,103 validators,104 hints,105 now: () => new Date().toISOString(),106107 log(level, msg, extra) {108 const line: LogLine = { t: new Date().toISOString(), level, msg: extra && Object.keys(extra).length ? `${msg} ${safeJson(extra)}` : msg };109 if (logs.length >= RUN_LOG_CAP) logs.splice(RUN_LOG_CAP / 2, 1); // keep head (discovery) and tail (latest)110 logs.push(line);111 dirtyLog = true;112 if (LEVEL_RANK[level] >= LEVEL_RANK[minLevel]) {113 if (opts.onLog) opts.onLog(line, extra);114 else console.error(`[${cfg.id}] ${level.toUpperCase()} ${line.msg}`);115 }116 },117118 async flushLog() {119 if (opts.dryRun || !dirtyLog) return;120 if (flushing) return flushing;121 flushing = (async () => {122 try {123 await getDb().update(connectorRuns).set({ log: logs.slice(-RUN_LOG_CAP) }).where(eq(connectorRuns.id, opts.runId));124 dirtyLog = false;125 } catch (e) {126 console.error(`[${cfg.id}] log flush failed: ${(e as Error).message}`);127 } finally { flushing = null; }128 })();129 return flushing;130 },131132 async getState<T>(key: string): Promise<T | null> {133 if (dryState.has(key)) return dryState.get(key) as T;134 const rows = await getDb().select({ value: connectorState.value }).from(connectorState).where(and(eq(connectorState.connectorId, cfg.id), eq(connectorState.key, key))).limit(1);135 return (rows[0]?.value as T | undefined) ?? null;136 },137138 async setState(key: string, value: unknown): Promise<void> {139 if (opts.dryRun) { dryState.set(key, value); return; }140 const now = new Date().toISOString();141 await getDb().insert(connectorState).values({ connectorId: cfg.id, key, value: value as never, updatedAt: now }).onConflictDoUpdate({ target: [connectorState.connectorId, connectorState.key], set: { value: value as never, updatedAt: now } });142 },143144 async isKnownUnchanged(url: string, contentHash: string): Promise<boolean> {145 const stored = await storedHashFor(url).catch(() => null);146 return stored != null && stored === contentHash;147 },148149 provenance(url: string, extra: Partial<Provenance> = {}): Provenance {150 const now = new Date().toISOString();151 return { sourceId: loaded.sourceId, connectorId: cfg.id, url, firstObserved: now, lastObserved: now, retrievedAt: now, confidence: SOURCE_KIND_BASE[cfg.kind] ?? "moderate", extractorVersion: loaded.extractorVersion, ...extra };152 },153154 async fetch(url: string, o: FetchOptions = {}): Promise<RawDocument> {155 const started = Date.now();156 let doc: RawDocument;157 let crawlDelay: number | null = null;158 let urlHost = host;159 try { urlHost = new URL(url).hostname; } catch { /* keep connector host */ }160161 // robots.txt unavailable (429/5xx/timeout, nothing cached) → "unknown": direct L1/L2 fetches only, never premium162 let robotsUnknown = false;163 if (cfg.fetch.respectRobots) {164 const r = await isAllowedByRobots(url).catch(() => ({ allowed: true, crawlDelay: null as number | null, sitemaps: [] as string[], unknown: true }));165 crawlDelay = r.crawlDelay;166 robotsUnknown = r.unknown;167 if (!r.allowed) {168 ctx.robotsBlocked++;169 ctx.log("warn", `robots.txt disallows ${url}`);170 doc = failDoc(url, 1, "direct", "robots_disallow", "disallowed by robots.txt", started);171 account(doc, o.group, started);172 return doc;173 }174 }175176 // ---- how high may this fetch escalate?177 let budgetMax = maxLevelForBudget(cfg.fetch.maxLevel, ctx.credits, cfg.fetch.maxCreditsPerRun);178 let capReason: string | null = budgetMax < cfg.fetch.maxLevel ? "run budget" : null;179 // premium credits this fetch may still spend (run budget, then connector-day budget): fetchWithEscalation skips180 // any level whose expected cost exceeds it, so one document never spends on L3 and L4 past the budget181 let creditsLeft = Math.max(0, cfg.fetch.maxCreditsPerRun - ctx.credits);182 if (budgetMax > 2 && robotsUnknown) { budgetMax = 2; capReason = "robots.txt unavailable — direct fetches only"; }183 if (budgetMax > 2 && isDiscoveryGroup(o.group)) { budgetMax = 2; capReason = "discovery is direct-only"; }184 if (budgetMax > 2 && cfg.fetch.maxCreditsPerDay !== undefined) {185 const usedToday = ((await budgetStore()?.connectorUsed(cfg.id)) ?? 0) + (budgetStore() ? 0 : dayCreditsThisRun);186 creditsLeft = Math.min(creditsLeft, Math.max(0, cfg.fetch.maxCreditsPerDay - usedToday));187 if (usedToday >= cfg.fetch.maxCreditsPerDay) {188 budgetMax = 2; capReason = `connector daily budget (${usedToday}/${cfg.fetch.maxCreditsPerDay})`;189 if (!dayCapLogged) { ctx.log("warn", `daily premium budget for ${cfg.id} exhausted (${usedToday}/${cfg.fetch.maxCreditsPerDay} credits) — direct fetches only until tomorrow (UTC)`); dayCapLogged = true; }190 }191 }192 const hint = hints.get(url);193 if (budgetMax > 2 && hint && !premiumAllowedAfterErrors(hint.errorCount)) { budgetMax = 2; capReason = `error gate (${hint.errorCount} consecutive failures)`; }194 if (budgetMax > 2 && scrapflyBudget.remaining() <= 0 && (budgetMax < 3 || firecrawlBudget.remaining() <= 0)) { budgetMax = 2; capReason = "provider daily budgets exhausted"; }195 if (capReason) ctx.budgetCapped++;196197 let level = Math.max(o.level ?? 1, cfg.fetch.level) as FetchLevel;198 if (level > budgetMax) { ctx.log("debug", `level ${level} requested for ${url} but capped at L${budgetMax} (${capReason ?? "connector maxLevel"})`); level = budgetMax; }199200 const delayMs = crawlDelay != null && Number.isFinite(crawlDelay) ? Math.min(30_000, Math.max(0, crawlDelay * 1000)) : 0;201 const v = validators.get(url);202 const release = await acquire(urlHost, delayMs);203 try {204 doc = await fetchWithEscalation(url, {205 ...o,206 etag: o.etag ?? v?.etag ?? null,207 lastModified: o.lastModified ?? v?.lastModified ?? null,208 level,209 maxLevel: budgetMax,210 creditsLeft,211 renderJs: o.renderJs ?? cfg.fetch.renderJs,212 country: o.country ?? cfg.fetch.country,213 waitForSelector: o.waitForSelector ?? cfg.fetch.waitForSelector,214 timeoutMs: o.timeoutMs ?? cfg.fetch.timeoutMs,215 maxBytes: o.maxBytes ?? cfg.fetch.maxBytes,216 accept: o.accept ?? cfg.fetch.accept,217 headers: { ...(cfg.fetch.userAgent === "browser" ? { "user-agent": BROWSER_UA } : {}), ...(cfg.fetch.headers ?? {}), ...(o.headers ?? {}) },218 });219 } catch (e) {220 doc = failDoc(url, level, "direct", "fetch_threw", (e as Error).message, started);221 } finally { release(); }222 if (robotsUnknown) doc.meta = { ...(doc.meta ?? {}), robotsUnknown: true };223 if (doc.credits > 0) { dayCreditsThisRun += doc.credits; budgetStore()?.spendConnector(cfg.id, doc.credits); }224 account(doc, o.group, started);225 ctx.log("debug", `fetch ${url} → ${doc.error?.code ?? doc.status}${doc.notModified ? " (304)" : ""} via ${doc.fetcher} L${doc.level} ${doc.durationMs}ms ${doc.body.length}b${doc.credits ? ` ${doc.credits}cr` : ""}`);226 return doc;227 },228 };229230 function account(doc: RawDocument, group: string | undefined, started: number): void {231 ctx.fetchCount++;232 ctx.credits += doc.credits;233 if (doc.notModified) ctx.notModifiedCount++;234 if (doc.error) ctx.fetchErrorCount++;235 const c = (ctx.costByFetcher[doc.fetcher] ??= { count: 0, credits: 0, ms: 0, errors: 0 });236 c.count++; c.credits += doc.credits; c.ms += doc.durationMs; if (doc.error) c.errors++;237 // Prometheus (contract: dci_crawl_fetches_total{connector,level,outcome}, dci_crawl_fetch_duration_seconds, dci_crawl_credits_total{provider})238 metrics.fetches.inc({ connector: cfg.id, level: `L${doc.level}`, outcome: fetchOutcome(doc) });239 metrics.fetchDuration.observe(Math.max(0, Date.now() - started) / 1000);240 if (doc.credits > 0 && (doc.fetcher === "scrapfly" || doc.fetcher === "firecrawl")) metrics.credits.inc({ provider: doc.fetcher }, doc.credits);241 let documentId = "";242 try { documentId = documentIdFor(doc.url); } catch { /* invalid url */ }243 void chInsert("crawl_log", [244 {245 ts: toChTs(doc.fetchedAt || new Date().toISOString()),246 connector_id: cfg.id,247 source_id: loaded.sourceId,248 document_id: documentId,249 url: doc.url.slice(0, 2048),250 fetch_level: doc.level,251 fetcher: doc.fetcher,252 status_code: Math.max(0, Math.min(65535, doc.status || 0)),253 duration_ms: Math.max(0, Math.round(doc.durationMs)),254 bytes: Math.max(0, doc.body.length),255 changed: 0,256 not_modified: doc.notModified ? 1 : 0,257 error: doc.error?.code ?? "",258 credits: doc.credits,259 run_id: opts.runId,260 },261 ]).catch(() => undefined);262 void group;263 }264265 return ctx;266}267268function safeJson(v: unknown): string {269 try { return JSON.stringify(v); } catch { return "[unserializable]"; }270}271272/** Snapshot of premium budgets for status endpoints / doctor (shared Redis view when installed, else this process). */273export function budgetSnapshot(): { scrapfly: { day: string; used: number; limit: number }; firecrawl: { day: string; used: number; limit: number } } {274 return { scrapfly: scrapflyBudget.snapshot(), firecrawl: firecrawlBudget.snapshot() };275}276277/** Sum of premium credits spent today across runs (from connector_runs.stats.credits). */278export async function creditsToday(): Promise<number> {279 const rows = await getDb().execute<{ credits: number | null }>(sql`select coalesce(sum((stats->>'credits')::float), 0) as credits from connector_runs where started_at >= date_trunc('day', now() at time zone 'utc')`);280 return Number(rows[0]?.credits ?? 0);281}282