SPB Git forge
38commits 1branches 0releases
338.7 MBsize
maindefault branch
4 h agolast push
HTML 53.9% TypeScript 44.5% JavaScript 0.6% SQL 0.5%
5.6 KB · 102 lines typescript
Raw Blame History
1/**2 * Shared plumbing for dataset/API connectors (PeeringDB, Overpass, Wikidata, World Bank).3 * Every network call goes through ctx.fetch so robots policy, per-host rate limits, SSRF checks4 * and crawl logging apply. Nothing here invents data: helpers only clean, parse and memoize.5 */6import type { ConnectorConfig, Connector, ConnectorContext, DiscoveredUrl, FetchOptions, RawDocument } from "@dci/connectors";7import type { NormalizedEntity, ValidationReport } from "@dci/core";8import { validateEntities } from "@dci/connectors";9import { cleanText } from "@dci/core";1011/** Base class: copies the YAML config into the Connector fields; subclasses implement the pipeline steps. */12export abstract class DatasetConnector implements Connector {13  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;14  constructor(readonly cfg: ConnectorConfig) {15    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;16  }17  abstract discover(ctx: ConnectorContext): Promise<DiscoveredUrl[]>;18  async fetch(ctx: ConnectorContext, u: DiscoveredUrl): Promise<RawDocument> {19    const doc = await ctx.fetch(u.url, { group: u.group, accept: "application/json" });20    doc.group = u.group; doc.pageType = u.pageType ?? "dataset"; doc.meta = { ...(doc.meta ?? {}), ...(u.meta ?? {}) };21    return doc;22  }23  abstract extract(ctx: ConnectorContext, doc: RawDocument): Promise<import("@dci/connectors").ExtractedRecord[]>;24  abstract normalize(ctx: ConnectorContext, records: import("@dci/connectors").ExtractedRecord[]): Promise<NormalizedEntity[]>;25  async validate(_ctx: ConnectorContext, entities: NormalizedEntity[]): Promise<ValidationReport> { return validateEntities(entities); }2627  /** typed access to YAML `params` */28  param<T>(key: string, fallback: T): T {29    const v = this.cfg.params?.[key];30    return v === undefined || v === null ? fallback : (v as T);31  }32}3334export function sleep(ms: number): Promise<void> { return new Promise((r) => setTimeout(r, Math.max(0, ms))); }3536/** Parse a JSON body defensively; returns null (and logs) instead of throwing on bad payloads. */37export function parseJson<T = unknown>(ctx: ConnectorContext, doc: RawDocument): T | null {38  if (doc.error || doc.status >= 400 || !doc.body.length) return null;39  try { return JSON.parse(doc.text) as T; } catch (e) { ctx.log("warn", `invalid JSON from ${doc.finalUrl}: ${(e as Error).message}`); return null; }40}4142/**43 * Per-run memo of fetched+parsed JSON documents so a composed dataset document (e.g. PeeringDB facilities44 * joined with ixfac/netfac) requests each upstream URL at most once per run, even when fetch() and45 * extract() both need it. Keyed by runId; entries of other runs are dropped to bound memory.46 */47const memo = new Map<string, Map<string, { doc: RawDocument; json: unknown }>>();48export async function fetchJsonOnce<T = unknown>(ctx: ConnectorContext, url: string, opts: FetchOptions & { pauseMs?: number } = {}): Promise<{ doc: RawDocument; json: T | null }> {49  for (const k of memo.keys()) if (k !== ctx.runId) memo.delete(k);50  let bucket = memo.get(ctx.runId);51  if (!bucket) { bucket = new Map(); memo.set(ctx.runId, bucket); }52  const hit = bucket.get(url);53  if (hit) return { doc: hit.doc, json: hit.json as T | null };54  if (opts.pauseMs && bucket.size > 0) await sleep(opts.pauseMs); // polite spacing between successive upstream calls of one run55  const { pauseMs: _p, ...fo } = opts;56  const doc = await ctx.fetch(url, { accept: "application/json", ...fo });57  const json = parseJson<T>(ctx, doc);58  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}`);59  bucket.set(url, { doc, json });60  return { doc, json };61}6263export function str(v: unknown): string | null {64  if (v == null) return null;65  const s = cleanText(String(v));66  return s && s !== "null" && s !== "undefined" ? s : null;67}6869export function num(v: unknown): number | null {70  if (v == null || v === "") return null;71  const n = typeof v === "number" ? v : Number(String(v).trim());72  return Number.isFinite(n) ? n : null;73}7475export function iso2(v: unknown): string | null {76  const s = str(v)?.toUpperCase() ?? null;77  return s && /^[A-Z]{2}$/.test(s) ? s : null;78}7980/** Keep only http(s) URLs; PeeringDB/OSM sometimes hold bare hostnames or junk. */81export function httpUrl(v: unknown): string | null {82  const s = str(v);83  if (!s) return null;84  const candidate = /^https?:\/\//i.test(s) ? s : /^[a-z0-9.-]+\.[a-z]{2,}(\/|$)/i.test(s) ? `https://${s}` : null;85  if (!candidate) return null;86  try { const u = new URL(candidate); return u.protocol === "http:" || u.protocol === "https:" ? u.toString() : null; } catch { return null; }87}8889/** Distinct, non-empty, case-insensitively deduplicated names, excluding `exclude` values. */90export function uniqNames(values: Array<string | null | undefined>, exclude: Array<string | null | undefined> = []): string[] {91  const seen = new Set(exclude.filter(Boolean).map((s) => s!.toLowerCase()));92  const out: string[] = [];93  for (const v of values) {94    const s = str(v);95    if (!s) continue;96    const k = s.toLowerCase();97    if (seen.has(k)) continue;98    seen.add(k); out.push(s);99  }100  return out;101}102