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%
18.4 KB · 223 lines typescript
Raw Blame History
1/**2 * PeeringDB connector — https://www.peeringdb.com/api/ (facilities, campuses, organizations, exchanges, carriers).3 *4 * One logical document per object type; the facility document is *composed*: fetch() also pulls the join5 * tables (ixfac, netfac, carrierfac) and the small lookup tables (campus, carrier, NSP networks) through6 * ctx.fetch, memoized per run so each upstream URL is requested once per run (PeeringDB throttles repeated7 * anonymous identical requests > 100 kB to 1/hour). Incremental runs use `?since=<unix ts>` on the primary8 * objects; join tables are always fetched in full because presence lists are recomputed from scratch.9 *10 * Legal: PeeringDB AUP (https://www.peeringdb.com/aup) restricts bulk use to "Internet operational issues"11 * (which the AUP says includes "Internet research and analysis"). Recorded in the YAML + docs; see docs/connectors/peeringdb.md.12 */13import type { ConnectorContext, DiscoveredUrl, ExtractedRecord, RawDocument } from "@dci/connectors";14import type { Fact, NormalizedCampus, NormalizedEntity, NormalizedFacility, NormalizedIxp, NormalizedOperator, Provenance } from "@dci/core";15import { validLatLng } from "@dci/core";16import { DatasetConnector, fetchJsonOnce, httpUrl, iso2, num, str, uniqNames } from "./shared.js";1718interface PdbList<T> { data: T[]; meta?: Record<string, unknown>; message?: string }19interface PdbBase { id: number; status: string; created?: string; updated?: string }20interface PdbFac extends PdbBase { org_id: number; org_name?: string; campus_id?: number | null; name: string; aka?: string; name_long?: string; website?: string; clli?: string; notes?: string; net_count?: number; ix_count?: number; carrier_count?: number; available_voltage_services?: string[]; diverse_serving_substations?: boolean | null; property?: string; region_continent?: string; address1?: string; address2?: string; city?: string; country?: string; state?: string; zipcode?: string; latitude?: number | null; longitude?: number | null }21interface PdbOrg extends PdbBase { name: string; aka?: string; name_long?: string; website?: string; city?: string; country?: string }22interface PdbCampus extends PdbBase { org_id: number; org_name?: string; name: string; name_long?: string; aka?: string; website?: string; country?: string; city?: string; zipcode?: string; state?: string; notes?: string }23interface PdbIx extends PdbBase { org_id: number; name: string; aka?: string; name_long?: string; city?: string; country?: string; region_continent?: string; website?: string; net_count?: number; fac_count?: number; notes?: string }24interface PdbIxFac extends PdbBase { ix_id: number; fac_id: number }25interface PdbNetFac extends PdbBase { net_id: number; fac_id: number }26interface PdbNet extends PdbBase { org_id: number; name: string; aka?: string; name_long?: string; asn: number; website?: string; info_type?: string; info_types?: string[] }27interface PdbCarrier extends PdbBase { org_id: number; org_name?: string; name: string; aka?: string; name_long?: string; website?: string; fac_count?: number }28interface PdbCarrierFac extends PdbBase { carrier_id: number; fac_id: number }2930type Group = "facilities" | "ixps" | "campuses" | "operators" | "carriers";31const GROUPS: Array<{ group: Group; object: string; priority: number }> = [32  { group: "facilities", object: "fac", priority: 90 },33  { group: "ixps", object: "ix", priority: 80 },34  { group: "campuses", object: "campus", priority: 70 },35  { group: "operators", object: "org", priority: 60 },36  { group: "carriers", object: "net", priority: 50 },37];3839const NET_FIELDS = "id,org_id,name,aka,name_long,asn,website,info_type,info_types,status";40const ORG_FIELDS = "id,name,aka,name_long,website,city,country,status";41const STATE_LAST_RUN = "peeringdb:lastRunAt";42const STATE_LAST_FULL = "peeringdb:lastFullAt";4344export class PeeringDbConnector extends DatasetConnector {45  private api(): string { return String(this.param("apiBase", "https://www.peeringdb.com/api")).replace(/\/+$/, ""); }46  private url(object: string, query: Record<string, string | number | undefined> = {}): string {47    const q = new URLSearchParams({ depth: "0" });48    for (const [k, v] of Object.entries(query)) if (v !== undefined && v !== "") q.set(k, String(v));49    return `${this.api()}/${object}?${q.toString()}`;50  }51  private headers(ctx: ConnectorContext): Record<string, string> {52    const key = ctx.env.PEERINGDB_API_KEY;53    return key ? { authorization: `Api-Key ${key}` } : {};54  }55  private async json<T>(ctx: ConnectorContext, url: string): Promise<{ doc: RawDocument; rows: T[] | null }> {56    const { doc, json } = await fetchJsonOnce<PdbList<T>>(ctx, url, { headers: this.headers(ctx), pauseMs: this.param<number>("minDelayMs", 2000), timeoutMs: 120_000, maxBytes: 96 * 1024 * 1024 });57    if (!json || !Array.isArray(json.data)) {58      if (json?.message) ctx.log("warn", `PeeringDB ${url} → ${json.message}`);59      return { doc, rows: null };60    }61    return { doc, rows: json.data };62  }6364  /** One URL per object type. Incremental `since=` after a successful full load; full refresh every `fullRefreshDays`. */65  async discover(ctx: ConnectorContext): Promise<DiscoveredUrl[]> {66    const nowMs = Date.now();67    const lastRun = await ctx.getState<string>(STATE_LAST_RUN);68    const lastFull = await ctx.getState<string>(STATE_LAST_FULL);69    const fullRefreshDays = this.param<number>("fullRefreshDays", 7);70    const incrementalWanted = this.param<boolean>("incremental", true);71    const full = !incrementalWanted || !lastRun || !lastFull || nowMs - Date.parse(lastFull) > fullRefreshDays * 86_400_000;72    const since = full ? undefined : Math.max(0, Math.floor(Date.parse(lastRun!) / 1000) - 3600); // 1 h overlap73    ctx.log("info", full ? "PeeringDB full load" : `PeeringDB incremental load since=${since}`, { lastRun, lastFull });74    if (!ctx.dryRun) { await ctx.setState(STATE_LAST_RUN, new Date(nowMs).toISOString()); if (full) await ctx.setState(STATE_LAST_FULL, new Date(nowMs).toISOString()); }75    return GROUPS.map((g) => ({ url: this.url(g.object, g.object === "net" ? { info_type: "NSP", fields: NET_FIELDS, since } : g.object === "org" ? { fields: ORG_FIELDS, since } : { since }), group: g.group, pageType: "dataset", priority: g.priority, meta: { object: g.object, mode: full ? "full" : "incremental", since: since ?? null } }));76  }7778  /** Fetch the primary object list, then pre-warm the join/lookup tables it needs (rate-limited, once per run). */79  async fetch(ctx: ConnectorContext, u: DiscoveredUrl): Promise<RawDocument> {80    const { doc } = await this.json<PdbBase>(ctx, u.url);81    doc.group = u.group; doc.pageType = "dataset"; doc.meta = { ...(doc.meta ?? {}), ...(u.meta ?? {}) };82    if (doc.error || doc.status >= 400) return doc;83    for (const dep of this.dependencies(u.group as Group)) await this.json(ctx, dep);84    return doc;85  }8687  private dependencies(group: Group): string[] {88    switch (group) {89      case "facilities": return [this.url("campus"), this.url("ixfac", { fields: "id,ix_id,fac_id,status" }), this.url("ix", { fields: "id,name,status" }), this.url("netfac", { fields: "id,net_id,fac_id,status" }), this.url("net", { info_type: "NSP", fields: NET_FIELDS }), this.url("carrier"), this.url("carrierfac", { fields: "id,carrier_id,fac_id,status" })];90      case "ixps": return [this.url("ixfac", { fields: "id,ix_id,fac_id,status" })];91      case "operators": return [this.url("fac", { fields: "id,org_id,status" })];92      case "carriers": return [this.url("carrier")];93      default: return [];94    }95  }9697  async extract(ctx: ConnectorContext, doc: RawDocument): Promise<ExtractedRecord[]> {98    if (doc.error || doc.notModified || doc.status >= 400 || !doc.body.length) return [];99    // group normally comes from discover(); for manual runs (`--url …`) infer it from the API path100    const fromUrl = GROUPS.find((g) => new RegExp(`/api/${g.object}(\\?|$)`).test(doc.finalUrl));101    const group = (GROUPS.some((g) => g.group === doc.group) ? doc.group : fromUrl?.group) as Group | undefined;102    const object = String(doc.meta?.object ?? fromUrl?.object ?? "");103    let parsed: PdbList<PdbBase> | null = null;104    try { parsed = JSON.parse(doc.text) as PdbList<PdbBase>; } catch { ctx.log("warn", `PeeringDB ${doc.finalUrl}: invalid JSON`); return []; }105    const rows = Array.isArray(parsed?.data) ? parsed.data : [];106    const live = rows.filter((r) => r.status === "ok");107    const deleted = rows.length - live.length;108    if (deleted) ctx.log("info", `PeeringDB ${object}: skipping ${deleted} non-"ok" records (deleted/pending) — removals are not propagated`, { url: doc.finalUrl });109    const url = doc.finalUrl;110    const rec = (kind: ExtractedRecord["kind"], key: string, data: Record<string, unknown>, methods: Record<string, string> = {}): ExtractedRecord => ({ kind, key, data, methods: { ...methods, _all: `dataset-field:${object}` }, certainty: 1, url, pageType: "dataset" });111112    if (group === "facilities") {113      // sequential on purpose: memo hits when fetch() pre-warmed them, otherwise politely spaced upstream calls114      const deps = this.dependencies(group);115      const campus = await this.json<PdbCampus>(ctx, deps[0]!);116      const ixfac = await this.json<PdbIxFac>(ctx, deps[1]!);117      const ix = await this.json<PdbIx>(ctx, deps[2]!);118      const netfac = await this.json<PdbNetFac>(ctx, deps[3]!);119      const net = await this.json<PdbNet>(ctx, deps[4]!);120      const carrier = await this.json<PdbCarrier>(ctx, deps[5]!);121      const carrierfac = await this.json<PdbCarrierFac>(ctx, deps[6]!);122      const missing = [["campus", campus], ["ixfac", ixfac], ["ix", ix], ["netfac", netfac], ["net", net], ["carrier", carrier], ["carrierfac", carrierfac]].filter(([, r]) => !(r as { rows: unknown[] | null }).rows).map(([n]) => n);123      if (missing.length) { ctx.log("error", `PeeringDB facilities: join tables unavailable (${missing.join(", ")}) — aborting this document to avoid partial presence lists`); return []; }124      const campusById = new Map(campus.rows!.filter((c) => c.status === "ok").map((c) => [c.id, c] as const));125      const ixName = new Map(ix.rows!.filter((x) => x.status === "ok").map((x) => [x.id, x.name] as const));126      const nspById = new Map(net.rows!.filter((n) => n.status === "ok" && (n.info_type === "NSP" || n.info_types?.includes("NSP"))).map((n) => [n.id, n] as const));127      const carrierById = new Map(carrier.rows!.filter((c) => c.status === "ok").map((c) => [c.id, c] as const));128      const ixAt = new Map<number, Set<string>>(); for (const r of ixfac.rows!) { if (r.status !== "ok") continue; const n = ixName.get(r.ix_id); if (!n) continue; (ixAt.get(r.fac_id) ?? ixAt.set(r.fac_id, new Set()).get(r.fac_id)!).add(n); }129      const nspAt = new Map<number, Set<string>>(); for (const r of netfac.rows!) { if (r.status !== "ok") continue; const n = nspById.get(r.net_id); if (!n) continue; (nspAt.get(r.fac_id) ?? nspAt.set(r.fac_id, new Set()).get(r.fac_id)!).add(n.name); }130      const carrierAt = new Map<number, Set<string>>(); for (const r of carrierfac.rows!) { if (r.status !== "ok") continue; const c = carrierById.get(r.carrier_id); if (!c) continue; (carrierAt.get(r.fac_id) ?? carrierAt.set(r.fac_id, new Set()).get(r.fac_id)!).add(c.name); }131      return (live as PdbFac[]).map((f) => rec("facility", `peeringdb:fac:${f.id}`, { ...f, _campus: f.campus_id ? campusById.get(f.campus_id) ?? null : null, _ixps: [...(ixAt.get(f.id) ?? [])].sort(), _carriers: [...new Set([...(carrierAt.get(f.id) ?? []), ...(nspAt.get(f.id) ?? [])])].sort() }, { geo: "dataset-field:latitude,longitude", carriers: "dataset-join:carrierfac,netfac(NSP)", ixps: "dataset-join:ixfac" }));132    }133    if (group === "ixps") {134      const ixfac = await this.json<PdbIxFac>(ctx, this.dependencies(group)[0]!);135      const facs = new Map<number, number[]>();136      for (const r of ixfac.rows ?? []) { if (r.status !== "ok") continue; (facs.get(r.ix_id) ?? facs.set(r.ix_id, []).get(r.ix_id)!).push(r.fac_id); }137      if (!ixfac.rows) ctx.log("warn", "PeeringDB ixfac unavailable — IXPs emitted without facility links");138      return (live as PdbIx[]).map((x) => rec("ixp", `peeringdb:ix:${x.id}`, { ...x, _facIds: (facs.get(x.id) ?? []).sort((a, b) => a - b) }, { facilityKeys: "dataset-join:ixfac" }));139    }140    if (group === "campuses") return (live as PdbCampus[]).map((c) => rec("campus", `peeringdb:campus:${c.id}`, { ...c }));141    if (group === "operators") {142      const fac = await this.json<PdbFac>(ctx, this.dependencies(group)[0]!);143      if (!fac.rows) { ctx.log("error", "PeeringDB org: facility list unavailable — cannot select facility-owning organizations"); return []; }144      const owning = new Set(fac.rows.filter((f) => f.status === "ok").map((f) => f.org_id));145      return (live as PdbOrg[]).filter((o) => owning.has(o.id)).map((o) => rec("operator", `peeringdb:org:${o.id}`, { ...o }));146    }147    if (group === "carriers") {148      const out = (live as PdbNet[]).filter((n) => n.info_type === "NSP" || n.info_types?.includes("NSP")).map((n) => rec("operator", `peeringdb:net:${n.id}`, { ...n, _kind: "carrier" }));149      const carrier = await this.json<PdbCarrier>(ctx, this.dependencies(group)[0]!);150      for (const c of carrier.rows ?? []) if (c.status === "ok") out.push(rec("operator", `peeringdb:carrier:${c.id}`, { ...c, _kind: "carrier", _object: "carrier" }));151      return out;152    }153    return [];154  }155156  async normalize(ctx: ConnectorContext, records: ExtractedRecord[]): Promise<NormalizedEntity[]> {157    const out: NormalizedEntity[] = [];158    const geoSource = "dataset:peeringdb";159    const carrierHotelIx = this.param<number>("carrierHotelMinIx", 3);160    const carrierHotelNets = this.param<number>("carrierHotelMinNets", 150);161    for (const r of records) {162      const x = r.data as Record<string, unknown>;163      const prov = (method: string, extra: Partial<Provenance> = {}): Provenance => ctx.provenance(r.url, { method, extractorVersion: this.parserVersion, ...extra });164      if (r.kind === "facility") {165        const f = x as unknown as PdbFac & { _campus: PdbCampus | null; _ixps: string[]; _carriers: string[] };166        const name = str(f.name); if (!name) continue;167        const lat = num(f.latitude), lng = num(f.longitude);168        const netCount = num(f.net_count), ixCount = num(f.ix_count), carrierCount = num(f.carrier_count);169        const facts: Fact[] = [];170        const fact = (field: string, value: unknown, method: string) => { if (value !== null && value !== undefined && value !== "" && !(Array.isArray(value) && !value.length)) facts.push({ field, value, provenance: prov(method) }); };171        fact("networkCount", netCount, "dataset-field:net_count"); fact("ixCount", ixCount, "dataset-field:ix_count"); fact("carrierCount", carrierCount, "dataset-field:carrier_count");172        fact("property", str(f.property), "dataset-field:property"); fact("diverseServingSubstations", typeof f.diverse_serving_substations === "boolean" ? f.diverse_serving_substations : null, "dataset-field:diverse_serving_substations");173        fact("availableVoltageServices", Array.isArray(f.available_voltage_services) ? f.available_voltage_services.filter((s) => typeof s === "string" && s) : null, "dataset-field:available_voltage_services");174        const externalIds: Record<string, string | number> = { peeringdb_fac: f.id, peeringdb_org: f.org_id };175        if (f.campus_id) externalIds.peeringdb_campus = f.campus_id;176        if (str(f.clli)) externalIds.clli = str(f.clli)!;177        const ent: NormalizedFacility = {178          entityType: "facility", key: r.key, name,179          aliases: uniqNames([f.aka, f.name_long], [name]),180          operatorName: str(f.org_name), operatorKey: `peeringdb:org:${f.org_id}`,181          campusName: f._campus ? str(f._campus.name) : null, campusKey: f.campus_id ? `peeringdb:campus:${f.campus_id}` : null,182          address: uniqNames([f.address1, f.address2]).join(", ") || null,183          city: str(f.city), regionName: str(f.state), countryIso2: iso2(f.country), postalCode: str(f.zipcode),184          geo: validLatLng(lat, lng) ? { lat: lat!, lng: lng!, precision: "exact", source: geoSource } : null,185          status: f.status === "ok" ? "operational" : null,186          facilityType: (ixCount ?? 0) >= carrierHotelIx && (netCount ?? 0) >= carrierHotelNets ? "carrier_hotel" : "colocation",187          website: httpUrl(f.website),188          carriers: f._carriers, ixps: f._ixps,189          externalIds, description: str(f.notes), facts,190          provenance: prov("dataset-field:fac"),191        };192        out.push(ent);193      } else if (r.kind === "ixp") {194        const x2 = x as unknown as PdbIx & { _facIds: number[] };195        const name = str(x2.name); if (!name) continue;196        const ent: NormalizedIxp = {197          entityType: "ixp", key: r.key, name, nameLong: str(x2.name_long), city: str(x2.city), countryIso2: iso2(x2.country), regionContinent: str(x2.region_continent), website: httpUrl(x2.website), networkCount: num(x2.net_count),198          facilityKeys: x2._facIds.map((id) => `peeringdb:fac:${id}`), externalIds: { peeringdb_ix: x2.id, peeringdb_org: x2.org_id }, provenance: prov("dataset-field:ix"),199        };200        out.push(ent);201      } else if (r.kind === "campus") {202        const c = x as unknown as PdbCampus;203        const name = str(c.name); if (!name) continue;204        const ent: NormalizedCampus = { entityType: "campus", key: r.key, name, operatorName: str(c.org_name), city: str(c.city), countryIso2: iso2(c.country), geo: null, externalIds: { peeringdb_campus: c.id, peeringdb_org: c.org_id }, provenance: prov("dataset-field:campus") };205        out.push(ent);206      } else if (r.kind === "operator") {207        const name = str(x.name); if (!name) continue;208        const isCarrier = x._kind === "carrier";209        const externalIds: Record<string, string | number> = {};210        if (x._object === "carrier") { externalIds.peeringdb_carrier = x.id as number; externalIds.peeringdb_org = x.org_id as number; }211        else if (isCarrier) { externalIds.peeringdb_net = x.id as number; if (num(x.asn) != null) externalIds.asn = num(x.asn)!; externalIds.peeringdb_org = x.org_id as number; }212        else externalIds.peeringdb_org = x.id as number;213        const ent: NormalizedOperator = {214          entityType: "operator", key: r.key, name, aliases: uniqNames([x.aka as string, x.name_long as string], [name]), kind: isCarrier ? "carrier" : null, website: httpUrl(x.website),215          hqCountryIso2: isCarrier ? null : iso2(x.country), hqCity: isCarrier ? null : str(x.city), externalIds, provenance: prov(`dataset-field:${x._object === "carrier" ? "carrier" : isCarrier ? "net" : "org"}`),216        };217        out.push(ent);218      }219    }220    return out;221  }222}223