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%
5.3 KB · 105 lines typescript
Raw Blame History
1/** Campus reconciliation: key → (operator, normalized name, country) → create. */2import { sql } from "@dci/db";3import { newId, normalizeName, type NormalizedCampus } from "@dci/core";4import { addRef, bump, safeCountry, uniqueSlug, type IngestContext, type Tx } from "./common.js";5import { validGeo } from "./geo.js";6import { assignMetro } from "./metros.js";7import { resolveOperator } from "./operators.js";8import { writeProvenance } from "./provenance.js";910export interface CampusInput {11  name: string;12  key?: string | null;13  operatorId?: string | null;14  countryIso2?: string | null;15  city?: string | null;16  lat?: number | null;17  lng?: number | null;18  externalIds?: Record<string, string | number>;19}2021export async function resolveCampus(tx: Tx, ctx: IngestContext, c: CampusInput): Promise<{ id: string; created: boolean }> {22  const name = c.name.trim();23  if (!name) throw new Error("campus name is empty");24  const norm = normalizeName(name);25  const country = await safeCountry(tx, ctx, c.countryIso2);26  let id: string | null = null;27  if (c.key) {28    const r = await tx.execute(sql`select entity_id from entity_keys where key = ${c.key} and entity_type = 'campus'`);29    if (r[0]) id = String(r[0].entity_id);30  }31  if (!id && c.externalIds) {32    for (const [k, v] of Object.entries(c.externalIds)) {33      // allowlist: only namespaces that name one campus may link records34      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;35      const r = await tx.execute(sql`select id from campuses where external_ids @> ${JSON.stringify({ [k]: v })}::jsonb limit 1`);36      if (r[0]) {37        id = String(r[0].id);38        break;39      }40    }41  }42  if (!id) {43    const r = await tx.execute(sql`44      select id from campuses where lower(regexp_replace(name, '[^a-zA-Z0-9]+', ' ', 'g')) = lower(regexp_replace(${name}, '[^a-zA-Z0-9]+', ' ', 'g'))45        and (${c.operatorId ?? null}::text is null or operator_id is null or operator_id = ${c.operatorId ?? null})46        and (${country}::text is null or country_iso2 is null or country_iso2 = ${country})47      order by (operator_id = ${c.operatorId ?? null}) desc nulls last limit 1`);48    if (r[0]) id = String(r[0].id);49    else {50      // normalized-name equality (drops "campus", "data center" noise)51      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`);52      const hit = all.find((row) => normalizeName(String(row.name)) === norm);53      if (hit) id = String(hit.id);54    }55  }56  let created = false;57  const geo = typeof c.lat === "number" && typeof c.lng === "number" && validGeo({ lat: c.lat, lng: c.lng, precision: "unknown", source: "" });58  const metro = await assignMetro(tx, { lat: geo ? c.lat : null, lng: geo ? c.lng : null, city: c.city, countryIso2: country });59  if (!id) {60    id = newId("campus");61    const slug = await uniqueSlug(tx, "campuses", name, c.city ?? country);62    await tx.execute(sql`insert into campuses (id, slug, name, operator_id, metro_id, country_iso2, city, lat, lng, external_ids)63      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)`);64    created = true;65  } else {66    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}),67      city = coalesce(city, ${c.city ?? null}), lat = coalesce(lat, ${geo ? c.lat : null}), lng = coalesce(lng, ${geo ? c.lng : null}),68      external_ids = external_ids || ${JSON.stringify(c.externalIds ?? {})}::jsonb, updated_at = now() where id = ${id}`);69  }70  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`);71  addRef(ctx, "campus", id);72  return { id, created };73}7475export async function ingestCampus(tx: Tx, ctx: IngestContext, c: NormalizedCampus): Promise<void> {76  const operatorId = c.operatorName ? (await resolveOperator(tx, ctx, { name: c.operatorName })).id : null;77  const { id, created } = await resolveCampus(tx, ctx, {78    name: c.name,79    key: c.key,80    operatorId,81    countryIso2: c.countryIso2,82    city: c.city,83    lat: validGeo(c.geo) ? c.geo.lat : null,84    lng: validGeo(c.geo) ? c.geo.lng : null,85    externalIds: c.externalIds,86  });87  if (created) ctx.stats.created++;88  else ctx.stats.updated++;89  bump(ctx, "campus");90  await writeProvenance(91    tx,92    ctx,93    "campus",94    id,95    [96      { field: "name", value: c.name },97      { field: "operatorName", value: c.operatorName ?? null },98      { field: "city", value: c.city ?? null },99      { field: "countryIso2", value: c.countryIso2 ?? null },100      { field: "geo", value: validGeo(c.geo) ? { lat: c.geo.lat, lng: c.geo.lng, precision: c.geo.precision } : null },101    ].map((f) => ({ ...f, provenance: c.provenance })),102    c.key,103  );104}105