TypeScript 61.9%
HTML 37.2%
SQL 0.7%
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