import { Agent, fetch as undiciFetch, type Dispatcher } from "undici"; import { gunzipSync } from "node:zlib"; import { assertUrlAllowed, safeLookup, UrlPolicyError } from "@dci/core"; import type { FetchLevel } from "@dci/core"; import type { Fetcher, FetchOptions, RawDocument } from "./types.js"; /** * Fetch abstraction with escalation levels: * 1 — direct HTTP (conditional GET, ETag / Last-Modified) * 2 — direct HTTP + parsing hints (browser UA, accept html) [same transport, different identity] * 3 — Firecrawl (markdown + html normalization) * 4 — Scrapfly rendered request (anti-bot + JS rendering) * Every URL passes the SSRF policy before any network activity, on every redirect hop. * * Identity: robots.txt is always evaluated for the bot token (DEFAULT_UA / DCI_USER_AGENT) — see robots.ts. The L2 * browser User-Agent is only ever sent for a URL the bot token was already allowed to fetch; it changes how a * server renders the page, never whether we are permitted to request it. HTTP 429 is never escalated past. */ export const DEFAULT_UA = process.env.DCI_USER_AGENT ?? "DataCenterIndexBot/0.1 (+https://www.datacenterindex.io/bot; contact@spboucher.ai)"; export const BROWSER_UA = "Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/605.1.15 (KHTML, like Gecko) Version/17.5 Safari/605.1.15"; const DEFAULT_TIMEOUT = 30_000; const DEFAULT_MAX_BYTES = 25 * 1024 * 1024; let agent: Dispatcher | null = null; let agentH1: Dispatcher | null = null; const h1Hosts = new Set(); function dispatcher(h1 = false): Dispatcher { if (h1) return (agentH1 ??= new Agent({ connect: { lookup: safeLookup as never, timeout: 10_000 }, connections: 32, pipelining: 1, keepAliveTimeout: 15_000, headersTimeout: 30_000, bodyTimeout: 45_000, allowH2: false })); return (agent ??= new Agent({ connect: { lookup: safeLookup as never, timeout: 10_000 }, connections: 64, pipelining: 1, keepAliveTimeout: 15_000, headersTimeout: 30_000, bodyTimeout: 45_000, allowH2: true })); } export async function closeFetchers(): Promise { await agent?.close(); await agentH1?.close(); agent = agentH1 = null; } function hdrs(h: Headers): Record { const out: Record = {}; h.forEach((v, k) => { if (/^(content-type|content-length|etag|last-modified|cache-control|server|retry-after|date|x-robots-tag|link|location|content-disposition|x-ratelimit-[a-z-]+)$/i.test(k)) out[k.toLowerCase()] = v; }); return out; } function decodeText(body: Buffer, contentType: string | null): string { const m = contentType?.match(/charset=([\w-]+)/i); const cs = (m?.[1] ?? "utf-8").toLowerCase(); try { if (cs === "utf-8" || cs === "utf8") return body.toString("utf8"); return new TextDecoder(cs as string).decode(body); } catch { return body.toString("utf8"); } } export function failDoc(url: string, level: FetchLevel, fetcher: RawDocument["fetcher"], code: string, message: string, started: number, status = 0): RawDocument { return { url, finalUrl: url, fetchedAt: new Date().toISOString(), status, contentType: null, body: Buffer.alloc(0), text: "", headers: {}, etag: null, lastModified: null, notModified: false, fetcher, level, durationMs: Date.now() - started, credits: 0, error: { code, message } }; } export class DirectFetcher implements Fetcher { readonly name = "direct" as const; constructor(readonly level: FetchLevel = 1) {} available(): boolean { return true; } async fetch(urlStr: string, opts: FetchOptions = {}): Promise { const started = Date.now(); const level = opts.level ?? this.level; const headers: Record = { "user-agent": level >= 2 ? BROWSER_UA : DEFAULT_UA, accept: opts.accept ?? "text/html,application/xhtml+xml,application/xml;q=0.9,application/json;q=0.9,application/pdf;q=0.8,*/*;q=0.5", "accept-language": "en-US,en;q=0.9,fr;q=0.6,de;q=0.4", "accept-encoding": "gzip, deflate, br", ...(level >= 2 ? { "sec-fetch-dest": "document", "sec-fetch-mode": "navigate", "sec-fetch-site": "none", "upgrade-insecure-requests": "1" } : {}), ...opts.headers, }; if (opts.etag) headers["if-none-match"] = opts.etag; if (opts.lastModified) headers["if-modified-since"] = opts.lastModified; let current = urlStr; let redirects = 0; for (;;) { try { await assertUrlAllowed(current); } catch (e) { return failDoc(urlStr, level, "direct", "ssrf_blocked", e instanceof UrlPolicyError ? e.message : String(e), started); } const ac = new AbortController(); const timer = setTimeout(() => ac.abort(), opts.timeoutMs ?? DEFAULT_TIMEOUT); let res: Response; const host = new URL(current).hostname; try { res = (await undiciFetch(current, { method: "GET", headers, redirect: "manual", signal: ac.signal, dispatcher: dispatcher(h1Hosts.has(host)) } as never)) as unknown as Response; } catch (e) { clearTimeout(timer); const msg = e instanceof Error ? `${e.name}: ${e.message}${(e as { cause?: Error }).cause ? " — " + String((e as { cause?: Error }).cause?.message) : ""}` : String(e); if (/NGHTTP2|HTTP\/2/i.test(msg) && !h1Hosts.has(host)) { h1Hosts.add(host); continue; } const code = ac.signal.aborted ? "timeout" : /ENOTFOUND|EAI_AGAIN|getaddrinfo/i.test(msg) ? "dns" : /blocked/i.test(msg) ? "ssrf_blocked" : /CERT|TLS|SSL|certificate/i.test(msg) ? "tls" : /ECONNREFUSED|ECONNRESET|EHOSTUNREACH|ETIMEDOUT|socket/i.test(msg) ? "connection" : "fetch_failed"; return failDoc(urlStr, level, "direct", code, msg, started); } if ([301, 302, 303, 307, 308].includes(res.status)) { clearTimeout(timer); const loc = res.headers.get("location"); if (!loc) return failDoc(urlStr, level, "direct", "redirect_without_location", `HTTP ${res.status}`, started, res.status); if (++redirects > 6) return failDoc(urlStr, level, "direct", "too_many_redirects", "more than 6 redirects", started, res.status); try { current = new URL(loc, current).toString(); } catch { return failDoc(urlStr, level, "direct", "bad_redirect", loc, started, res.status); } delete headers["if-none-match"]; delete headers["if-modified-since"]; continue; } const h = hdrs(res.headers); const base = { url: urlStr, finalUrl: current, fetchedAt: new Date().toISOString(), status: res.status, contentType: res.headers.get("content-type"), headers: h, etag: res.headers.get("etag"), lastModified: res.headers.get("last-modified"), fetcher: "direct" as const, level, credits: 0 }; if (res.status === 304) { clearTimeout(timer); return { ...base, body: Buffer.alloc(0), text: "", notModified: true, durationMs: Date.now() - started }; } const maxBytes = opts.maxBytes ?? DEFAULT_MAX_BYTES; const declared = Number(res.headers.get("content-length") ?? 0); if (declared > maxBytes) { clearTimeout(timer); return failDoc(urlStr, level, "direct", "too_large", `content-length ${declared}`, started, res.status); } const chunks: Uint8Array[] = []; let total = 0; try { const reader = res.body?.getReader(); if (reader) for (;;) { const { done, value } = await reader.read(); if (done) break; total += value.byteLength; if (total > maxBytes) { await reader.cancel(); clearTimeout(timer); return failDoc(urlStr, level, "direct", "too_large", `body > ${maxBytes}`, started, res.status); } chunks.push(value); } } catch (e) { clearTimeout(timer); return failDoc(urlStr, level, "direct", ac.signal.aborted ? "timeout" : "body_read_failed", (e as Error).message, started, res.status); } clearTimeout(timer); let body: Buffer = Buffer.concat(chunks); if (/\.gz$/i.test(current) || (h["content-type"] ?? "").includes("gzip")) { const r = safeGunzip(body, maxBytes); if (r.error) return failDoc(urlStr, level, "direct", "too_large", r.error, started, res.status); body = r.body; } const ct = base.contentType; const isText = !ct || /text|json|xml|javascript|html|csv|markdown/i.test(ct); return applyRobotsDirectives({ ...base, body, text: isText ? decodeText(body, ct) : "", notModified: false, durationMs: Date.now() - started }); } } } /** Decompression cap: an archive may expand to at most max(4 × maxBytes, 64 MB) — anything bigger is a bomb. */ export function gunzipLimit(maxBytes: number): number { return Math.max(maxBytes * 4, 64 * 1024 * 1024); } export function safeGunzip(body: Buffer, maxBytes = DEFAULT_MAX_BYTES): { body: Buffer; error?: string } { const maxOutputLength = gunzipLimit(maxBytes); try { return { body: gunzipSync(body, { maxOutputLength }) }; } catch (e) { const code = (e as { code?: string }).code; if (code === "ERR_BUFFER_TOO_LARGE" || e instanceof RangeError) return { body, error: `gzip output exceeds ${maxOutputLength} bytes` }; return { body }; // not gzip after all: keep the raw bytes } } /** * `X-Robots-Tag` / `` directives we honour: `noarchive` and `noindex` both mean "do not keep a * copy of this page" — the pipeline still extracts facts (facts are not copyrightable) but skips the raw archive. * Sets `doc.meta.noarchive = true` (and `doc.meta.robotsDirective` with the matched token). Idempotent. */ export function applyRobotsDirectives(doc: RawDocument): RawDocument { const dir = robotsDirective(doc); if (dir) doc.meta = { ...(doc.meta ?? {}), noarchive: true, robotsDirective: dir }; return doc; } export function robotsDirective(doc: RawDocument): string | null { const header = doc.headers["x-robots-tag"]; if (header && /\b(noarchive|noindex|none)\b/i.test(header)) return `header:${header.match(/\b(noarchive|noindex|none)\b/i)![1]!.toLowerCase()}`; if (doc.text && /html/i.test(doc.contentType ?? "")) { const head = doc.text.slice(0, 200_000); const re = /]*\bname\s*=\s*["']?(?:robots|datacenterindexbot)["']?[^>]*>/gi; for (const m of head.matchAll(re)) { const content = m[0].match(/\bcontent\s*=\s*["']([^"']*)["']/i)?.[1] ?? ""; const tok = content.match(/\b(noarchive|noindex|none)\b/i); if (tok) return `meta:${tok[1]!.toLowerCase()}`; } } return null; } /** `Retry-After` header → milliseconds (delta-seconds or HTTP-date), bounded to [0, 7 days]; null when absent/invalid. */ export function parseRetryAfter(value: string | null | undefined, now = Date.now()): number | null { if (!value) return null; const v = value.trim(); const max = 7 * 86_400_000; if (/^\d+$/.test(v)) return Math.min(max, Number(v) * 1000); const t = Date.parse(v); if (Number.isFinite(t)) return Math.min(max, Math.max(0, t - now)); return null; } export type PremiumProvider = "scrapfly" | "firecrawl"; /** * Pluggable shared store for the daily premium budgets. The in-process counter is always kept; a store lets the * runtime persist credits (e.g. Redis `dci:budget::`) so budgets survive restarts and are shared by * several worker processes. `used()` returns the latest known shared value (may lag) or null when unknown; * `spend()` must never throw (fire-and-forget is fine). */ export interface BudgetStore { used(provider: PremiumProvider, day: string): number | null; spend(provider: PremiumProvider, day: string, credits: number): void; } let budgetStore: BudgetStore | null = null; export function setBudgetStore(store: BudgetStore | null): void { budgetStore = store; } /** UTC day key used by the budgets and the shared store. */ export function budgetDay(now = new Date()): string { return now.toISOString().slice(0, 10); } /** Daily budget bookkeeping shared by premium fetchers. */ export class Budget { private day = ""; private local = 0; constructor(readonly provider: PremiumProvider, private readonly envKey: string, private readonly fallback: number) {} private roll() { const d = budgetDay(); if (d !== this.day) { this.day = d; this.local = 0; } } limit(): number { return Number(process.env[this.envKey] ?? this.fallback); } /** credits used today: the shared store's view when available, never below what this process spent itself */ used(): number { this.roll(); return Math.max(this.local, budgetStore?.used(this.provider, this.day) ?? 0); } remaining(): number { return Math.max(0, this.limit() - this.used()); } spend(n = 1) { this.roll(); this.local += n; try { budgetStore?.spend(this.provider, this.day, n); } catch { /* store must not break fetching */ } } snapshot() { this.roll(); return { day: this.day, used: this.used(), local: this.local, limit: this.limit() }; } } export const scrapflyBudget = new Budget("scrapfly", "DCI_SCRAPFLY_DAILY_BUDGET", 400); export const firecrawlBudget = new Budget("firecrawl", "DCI_FIRECRAWL_DAILY_BUDGET", 200); export class FirecrawlFetcher implements Fetcher { readonly name = "firecrawl" as const; readonly level: FetchLevel = 3; available(): boolean { return Boolean(process.env.FIRECRAWL_API_KEY) && firecrawlBudget.remaining() > 0; } async fetch(url: string, opts: FetchOptions = {}): Promise { const started = Date.now(); if (!this.available()) return failDoc(url, 3, "firecrawl", "firecrawl_unavailable", "no key or daily budget exhausted", started); try { await assertUrlAllowed(url); } catch (e) { return failDoc(url, 3, "firecrawl", "ssrf_blocked", (e as Error).message, started); } const ac = new AbortController(); const timer = setTimeout(() => ac.abort(), opts.timeoutMs ?? 90_000); try { const res = await fetch("https://api.firecrawl.dev/v1/scrape", { method: "POST", signal: ac.signal, headers: { "content-type": "application/json", authorization: `Bearer ${process.env.FIRECRAWL_API_KEY}` }, body: JSON.stringify({ url, formats: ["markdown", "html"], onlyMainContent: false, waitFor: opts.waitForSelector ? 3000 : 0 }) }); const json = (await res.json()) as { success?: boolean; data?: { markdown?: string; html?: string; metadata?: { statusCode?: number; sourceURL?: string; url?: string; title?: string } }; error?: string }; // the daily budget is debited only when Firecrawl actually delivered a document (failed calls are not billed as a scrape) if (!res.ok || !json.success || !json.data) return failDoc(url, 3, "firecrawl", "firecrawl_error", `${res.status} ${json.error ?? "no data"}`, started); firecrawlBudget.spend(1); const html = json.data.html ?? ""; const body = Buffer.from(html || json.data.markdown || "", "utf8"); return { url, finalUrl: json.data.metadata?.url ?? json.data.metadata?.sourceURL ?? url, fetchedAt: new Date().toISOString(), status: json.data.metadata?.statusCode ?? 200, contentType: html ? "text/html; charset=utf-8" : "text/markdown", body, text: body.toString("utf8"), headers: {}, etag: null, lastModified: null, notModified: false, fetcher: "firecrawl", level: 3, durationMs: Date.now() - started, credits: 1, markdown: json.data.markdown ?? null }; } catch (e) { return failDoc(url, 3, "firecrawl", ac.signal.aborted ? "timeout" : "firecrawl_error", (e as Error).message, started); } finally { clearTimeout(timer); } } } /** Scrapfly: https://scrapfly.io/docs/scrape-api/getting-started — GET /scrape?key&url&asp&render_js&country&cache */ export class ScrapflyFetcher implements Fetcher { readonly name = "scrapfly" as const; readonly level: FetchLevel = 4; available(): boolean { return Boolean(process.env.SCRAPFLY_API_KEY) && scrapflyBudget.remaining() > 0; } async fetch(url: string, opts: FetchOptions = {}): Promise { const started = Date.now(); if (!this.available()) return failDoc(url, 4, "scrapfly", "scrapfly_unavailable", "no key or daily budget exhausted", started); try { await assertUrlAllowed(url); } catch (e) { return failDoc(url, 4, "scrapfly", "ssrf_blocked", (e as Error).message, started); } const q = new URLSearchParams({ key: process.env.SCRAPFLY_API_KEY!, url, asp: "true", render_js: opts.renderJs === false ? "false" : "true", country: opts.country ?? "us", retry: "true", cache: "true", cache_ttl: "3600" }); if (opts.waitForSelector) q.set("wait_for_selector", opts.waitForSelector); if (opts.renderJs !== false) q.set("rendering_wait", "1500"); const ac = new AbortController(); const timer = setTimeout(() => ac.abort(), opts.timeoutMs ?? 120_000); try { const res = await fetch(`https://api.scrapfly.io/scrape?${q}`, { signal: ac.signal, headers: { accept: "application/json" } }); const json = (await res.json()) as { result?: { content?: string; status_code?: number; success?: boolean; reason?: string; response_headers?: Record; url?: string; error?: { message?: string }; format?: string }; context?: { cost?: { total?: number } }; message?: string; code?: string }; const r = json.result; const cost = Number(json.context?.cost?.total ?? 1) || 1; scrapflyBudget.spend(cost); if (!res.ok || !r) return failDoc(url, 4, "scrapfly", "scrapfly_error", `${res.status} ${json.message ?? json.code ?? "no result"}`, started); if (!r.success) return { ...failDoc(url, 4, "scrapfly", "scrapfly_failed", `${r.status_code ?? 0} ${r.reason ?? r.error?.message ?? "upstream failed"}`, started, r.status_code ?? 0), credits: cost }; const h = Object.fromEntries(Object.entries(r.response_headers ?? {}).map(([k, v]) => [k.toLowerCase(), String(v)])); const body = Buffer.from(r.content ?? "", "utf8"); return applyRobotsDirectives({ url, finalUrl: r.url ?? url, fetchedAt: new Date().toISOString(), status: r.status_code ?? 200, contentType: h["content-type"] ?? "text/html; charset=utf-8", body, text: body.toString("utf8"), headers: h, etag: null, lastModified: h["last-modified"] ?? null, notModified: false, fetcher: "scrapfly", level: 4, durationMs: Date.now() - started, credits: cost }); } catch (e) { return failDoc(url, 4, "scrapfly", ac.signal.aborted ? "timeout" : "scrapfly_error", (e as Error).message, started); } finally { clearTimeout(timer); } } } export const fetchers: Record = { 1: new DirectFetcher(1), 2: new DirectFetcher(2), 3: new FirecrawlFetcher(), 4: new ScrapflyFetcher() }; /** Heuristics: a 200 that is really an anti-bot interstitial or an empty JS shell. */ export function looksBlockedOrEmpty(doc: RawDocument): boolean { if (doc.status === 403 || doc.status === 429 || doc.status === 503 || doc.status === 401) return true; if (doc.status !== 200 || !doc.contentType || !/html/i.test(doc.contentType)) return false; const t = doc.text; const textLen = t.replace(//gi, "").replace(//gi, "").replace(/<[^>]+>/g, " ").replace(/\s+/g, " ").trim().length; // Hard interstitial markers (the page IS the challenge) — only trusted when the page carries little real text, // because Cloudflare-fronted sites embed the challenge-platform beacon on perfectly normal pages. const interstitial = /([^<]*(Just a moment|Attention Required|Access Denied|Pardon Our Interruption|Security Checkpoint|Please Wait)[^<]*<\/title>|cf-browser-verification|_Incapsula_Resource|Request unsuccessful\. Incapsula|challenge-platform|Enable JavaScript and cookies to continue|are you a robot|verify you are human|Reference #\d+\.[0-9a-f]+)/i.test(t); if (interstitial && textLen < 1500) return true; return textLen < 400 && /<(div id="(root|app|__next|__nuxt)"|app-root|noscript)/i.test(t); } /** Expected premium cost of one fetch at a level (credits): Firecrawl 1 scrape; Scrapfly ≈1 with ASP only, ≈6 with JS rendering. */ export function expectedCredits(level: FetchLevel, opts: Pick<FetchOptions, "renderJs"> = {}): number { if (level === 3) return 1; if (level === 4) return opts.renderJs === false ? 1 : 6; return 0; } /** * Escalating fetch: start at `level`, escalate up to `maxLevel` when blocked/empty. Premium fetchers are * skipped when unavailable, and when their expected cost exceeds `creditsLeft` (the run's remaining premium budget — * so one call never spends on L3 and L4 together when the run budget cannot cover both). HTTP 429 never escalates: * the server asked us to slow down, the document is returned with `meta.retryAfterMs` from `Retry-After` when sent. * Returns the best document obtained (never throws). */ export async function fetchWithEscalation(url: string, opts: FetchOptions & { maxLevel?: FetchLevel; creditsLeft?: number } = {}): Promise<RawDocument> { let level: FetchLevel = opts.level ?? 1; const max: FetchLevel = opts.maxLevel ?? 2; let creditsLeft = opts.creditsLeft; let best: RawDocument | null = null; const attempts: string[] = []; while (level <= max) { const f = fetchers[level]; if (!f.available()) { attempts.push(`L${level}:unavailable`); level = (level + 1) as FetchLevel; continue; } const cost = expectedCredits(level, opts); if (creditsLeft !== undefined && cost > 0 && cost > creditsLeft) { attempts.push(`L${level}:over_budget(${cost}>${creditsLeft})`); level = (level + 1) as FetchLevel; continue; } const doc = applyRobotsDirectives(await f.fetch(url, { ...opts, level })); if (creditsLeft !== undefined) creditsLeft = Math.max(0, creditsLeft - doc.credits); doc.meta = { ...(doc.meta ?? {}), attempts: [...attempts, `L${level}:${doc.error?.code ?? doc.status}`] }; if (doc.notModified) return doc; if (doc.status === 429) { const retryAfterMs = parseRetryAfter(doc.headers["retry-after"]); doc.meta = { ...doc.meta, rateLimited: true, ...(retryAfterMs != null ? { retryAfterMs } : {}) }; return doc; // never buy our way past a rate limit } if (!doc.error && !looksBlockedOrEmpty(doc)) return doc; // hard failures that escalation cannot fix if (doc.error && ["ssrf_blocked", "dns", "too_large"].includes(doc.error.code)) return doc; if (!doc.error && doc.status === 404) return doc; if (!doc.error && doc.status === 410) return doc; best = best ?? doc; if (!doc.error && doc.status >= 200 && doc.status < 300) best = doc; // blocked-looking 200 still better than nothing attempts.push(`L${level}:${doc.error?.code ?? doc.status}`); level = (level + 1) as FetchLevel; } // the fallback document carries the complete attempt list (including levels skipped as unavailable / over budget) if (best) best.meta = { ...(best.meta ?? {}), attempts: [...attempts] }; return best ?? failDoc(url, max, "direct", "no_fetcher", attempts.join(","), Date.now()); }