/** * Shared plumbing for dataset/API connectors (PeeringDB, Overpass, Wikidata, World Bank). * Every network call goes through ctx.fetch so robots policy, per-host rate limits, SSRF checks * and crawl logging apply. Nothing here invents data: helpers only clean, parse and memoize. */ import type { ConnectorConfig, Connector, ConnectorContext, DiscoveredUrl, FetchOptions, RawDocument } from "@dci/connectors"; import type { NormalizedEntity, ValidationReport } from "@dci/core"; import { validateEntities } from "@dci/connectors"; import { cleanText } from "@dci/core"; /** Base class: copies the YAML config into the Connector fields; subclasses implement the pipeline steps. */ export abstract class DatasetConnector implements Connector { id: string; sourceName: string; sourceDomain: string; sourceKind: Connector["sourceKind"]; type: Connector["type"]; parserVersion: string; schedule: Connector["schedule"]; license: string | null; attribution: string | null; priority: number; constructor(readonly cfg: ConnectorConfig) { this.id = cfg.id; this.sourceName = cfg.name; this.sourceDomain = cfg.domain; this.sourceKind = cfg.kind; this.type = cfg.mode; this.parserVersion = cfg.parserVersion; this.schedule = cfg.schedule; this.license = cfg.license ?? null; this.attribution = cfg.attribution ?? null; this.priority = cfg.priority; } abstract discover(ctx: ConnectorContext): Promise; async fetch(ctx: ConnectorContext, u: DiscoveredUrl): Promise { const doc = await ctx.fetch(u.url, { group: u.group, accept: "application/json" }); doc.group = u.group; doc.pageType = u.pageType ?? "dataset"; doc.meta = { ...(doc.meta ?? {}), ...(u.meta ?? {}) }; return doc; } abstract extract(ctx: ConnectorContext, doc: RawDocument): Promise; abstract normalize(ctx: ConnectorContext, records: import("@dci/connectors").ExtractedRecord[]): Promise; async validate(_ctx: ConnectorContext, entities: NormalizedEntity[]): Promise { return validateEntities(entities); } /** typed access to YAML `params` */ param(key: string, fallback: T): T { const v = this.cfg.params?.[key]; return v === undefined || v === null ? fallback : (v as T); } } export function sleep(ms: number): Promise { return new Promise((r) => setTimeout(r, Math.max(0, ms))); } /** Parse a JSON body defensively; returns null (and logs) instead of throwing on bad payloads. */ export function parseJson(ctx: ConnectorContext, doc: RawDocument): T | null { if (doc.error || doc.status >= 400 || !doc.body.length) return null; try { return JSON.parse(doc.text) as T; } catch (e) { ctx.log("warn", `invalid JSON from ${doc.finalUrl}: ${(e as Error).message}`); return null; } } /** * Per-run memo of fetched+parsed JSON documents so a composed dataset document (e.g. PeeringDB facilities * joined with ixfac/netfac) requests each upstream URL at most once per run, even when fetch() and * extract() both need it. Keyed by runId; entries of other runs are dropped to bound memory. */ const memo = new Map>(); export async function fetchJsonOnce(ctx: ConnectorContext, url: string, opts: FetchOptions & { pauseMs?: number } = {}): Promise<{ doc: RawDocument; json: T | null }> { for (const k of memo.keys()) if (k !== ctx.runId) memo.delete(k); let bucket = memo.get(ctx.runId); if (!bucket) { bucket = new Map(); memo.set(ctx.runId, bucket); } const hit = bucket.get(url); if (hit) return { doc: hit.doc, json: hit.json as T | null }; if (opts.pauseMs && bucket.size > 0) await sleep(opts.pauseMs); // polite spacing between successive upstream calls of one run const { pauseMs: _p, ...fo } = opts; const doc = await ctx.fetch(url, { accept: "application/json", ...fo }); const json = parseJson(ctx, doc); if (doc.status === 429 || doc.status === 503) ctx.log("warn", `rate limited by ${new URL(url).hostname} (HTTP ${doc.status}${doc.headers["retry-after"] ? ", retry-after " + doc.headers["retry-after"] : ""}) → ${url}`); bucket.set(url, { doc, json }); return { doc, json }; } export function str(v: unknown): string | null { if (v == null) return null; const s = cleanText(String(v)); return s && s !== "null" && s !== "undefined" ? s : null; } export function num(v: unknown): number | null { if (v == null || v === "") return null; const n = typeof v === "number" ? v : Number(String(v).trim()); return Number.isFinite(n) ? n : null; } export function iso2(v: unknown): string | null { const s = str(v)?.toUpperCase() ?? null; return s && /^[A-Z]{2}$/.test(s) ? s : null; } /** Keep only http(s) URLs; PeeringDB/OSM sometimes hold bare hostnames or junk. */ export function httpUrl(v: unknown): string | null { const s = str(v); if (!s) return null; const candidate = /^https?:\/\//i.test(s) ? s : /^[a-z0-9.-]+\.[a-z]{2,}(\/|$)/i.test(s) ? `https://${s}` : null; if (!candidate) return null; try { const u = new URL(candidate); return u.protocol === "http:" || u.protocol === "https:" ? u.toString() : null; } catch { return null; } } /** Distinct, non-empty, case-insensitively deduplicated names, excluding `exclude` values. */ export function uniqNames(values: Array, exclude: Array = []): string[] { const seen = new Set(exclude.filter(Boolean).map((s) => s!.toLowerCase())); const out: string[] = []; for (const v of values) { const s = str(v); if (!s) continue; const k = s.toLowerCase(); if (seen.has(k)) continue; seen.add(k); out.push(s); } return out; }