/** * Operator resolution: entity key → external ids → curated canonical table → exact normalized name → alias → * trigram similarity ≥ 0.92 with the same website domain → create. Aliases are accumulated on the record. */ import { sql, textArray } from "@dci/db"; import { newId, normalizeName, type NormalizedOperator, type OperatorKind, type Provenance } from "@dci/core"; import { addRef, bump, safeCountry, uniqueSlug, uniqStrings, websiteDomain, normalizeWebsite, type IngestContext, type Tx } from "./common.js"; import { findCanonicalOperator, HYPERSCALERS, type CanonicalOperator } from "./canonical-operators.js"; import { writeProvenance } from "./provenance.js"; export interface OperatorInput { name: string; key?: string | null; aliases?: string[]; kind?: OperatorKind | null; website?: string | null; hqCountryIso2?: string | null; hqCity?: string | null; parentName?: string | null; description?: string | null; externalIds?: Record; /** hint from the caller (e.g. facility.carriers → carrier, cloudProviders → cloud) */ roleHint?: "carrier" | "cloud" | null; provenance?: Provenance | null; } export interface OperatorRow { id: string; name: string; normalizedName: string; kind: string | null; website: string | null; aliases: string[]; externalIds: Record; isCloudProvider: boolean; isCarrier: boolean; hqCountryIso2: string | null; parentId: string | null; } const CARRIER_HINTS = /\b(telecom|telekom|telecommunications|communications|networks?|fiber|fibre|broadband|carrier|telco)\b/i; /** Kind inference: curated table first, then role hint, then name heuristics. */ export function inferOperatorKind(name: string, canonical: CanonicalOperator | null, hint: OperatorInput["roleHint"], provided: OperatorKind | null | undefined): OperatorKind | null { if (canonical) return canonical.kind; if (provided) return provided; if (hint === "carrier") return "carrier"; if (hint === "cloud") return "cloud"; if (CARRIER_HINTS.test(name)) return "carrier"; if (/\b(university|institute|research|laboratory|cern)\b/i.test(name)) return "research"; if (/\b(ministry|government|federal|state of|city of|county)\b/i.test(name)) return "government"; if (/\b(internet exchange|-ix\b|ixp)\b/i.test(name)) return "ixp_operator"; return null; } function rowFrom(r: Record): OperatorRow { return { id: String(r.id), name: String(r.name), normalizedName: String(r.normalized_name), kind: r.kind == null ? null : String(r.kind), website: r.website == null ? null : String(r.website), aliases: Array.isArray(r.aliases) ? (r.aliases as string[]) : [], externalIds: (r.external_ids as Record) ?? {}, isCloudProvider: Boolean(r.is_cloud_provider), isCarrier: Boolean(r.is_carrier), hqCountryIso2: r.hq_country_iso2 == null ? null : String(r.hq_country_iso2), parentId: r.parent_id == null ? null : String(r.parent_id), }; } const COLS = sql`id, name, normalized_name, kind, website, aliases, external_ids, is_cloud_provider, is_carrier, hq_country_iso2, parent_id`; async function byKey(tx: Tx, key: string): Promise { 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)`); return rows[0] ? rowFrom(rows[0]) : null; } /** External-id namespaces that identify ONE operator (allowlist — a shared `hq_country` or `stock_exchange` must never fold two companies). */ export const OPERATOR_IDENTIFYING_KEYS: ReadonlySet = new Set(["wikidata", "wikipedia_en", "peeringdb_org", "peeringdb_net", "lei", "cik", "asn", "crunchbase", "linkedin", "gleif"]); async function byExternalIds(tx: Tx, ext: Record | undefined): Promise { if (!ext) return null; for (const [k, v] of Object.entries(ext)) { if (v == null || v === "" || !OPERATOR_IDENTIFYING_KEYS.has(k)) continue; const rows = await tx.execute(sql`select ${COLS} from operators where external_ids @> ${JSON.stringify({ [k]: v })}::jsonb limit 1`); if (rows[0]) return rowFrom(rows[0]); } return null; } async function byNormalizedOrAlias(tx: Tx, names: string[]): Promise { const norms = uniqStrings(names.map((n) => normalizeName(n))).filter(Boolean); const lowers = uniqStrings(names).map((n) => n.toLowerCase()); if (!norms.length) return null; const rows = await tx.execute(sql` select ${COLS} from operators where normalized_name in ${norms} or exists (select 1 from unnest(aliases) a where lower(a) in ${lowers}) order by (normalized_name in ${norms}) desc, created_at asc limit 1`); return rows[0] ? rowFrom(rows[0]) : null; } async function byTrigramSameDomain(tx: Tx, name: string, website: string | null | undefined): Promise { const dom = websiteDomain(website); const n = normalizeName(name); if (!dom || !n) return null; const rows = await tx.execute(sql` select ${COLS}, similarity(normalized_name, ${n}) as sim from operators where website is not null and similarity(normalized_name, ${n}) >= 0.92 order by sim desc limit 5`); for (const r of rows) if (websiteDomain(String(r.website)) === dom) return rowFrom(r); return null; } /** * Resolve (or create) an operator. Returns the id and whether it was created. Results are cached per batch by * normalized name so a page listing 200 Equinix sites hits the database once. */ export async function resolveOperator(tx: Tx, ctx: IngestContext, input: OperatorInput): Promise<{ id: string; created: boolean; name: string }> { const rawName = input.name.trim(); if (!rawName) throw new Error("operator name is empty"); const canonical = findCanonicalOperator(rawName); const displayName = canonical?.name ?? rawName; const norm = normalizeName(displayName); const cacheKey = input.key ? `key:${input.key}` : `n:${norm}`; const cached = ctx.caches.operatorsByNorm.get(cacheKey) ?? ctx.caches.operatorsByNorm.get(`n:${norm}`); if (cached) return { id: cached, created: false, name: ctx.caches.operatorNamesById.get(cached) ?? displayName }; // Serialize concurrent batches resolving the same operator (two connectors seeing "Digital Realty" at once would // otherwise both miss the lookup and both insert). Transaction-scoped, released at batch commit. await tx.execute(sql`select pg_advisory_xact_lock(hashtext(${`dci:operator:${norm}`}))`); let row: OperatorRow | null = null; if (input.key) row = await byKey(tx, input.key); if (!row) row = await byExternalIds(tx, input.externalIds); if (!row) row = await byNormalizedOrAlias(tx, uniqStrings([displayName, rawName, ...(canonical?.aliases ?? []), ...(input.aliases ?? [])])); if (!row) row = await byTrigramSameDomain(tx, displayName, input.website ?? canonical?.website); const kind = inferOperatorKind(displayName, canonical, input.roleHint ?? null, input.kind); const isHyper = !!canonical && HYPERSCALERS.has(canonical.name); const isCloud = !!canonical?.isCloudProvider || kind === "cloud" || input.roleHint === "cloud" || (isHyper && canonical?.isCloudProvider !== false && kind === "hyperscaler" && !!canonical?.isCloudProvider); const isCarrier = !!canonical?.isCarrier || kind === "carrier" || input.roleHint === "carrier"; const website = normalizeWebsite(input.website) ?? canonical?.website ?? null; const hq = await safeCountry(tx, ctx, input.hqCountryIso2 ?? canonical?.hqCountryIso2 ?? null); const newAliases = uniqStrings([rawName !== displayName ? rawName : null, ...(input.aliases ?? []), ...(canonical?.aliases ?? [])]).filter((a) => normalizeName(a) !== norm || a !== displayName); let created = false; if (!row) { const id = newId("operator"); const slug = await uniqueSlug(tx, "operators", displayName, hq); await tx.execute(sql` insert into operators (id, slug, name, normalized_name, kind, website, hq_country_iso2, hq_city, description, aliases, external_ids, is_cloud_provider, is_carrier) values (${id}, ${slug}, ${displayName}, ${norm}, ${kind}, ${website}, ${hq}, ${input.hqCity ?? null}, ${input.description ?? null}, ${textArray(newAliases)}::text[], ${JSON.stringify(input.externalIds ?? {})}::jsonb, ${isCloud}, ${isCarrier})`); row = { id, name: displayName, normalizedName: norm, kind, website, aliases: newAliases, externalIds: input.externalIds ?? {}, isCloudProvider: isCloud, isCarrier, hqCountryIso2: hq, parentId: null }; created = true; ctx.stats.created++; bump(ctx, "operator"); } else { // enrich: fill blanks, merge aliases + external ids, upgrade flags/kind when the curated table knows better const mergedAliases = uniqStrings([...row.aliases, ...newAliases, row.name !== displayName && canonical ? row.name : null]).filter((a) => a !== (canonical?.name ?? row!.name)); const mergedExt = { ...row.externalIds, ...(input.externalIds ?? {}) }; const nextKind = row.kind ?? kind; const nextName = canonical ? canonical.name : row.name; 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); if (changed) { await tx.execute(sql` update operators set name = ${nextName}, normalized_name = ${normalizeName(nextName)}, kind = ${nextKind}, website = coalesce(website, ${website}), hq_country_iso2 = coalesce(hq_country_iso2, ${hq}), hq_city = coalesce(hq_city, ${input.hqCity ?? null}), description = coalesce(description, ${input.description ?? null}), aliases = ${textArray(mergedAliases)}::text[], external_ids = ${JSON.stringify(mergedExt)}::jsonb, is_cloud_provider = is_cloud_provider or ${isCloud}, is_carrier = is_carrier or ${isCarrier}, updated_at = now() where id = ${row.id}`); row = { ...row, name: nextName, aliases: mergedAliases, externalIds: mergedExt, kind: nextKind }; } } if (input.key) { 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`); } // parent linkage (curated or provided) const parentName = input.parentName ?? canonical?.parent ?? null; if (parentName && !row.parentId && normalizeName(parentName) !== row.normalizedName) { const parent = await resolveOperator(tx, ctx, { name: parentName }); 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`); } ctx.caches.operatorsByNorm.set(cacheKey, row.id); ctx.caches.operatorsByNorm.set(`n:${norm}`, row.id); ctx.caches.operatorsByNorm.set(`n:${normalizeName(rawName)}`, row.id); ctx.caches.operatorNamesById.set(row.id, row.name); addRef(ctx, "operator", row.id); return { id: row.id, created, name: row.name }; } /** Ingest a NormalizedOperator entity (full record with provenance). */ export async function ingestOperator(tx: Tx, ctx: IngestContext, op: NormalizedOperator): Promise { const before = await tx.execute(sql`select entity_id from entity_keys where key = ${op.key} and entity_type = 'operator'`); const { id, created } = await resolveOperator(tx, ctx, { name: op.name, key: op.key, aliases: op.aliases, kind: op.kind ?? null, website: op.website, hqCountryIso2: op.hqCountryIso2, hqCity: op.hqCity, parentName: op.parentName, description: op.description, externalIds: op.externalIds, provenance: op.provenance, }); if (!created) { if (before.length) ctx.stats.updated++; else ctx.stats.merged++; bump(ctx, "operator"); } const fields: Array<{ field: string; value: unknown }> = [ { field: "name", value: op.name }, { field: "kind", value: op.kind ?? null }, { field: "website", value: op.website ?? null }, { field: "hqCountryIso2", value: op.hqCountryIso2 ?? null }, { field: "hqCity", value: op.hqCity ?? null }, { field: "parentName", value: op.parentName ?? null }, { field: "description", value: op.description ?? null }, { field: "externalIds", value: op.externalIds && Object.keys(op.externalIds).length ? op.externalIds : null }, ]; await writeProvenance(tx, ctx, "operator", id, fields.map((f) => ({ ...f, provenance: op.provenance })), op.key); } /** Load `id → name` for a set of operator ids (event titles). */ export async function operatorNames(tx: Tx, ctx: IngestContext, ids: Array): Promise> { const out = new Map(); const missing: string[] = []; for (const id of ids) { if (!id) continue; const cached = ctx.caches.operatorNamesById.get(id); if (cached) out.set(id, cached); else missing.push(id); } if (missing.length) { const rows = await tx.execute(sql`select id, name from operators where id in ${missing}`); for (const r of rows) { out.set(String(r.id), String(r.name)); ctx.caches.operatorNamesById.set(String(r.id), String(r.name)); } } return out; } /** Look an operator up without creating it (canonical name/alias, exact normalized name or alias). */ export async function lookupOperatorId(tx: Tx, name: string): Promise { const canonical = findCanonicalOperator(name); const row = await byNormalizedOrAlias(tx, uniqStrings([canonical?.name ?? null, name, ...(canonical?.aliases ?? [])])); return row?.id ?? null; } /** Fold operator `fromId` into `intoId` (duplicates created by a race or reviewed by an admin). */ export async function mergeOperators(tx: Tx, fromId: string, intoId: string): Promise { if (fromId === intoId) return; const from = (await tx.execute(sql`select name, aliases, external_ids from operators where id = ${fromId}`))[0]; if (!from) throw new Error(`operator ${fromId} not found`); 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"]) { const [table, col] = t.split(":") as [string, string]; await tx.execute(sql`update ${sql.identifier(table)} set ${sql.identifier(col)} = ${intoId} where ${sql.identifier(col)} = ${fromId}`); } 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`); await tx.execute(sql`delete from facility_tenants where operator_id = ${fromId}`); await tx.execute(sql`update entity_keys set entity_id = ${intoId} where entity_type = 'operator' and entity_id = ${fromId}`); await tx.execute(sql`update provenance set entity_id = ${intoId} where entity_type = 'operator' and entity_id = ${fromId} 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)`); await tx.execute(sql`delete from provenance where entity_type = 'operator' and entity_id = ${fromId}`); await tx.execute(sql`update news_items set operator_ids = array_replace(operator_ids, ${fromId}, ${intoId}) where ${fromId} = any(operator_ids)`); const aliases = Array.isArray(from.aliases) ? (from.aliases as string[]) : []; await tx.execute(sql`update operators set aliases = (select coalesce(array_agg(distinct a), '{}'::text[]) from unnest(aliases || ${textArray([...aliases, String(from.name)])}::text[]) a where a <> operators.name), external_ids = ${JSON.stringify(from.external_ids ?? {})}::jsonb || external_ids, updated_at = now() where id = ${intoId}`); await tx.execute(sql`delete from operators where id = ${fromId}`); }