SPB Git forge

spb/rareindex

Public
54commits 1branches 0releases
7.1 MBsize
maindefault branch
10 days agolast push
TypeScript 61.9% HTML 37.2% SQL 0.7%
7.3 KB · 176 lines typescript
Raw Blame History
1import type { Engine, ExtractionResult, Logger } from '@rareindex/shared';2import { createHttpEngine, type HttpEngine } from './engines/http.js';3import { createFirecrawlEngine, type FirecrawlEngine } from './engines/firecrawl.js';4import { createScrapflyEngine, type ScrapflyEngine } from './engines/scrapfly.js';5import { documentQuality, isChallengePage, qualityScore, MIN_QUALITY } from './quality.js';6import { policyFor } from './domains.js';7import { CircuitOpenError, HostGates } from './circuit.js';8import { domainOf, type EngineStats, type FetchOptions } from './types.js';910export interface RouterAttempt {11  engine: Engine;12  url: string;13  host: string;14  success: boolean;15  credits: number;16  ms: number;17  quality: number;18  error: string | null;19  httpStatus: number | null;20  blocked: boolean;21}2223export interface RouterOptions {24  firecrawlApiKey?: string;25  scrapflyApiKey?: string;26  log?: Logger;27  /** called for every engine attempt (cost ledger, metrics) */28  onAttempt?: (info: RouterAttempt) => void;29  /** shared gates (one per process) so several connectors on the same host respect one limit */30  gates?: HostGates;31}3233export interface Router {34  fetch(url: string, opts?: FetchOptions, defaults?: { enginePriority: Engine[]; stats?: Record<string, EngineStats> }): Promise<ExtractionResult>;35  engines: { http: HttpEngine; firecrawl: FirecrawlEngine | null; scrapfly: ScrapflyEngine | null };36  available(): Engine[];37  gates: HostGates;38}3940/** Process-wide gates keyed by host; policies come from connectors/domains.json. */41export const sharedGates = new HostGates((host) => {42  const p = policyFor(host);43  return { concurrency: p.concurrency, minIntervalMs: p.minIntervalMs, circuitFailures: p.circuitFailures, circuitCooldownMs: p.circuitCooldownMs };44});4546/**47 * Firecrawl → Scrapfly router (§198). Tries engines in priority order, scores each result with the48 * connector's parser (or a document heuristic), and falls through until quality ≥ minQuality.49 * A result that fails everywhere is returned with requiresReview=true — never silently dropped.50 * Every attempt passes through the per-host gate (concurrency, interval, circuit breaker).51 */52export function createRouter(opts: RouterOptions = {}): Router {53  const http = createHttpEngine();54  const firecrawl = opts.firecrawlApiKey ? createFirecrawlEngine({ apiKey: opts.firecrawlApiKey }) : null;55  const scrapfly = opts.scrapflyApiKey ? createScrapflyEngine({ apiKey: opts.scrapflyApiKey }) : null;56  const gates = opts.gates ?? sharedGates;5758  function available(): Engine[] {59    const out: Engine[] = ['api', 'feed'];60    if (firecrawl) out.push('firecrawl');61    if (scrapfly) out.push('scrapfly');62    return out;63  }6465  async function runEngine(engine: Engine, url: string, o: FetchOptions): Promise<ExtractionResult | null> {66    switch (engine) {67      case 'api':68      case 'feed':69        return http.fetch(url, o);70      case 'firecrawl':71        return firecrawl ? firecrawl.scrape(url, o) : null;72      case 'scrapfly':73        return scrapfly ? scrapfly.scrape(url, o) : null;74      default:75        return null;76    }77  }7879  function empty(url: string, engine: Engine, error: string): ExtractionResult {80    return { success: false, engine, url, finalUrl: null, httpStatus: null, html: null, markdown: null, json: null, costCredits: 0, durationMs: 0, fetchedAt: new Date(), qualityScore: 0, requiresReview: true, error };81  }8283  async function routedFetch(url: string, o: FetchOptions = {}, defaults?: { enginePriority: Engine[]; stats?: Record<string, EngineStats> }): Promise<ExtractionResult> {84    const host = domainOf(url);85    const policy = policyFor(host);86    let engines = (o.engines ?? defaults?.enginePriority ?? ['api', 'firecrawl', 'scrapfly']).filter((e) => available().includes(e));87    if (policy.engines) engines = engines.filter((e) => policy.engines!.includes(e));88    const merged: FetchOptions = {89      ...o,90      timeoutMs: o.timeoutMs ?? policy.timeoutMs,91      country: o.country ?? policy.scrapfly.country ?? policy.country,92      waitForMs: o.waitForMs ?? policy.firecrawl.waitForMs,93      renderJs: o.renderJs ?? policy.scrapfly.renderJs,94    };95    const minQuality = o.minQuality ?? MIN_QUALITY;96    let last: ExtractionResult | null = null;97    const errors: string[] = [];98    for (const engine of engines) {99      let release: (() => void) | null = null;100      if (!o.ignoreCircuit) {101        try {102          release = await gates.acquire(host);103        } catch (err) {104          if (err instanceof CircuitOpenError) {105            const stats = defaults?.stats;106            if (stats) {107              const s = (stats[engine] ??= { attempts: 0, success: 0, credits: 0, ms: 0 });108              s.circuitOpen = (s.circuitOpen ?? 0) + 1;109            }110            opts.log?.warn({ host, url, retryAt: err.retryAt }, 'circuit open; request refused');111            return empty(url, engine, err.message);112          }113          throw err;114        }115      }116      let res: ExtractionResult | null;117      try {118        res = await runEngine(engine, url, merged);119      } finally {120        release?.();121      }122      if (!res) continue;123      let quality = 0;124      const blocked = res.success ? isChallengePage(res) : false;125      if (res.success && !blocked) {126        if (o.parse) {127          try {128            const found = o.parse(res);129            quality = qualityScore(found, o.expect);130          } catch (err) {131            quality = 0;132            res.error = `parse: ${err instanceof Error ? err.message : String(err)}`;133          }134        } else {135          quality = documentQuality(res);136        }137      } else if (blocked) {138        res.error = res.error ?? 'anti-bot challenge page';139      }140      res.qualityScore = quality;141      const ok = res.success && !blocked && quality >= minQuality;142      const stats = defaults?.stats;143      if (stats) {144        const s = (stats[engine] ??= { attempts: 0, success: 0, credits: 0, ms: 0 });145        s.attempts++;146        s.credits += res.costCredits ?? 0;147        s.ms += res.durationMs ?? 0;148        if (ok) s.success++;149        if (res.httpStatus !== null && res.httpStatus !== undefined) {150          s.statuses ??= {};151          s.statuses[String(res.httpStatus)] = (s.statuses[String(res.httpStatus)] ?? 0) + 1;152        }153        if (blocked) s.blocked = (s.blocked ?? 0) + 1;154      }155      // A 404/410 is a definitive answer, not a host failure.156      const hostFailure = !ok && res.httpStatus !== 404 && res.httpStatus !== 410;157      if (!o.ignoreCircuit) gates.report(host, !hostFailure);158      opts.onAttempt?.({ engine, url, host, success: ok, credits: res.costCredits ?? 0, ms: res.durationMs ?? 0, quality, error: res.error, httpStatus: res.httpStatus ?? null, blocked });159      if (ok) return res;160      errors.push(`${engine}: ${res.error ?? `quality ${quality}`}`);161      last = res;162      opts.log?.debug({ url, engine, quality, error: res.error }, 'engine fallback');163      // Do not escalate on hard "not found"/"gone" statuses: the page does not exist.164      if (res.httpStatus === 404 || res.httpStatus === 410) break;165    }166    return {167      ...(last ?? empty(url, engines[0] ?? 'api', 'no engine available')),168      success: false,169      requiresReview: true,170      error: errors.join(' | ') || 'no engine available',171    };172  }173174  return { fetch: routedFetch, engines: { http, firecrawl, scrapfly }, available, gates };175}176