/** * PeeringDB connector — https://www.peeringdb.com/api/ (facilities, campuses, organizations, exchanges, carriers). * * One logical document per object type; the facility document is *composed*: fetch() also pulls the join * tables (ixfac, netfac, carrierfac) and the small lookup tables (campus, carrier, NSP networks) through * ctx.fetch, memoized per run so each upstream URL is requested once per run (PeeringDB throttles repeated * anonymous identical requests > 100 kB to 1/hour). Incremental runs use `?since=` on the primary * objects; join tables are always fetched in full because presence lists are recomputed from scratch. * * Legal: PeeringDB AUP (https://www.peeringdb.com/aup) restricts bulk use to "Internet operational issues" * (which the AUP says includes "Internet research and analysis"). Recorded in the YAML + docs; see docs/connectors/peeringdb.md. */ import type { ConnectorContext, DiscoveredUrl, ExtractedRecord, RawDocument } from "@dci/connectors"; import type { Fact, NormalizedCampus, NormalizedEntity, NormalizedFacility, NormalizedIxp, NormalizedOperator, Provenance } from "@dci/core"; import { validLatLng } from "@dci/core"; import { DatasetConnector, fetchJsonOnce, httpUrl, iso2, num, str, uniqNames } from "./shared.js"; interface PdbList { data: T[]; meta?: Record; message?: string } interface PdbBase { id: number; status: string; created?: string; updated?: string } interface 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 } interface PdbOrg extends PdbBase { name: string; aka?: string; name_long?: string; website?: string; city?: string; country?: string } interface 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 } interface 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 } interface PdbIxFac extends PdbBase { ix_id: number; fac_id: number } interface PdbNetFac extends PdbBase { net_id: number; fac_id: number } interface PdbNet extends PdbBase { org_id: number; name: string; aka?: string; name_long?: string; asn: number; website?: string; info_type?: string; info_types?: string[] } interface PdbCarrier extends PdbBase { org_id: number; org_name?: string; name: string; aka?: string; name_long?: string; website?: string; fac_count?: number } interface PdbCarrierFac extends PdbBase { carrier_id: number; fac_id: number } type Group = "facilities" | "ixps" | "campuses" | "operators" | "carriers"; const GROUPS: Array<{ group: Group; object: string; priority: number }> = [ { group: "facilities", object: "fac", priority: 90 }, { group: "ixps", object: "ix", priority: 80 }, { group: "campuses", object: "campus", priority: 70 }, { group: "operators", object: "org", priority: 60 }, { group: "carriers", object: "net", priority: 50 }, ]; const NET_FIELDS = "id,org_id,name,aka,name_long,asn,website,info_type,info_types,status"; const ORG_FIELDS = "id,name,aka,name_long,website,city,country,status"; const STATE_LAST_RUN = "peeringdb:lastRunAt"; const STATE_LAST_FULL = "peeringdb:lastFullAt"; export class PeeringDbConnector extends DatasetConnector { private api(): string { return String(this.param("apiBase", "https://www.peeringdb.com/api")).replace(/\/+$/, ""); } private url(object: string, query: Record = {}): string { const q = new URLSearchParams({ depth: "0" }); for (const [k, v] of Object.entries(query)) if (v !== undefined && v !== "") q.set(k, String(v)); return `${this.api()}/${object}?${q.toString()}`; } private headers(ctx: ConnectorContext): Record { const key = ctx.env.PEERINGDB_API_KEY; return key ? { authorization: `Api-Key ${key}` } : {}; } private async json(ctx: ConnectorContext, url: string): Promise<{ doc: RawDocument; rows: T[] | null }> { const { doc, json } = await fetchJsonOnce>(ctx, url, { headers: this.headers(ctx), pauseMs: this.param("minDelayMs", 2000), timeoutMs: 120_000, maxBytes: 96 * 1024 * 1024 }); if (!json || !Array.isArray(json.data)) { if (json?.message) ctx.log("warn", `PeeringDB ${url} → ${json.message}`); return { doc, rows: null }; } return { doc, rows: json.data }; } /** One URL per object type. Incremental `since=` after a successful full load; full refresh every `fullRefreshDays`. */ async discover(ctx: ConnectorContext): Promise { const nowMs = Date.now(); const lastRun = await ctx.getState(STATE_LAST_RUN); const lastFull = await ctx.getState(STATE_LAST_FULL); const fullRefreshDays = this.param("fullRefreshDays", 7); const incrementalWanted = this.param("incremental", true); const full = !incrementalWanted || !lastRun || !lastFull || nowMs - Date.parse(lastFull) > fullRefreshDays * 86_400_000; const since = full ? undefined : Math.max(0, Math.floor(Date.parse(lastRun!) / 1000) - 3600); // 1 h overlap ctx.log("info", full ? "PeeringDB full load" : `PeeringDB incremental load since=${since}`, { lastRun, lastFull }); 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()); } 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 } })); } /** Fetch the primary object list, then pre-warm the join/lookup tables it needs (rate-limited, once per run). */ async fetch(ctx: ConnectorContext, u: DiscoveredUrl): Promise { const { doc } = await this.json(ctx, u.url); doc.group = u.group; doc.pageType = "dataset"; doc.meta = { ...(doc.meta ?? {}), ...(u.meta ?? {}) }; if (doc.error || doc.status >= 400) return doc; for (const dep of this.dependencies(u.group as Group)) await this.json(ctx, dep); return doc; } private dependencies(group: Group): string[] { switch (group) { 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" })]; case "ixps": return [this.url("ixfac", { fields: "id,ix_id,fac_id,status" })]; case "operators": return [this.url("fac", { fields: "id,org_id,status" })]; case "carriers": return [this.url("carrier")]; default: return []; } } async extract(ctx: ConnectorContext, doc: RawDocument): Promise { if (doc.error || doc.notModified || doc.status >= 400 || !doc.body.length) return []; // group normally comes from discover(); for manual runs (`--url …`) infer it from the API path const fromUrl = GROUPS.find((g) => new RegExp(`/api/${g.object}(\\?|$)`).test(doc.finalUrl)); const group = (GROUPS.some((g) => g.group === doc.group) ? doc.group : fromUrl?.group) as Group | undefined; const object = String(doc.meta?.object ?? fromUrl?.object ?? ""); let parsed: PdbList | null = null; try { parsed = JSON.parse(doc.text) as PdbList; } catch { ctx.log("warn", `PeeringDB ${doc.finalUrl}: invalid JSON`); return []; } const rows = Array.isArray(parsed?.data) ? parsed.data : []; const live = rows.filter((r) => r.status === "ok"); const deleted = rows.length - live.length; if (deleted) ctx.log("info", `PeeringDB ${object}: skipping ${deleted} non-"ok" records (deleted/pending) — removals are not propagated`, { url: doc.finalUrl }); const url = doc.finalUrl; const rec = (kind: ExtractedRecord["kind"], key: string, data: Record, methods: Record = {}): ExtractedRecord => ({ kind, key, data, methods: { ...methods, _all: `dataset-field:${object}` }, certainty: 1, url, pageType: "dataset" }); if (group === "facilities") { // sequential on purpose: memo hits when fetch() pre-warmed them, otherwise politely spaced upstream calls const deps = this.dependencies(group); const campus = await this.json(ctx, deps[0]!); const ixfac = await this.json(ctx, deps[1]!); const ix = await this.json(ctx, deps[2]!); const netfac = await this.json(ctx, deps[3]!); const net = await this.json(ctx, deps[4]!); const carrier = await this.json(ctx, deps[5]!); const carrierfac = await this.json(ctx, deps[6]!); 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); if (missing.length) { ctx.log("error", `PeeringDB facilities: join tables unavailable (${missing.join(", ")}) — aborting this document to avoid partial presence lists`); return []; } const campusById = new Map(campus.rows!.filter((c) => c.status === "ok").map((c) => [c.id, c] as const)); const ixName = new Map(ix.rows!.filter((x) => x.status === "ok").map((x) => [x.id, x.name] as const)); 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)); const carrierById = new Map(carrier.rows!.filter((c) => c.status === "ok").map((c) => [c.id, c] as const)); const ixAt = new Map>(); 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); } const nspAt = new Map>(); 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); } const carrierAt = new Map>(); 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); } 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" })); } if (group === "ixps") { const ixfac = await this.json(ctx, this.dependencies(group)[0]!); const facs = new Map(); 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); } if (!ixfac.rows) ctx.log("warn", "PeeringDB ixfac unavailable — IXPs emitted without facility links"); 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" })); } if (group === "campuses") return (live as PdbCampus[]).map((c) => rec("campus", `peeringdb:campus:${c.id}`, { ...c })); if (group === "operators") { const fac = await this.json(ctx, this.dependencies(group)[0]!); if (!fac.rows) { ctx.log("error", "PeeringDB org: facility list unavailable — cannot select facility-owning organizations"); return []; } const owning = new Set(fac.rows.filter((f) => f.status === "ok").map((f) => f.org_id)); return (live as PdbOrg[]).filter((o) => owning.has(o.id)).map((o) => rec("operator", `peeringdb:org:${o.id}`, { ...o })); } if (group === "carriers") { 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" })); const carrier = await this.json(ctx, this.dependencies(group)[0]!); for (const c of carrier.rows ?? []) if (c.status === "ok") out.push(rec("operator", `peeringdb:carrier:${c.id}`, { ...c, _kind: "carrier", _object: "carrier" })); return out; } return []; } async normalize(ctx: ConnectorContext, records: ExtractedRecord[]): Promise { const out: NormalizedEntity[] = []; const geoSource = "dataset:peeringdb"; const carrierHotelIx = this.param("carrierHotelMinIx", 3); const carrierHotelNets = this.param("carrierHotelMinNets", 150); for (const r of records) { const x = r.data as Record; const prov = (method: string, extra: Partial = {}): Provenance => ctx.provenance(r.url, { method, extractorVersion: this.parserVersion, ...extra }); if (r.kind === "facility") { const f = x as unknown as PdbFac & { _campus: PdbCampus | null; _ixps: string[]; _carriers: string[] }; const name = str(f.name); if (!name) continue; const lat = num(f.latitude), lng = num(f.longitude); const netCount = num(f.net_count), ixCount = num(f.ix_count), carrierCount = num(f.carrier_count); const facts: Fact[] = []; 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) }); }; fact("networkCount", netCount, "dataset-field:net_count"); fact("ixCount", ixCount, "dataset-field:ix_count"); fact("carrierCount", carrierCount, "dataset-field:carrier_count"); 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"); fact("availableVoltageServices", Array.isArray(f.available_voltage_services) ? f.available_voltage_services.filter((s) => typeof s === "string" && s) : null, "dataset-field:available_voltage_services"); const externalIds: Record = { peeringdb_fac: f.id, peeringdb_org: f.org_id }; if (f.campus_id) externalIds.peeringdb_campus = f.campus_id; if (str(f.clli)) externalIds.clli = str(f.clli)!; const ent: NormalizedFacility = { entityType: "facility", key: r.key, name, aliases: uniqNames([f.aka, f.name_long], [name]), operatorName: str(f.org_name), operatorKey: `peeringdb:org:${f.org_id}`, campusName: f._campus ? str(f._campus.name) : null, campusKey: f.campus_id ? `peeringdb:campus:${f.campus_id}` : null, address: uniqNames([f.address1, f.address2]).join(", ") || null, city: str(f.city), regionName: str(f.state), countryIso2: iso2(f.country), postalCode: str(f.zipcode), geo: validLatLng(lat, lng) ? { lat: lat!, lng: lng!, precision: "exact", source: geoSource } : null, status: f.status === "ok" ? "operational" : null, facilityType: (ixCount ?? 0) >= carrierHotelIx && (netCount ?? 0) >= carrierHotelNets ? "carrier_hotel" : "colocation", website: httpUrl(f.website), carriers: f._carriers, ixps: f._ixps, externalIds, description: str(f.notes), facts, provenance: prov("dataset-field:fac"), }; out.push(ent); } else if (r.kind === "ixp") { const x2 = x as unknown as PdbIx & { _facIds: number[] }; const name = str(x2.name); if (!name) continue; const ent: NormalizedIxp = { 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), facilityKeys: x2._facIds.map((id) => `peeringdb:fac:${id}`), externalIds: { peeringdb_ix: x2.id, peeringdb_org: x2.org_id }, provenance: prov("dataset-field:ix"), }; out.push(ent); } else if (r.kind === "campus") { const c = x as unknown as PdbCampus; const name = str(c.name); if (!name) continue; 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") }; out.push(ent); } else if (r.kind === "operator") { const name = str(x.name); if (!name) continue; const isCarrier = x._kind === "carrier"; const externalIds: Record = {}; if (x._object === "carrier") { externalIds.peeringdb_carrier = x.id as number; externalIds.peeringdb_org = x.org_id as number; } 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; } else externalIds.peeringdb_org = x.id as number; const ent: NormalizedOperator = { 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), hqCountryIso2: isCarrier ? null : iso2(x.country), hqCity: isCarrier ? null : str(x.city), externalIds, provenance: prov(`dataset-field:${x._object === "carrier" ? "carrier" : isCarrier ? "net" : "org"}`), }; out.push(ent); } } return out; } }