/** IXP reconciliation (key → external ids → normalized name + country) and facility_ixps links. */ import { sql } from "@dci/db"; import { newId, normalizeName, type NormalizedIxp } from "@dci/core"; import { addRef, bump, safeCountry, uniqueSlug, type IngestContext, type Tx } from "./common.js"; import { assignMetro } from "./metros.js"; import { writeProvenance } from "./provenance.js"; import { facilityIdForKey } from "./keys.js"; export interface IxpInput { name: string; key?: string | null; nameLong?: string | null; city?: string | null; countryIso2?: string | null; regionContinent?: string | null; website?: string | null; networkCount?: number | null; externalIds?: Record; } export async function resolveIxp(tx: Tx, ctx: IngestContext, i: IxpInput): Promise<{ id: string; created: boolean }> { const name = i.name.trim(); if (!name) throw new Error("ixp name is empty"); const country = await safeCountry(tx, ctx, i.countryIso2); let id: string | null = null; if (i.key) { const r = await tx.execute(sql`select entity_id from entity_keys where key = ${i.key} and entity_type = 'ixp'`); if (r[0]) id = String(r[0].entity_id); } if (!id && i.externalIds) { for (const [k, v] of Object.entries(i.externalIds)) { if (v == null || v === "" || !/^(peeringdb_ix|peeringdb|pdb_ix|wikidata|ixpdb|pch_id|euro_ix)$/.test(k)) continue; const r = await tx.execute(sql`select id from ixps where external_ids @> ${JSON.stringify({ [k]: v })}::jsonb limit 1`); if (r[0]) { id = String(r[0].id); break; } } } if (!id) { const norm = normalizeName(name); const rows = await tx.execute(sql`select id, name, name_long from ixps where (${country}::text is null or country_iso2 is null or country_iso2 = ${country}) limit 2000`); const hit = rows.find((r) => normalizeName(String(r.name)) === norm || (r.name_long && normalizeName(String(r.name_long)) === norm) || (i.nameLong && normalizeName(i.nameLong) === normalizeName(String(r.name)))); if (hit) id = String(hit.id); } const metro = await assignMetro(tx, { city: i.city, countryIso2: country }); let created = false; if (!id) { id = newId("ixp"); const slug = await uniqueSlug(tx, "ixps", name, i.city ?? country); await tx.execute(sql`insert into ixps (id, slug, name, name_long, city, country_iso2, metro_id, region_continent, website, network_count, external_ids) values (${id}, ${slug}, ${name}, ${i.nameLong ?? null}, ${i.city ?? null}, ${metro.countryIso2}, ${metro.metroId}, ${i.regionContinent ?? null}, ${i.website ?? null}, ${i.networkCount ?? null}, ${JSON.stringify(i.externalIds ?? {})}::jsonb)`); created = true; } else { await tx.execute(sql`update ixps set name_long = coalesce(${i.nameLong ?? null}, name_long), city = coalesce(city, ${i.city ?? null}), country_iso2 = coalesce(country_iso2, ${metro.countryIso2}), metro_id = coalesce(metro_id, ${metro.metroId}), region_continent = coalesce(region_continent, ${i.regionContinent ?? null}), website = coalesce(${i.website ?? null}, website), network_count = coalesce(${i.networkCount ?? null}, network_count), external_ids = external_ids || ${JSON.stringify(i.externalIds ?? {})}::jsonb, updated_at = now() where id = ${id}`); } if (i.key) await tx.execute(sql`insert into entity_keys (key, entity_type, entity_id, connector_id) values (${i.key}, 'ixp', ${id}, ${ctx.run.connectorId}) on conflict (key) do update set entity_id = excluded.entity_id`); addRef(ctx, "ixp", id); return { id, created }; } export async function linkFacilityIxp(tx: Tx, ctx: IngestContext, facilityId: string, ixpId: string): Promise { await tx.execute(sql`insert into facility_ixps (facility_id, ixp_id, source_id) values (${facilityId}, ${ixpId}, ${ctx.run.sourceId}) on conflict do nothing`); } export async function refreshIxpCount(tx: Tx, facilityId: string): Promise { await tx.execute(sql`update facilities f set ixp_count = (select count(*) from facility_ixps x where x.facility_id = f.id) where f.id = ${facilityId}`); } /** Attach IXPs listed on a facility page by name (creates the IXP when unknown, country from the facility). */ export async function attachIxpsByName(tx: Tx, ctx: IngestContext, facilityId: string, names: string[], countryIso2: string | null): Promise { let n = 0; for (const raw of names) { const name = raw?.trim(); if (!name || name.length < 2 || name.length > 120) continue; try { const { id } = await resolveIxp(tx, ctx, { name, countryIso2 }); await linkFacilityIxp(tx, ctx, facilityId, id); n++; } catch (e) { console.warn(`[ingest] ixp "${name}" skipped: ${(e as Error).message}`); } } if (n) await refreshIxpCount(tx, facilityId); return n; } export async function ingestIxp(tx: Tx, ctx: IngestContext, ix: NormalizedIxp): Promise { const { id, created } = await resolveIxp(tx, ctx, { name: ix.name, key: ix.key, nameLong: ix.nameLong, city: ix.city, countryIso2: ix.countryIso2, regionContinent: ix.regionContinent, website: ix.website, networkCount: ix.networkCount, externalIds: ix.externalIds, }); if (created) ctx.stats.created++; else ctx.stats.updated++; bump(ctx, "ixp"); const touched = new Set(); for (const fk of ix.facilityKeys ?? []) { const fid = await facilityIdForKey(tx, ctx, fk); if (!fid) continue; await linkFacilityIxp(tx, ctx, fid, id); touched.add(fid); } for (const fid of touched) await refreshIxpCount(tx, fid); await writeProvenance( tx, ctx, "ixp", id, [ { field: "name", value: ix.name }, { field: "nameLong", value: ix.nameLong ?? null }, { field: "city", value: ix.city ?? null }, { field: "countryIso2", value: ix.countryIso2 ?? null }, { field: "website", value: ix.website ?? null }, { field: "networkCount", value: ix.networkCount ?? null }, { field: "facilityKeys", value: ix.facilityKeys?.length ? ix.facilityKeys : null }, ].map((f) => ({ ...f, provenance: ix.provenance })), ix.key, ); }