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%
6.1 KB · 130 lines typescript
Raw Blame History
1/** IXP reconciliation (key → external ids → normalized name + country) and facility_ixps links. */2import { sql } from "@dci/db";3import { newId, normalizeName, type NormalizedIxp } from "@dci/core";4import { addRef, bump, safeCountry, uniqueSlug, type IngestContext, type Tx } from "./common.js";5import { assignMetro } from "./metros.js";6import { writeProvenance } from "./provenance.js";7import { facilityIdForKey } from "./keys.js";89export interface IxpInput {10  name: string;11  key?: string | null;12  nameLong?: string | null;13  city?: string | null;14  countryIso2?: string | null;15  regionContinent?: string | null;16  website?: string | null;17  networkCount?: number | null;18  externalIds?: Record<string, string | number>;19}2021export async function resolveIxp(tx: Tx, ctx: IngestContext, i: IxpInput): Promise<{ id: string; created: boolean }> {22  const name = i.name.trim();23  if (!name) throw new Error("ixp name is empty");24  const country = await safeCountry(tx, ctx, i.countryIso2);25  let id: string | null = null;26  if (i.key) {27    const r = await tx.execute(sql`select entity_id from entity_keys where key = ${i.key} and entity_type = 'ixp'`);28    if (r[0]) id = String(r[0].entity_id);29  }30  if (!id && i.externalIds) {31    for (const [k, v] of Object.entries(i.externalIds)) {32      if (v == null || v === "" || !/^(peeringdb_ix|peeringdb|pdb_ix|wikidata|ixpdb|pch_id|euro_ix)$/.test(k)) continue;33      const r = await tx.execute(sql`select id from ixps where external_ids @> ${JSON.stringify({ [k]: v })}::jsonb limit 1`);34      if (r[0]) {35        id = String(r[0].id);36        break;37      }38    }39  }40  if (!id) {41    const norm = normalizeName(name);42    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`);43    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))));44    if (hit) id = String(hit.id);45  }46  const metro = await assignMetro(tx, { city: i.city, countryIso2: country });47  let created = false;48  if (!id) {49    id = newId("ixp");50    const slug = await uniqueSlug(tx, "ixps", name, i.city ?? country);51    await tx.execute(sql`insert into ixps (id, slug, name, name_long, city, country_iso2, metro_id, region_continent, website, network_count, external_ids)52      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)`);53    created = true;54  } else {55    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}),56      region_continent = coalesce(region_continent, ${i.regionContinent ?? null}), website = coalesce(${i.website ?? null}, website), network_count = coalesce(${i.networkCount ?? null}, network_count),57      external_ids = external_ids || ${JSON.stringify(i.externalIds ?? {})}::jsonb, updated_at = now() where id = ${id}`);58  }59  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`);60  addRef(ctx, "ixp", id);61  return { id, created };62}6364export async function linkFacilityIxp(tx: Tx, ctx: IngestContext, facilityId: string, ixpId: string): Promise<void> {65  await tx.execute(sql`insert into facility_ixps (facility_id, ixp_id, source_id) values (${facilityId}, ${ixpId}, ${ctx.run.sourceId}) on conflict do nothing`);66}6768export async function refreshIxpCount(tx: Tx, facilityId: string): Promise<void> {69  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}`);70}7172/** Attach IXPs listed on a facility page by name (creates the IXP when unknown, country from the facility). */73export async function attachIxpsByName(tx: Tx, ctx: IngestContext, facilityId: string, names: string[], countryIso2: string | null): Promise<number> {74  let n = 0;75  for (const raw of names) {76    const name = raw?.trim();77    if (!name || name.length < 2 || name.length > 120) continue;78    try {79      const { id } = await resolveIxp(tx, ctx, { name, countryIso2 });80      await linkFacilityIxp(tx, ctx, facilityId, id);81      n++;82    } catch (e) {83      console.warn(`[ingest] ixp "${name}" skipped: ${(e as Error).message}`);84    }85  }86  if (n) await refreshIxpCount(tx, facilityId);87  return n;88}8990export async function ingestIxp(tx: Tx, ctx: IngestContext, ix: NormalizedIxp): Promise<void> {91  const { id, created } = await resolveIxp(tx, ctx, {92    name: ix.name,93    key: ix.key,94    nameLong: ix.nameLong,95    city: ix.city,96    countryIso2: ix.countryIso2,97    regionContinent: ix.regionContinent,98    website: ix.website,99    networkCount: ix.networkCount,100    externalIds: ix.externalIds,101  });102  if (created) ctx.stats.created++;103  else ctx.stats.updated++;104  bump(ctx, "ixp");105  const touched = new Set<string>();106  for (const fk of ix.facilityKeys ?? []) {107    const fid = await facilityIdForKey(tx, ctx, fk);108    if (!fid) continue;109    await linkFacilityIxp(tx, ctx, fid, id);110    touched.add(fid);111  }112  for (const fid of touched) await refreshIxpCount(tx, fid);113  await writeProvenance(114    tx,115    ctx,116    "ixp",117    id,118    [119      { field: "name", value: ix.name },120      { field: "nameLong", value: ix.nameLong ?? null },121      { field: "city", value: ix.city ?? null },122      { field: "countryIso2", value: ix.countryIso2 ?? null },123      { field: "website", value: ix.website ?? null },124      { field: "networkCount", value: ix.networkCount ?? null },125      { field: "facilityKeys", value: ix.facilityKeys?.length ? ix.facilityKeys : null },126    ].map((f) => ({ ...f, provenance: ix.provenance })),127    ix.key,128  );129}130