/** * Production ConnectorContext: robots.txt, per-host rate limit + crawl-delay, escalating fetch bounded by the * connector's levels and premium credit budgets (per run, per connector per day, per provider per day), * connector_state persistence, content-hash lookups, structured logging mirrored into connector_runs.log, * Prometheus fetch metrics and a ClickHouse crawl_log row per fetch (never fatal). * * Cost-control rules enforced here (see docs/CRAWL-OPERATIONS.md): * - discovery fetches (group `sitemap` / `rss`) are direct-only: L1 → L2, never Firecrawl / Scrapfly; * - premium levels stop once `fetch.maxCreditsPerRun` is spent in this run, once `fetch.maxCreditsPerDay` is * spent by this connector today (Redis `dci:budget:connector::`), or once a provider's daily budget * (`DCI_*_DAILY_BUDGET`, shared via Redis `dci:budget::`) is exhausted; * - a document that failed twice in a row is fetched direct-only until it either succeeds or reaches every 4th * attempt (`premiumAllowedAfterErrors`), so a hard 403 does not burn credits on every backoff cycle. */ import type { ConnectorContext, FetchOptions, RawDocument } from "@dci/connectors"; import { fetchWithEscalation, isAllowedByRobots, acquire, configureHost, failDoc, BROWSER_UA, scrapflyBudget, firecrawlBudget } from "@dci/connectors"; import type { FetchLevel, Provenance } from "@dci/core"; import { SOURCE_KIND_BASE } from "@dci/core"; import { getDb, connectorState, connectorRuns, sql, eq, and } from "@dci/db"; import { chInsert } from "@dci/db/clickhouse"; import { budgetStore } from "./budget.js"; import type { LoadedConnector } from "./configs.js"; import { getEnv } from "./env.js"; import { documentIdFor, storedHashFor } from "./documents.js"; import { fetchOutcome, metrics } from "./prom.js"; import { isDiscoveryGroup, maxLevelForBudget, premiumAllowedAfterErrors } from "./scheduling.js"; export type LogLevel = "debug" | "info" | "warn" | "error"; export interface LogLine { t: string; level: LogLevel; msg: string } const LEVEL_RANK: Record = { debug: 0, info: 1, warn: 2, error: 3 }; export const RUN_LOG_CAP = 500; export interface FetcherCost { count: number; credits: number; ms: number; errors: number } /** Per-URL hints the pipeline injects before calling connector.fetch (connectors call ctx.fetch without them). */ export interface FetchHint { errorCount: number } export interface RunContext extends ConnectorContext { readonly loaded: LoadedConnector; /** premium credits spent in this run */ credits: number; fetchCount: number; notModifiedCount: number; fetchErrorCount: number; robotsBlocked: number; /** fetches whose level was capped by a budget rule (run / connector-day / provider-day / error gate) */ budgetCapped: number; costByFetcher: Record; logs: LogLine[]; /** persist the (capped) log into connector_runs.log — no-op in dry run */ flushLog(): Promise; /** connector_state snapshot written during a dry run (kept in memory only) */ readonly dryState: Map; /** conditional-GET validators injected by the pipeline per URL */ readonly validators: Map; /** document scheduling hints injected by the pipeline per URL (error gate) */ readonly hints: Map; } export interface CreateContextOptions { runId: string; dryRun: boolean; /** forward log lines (CLI pretty printing) */ onLog?: (line: LogLine, extra?: Record) => void; /** minimum level printed to stderr */ logLevel?: LogLevel; } function toChTs(iso: string): string { return iso.replace("T", " ").replace("Z", ""); } export function createRunContext(loaded: LoadedConnector, opts: CreateContextOptions): RunContext { const { cfg } = loaded; const env = getEnv(); const minLevel = opts.logLevel ?? env.logLevel; const host = cfg.domain.replace(/^https?:\/\//, "").split("/")[0]!; configureHost(host, cfg.fetch.rpm, cfg.fetch.concurrency); const dryState = new Map(); const validators = new Map(); const hints = new Map(); const logs: LogLine[] = []; let dirtyLog = false; let flushing: Promise | null = null; // credits this run added to the connector's daily counter (the Redis read may lag a few seconds) let dayCreditsThisRun = 0; let dayCapLogged = false; const ctx: RunContext = { loaded, connectorId: cfg.id, sourceId: loaded.sourceId, runId: opts.runId, dryRun: opts.dryRun, env: process.env, credits: 0, fetchCount: 0, notModifiedCount: 0, fetchErrorCount: 0, robotsBlocked: 0, budgetCapped: 0, costByFetcher: {}, logs, dryState, validators, hints, now: () => new Date().toISOString(), log(level, msg, extra) { const line: LogLine = { t: new Date().toISOString(), level, msg: extra && Object.keys(extra).length ? `${msg} ${safeJson(extra)}` : msg }; if (logs.length >= RUN_LOG_CAP) logs.splice(RUN_LOG_CAP / 2, 1); // keep head (discovery) and tail (latest) logs.push(line); dirtyLog = true; if (LEVEL_RANK[level] >= LEVEL_RANK[minLevel]) { if (opts.onLog) opts.onLog(line, extra); else console.error(`[${cfg.id}] ${level.toUpperCase()} ${line.msg}`); } }, async flushLog() { if (opts.dryRun || !dirtyLog) return; if (flushing) return flushing; flushing = (async () => { try { await getDb().update(connectorRuns).set({ log: logs.slice(-RUN_LOG_CAP) }).where(eq(connectorRuns.id, opts.runId)); dirtyLog = false; } catch (e) { console.error(`[${cfg.id}] log flush failed: ${(e as Error).message}`); } finally { flushing = null; } })(); return flushing; }, async getState(key: string): Promise { if (dryState.has(key)) return dryState.get(key) as T; const rows = await getDb().select({ value: connectorState.value }).from(connectorState).where(and(eq(connectorState.connectorId, cfg.id), eq(connectorState.key, key))).limit(1); return (rows[0]?.value as T | undefined) ?? null; }, async setState(key: string, value: unknown): Promise { if (opts.dryRun) { dryState.set(key, value); return; } const now = new Date().toISOString(); 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 } }); }, async isKnownUnchanged(url: string, contentHash: string): Promise { const stored = await storedHashFor(url).catch(() => null); return stored != null && stored === contentHash; }, provenance(url: string, extra: Partial = {}): Provenance { const now = new Date().toISOString(); 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 }; }, async fetch(url: string, o: FetchOptions = {}): Promise { const started = Date.now(); let doc: RawDocument; let crawlDelay: number | null = null; let urlHost = host; try { urlHost = new URL(url).hostname; } catch { /* keep connector host */ } // robots.txt unavailable (429/5xx/timeout, nothing cached) → "unknown": direct L1/L2 fetches only, never premium let robotsUnknown = false; if (cfg.fetch.respectRobots) { const r = await isAllowedByRobots(url).catch(() => ({ allowed: true, crawlDelay: null as number | null, sitemaps: [] as string[], unknown: true })); crawlDelay = r.crawlDelay; robotsUnknown = r.unknown; if (!r.allowed) { ctx.robotsBlocked++; ctx.log("warn", `robots.txt disallows ${url}`); doc = failDoc(url, 1, "direct", "robots_disallow", "disallowed by robots.txt", started); account(doc, o.group, started); return doc; } } // ---- how high may this fetch escalate? let budgetMax = maxLevelForBudget(cfg.fetch.maxLevel, ctx.credits, cfg.fetch.maxCreditsPerRun); let capReason: string | null = budgetMax < cfg.fetch.maxLevel ? "run budget" : null; // premium credits this fetch may still spend (run budget, then connector-day budget): fetchWithEscalation skips // any level whose expected cost exceeds it, so one document never spends on L3 and L4 past the budget let creditsLeft = Math.max(0, cfg.fetch.maxCreditsPerRun - ctx.credits); if (budgetMax > 2 && robotsUnknown) { budgetMax = 2; capReason = "robots.txt unavailable — direct fetches only"; } if (budgetMax > 2 && isDiscoveryGroup(o.group)) { budgetMax = 2; capReason = "discovery is direct-only"; } if (budgetMax > 2 && cfg.fetch.maxCreditsPerDay !== undefined) { const usedToday = ((await budgetStore()?.connectorUsed(cfg.id)) ?? 0) + (budgetStore() ? 0 : dayCreditsThisRun); creditsLeft = Math.min(creditsLeft, Math.max(0, cfg.fetch.maxCreditsPerDay - usedToday)); if (usedToday >= cfg.fetch.maxCreditsPerDay) { budgetMax = 2; capReason = `connector daily budget (${usedToday}/${cfg.fetch.maxCreditsPerDay})`; 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; } } } const hint = hints.get(url); if (budgetMax > 2 && hint && !premiumAllowedAfterErrors(hint.errorCount)) { budgetMax = 2; capReason = `error gate (${hint.errorCount} consecutive failures)`; } if (budgetMax > 2 && scrapflyBudget.remaining() <= 0 && (budgetMax < 3 || firecrawlBudget.remaining() <= 0)) { budgetMax = 2; capReason = "provider daily budgets exhausted"; } if (capReason) ctx.budgetCapped++; let level = Math.max(o.level ?? 1, cfg.fetch.level) as FetchLevel; if (level > budgetMax) { ctx.log("debug", `level ${level} requested for ${url} but capped at L${budgetMax} (${capReason ?? "connector maxLevel"})`); level = budgetMax; } const delayMs = crawlDelay != null && Number.isFinite(crawlDelay) ? Math.min(30_000, Math.max(0, crawlDelay * 1000)) : 0; const v = validators.get(url); const release = await acquire(urlHost, delayMs); try { doc = await fetchWithEscalation(url, { ...o, etag: o.etag ?? v?.etag ?? null, lastModified: o.lastModified ?? v?.lastModified ?? null, level, maxLevel: budgetMax, creditsLeft, renderJs: o.renderJs ?? cfg.fetch.renderJs, country: o.country ?? cfg.fetch.country, waitForSelector: o.waitForSelector ?? cfg.fetch.waitForSelector, timeoutMs: o.timeoutMs ?? cfg.fetch.timeoutMs, maxBytes: o.maxBytes ?? cfg.fetch.maxBytes, accept: o.accept ?? cfg.fetch.accept, headers: { ...(cfg.fetch.userAgent === "browser" ? { "user-agent": BROWSER_UA } : {}), ...(cfg.fetch.headers ?? {}), ...(o.headers ?? {}) }, }); } catch (e) { doc = failDoc(url, level, "direct", "fetch_threw", (e as Error).message, started); } finally { release(); } if (robotsUnknown) doc.meta = { ...(doc.meta ?? {}), robotsUnknown: true }; if (doc.credits > 0) { dayCreditsThisRun += doc.credits; budgetStore()?.spendConnector(cfg.id, doc.credits); } account(doc, o.group, started); 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` : ""}`); return doc; }, }; function account(doc: RawDocument, group: string | undefined, started: number): void { ctx.fetchCount++; ctx.credits += doc.credits; if (doc.notModified) ctx.notModifiedCount++; if (doc.error) ctx.fetchErrorCount++; const c = (ctx.costByFetcher[doc.fetcher] ??= { count: 0, credits: 0, ms: 0, errors: 0 }); c.count++; c.credits += doc.credits; c.ms += doc.durationMs; if (doc.error) c.errors++; // Prometheus (contract: dci_crawl_fetches_total{connector,level,outcome}, dci_crawl_fetch_duration_seconds, dci_crawl_credits_total{provider}) metrics.fetches.inc({ connector: cfg.id, level: `L${doc.level}`, outcome: fetchOutcome(doc) }); metrics.fetchDuration.observe(Math.max(0, Date.now() - started) / 1000); if (doc.credits > 0 && (doc.fetcher === "scrapfly" || doc.fetcher === "firecrawl")) metrics.credits.inc({ provider: doc.fetcher }, doc.credits); let documentId = ""; try { documentId = documentIdFor(doc.url); } catch { /* invalid url */ } void chInsert("crawl_log", [ { ts: toChTs(doc.fetchedAt || new Date().toISOString()), connector_id: cfg.id, source_id: loaded.sourceId, document_id: documentId, url: doc.url.slice(0, 2048), fetch_level: doc.level, fetcher: doc.fetcher, status_code: Math.max(0, Math.min(65535, doc.status || 0)), duration_ms: Math.max(0, Math.round(doc.durationMs)), bytes: Math.max(0, doc.body.length), changed: 0, not_modified: doc.notModified ? 1 : 0, error: doc.error?.code ?? "", credits: doc.credits, run_id: opts.runId, }, ]).catch(() => undefined); void group; } return ctx; } function safeJson(v: unknown): string { try { return JSON.stringify(v); } catch { return "[unserializable]"; } } /** Snapshot of premium budgets for status endpoints / doctor (shared Redis view when installed, else this process). */ export function budgetSnapshot(): { scrapfly: { day: string; used: number; limit: number }; firecrawl: { day: string; used: number; limit: number } } { return { scrapfly: scrapflyBudget.snapshot(), firecrawl: firecrawlBudget.snapshot() }; } /** Sum of premium credits spent today across runs (from connector_runs.stats.credits). */ export async function creditsToday(): Promise { 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')`); return Number(rows[0]?.credits ?? 0); }