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%
14.4 KB · 282 lines typescript
Raw Blame History
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