spb/datacenterindex
Public
HTML 53.9%
TypeScript 44.5%
JavaScript 0.6%
SQL 0.5%
1import { Agent, fetch as undiciFetch, type Dispatcher } from "undici";2import { gunzipSync } from "node:zlib";3import { assertUrlAllowed, safeLookup, UrlPolicyError } from "@dci/core";4import type { FetchLevel } from "@dci/core";5import type { Fetcher, FetchOptions, RawDocument } from "./types.js";67/**8 * Fetch abstraction with escalation levels:9 * 1 — direct HTTP (conditional GET, ETag / Last-Modified)10 * 2 — direct HTTP + parsing hints (browser UA, accept html) [same transport, different identity]11 * 3 — Firecrawl (markdown + html normalization)12 * 4 — Scrapfly rendered request (anti-bot + JS rendering)13 * Every URL passes the SSRF policy before any network activity, on every redirect hop.14 *15 * Identity: robots.txt is always evaluated for the bot token (DEFAULT_UA / DCI_USER_AGENT) — see robots.ts. The L216 * browser User-Agent is only ever sent for a URL the bot token was already allowed to fetch; it changes how a17 * server renders the page, never whether we are permitted to request it. HTTP 429 is never escalated past.18 */1920export const DEFAULT_UA = process.env.DCI_USER_AGENT ?? "DataCenterIndexBot/0.1 (+https://www.datacenterindex.io/bot; contact@spboucher.ai)";21export 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";22const DEFAULT_TIMEOUT = 30_000;23const DEFAULT_MAX_BYTES = 25 * 1024 * 1024;2425let agent: Dispatcher | null = null;26let agentH1: Dispatcher | null = null;27const h1Hosts = new Set<string>();28function dispatcher(h1 = false): Dispatcher {29 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 }));30 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 }));31}32export async function closeFetchers(): Promise<void> {33 await agent?.close();34 await agentH1?.close();35 agent = agentH1 = null;36}3738function hdrs(h: Headers): Record<string, string> {39 const out: Record<string, string> = {};40 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; });41 return out;42}4344function decodeText(body: Buffer, contentType: string | null): string {45 const m = contentType?.match(/charset=([\w-]+)/i);46 const cs = (m?.[1] ?? "utf-8").toLowerCase();47 try {48 if (cs === "utf-8" || cs === "utf8") return body.toString("utf8");49 return new TextDecoder(cs as string).decode(body);50 } catch {51 return body.toString("utf8");52 }53}5455export function failDoc(url: string, level: FetchLevel, fetcher: RawDocument["fetcher"], code: string, message: string, started: number, status = 0): RawDocument {56 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 } };57}5859export class DirectFetcher implements Fetcher {60 readonly name = "direct" as const;61 constructor(readonly level: FetchLevel = 1) {}62 available(): boolean { return true; }63 async fetch(urlStr: string, opts: FetchOptions = {}): Promise<RawDocument> {64 const started = Date.now();65 const level = opts.level ?? this.level;66 const headers: Record<string, string> = {67 "user-agent": level >= 2 ? BROWSER_UA : DEFAULT_UA,68 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",69 "accept-language": "en-US,en;q=0.9,fr;q=0.6,de;q=0.4",70 "accept-encoding": "gzip, deflate, br",71 ...(level >= 2 ? { "sec-fetch-dest": "document", "sec-fetch-mode": "navigate", "sec-fetch-site": "none", "upgrade-insecure-requests": "1" } : {}),72 ...opts.headers,73 };74 if (opts.etag) headers["if-none-match"] = opts.etag;75 if (opts.lastModified) headers["if-modified-since"] = opts.lastModified;76 let current = urlStr;77 let redirects = 0;78 for (;;) {79 try { await assertUrlAllowed(current); } catch (e) { return failDoc(urlStr, level, "direct", "ssrf_blocked", e instanceof UrlPolicyError ? e.message : String(e), started); }80 const ac = new AbortController();81 const timer = setTimeout(() => ac.abort(), opts.timeoutMs ?? DEFAULT_TIMEOUT);82 let res: Response;83 const host = new URL(current).hostname;84 try {85 res = (await undiciFetch(current, { method: "GET", headers, redirect: "manual", signal: ac.signal, dispatcher: dispatcher(h1Hosts.has(host)) } as never)) as unknown as Response;86 } catch (e) {87 clearTimeout(timer);88 const msg = e instanceof Error ? `${e.name}: ${e.message}${(e as { cause?: Error }).cause ? " — " + String((e as { cause?: Error }).cause?.message) : ""}` : String(e);89 if (/NGHTTP2|HTTP\/2/i.test(msg) && !h1Hosts.has(host)) { h1Hosts.add(host); continue; }90 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";91 return failDoc(urlStr, level, "direct", code, msg, started);92 }93 if ([301, 302, 303, 307, 308].includes(res.status)) {94 clearTimeout(timer);95 const loc = res.headers.get("location");96 if (!loc) return failDoc(urlStr, level, "direct", "redirect_without_location", `HTTP ${res.status}`, started, res.status);97 if (++redirects > 6) return failDoc(urlStr, level, "direct", "too_many_redirects", "more than 6 redirects", started, res.status);98 try { current = new URL(loc, current).toString(); } catch { return failDoc(urlStr, level, "direct", "bad_redirect", loc, started, res.status); }99 delete headers["if-none-match"]; delete headers["if-modified-since"];100 continue;101 }102 const h = hdrs(res.headers);103 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 };104 if (res.status === 304) { clearTimeout(timer); return { ...base, body: Buffer.alloc(0), text: "", notModified: true, durationMs: Date.now() - started }; }105 const maxBytes = opts.maxBytes ?? DEFAULT_MAX_BYTES;106 const declared = Number(res.headers.get("content-length") ?? 0);107 if (declared > maxBytes) { clearTimeout(timer); return failDoc(urlStr, level, "direct", "too_large", `content-length ${declared}`, started, res.status); }108 const chunks: Uint8Array[] = [];109 let total = 0;110 try {111 const reader = res.body?.getReader();112 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); }113 } catch (e) {114 clearTimeout(timer);115 return failDoc(urlStr, level, "direct", ac.signal.aborted ? "timeout" : "body_read_failed", (e as Error).message, started, res.status);116 }117 clearTimeout(timer);118 let body: Buffer = Buffer.concat(chunks);119 if (/\.gz$/i.test(current) || (h["content-type"] ?? "").includes("gzip")) {120 const r = safeGunzip(body, maxBytes);121 if (r.error) return failDoc(urlStr, level, "direct", "too_large", r.error, started, res.status);122 body = r.body;123 }124 const ct = base.contentType;125 const isText = !ct || /text|json|xml|javascript|html|csv|markdown/i.test(ct);126 return applyRobotsDirectives({ ...base, body, text: isText ? decodeText(body, ct) : "", notModified: false, durationMs: Date.now() - started });127 }128 }129}130131/** Decompression cap: an archive may expand to at most max(4 × maxBytes, 64 MB) — anything bigger is a bomb. */132export function gunzipLimit(maxBytes: number): number { return Math.max(maxBytes * 4, 64 * 1024 * 1024); }133export function safeGunzip(body: Buffer, maxBytes = DEFAULT_MAX_BYTES): { body: Buffer; error?: string } {134 const maxOutputLength = gunzipLimit(maxBytes);135 try {136 return { body: gunzipSync(body, { maxOutputLength }) };137 } catch (e) {138 const code = (e as { code?: string }).code;139 if (code === "ERR_BUFFER_TOO_LARGE" || e instanceof RangeError) return { body, error: `gzip output exceeds ${maxOutputLength} bytes` };140 return { body }; // not gzip after all: keep the raw bytes141 }142}143144/**145 * `X-Robots-Tag` / `<meta name="robots">` directives we honour: `noarchive` and `noindex` both mean "do not keep a146 * copy of this page" — the pipeline still extracts facts (facts are not copyrightable) but skips the raw archive.147 * Sets `doc.meta.noarchive = true` (and `doc.meta.robotsDirective` with the matched token). Idempotent.148 */149export function applyRobotsDirectives(doc: RawDocument): RawDocument {150 const dir = robotsDirective(doc);151 if (dir) doc.meta = { ...(doc.meta ?? {}), noarchive: true, robotsDirective: dir };152 return doc;153}154export function robotsDirective(doc: RawDocument): string | null {155 const header = doc.headers["x-robots-tag"];156 if (header && /\b(noarchive|noindex|none)\b/i.test(header)) return `header:${header.match(/\b(noarchive|noindex|none)\b/i)![1]!.toLowerCase()}`;157 if (doc.text && /html/i.test(doc.contentType ?? "")) {158 const head = doc.text.slice(0, 200_000);159 const re = /<meta\b[^>]*\bname\s*=\s*["']?(?:robots|datacenterindexbot)["']?[^>]*>/gi;160 for (const m of head.matchAll(re)) {161 const content = m[0].match(/\bcontent\s*=\s*["']([^"']*)["']/i)?.[1] ?? "";162 const tok = content.match(/\b(noarchive|noindex|none)\b/i);163 if (tok) return `meta:${tok[1]!.toLowerCase()}`;164 }165 }166 return null;167}168169/** `Retry-After` header → milliseconds (delta-seconds or HTTP-date), bounded to [0, 7 days]; null when absent/invalid. */170export function parseRetryAfter(value: string | null | undefined, now = Date.now()): number | null {171 if (!value) return null;172 const v = value.trim();173 const max = 7 * 86_400_000;174 if (/^\d+$/.test(v)) return Math.min(max, Number(v) * 1000);175 const t = Date.parse(v);176 if (Number.isFinite(t)) return Math.min(max, Math.max(0, t - now));177 return null;178}179180export type PremiumProvider = "scrapfly" | "firecrawl";181182/**183 * Pluggable shared store for the daily premium budgets. The in-process counter is always kept; a store lets the184 * runtime persist credits (e.g. Redis `dci:budget:<provider>:<day>`) so budgets survive restarts and are shared by185 * several worker processes. `used()` returns the latest known shared value (may lag) or null when unknown;186 * `spend()` must never throw (fire-and-forget is fine).187 */188export interface BudgetStore {189 used(provider: PremiumProvider, day: string): number | null;190 spend(provider: PremiumProvider, day: string, credits: number): void;191}192let budgetStore: BudgetStore | null = null;193export function setBudgetStore(store: BudgetStore | null): void { budgetStore = store; }194/** UTC day key used by the budgets and the shared store. */195export function budgetDay(now = new Date()): string { return now.toISOString().slice(0, 10); }196197/** Daily budget bookkeeping shared by premium fetchers. */198export class Budget {199 private day = "";200 private local = 0;201 constructor(readonly provider: PremiumProvider, private readonly envKey: string, private readonly fallback: number) {}202 private roll() { const d = budgetDay(); if (d !== this.day) { this.day = d; this.local = 0; } }203 limit(): number { return Number(process.env[this.envKey] ?? this.fallback); }204 /** credits used today: the shared store's view when available, never below what this process spent itself */205 used(): number { this.roll(); return Math.max(this.local, budgetStore?.used(this.provider, this.day) ?? 0); }206 remaining(): number { return Math.max(0, this.limit() - this.used()); }207 spend(n = 1) { this.roll(); this.local += n; try { budgetStore?.spend(this.provider, this.day, n); } catch { /* store must not break fetching */ } }208 snapshot() { this.roll(); return { day: this.day, used: this.used(), local: this.local, limit: this.limit() }; }209}210export const scrapflyBudget = new Budget("scrapfly", "DCI_SCRAPFLY_DAILY_BUDGET", 400);211export const firecrawlBudget = new Budget("firecrawl", "DCI_FIRECRAWL_DAILY_BUDGET", 200);212213export class FirecrawlFetcher implements Fetcher {214 readonly name = "firecrawl" as const;215 readonly level: FetchLevel = 3;216 available(): boolean { return Boolean(process.env.FIRECRAWL_API_KEY) && firecrawlBudget.remaining() > 0; }217 async fetch(url: string, opts: FetchOptions = {}): Promise<RawDocument> {218 const started = Date.now();219 if (!this.available()) return failDoc(url, 3, "firecrawl", "firecrawl_unavailable", "no key or daily budget exhausted", started);220 try { await assertUrlAllowed(url); } catch (e) { return failDoc(url, 3, "firecrawl", "ssrf_blocked", (e as Error).message, started); }221 const ac = new AbortController();222 const timer = setTimeout(() => ac.abort(), opts.timeoutMs ?? 90_000);223 try {224 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 }) });225 const json = (await res.json()) as { success?: boolean; data?: { markdown?: string; html?: string; metadata?: { statusCode?: number; sourceURL?: string; url?: string; title?: string } }; error?: string };226 // the daily budget is debited only when Firecrawl actually delivered a document (failed calls are not billed as a scrape)227 if (!res.ok || !json.success || !json.data) return failDoc(url, 3, "firecrawl", "firecrawl_error", `${res.status} ${json.error ?? "no data"}`, started);228 firecrawlBudget.spend(1);229 const html = json.data.html ?? "";230 const body = Buffer.from(html || json.data.markdown || "", "utf8");231 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 };232 } catch (e) {233 return failDoc(url, 3, "firecrawl", ac.signal.aborted ? "timeout" : "firecrawl_error", (e as Error).message, started);234 } finally { clearTimeout(timer); }235 }236}237238/** Scrapfly: https://scrapfly.io/docs/scrape-api/getting-started — GET /scrape?key&url&asp&render_js&country&cache */239export class ScrapflyFetcher implements Fetcher {240 readonly name = "scrapfly" as const;241 readonly level: FetchLevel = 4;242 available(): boolean { return Boolean(process.env.SCRAPFLY_API_KEY) && scrapflyBudget.remaining() > 0; }243 async fetch(url: string, opts: FetchOptions = {}): Promise<RawDocument> {244 const started = Date.now();245 if (!this.available()) return failDoc(url, 4, "scrapfly", "scrapfly_unavailable", "no key or daily budget exhausted", started);246 try { await assertUrlAllowed(url); } catch (e) { return failDoc(url, 4, "scrapfly", "ssrf_blocked", (e as Error).message, started); }247 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" });248 if (opts.waitForSelector) q.set("wait_for_selector", opts.waitForSelector);249 if (opts.renderJs !== false) q.set("rendering_wait", "1500");250 const ac = new AbortController();251 const timer = setTimeout(() => ac.abort(), opts.timeoutMs ?? 120_000);252 try {253 const res = await fetch(`https://api.scrapfly.io/scrape?${q}`, { signal: ac.signal, headers: { accept: "application/json" } });254 const json = (await res.json()) as { result?: { content?: string; status_code?: number; success?: boolean; reason?: string; response_headers?: Record<string, string>; url?: string; error?: { message?: string }; format?: string }; context?: { cost?: { total?: number } }; message?: string; code?: string };255 const r = json.result;256 const cost = Number(json.context?.cost?.total ?? 1) || 1;257 scrapflyBudget.spend(cost);258 if (!res.ok || !r) return failDoc(url, 4, "scrapfly", "scrapfly_error", `${res.status} ${json.message ?? json.code ?? "no result"}`, started);259 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 };260 const h = Object.fromEntries(Object.entries(r.response_headers ?? {}).map(([k, v]) => [k.toLowerCase(), String(v)]));261 const body = Buffer.from(r.content ?? "", "utf8");262 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 });263 } catch (e) {264 return failDoc(url, 4, "scrapfly", ac.signal.aborted ? "timeout" : "scrapfly_error", (e as Error).message, started);265 } finally { clearTimeout(timer); }266 }267}268269export const fetchers: Record<FetchLevel, Fetcher> = { 1: new DirectFetcher(1), 2: new DirectFetcher(2), 3: new FirecrawlFetcher(), 4: new ScrapflyFetcher() };270271/** Heuristics: a 200 that is really an anti-bot interstitial or an empty JS shell. */272export function looksBlockedOrEmpty(doc: RawDocument): boolean {273 if (doc.status === 403 || doc.status === 429 || doc.status === 503 || doc.status === 401) return true;274 if (doc.status !== 200 || !doc.contentType || !/html/i.test(doc.contentType)) return false;275 const t = doc.text;276 const textLen = t.replace(/<script[\s\S]*?<\/script>/gi, "").replace(/<style[\s\S]*?<\/style>/gi, "").replace(/<[^>]+>/g, " ").replace(/\s+/g, " ").trim().length;277 // Hard interstitial markers (the page IS the challenge) — only trusted when the page carries little real text,278 // because Cloudflare-fronted sites embed the challenge-platform beacon on perfectly normal pages.279 const interstitial = /(<title>[^<]*(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);280 if (interstitial && textLen < 1500) return true;281 return textLen < 400 && /<(div id="(root|app|__next|__nuxt)"|app-root|noscript)/i.test(t);282}283284/** Expected premium cost of one fetch at a level (credits): Firecrawl 1 scrape; Scrapfly ≈1 with ASP only, ≈6 with JS rendering. */285export function expectedCredits(level: FetchLevel, opts: Pick<FetchOptions, "renderJs"> = {}): number {286 if (level === 3) return 1;287 if (level === 4) return opts.renderJs === false ? 1 : 6;288 return 0;289}290291/**292 * Escalating fetch: start at `level`, escalate up to `maxLevel` when blocked/empty. Premium fetchers are293 * skipped when unavailable, and when their expected cost exceeds `creditsLeft` (the run's remaining premium budget —294 * so one call never spends on L3 and L4 together when the run budget cannot cover both). HTTP 429 never escalates:295 * the server asked us to slow down, the document is returned with `meta.retryAfterMs` from `Retry-After` when sent.296 * Returns the best document obtained (never throws).297 */298export async function fetchWithEscalation(url: string, opts: FetchOptions & { maxLevel?: FetchLevel; creditsLeft?: number } = {}): Promise<RawDocument> {299 let level: FetchLevel = opts.level ?? 1;300 const max: FetchLevel = opts.maxLevel ?? 2;301 let creditsLeft = opts.creditsLeft;302 let best: RawDocument | null = null;303 const attempts: string[] = [];304 while (level <= max) {305 const f = fetchers[level];306 if (!f.available()) { attempts.push(`L${level}:unavailable`); level = (level + 1) as FetchLevel; continue; }307 const cost = expectedCredits(level, opts);308 if (creditsLeft !== undefined && cost > 0 && cost > creditsLeft) { attempts.push(`L${level}:over_budget(${cost}>${creditsLeft})`); level = (level + 1) as FetchLevel; continue; }309 const doc = applyRobotsDirectives(await f.fetch(url, { ...opts, level }));310 if (creditsLeft !== undefined) creditsLeft = Math.max(0, creditsLeft - doc.credits);311 doc.meta = { ...(doc.meta ?? {}), attempts: [...attempts, `L${level}:${doc.error?.code ?? doc.status}`] };312 if (doc.notModified) return doc;313 if (doc.status === 429) {314 const retryAfterMs = parseRetryAfter(doc.headers["retry-after"]);315 doc.meta = { ...doc.meta, rateLimited: true, ...(retryAfterMs != null ? { retryAfterMs } : {}) };316 return doc; // never buy our way past a rate limit317 }318 if (!doc.error && !looksBlockedOrEmpty(doc)) return doc;319 // hard failures that escalation cannot fix320 if (doc.error && ["ssrf_blocked", "dns", "too_large"].includes(doc.error.code)) return doc;321 if (!doc.error && doc.status === 404) return doc;322 if (!doc.error && doc.status === 410) return doc;323 best = best ?? doc;324 if (!doc.error && doc.status >= 200 && doc.status < 300) best = doc; // blocked-looking 200 still better than nothing325 attempts.push(`L${level}:${doc.error?.code ?? doc.status}`);326 level = (level + 1) as FetchLevel;327 }328 // the fallback document carries the complete attempt list (including levels skipped as unavailable / over budget)329 if (best) best.meta = { ...(best.meta ?? {}), attempts: [...attempts] };330 return best ?? failDoc(url, max, "direct", "no_fetcher", attempts.join(","), Date.now());331}332