/** Campus reconciliation: key → (operator, normalized name, country) → create. */ import { sql } from "@dci/db"; import { newId, normalizeName, type NormalizedCampus } from "@dci/core"; import { addRef, bump, safeCountry, uniqueSlug, type IngestContext, type Tx } from "./common.js"; import { validGeo } from "./geo.js"; import { assignMetro } from "./metros.js"; import { resolveOperator } from "./operators.js"; import { writeProvenance } from "./provenance.js"; export interface CampusInput { name: string; key?: string | null; operatorId?: string | null; countryIso2?: string | null; city?: string | null; lat?: number | null; lng?: number | null; externalIds?: Record; } export async function resolveCampus(tx: Tx, ctx: IngestContext, c: CampusInput): Promise<{ id: string; created: boolean }> { const name = c.name.trim(); if (!name) throw new Error("campus name is empty"); const norm = normalizeName(name); const country = await safeCountry(tx, ctx, c.countryIso2); let id: string | null = null; if (c.key) { const r = await tx.execute(sql`select entity_id from entity_keys where key = ${c.key} and entity_type = 'campus'`); if (r[0]) id = String(r[0].entity_id); } if (!id && c.externalIds) { for (const [k, v] of Object.entries(c.externalIds)) { // allowlist: only namespaces that name one campus may link records if (v == null || v === "" || !/^(osm|wikidata|peeringdb_campus|[a-z0-9]+_campus(_slug|_code)?)$/.test(k) || /^(operator|owner|state|region|country|postal|market|source)_/.test(k)) continue; const r = await tx.execute(sql`select id from campuses where external_ids @> ${JSON.stringify({ [k]: v })}::jsonb limit 1`); if (r[0]) { id = String(r[0].id); break; } } } if (!id) { const r = await tx.execute(sql` select id from campuses where lower(regexp_replace(name, '[^a-zA-Z0-9]+', ' ', 'g')) = lower(regexp_replace(${name}, '[^a-zA-Z0-9]+', ' ', 'g')) and (${c.operatorId ?? null}::text is null or operator_id is null or operator_id = ${c.operatorId ?? null}) and (${country}::text is null or country_iso2 is null or country_iso2 = ${country}) order by (operator_id = ${c.operatorId ?? null}) desc nulls last limit 1`); if (r[0]) id = String(r[0].id); else { // normalized-name equality (drops "campus", "data center" noise) const all = await tx.execute(sql`select id, name from campuses where (${country}::text is null or country_iso2 is null or country_iso2 = ${country}) and (${c.operatorId ?? null}::text is null or operator_id is null or operator_id = ${c.operatorId ?? null}) limit 500`); const hit = all.find((row) => normalizeName(String(row.name)) === norm); if (hit) id = String(hit.id); } } let created = false; const geo = typeof c.lat === "number" && typeof c.lng === "number" && validGeo({ lat: c.lat, lng: c.lng, precision: "unknown", source: "" }); const metro = await assignMetro(tx, { lat: geo ? c.lat : null, lng: geo ? c.lng : null, city: c.city, countryIso2: country }); if (!id) { id = newId("campus"); const slug = await uniqueSlug(tx, "campuses", name, c.city ?? country); await tx.execute(sql`insert into campuses (id, slug, name, operator_id, metro_id, country_iso2, city, lat, lng, external_ids) values (${id}, ${slug}, ${name}, ${c.operatorId ?? null}, ${metro.metroId}, ${metro.countryIso2}, ${c.city ?? null}, ${geo ? c.lat : null}, ${geo ? c.lng : null}, ${JSON.stringify(c.externalIds ?? {})}::jsonb)`); created = true; } else { await tx.execute(sql`update campuses set operator_id = coalesce(operator_id, ${c.operatorId ?? null}), metro_id = coalesce(metro_id, ${metro.metroId}), country_iso2 = coalesce(country_iso2, ${metro.countryIso2}), city = coalesce(city, ${c.city ?? null}), lat = coalesce(lat, ${geo ? c.lat : null}), lng = coalesce(lng, ${geo ? c.lng : null}), external_ids = external_ids || ${JSON.stringify(c.externalIds ?? {})}::jsonb, updated_at = now() where id = ${id}`); } if (c.key) await tx.execute(sql`insert into entity_keys (key, entity_type, entity_id, connector_id) values (${c.key}, 'campus', ${id}, ${ctx.run.connectorId}) on conflict (key) do update set entity_id = excluded.entity_id`); addRef(ctx, "campus", id); return { id, created }; } export async function ingestCampus(tx: Tx, ctx: IngestContext, c: NormalizedCampus): Promise { const operatorId = c.operatorName ? (await resolveOperator(tx, ctx, { name: c.operatorName })).id : null; const { id, created } = await resolveCampus(tx, ctx, { name: c.name, key: c.key, operatorId, countryIso2: c.countryIso2, city: c.city, lat: validGeo(c.geo) ? c.geo.lat : null, lng: validGeo(c.geo) ? c.geo.lng : null, externalIds: c.externalIds, }); if (created) ctx.stats.created++; else ctx.stats.updated++; bump(ctx, "campus"); await writeProvenance( tx, ctx, "campus", id, [ { field: "name", value: c.name }, { field: "operatorName", value: c.operatorName ?? null }, { field: "city", value: c.city ?? null }, { field: "countryIso2", value: c.countryIso2 ?? null }, { field: "geo", value: validGeo(c.geo) ? { lat: c.geo.lat, lng: c.geo.lng, precision: c.geo.precision } : null }, ].map((f) => ({ ...f, provenance: c.provenance })), c.key, ); }