import type { Engine, ExtractionResult, Logger } from '@rareindex/shared'; import { createHttpEngine, type HttpEngine } from './engines/http.js'; import { createFirecrawlEngine, type FirecrawlEngine } from './engines/firecrawl.js'; import { createScrapflyEngine, type ScrapflyEngine } from './engines/scrapfly.js'; import { documentQuality, isChallengePage, qualityScore, MIN_QUALITY } from './quality.js'; import { policyFor } from './domains.js'; import { CircuitOpenError, HostGates } from './circuit.js'; import { domainOf, type EngineStats, type FetchOptions } from './types.js'; export interface RouterAttempt { engine: Engine; url: string; host: string; success: boolean; credits: number; ms: number; quality: number; error: string | null; httpStatus: number | null; blocked: boolean; } export interface RouterOptions { firecrawlApiKey?: string; scrapflyApiKey?: string; log?: Logger; /** called for every engine attempt (cost ledger, metrics) */ onAttempt?: (info: RouterAttempt) => void; /** shared gates (one per process) so several connectors on the same host respect one limit */ gates?: HostGates; } export interface Router { fetch(url: string, opts?: FetchOptions, defaults?: { enginePriority: Engine[]; stats?: Record }): Promise; engines: { http: HttpEngine; firecrawl: FirecrawlEngine | null; scrapfly: ScrapflyEngine | null }; available(): Engine[]; gates: HostGates; } /** Process-wide gates keyed by host; policies come from connectors/domains.json. */ export const sharedGates = new HostGates((host) => { const p = policyFor(host); return { concurrency: p.concurrency, minIntervalMs: p.minIntervalMs, circuitFailures: p.circuitFailures, circuitCooldownMs: p.circuitCooldownMs }; }); /** * Firecrawl → Scrapfly router (§198). Tries engines in priority order, scores each result with the * connector's parser (or a document heuristic), and falls through until quality ≥ minQuality. * A result that fails everywhere is returned with requiresReview=true — never silently dropped. * Every attempt passes through the per-host gate (concurrency, interval, circuit breaker). */ export function createRouter(opts: RouterOptions = {}): Router { const http = createHttpEngine(); const firecrawl = opts.firecrawlApiKey ? createFirecrawlEngine({ apiKey: opts.firecrawlApiKey }) : null; const scrapfly = opts.scrapflyApiKey ? createScrapflyEngine({ apiKey: opts.scrapflyApiKey }) : null; const gates = opts.gates ?? sharedGates; function available(): Engine[] { const out: Engine[] = ['api', 'feed']; if (firecrawl) out.push('firecrawl'); if (scrapfly) out.push('scrapfly'); return out; } async function runEngine(engine: Engine, url: string, o: FetchOptions): Promise { switch (engine) { case 'api': case 'feed': return http.fetch(url, o); case 'firecrawl': return firecrawl ? firecrawl.scrape(url, o) : null; case 'scrapfly': return scrapfly ? scrapfly.scrape(url, o) : null; default: return null; } } function empty(url: string, engine: Engine, error: string): ExtractionResult { 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 }; } async function routedFetch(url: string, o: FetchOptions = {}, defaults?: { enginePriority: Engine[]; stats?: Record }): Promise { const host = domainOf(url); const policy = policyFor(host); let engines = (o.engines ?? defaults?.enginePriority ?? ['api', 'firecrawl', 'scrapfly']).filter((e) => available().includes(e)); if (policy.engines) engines = engines.filter((e) => policy.engines!.includes(e)); const merged: FetchOptions = { ...o, timeoutMs: o.timeoutMs ?? policy.timeoutMs, country: o.country ?? policy.scrapfly.country ?? policy.country, waitForMs: o.waitForMs ?? policy.firecrawl.waitForMs, renderJs: o.renderJs ?? policy.scrapfly.renderJs, }; const minQuality = o.minQuality ?? MIN_QUALITY; let last: ExtractionResult | null = null; const errors: string[] = []; for (const engine of engines) { let release: (() => void) | null = null; if (!o.ignoreCircuit) { try { release = await gates.acquire(host); } catch (err) { if (err instanceof CircuitOpenError) { const stats = defaults?.stats; if (stats) { const s = (stats[engine] ??= { attempts: 0, success: 0, credits: 0, ms: 0 }); s.circuitOpen = (s.circuitOpen ?? 0) + 1; } opts.log?.warn({ host, url, retryAt: err.retryAt }, 'circuit open; request refused'); return empty(url, engine, err.message); } throw err; } } let res: ExtractionResult | null; try { res = await runEngine(engine, url, merged); } finally { release?.(); } if (!res) continue; let quality = 0; const blocked = res.success ? isChallengePage(res) : false; if (res.success && !blocked) { if (o.parse) { try { const found = o.parse(res); quality = qualityScore(found, o.expect); } catch (err) { quality = 0; res.error = `parse: ${err instanceof Error ? err.message : String(err)}`; } } else { quality = documentQuality(res); } } else if (blocked) { res.error = res.error ?? 'anti-bot challenge page'; } res.qualityScore = quality; const ok = res.success && !blocked && quality >= minQuality; const stats = defaults?.stats; if (stats) { const s = (stats[engine] ??= { attempts: 0, success: 0, credits: 0, ms: 0 }); s.attempts++; s.credits += res.costCredits ?? 0; s.ms += res.durationMs ?? 0; if (ok) s.success++; if (res.httpStatus !== null && res.httpStatus !== undefined) { s.statuses ??= {}; s.statuses[String(res.httpStatus)] = (s.statuses[String(res.httpStatus)] ?? 0) + 1; } if (blocked) s.blocked = (s.blocked ?? 0) + 1; } // A 404/410 is a definitive answer, not a host failure. const hostFailure = !ok && res.httpStatus !== 404 && res.httpStatus !== 410; if (!o.ignoreCircuit) gates.report(host, !hostFailure); opts.onAttempt?.({ engine, url, host, success: ok, credits: res.costCredits ?? 0, ms: res.durationMs ?? 0, quality, error: res.error, httpStatus: res.httpStatus ?? null, blocked }); if (ok) return res; errors.push(`${engine}: ${res.error ?? `quality ${quality}`}`); last = res; opts.log?.debug({ url, engine, quality, error: res.error }, 'engine fallback'); // Do not escalate on hard "not found"/"gone" statuses: the page does not exist. if (res.httpStatus === 404 || res.httpStatus === 410) break; } return { ...(last ?? empty(url, engines[0] ?? 'api', 'no engine available')), success: false, requiresReview: true, error: errors.join(' | ') || 'no engine available', }; } return { fetch: routedFetch, engines: { http, firecrawl, scrapfly }, available, gates }; }