SPB Git forge
38commits 1branches 0releases
338.7 MBsize
maindefault branch
3 h agolast push
HTML 53.9% TypeScript 44.5% JavaScript 0.6% SQL 0.5%
22.2 KB · 332 lines typescript
Raw Blame History
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