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%
15.8 KB · 281 lines typescript
Raw Blame History
1/**2 * Operator resolution: entity key → external ids → curated canonical table → exact normalized name → alias →3 * trigram similarity ≥ 0.92 with the same website domain → create. Aliases are accumulated on the record.4 */5import { sql, textArray } from "@dci/db";6import { newId, normalizeName, type NormalizedOperator, type OperatorKind, type Provenance } from "@dci/core";7import { addRef, bump, safeCountry, uniqueSlug, uniqStrings, websiteDomain, normalizeWebsite, type IngestContext, type Tx } from "./common.js";8import { findCanonicalOperator, HYPERSCALERS, type CanonicalOperator } from "./canonical-operators.js";9import { writeProvenance } from "./provenance.js";1011export interface OperatorInput {12  name: string;13  key?: string | null;14  aliases?: string[];15  kind?: OperatorKind | null;16  website?: string | null;17  hqCountryIso2?: string | null;18  hqCity?: string | null;19  parentName?: string | null;20  description?: string | null;21  externalIds?: Record<string, string | number>;22  /** hint from the caller (e.g. facility.carriers → carrier, cloudProviders → cloud) */23  roleHint?: "carrier" | "cloud" | null;24  provenance?: Provenance | null;25}2627export interface OperatorRow {28  id: string;29  name: string;30  normalizedName: string;31  kind: string | null;32  website: string | null;33  aliases: string[];34  externalIds: Record<string, string | number>;35  isCloudProvider: boolean;36  isCarrier: boolean;37  hqCountryIso2: string | null;38  parentId: string | null;39}4041const CARRIER_HINTS = /\b(telecom|telekom|telecommunications|communications|networks?|fiber|fibre|broadband|carrier|telco)\b/i;4243/** Kind inference: curated table first, then role hint, then name heuristics. */44export function inferOperatorKind(name: string, canonical: CanonicalOperator | null, hint: OperatorInput["roleHint"], provided: OperatorKind | null | undefined): OperatorKind | null {45  if (canonical) return canonical.kind;46  if (provided) return provided;47  if (hint === "carrier") return "carrier";48  if (hint === "cloud") return "cloud";49  if (CARRIER_HINTS.test(name)) return "carrier";50  if (/\b(university|institute|research|laboratory|cern)\b/i.test(name)) return "research";51  if (/\b(ministry|government|federal|state of|city of|county)\b/i.test(name)) return "government";52  if (/\b(internet exchange|-ix\b|ixp)\b/i.test(name)) return "ixp_operator";53  return null;54}5556function rowFrom(r: Record<string, unknown>): OperatorRow {57  return {58    id: String(r.id),59    name: String(r.name),60    normalizedName: String(r.normalized_name),61    kind: r.kind == null ? null : String(r.kind),62    website: r.website == null ? null : String(r.website),63    aliases: Array.isArray(r.aliases) ? (r.aliases as string[]) : [],64    externalIds: (r.external_ids as Record<string, string | number>) ?? {},65    isCloudProvider: Boolean(r.is_cloud_provider),66    isCarrier: Boolean(r.is_carrier),67    hqCountryIso2: r.hq_country_iso2 == null ? null : String(r.hq_country_iso2),68    parentId: r.parent_id == null ? null : String(r.parent_id),69  };70}7172const COLS = sql`id, name, normalized_name, kind, website, aliases, external_ids, is_cloud_provider, is_carrier, hq_country_iso2, parent_id`;7374async function byKey(tx: Tx, key: string): Promise<OperatorRow | null> {75  const rows = await tx.execute(sql`select ${COLS} from operators o where o.id = (select entity_id from entity_keys where key = ${key} and entity_type = 'operator' limit 1)`);76  return rows[0] ? rowFrom(rows[0]) : null;77}7879/** External-id namespaces that identify ONE operator (allowlist — a shared `hq_country` or `stock_exchange` must never fold two companies). */80export const OPERATOR_IDENTIFYING_KEYS: ReadonlySet<string> = new Set(["wikidata", "wikipedia_en", "peeringdb_org", "peeringdb_net", "lei", "cik", "asn", "crunchbase", "linkedin", "gleif"]);81async function byExternalIds(tx: Tx, ext: Record<string, string | number> | undefined): Promise<OperatorRow | null> {82  if (!ext) return null;83  for (const [k, v] of Object.entries(ext)) {84    if (v == null || v === "" || !OPERATOR_IDENTIFYING_KEYS.has(k)) continue;85    const rows = await tx.execute(sql`select ${COLS} from operators where external_ids @> ${JSON.stringify({ [k]: v })}::jsonb limit 1`);86    if (rows[0]) return rowFrom(rows[0]);87  }88  return null;89}9091async function byNormalizedOrAlias(tx: Tx, names: string[]): Promise<OperatorRow | null> {92  const norms = uniqStrings(names.map((n) => normalizeName(n))).filter(Boolean);93  const lowers = uniqStrings(names).map((n) => n.toLowerCase());94  if (!norms.length) return null;95  const rows = await tx.execute(sql`96    select ${COLS} from operators97    where normalized_name in ${norms}98       or exists (select 1 from unnest(aliases) a where lower(a) in ${lowers})99    order by (normalized_name in ${norms}) desc, created_at asc100    limit 1`);101  return rows[0] ? rowFrom(rows[0]) : null;102}103104async function byTrigramSameDomain(tx: Tx, name: string, website: string | null | undefined): Promise<OperatorRow | null> {105  const dom = websiteDomain(website);106  const n = normalizeName(name);107  if (!dom || !n) return null;108  const rows = await tx.execute(sql`109    select ${COLS}, similarity(normalized_name, ${n}) as sim from operators110    where website is not null and similarity(normalized_name, ${n}) >= 0.92111    order by sim desc limit 5`);112  for (const r of rows) if (websiteDomain(String(r.website)) === dom) return rowFrom(r);113  return null;114}115116/**117 * Resolve (or create) an operator. Returns the id and whether it was created. Results are cached per batch by118 * normalized name so a page listing 200 Equinix sites hits the database once.119 */120export async function resolveOperator(tx: Tx, ctx: IngestContext, input: OperatorInput): Promise<{ id: string; created: boolean; name: string }> {121  const rawName = input.name.trim();122  if (!rawName) throw new Error("operator name is empty");123  const canonical = findCanonicalOperator(rawName);124  const displayName = canonical?.name ?? rawName;125  const norm = normalizeName(displayName);126  const cacheKey = input.key ? `key:${input.key}` : `n:${norm}`;127  const cached = ctx.caches.operatorsByNorm.get(cacheKey) ?? ctx.caches.operatorsByNorm.get(`n:${norm}`);128  if (cached) return { id: cached, created: false, name: ctx.caches.operatorNamesById.get(cached) ?? displayName };129130  // Serialize concurrent batches resolving the same operator (two connectors seeing "Digital Realty" at once would131  // otherwise both miss the lookup and both insert). Transaction-scoped, released at batch commit.132  await tx.execute(sql`select pg_advisory_xact_lock(hashtext(${`dci:operator:${norm}`}))`);133134  let row: OperatorRow | null = null;135  if (input.key) row = await byKey(tx, input.key);136  if (!row) row = await byExternalIds(tx, input.externalIds);137  if (!row) row = await byNormalizedOrAlias(tx, uniqStrings([displayName, rawName, ...(canonical?.aliases ?? []), ...(input.aliases ?? [])]));138  if (!row) row = await byTrigramSameDomain(tx, displayName, input.website ?? canonical?.website);139140  const kind = inferOperatorKind(displayName, canonical, input.roleHint ?? null, input.kind);141  const isHyper = !!canonical && HYPERSCALERS.has(canonical.name);142  const isCloud = !!canonical?.isCloudProvider || kind === "cloud" || input.roleHint === "cloud" || (isHyper && canonical?.isCloudProvider !== false && kind === "hyperscaler" && !!canonical?.isCloudProvider);143  const isCarrier = !!canonical?.isCarrier || kind === "carrier" || input.roleHint === "carrier";144  const website = normalizeWebsite(input.website) ?? canonical?.website ?? null;145  const hq = await safeCountry(tx, ctx, input.hqCountryIso2 ?? canonical?.hqCountryIso2 ?? null);146  const newAliases = uniqStrings([rawName !== displayName ? rawName : null, ...(input.aliases ?? []), ...(canonical?.aliases ?? [])]).filter((a) => normalizeName(a) !== norm || a !== displayName);147148  let created = false;149  if (!row) {150    const id = newId("operator");151    const slug = await uniqueSlug(tx, "operators", displayName, hq);152    await tx.execute(sql`153      insert into operators (id, slug, name, normalized_name, kind, website, hq_country_iso2, hq_city, description, aliases, external_ids, is_cloud_provider, is_carrier)154      values (${id}, ${slug}, ${displayName}, ${norm}, ${kind}, ${website}, ${hq}, ${input.hqCity ?? null}, ${input.description ?? null},155        ${textArray(newAliases)}::text[], ${JSON.stringify(input.externalIds ?? {})}::jsonb, ${isCloud}, ${isCarrier})`);156    row = { id, name: displayName, normalizedName: norm, kind, website, aliases: newAliases, externalIds: input.externalIds ?? {}, isCloudProvider: isCloud, isCarrier, hqCountryIso2: hq, parentId: null };157    created = true;158    ctx.stats.created++;159    bump(ctx, "operator");160  } else {161    // enrich: fill blanks, merge aliases + external ids, upgrade flags/kind when the curated table knows better162    const mergedAliases = uniqStrings([...row.aliases, ...newAliases, row.name !== displayName && canonical ? row.name : null]).filter((a) => a !== (canonical?.name ?? row!.name));163    const mergedExt = { ...row.externalIds, ...(input.externalIds ?? {}) };164    const nextKind = row.kind ?? kind;165    const nextName = canonical ? canonical.name : row.name;166    const changed = mergedAliases.length !== row.aliases.length || JSON.stringify(mergedExt) !== JSON.stringify(row.externalIds) || nextKind !== row.kind || (!row.website && website) || (!row.hqCountryIso2 && hq) || nextName !== row.name || (isCloud && !row.isCloudProvider) || (isCarrier && !row.isCarrier);167    if (changed) {168      await tx.execute(sql`169        update operators set170          name = ${nextName}, normalized_name = ${normalizeName(nextName)}, kind = ${nextKind}, website = coalesce(website, ${website}), hq_country_iso2 = coalesce(hq_country_iso2, ${hq}),171          hq_city = coalesce(hq_city, ${input.hqCity ?? null}), description = coalesce(description, ${input.description ?? null}),172          aliases = ${textArray(mergedAliases)}::text[],173          external_ids = ${JSON.stringify(mergedExt)}::jsonb, is_cloud_provider = is_cloud_provider or ${isCloud}, is_carrier = is_carrier or ${isCarrier}, updated_at = now()174        where id = ${row.id}`);175      row = { ...row, name: nextName, aliases: mergedAliases, externalIds: mergedExt, kind: nextKind };176    }177  }178179  if (input.key) {180    await tx.execute(sql`insert into entity_keys (key, entity_type, entity_id, connector_id) values (${input.key}, 'operator', ${row.id}, ${ctx.run.connectorId}) on conflict (key) do update set entity_id = excluded.entity_id`);181  }182  // parent linkage (curated or provided)183  const parentName = input.parentName ?? canonical?.parent ?? null;184  if (parentName && !row.parentId && normalizeName(parentName) !== row.normalizedName) {185    const parent = await resolveOperator(tx, ctx, { name: parentName });186    if (parent.id !== row.id) await tx.execute(sql`update operators set parent_id = ${parent.id} where id = ${row.id} and parent_id is null`);187  }188189  ctx.caches.operatorsByNorm.set(cacheKey, row.id);190  ctx.caches.operatorsByNorm.set(`n:${norm}`, row.id);191  ctx.caches.operatorsByNorm.set(`n:${normalizeName(rawName)}`, row.id);192  ctx.caches.operatorNamesById.set(row.id, row.name);193  addRef(ctx, "operator", row.id);194  return { id: row.id, created, name: row.name };195}196197/** Ingest a NormalizedOperator entity (full record with provenance). */198export async function ingestOperator(tx: Tx, ctx: IngestContext, op: NormalizedOperator): Promise<void> {199  const before = await tx.execute(sql`select entity_id from entity_keys where key = ${op.key} and entity_type = 'operator'`);200  const { id, created } = await resolveOperator(tx, ctx, {201    name: op.name,202    key: op.key,203    aliases: op.aliases,204    kind: op.kind ?? null,205    website: op.website,206    hqCountryIso2: op.hqCountryIso2,207    hqCity: op.hqCity,208    parentName: op.parentName,209    description: op.description,210    externalIds: op.externalIds,211    provenance: op.provenance,212  });213  if (!created) {214    if (before.length) ctx.stats.updated++;215    else ctx.stats.merged++;216    bump(ctx, "operator");217  }218  const fields: Array<{ field: string; value: unknown }> = [219    { field: "name", value: op.name },220    { field: "kind", value: op.kind ?? null },221    { field: "website", value: op.website ?? null },222    { field: "hqCountryIso2", value: op.hqCountryIso2 ?? null },223    { field: "hqCity", value: op.hqCity ?? null },224    { field: "parentName", value: op.parentName ?? null },225    { field: "description", value: op.description ?? null },226    { field: "externalIds", value: op.externalIds && Object.keys(op.externalIds).length ? op.externalIds : null },227  ];228  await writeProvenance(tx, ctx, "operator", id, fields.map((f) => ({ ...f, provenance: op.provenance })), op.key);229}230231/** Load `id → name` for a set of operator ids (event titles). */232export async function operatorNames(tx: Tx, ctx: IngestContext, ids: Array<string | null | undefined>): Promise<Map<string, string>> {233  const out = new Map<string, string>();234  const missing: string[] = [];235  for (const id of ids) {236    if (!id) continue;237    const cached = ctx.caches.operatorNamesById.get(id);238    if (cached) out.set(id, cached);239    else missing.push(id);240  }241  if (missing.length) {242    const rows = await tx.execute(sql`select id, name from operators where id in ${missing}`);243    for (const r of rows) {244      out.set(String(r.id), String(r.name));245      ctx.caches.operatorNamesById.set(String(r.id), String(r.name));246    }247  }248  return out;249}250251/** Look an operator up without creating it (canonical name/alias, exact normalized name or alias). */252export async function lookupOperatorId(tx: Tx, name: string): Promise<string | null> {253  const canonical = findCanonicalOperator(name);254  const row = await byNormalizedOrAlias(tx, uniqStrings([canonical?.name ?? null, name, ...(canonical?.aliases ?? [])]));255  return row?.id ?? null;256}257258/** Fold operator `fromId` into `intoId` (duplicates created by a race or reviewed by an admin). */259export async function mergeOperators(tx: Tx, fromId: string, intoId: string): Promise<void> {260  if (fromId === intoId) return;261  const from = (await tx.execute(sql`select name, aliases, external_ids from operators where id = ${fromId}`))[0];262  if (!from) throw new Error(`operator ${fromId} not found`);263  for (const t of ["facilities:operator_id", "facilities:owner_id", "campuses:operator_id", "cloud_regions:provider_id", "projects:operator_id", "events:operator_id", "operators:parent_id"]) {264    const [table, col] = t.split(":") as [string, string];265    await tx.execute(sql`update ${sql.identifier(table)} set ${sql.identifier(col)} = ${intoId} where ${sql.identifier(col)} = ${fromId}`);266  }267  await tx.execute(sql`insert into facility_tenants (facility_id, operator_id, role, asn, source_id) select facility_id, ${intoId}, role, asn, source_id from facility_tenants where operator_id = ${fromId} on conflict do nothing`);268  await tx.execute(sql`delete from facility_tenants where operator_id = ${fromId}`);269  await tx.execute(sql`update entity_keys set entity_id = ${intoId} where entity_type = 'operator' and entity_id = ${fromId}`);270  await tx.execute(sql`update provenance set entity_id = ${intoId} where entity_type = 'operator' and entity_id = ${fromId}271    and not exists (select 1 from provenance q where q.entity_type = 'operator' and q.entity_id = ${intoId} and q.field = provenance.field and q.source_id = provenance.source_id and q.url = provenance.url)`);272  await tx.execute(sql`delete from provenance where entity_type = 'operator' and entity_id = ${fromId}`);273  await tx.execute(sql`update news_items set operator_ids = array_replace(operator_ids, ${fromId}, ${intoId}) where ${fromId} = any(operator_ids)`);274  const aliases = Array.isArray(from.aliases) ? (from.aliases as string[]) : [];275  await tx.execute(sql`update operators set276      aliases = (select coalesce(array_agg(distinct a), '{}'::text[]) from unnest(aliases || ${textArray([...aliases, String(from.name)])}::text[]) a where a <> operators.name),277      external_ids = ${JSON.stringify(from.external_ids ?? {})}::jsonb || external_ids, updated_at = now()278    where id = ${intoId}`);279  await tx.execute(sql`delete from operators where id = ${fromId}`);280}281