import type { Provenance } from "@dci/core"; import { contentFingerprint } from "@dci/core"; import type { ConnectorContext, FetchOptions, RawDocument } from "./types.js"; import type { ConnectorConfig } from "./config.js"; import { fetchWithEscalation, BROWSER_UA } from "./fetchers.js"; import { isAllowedByRobots } from "./robots.js"; import { acquire, configureHost } from "./ratelimit.js"; import type { FetchLevel } from "@dci/core"; /** * Minimal in-memory ConnectorContext for dry runs / tests / connector authoring. * Applies robots, rate limits and escalation exactly like the production runtime, but persists nothing. */ export function createTestContext(cfg: ConnectorConfig, opts: { verbose?: boolean; sourceId?: string } = {}): ConnectorContext & { fetched: RawDocument[]; credits: number } { const state = new Map(); const hashes = new Map(); const host = cfg.domain.replace(/^https?:\/\//, "").split("/")[0]!; configureHost(host, cfg.fetch.rpm, cfg.fetch.concurrency); const fetched: RawDocument[] = []; const ctx = { connectorId: cfg.id, sourceId: opts.sourceId ?? `src_${cfg.id}`, runId: `run_test_${Date.now().toString(36)}`, fetched, credits: 0, dryRun: true, env: process.env, now: () => new Date().toISOString(), log(level: "debug" | "info" | "warn" | "error", msg: string, extra?: Record) { if (level === "debug" && !opts.verbose) return; console.error(`[${cfg.id}] ${level.toUpperCase()} ${msg}${extra ? " " + JSON.stringify(extra) : ""}`); }, async getState(key: string): Promise { return (state.get(key) as T) ?? null; }, async setState(key: string, value: unknown) { state.set(key, value); }, async isKnownUnchanged(url: string, hash: string) { return hashes.get(url) === hash; }, provenance(url: string, extra: Partial = {}): Provenance { const now = new Date().toISOString(); return { sourceId: ctx.sourceId, connectorId: cfg.id, url, firstObserved: now, lastObserved: now, retrievedAt: now, confidence: cfg.kind === "operator" || cfg.kind === "government" || cfg.kind === "filing" || cfg.kind === "utility" || cfg.kind === "cloud_provider" || cfg.kind === "registry" ? "high" : "moderate", extractorVersion: cfg.parserVersion, ...extra }; }, async fetch(url: string, o: FetchOptions = {}): Promise { let robotsUnknown = false; let crawlDelay: number | null = null; if (cfg.fetch.respectRobots) { const r = await isAllowedByRobots(url).catch(() => ({ allowed: true, crawlDelay: null as number | null, sitemaps: [] as string[], unknown: true })); robotsUnknown = r.unknown; crawlDelay = r.crawlDelay; if (!r.allowed) { ctx.log("warn", `robots.txt disallows ${url}`); return { url, finalUrl: url, fetchedAt: new Date().toISOString(), status: 0, contentType: null, body: Buffer.alloc(0), text: "", headers: {}, etag: null, lastModified: null, notModified: false, fetcher: "direct", level: 1, durationMs: 0, credits: 0, error: { code: "robots_disallow", message: "disallowed by robots.txt" } }; } } const release = await acquire(new URL(url).hostname, crawlDelay != null ? Math.min(30_000, Math.max(0, crawlDelay * 1000)) : 0); try { const level = Math.max(o.level ?? 1, cfg.fetch.level) as FetchLevel; // robots.txt unavailable → direct fetches only (same rule as the production context); run budget bounds premium spend const maxLevel = (robotsUnknown ? Math.min(2, cfg.fetch.maxLevel) : cfg.fetch.maxLevel) as FetchLevel; const creditsLeft = Math.max(0, cfg.fetch.maxCreditsPerRun - ctx.credits); const doc = await fetchWithEscalation(url, { ...o, level: Math.min(level, maxLevel) as FetchLevel, maxLevel, creditsLeft, renderJs: o.renderJs ?? cfg.fetch.renderJs, country: o.country ?? cfg.fetch.country, waitForSelector: o.waitForSelector ?? cfg.fetch.waitForSelector, timeoutMs: o.timeoutMs ?? cfg.fetch.timeoutMs, maxBytes: o.maxBytes ?? cfg.fetch.maxBytes, accept: o.accept ?? cfg.fetch.accept, headers: { ...(cfg.fetch.userAgent === "browser" ? { "user-agent": BROWSER_UA } : {}), ...(cfg.fetch.headers ?? {}), ...(o.headers ?? {}) } }); if (robotsUnknown) doc.meta = { ...(doc.meta ?? {}), robotsUnknown: true }; ctx.credits += doc.credits; if (!doc.error && doc.body.length) hashes.set(url, contentFingerprint(doc.text, doc.contentType)); fetched.push(doc); ctx.log("debug", `fetch ${url} → ${doc.error?.code ?? doc.status} via ${doc.fetcher} L${doc.level} ${doc.durationMs}ms ${doc.body.length}b`); return doc; } finally { release(); } }, }; return ctx; }