spb/datacenterindex
Public
HTML 53.9%
TypeScript 44.5%
JavaScript 0.6%
SQL 0.5%
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