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%
47.7 KB · 663 lines typescript
Raw Blame History
1/**2 * Facility reconciliation + persistence (docs/RECONCILIATION.md):3 *  1. entity_keys exact key → same facility (update)4 *  2. shared external ids (peeringdb_fac, osm, wikidata, …) → same facility5 *  3. candidates within the match radius of the coordinates, or same normalized name in the same country, or same6 *     operator in the same country (facility-code matching) → weighted score7 *     ≥ 0.92 auto-merge · 0.60–0.92 create-as-new + pending entity_matches row (nothing is lost, admin merges) · < 0.60 create8 * Field merge follows the authority policy in match.ts; coordinates never lose precision.9 */10import { sql, textArray } from "@dci/db";11import { cleanText, geohash, newId, normalizeName, sha256, validLatLng, type ConfidenceLevel, type NormalizedFacility, type Provenance } from "@dci/core";12import { addRef, bump, isPrimaryKind, provenanceFor, safeCountry, uniqueSlug, uniqStrings, normalizeWebsite, type IngestContext, type Tx } from "./common.js";13import { isHyperscalerName } from "./canonical-operators.js";14import { inferOperatorFromName } from "./operator-inference.js";1516/** Country pairs whose city-level points legitimately straddle a border / enclave (stated:point-in). */17export const BORDER_TOLERANT: ReadonlySet<string> = new Set(["HK:CN", "MO:CN", "SG:MY", "MY:SG", "LU:DE", "LU:FR", "LU:BE", "MC:FR", "VA:IT", "SM:IT", "PS:IL", "IL:PS", "XK:RS", "RS:XK", "LI:CH", "LI:AT", "AD:FR", "AD:ES", "GI:ES", "BH:SA", "CH:DE", "DE:CH", "NL:BE", "BE:NL", "NL:DE", "DE:NL", "AT:DE", "DE:AT", "IE:GB", "GB:IE", "US:CA", "CA:US", "US:MX", "MX:US", "DK:SE", "SE:DK", "FR:BE", "BE:FR", "FR:CH", "CH:FR", "FR:DE", "DE:FR", "CZ:DE", "DE:CZ", "PL:DE", "DE:PL", "SK:AT", "AT:SK", "HU:AT", "AT:HU", "PT:ES", "ES:PT", "NO:SE", "SE:NO", "FI:SE", "SE:FI", "AE:OM", "OM:AE", "CN:HK", "CN:MO"]);18import { resolveCampus } from "./campuses.js";19import { discoverySignificance, emitDiffEvents, recordEvent, TRACKED_FACILITY_FIELDS } from "./events.js";20import { bboxAround, validGeo } from "./geo.js";21import { countryFromPoint } from "./country-lookup.js";22import { facilityIdForKey, upsertKey } from "./keys.js";23import { attachIxpsByName } from "./ixps.js";24import { bestMatch, bestMw, completenessScore, decide, facilityConfidence, isPipelineStatus, scoreFacilityMatch, shouldReplace, shouldReplaceGeo, shouldReplaceMw, type FacilityCandidate, type FieldObservation, type MatchDecision, type MatchScore } from "./match.js";25import { CAPACITY_COLUMN, classifyAiEvidence, classifyCapacitySemantics, classifyScope, findEvidence, type CapacityPredicate, type ClaimScope } from "@dci/core";26import { markWinners, recordCapacityClaim, writeQualityFlags } from "./claims.js";27import { assignMetro } from "./metros.js";28import { operatorNames, resolveOperator } from "./operators.js";29import { backingObservation, loadCurrentProvenance, summarizeSources, writeProvenance, type CurrentProvenance, type ObservedField } from "./provenance.js";30import { attachTenants, refreshTenantCounts } from "./tenants.js";3132/** Stored facility columns we merge into (camelCase keys ↔ snake_case columns). */33const COLUMNS: Record<string, string> = {34  name: "name",35  operatorId: "operator_id",36  ownerId: "owner_id",37  campusId: "campus_id",38  metroId: "metro_id",39  countryIso2: "country_iso2",40  city: "city",41  regionName: "region_name",42  address: "address",43  postalCode: "postal_code",44  lat: "lat",45  lng: "lng",46  geoPrecision: "geo_precision",47  geoSource: "geo_source",48  geohash: "geohash",49  status: "status",50  facilityType: "facility_type",51  tier: "tier",52  buildingSqm: "building_sqm",53  siteAreaHa: "site_area_ha",54  itCapacityMw: "it_capacity_mw",55  totalPowerMw: "total_power_mw",56  plannedPowerMw: "planned_power_mw",57  mwIsEstimate: "mw_is_estimate",58  rackCount: "rack_count",59  pue: "pue",60  coolingType: "cooling_type",61  renewableClaim: "renewable_claim",62  openedOn: "opened_on",63  constructionStartedOn: "construction_started_on",64  announcedOn: "announced_on",65  website: "website",66  description: "description",67  isAi: "is_ai",68  isHyperscale: "is_hyperscale",69  certifications: "certifications",70  confidence: "confidence",71  completeness: "completeness",72  externalIds: "external_ids",73  sourceCount: "source_count",74  lastVerified: "last_verified",75  parentFacilityId: "parent_facility_id",76  recordScope: "record_scope",77  aiEvidence: "ai_evidence",78  utilityCapacityMw: "utility_capacity_mw",79  gridConnectionMw: "grid_connection_mw",80  ultimateCampusMw: "ultimate_campus_mw",81  capacityScope: "capacity_scope",82  capacitySemantics: "capacity_semantics",83  developerId: "developer_id",84  landownerId: "landowner_id",85};8687const SCALAR_FIELDS = ["name", "city", "regionName", "address", "postalCode", "status", "facilityType", "tier", "buildingSqm", "siteAreaHa", "rackCount", "pue", "coolingType", "renewableClaim", "openedOn", "constructionStartedOn", "announcedOn", "website", "description"] as const;88const MW_FIELDS = ["itCapacityMw", "totalPowerMw", "plannedPowerMw"] as const;8990type Row = Record<string, unknown>;9192async function loadFacility(tx: Tx, id: string): Promise<Row | null> {93  const rows = await tx.execute(sql`select f.*, (select coalesce(array_agg(alias), '{}'::text[]) from facility_aliases a where a.facility_id = f.id) as alias_list from facilities f where f.id = ${id}`);94  if (!rows[0]) return null;95  const r = rows[0];96  // camelCase view of the row97  const out: Row = {};98  for (const [camel, col] of Object.entries(COLUMNS)) out[camel] = r[col] ?? null;99  out.id = r.id;100  out.slug = r.slug;101  out.mergedInto = r.merged_into ?? null;102  out.aliases = Array.isArray(r.alias_list) ? r.alias_list : [];103  out.normalizedName = r.normalized_name;104  return out;105}106107async function followMerged(tx: Tx, id: string): Promise<string> {108  let cur = id;109  for (let i = 0; i < 5; i++) {110    const r = await tx.execute(sql`select merged_into from facilities where id = ${cur}`);111    const m = r[0]?.merged_into;112    if (!m) break;113    cur = String(m);114  }115  return cur;116}117118/**119 * External-id namespaces that identify ONE facility. Anything else stored in `externalIds` (an operator's Wikidata120 * id, a state name, an investment figure, a shared `ref` tag…) is informative but must never link records — a121 * shared `operator_wikidata` once folded 121 AWS OpenStreetMap features into a single facility.122 */123export const IDENTIFYING_EXTERNAL_ID_KEYS: ReadonlySet<string> = new Set([124  "osm", "wikidata", "wikipedia_en", "peeringdb_fac", "peeringdb", "pdb_fac", "dcmap", "geonames",125  "facebook_page", "meta_info_sheet", "meta_location", "google_detail_page", "google_location",126  "equinix_ibx", "digitalrealty_node", "digitalrealty_site_code", "ntt_slug", "cyrusone_slug", "stack_slug", "qts_slug", "vantage_slug", "coresite_code", "switch_slug",127  "databank_code", "edgeconnex_slug", "tierpoint_slug", "airtrunk_code", "nextdc_code", "atnorth_code", "cloudhq_slug", "skybox_slug", "odata_code", "telehouse_slug", "coltdcs_slug", "globalswitch_slug", "virtus_slug", "digitaledge_code",128]);129130/**131 * ALLOWLIST ONLY. `state_code`, `region_code`, `postal_code`, `source_url`, `osm_ref` (a shared `ref` tag), `stack_campus`132 * (one campus, many buildings) or `investment_currency` all end in an id-looking suffix and would fold unrelated133 * facilities into one row (the `*_code` heuristic once folded 121 AWS OpenStreetMap features). A new connector that134 * needs its key to link records adds it here explicitly.135 */136export function isIdentifyingExternalId(key: string): boolean {137  return IDENTIFYING_EXTERNAL_ID_KEYS.has(key);138}139140async function byExternalIds(tx: Tx, ext: Record<string, string | number> | undefined): Promise<string | null> {141  if (!ext) return null;142  for (const [k, v] of Object.entries(ext)) {143    if (v == null || v === "" || !isIdentifyingExternalId(k)) continue;144    const rows = await tx.execute(sql`select id from facilities where merged_into is null and external_ids @> ${JSON.stringify({ [k]: v })}::jsonb order by created_at asc limit 1`);145    if (rows[0]) return String(rows[0].id);146  }147  return null;148}149150async function loadCandidates(tx: Tx, nf: NormalizedFacility, operatorId: string | null, countryIso2: string | null): Promise<FacilityCandidate[]> {151  const norm = normalizeName(nf.name);152  const aliasNorms = uniqStrings([norm, ...(nf.aliases ?? []).map((a) => normalizeName(a))]).filter(Boolean);153  const geo = validGeo(nf.geo) ? nf.geo : null;154  const box = geo ? bboxAround(geo.lat, geo.lng, 40) : null;155  const conds = [];156  if (box) conds.push(sql`(f.lat between ${box.minLat} and ${box.maxLat} and f.lng between ${box.minLng} and ${box.maxLng})`);157  if (aliasNorms.length) {158    conds.push(countryIso2 ? sql`(f.normalized_name in ${aliasNorms} and (f.country_iso2 = ${countryIso2} or f.country_iso2 is null))` : sql`(f.normalized_name in ${aliasNorms})`);159    conds.push(sql`exists (select 1 from facility_aliases a where a.facility_id = f.id and a.normalized in ${aliasNorms})`);160  }161  if (operatorId && countryIso2) conds.push(sql`(f.operator_id = ${operatorId} and f.country_iso2 = ${countryIso2})`);162  if (!conds.length) return [];163  // deterministic and relevance-ordered: same operator first, then same normalized name, then nearest — a dense metro164  // (Ashburn, Dallas, Singapore) has far more than 400 rows in a 40 km box and the true duplicate must be on the first page165  const orderGeo = geo ? sql`, ((f.lat - ${geo.lat}) * (f.lat - ${geo.lat}) + (f.lng - ${geo.lng}) * (f.lng - ${geo.lng})) asc nulls last` : sql``;166  const rows = await tx.execute(sql`167    select f.id, f.name, f.normalized_name, f.operator_id, o.name as operator_name, f.country_iso2, f.city, f.address, f.lat, f.lng, f.geo_precision, f.external_ids,168      (select coalesce(array_agg(alias), '{}'::text[]) from facility_aliases a where a.facility_id = f.id) as aliases169    from facilities f left join operators o on o.id = f.operator_id where f.merged_into is null and (${sql.join(conds, sql` or `)})170    order by (${operatorId ?? null}::text is not null and f.operator_id = ${operatorId ?? null}) desc, (f.normalized_name in ${aliasNorms.length ? aliasNorms : ["__none__"]}) desc${orderGeo}, f.created_at asc171    limit 400`);172  return rows.map((r) => ({173    id: String(r.id),174    name: String(r.name),175    normalizedName: String(r.normalized_name),176    aliases: Array.isArray(r.aliases) ? (r.aliases as string[]) : [],177    operatorId: r.operator_id == null ? null : String(r.operator_id),178    operatorName: r.operator_name == null ? null : String(r.operator_name),179    countryIso2: r.country_iso2 == null ? null : String(r.country_iso2),180    city: r.city == null ? null : String(r.city),181    address: r.address == null ? null : String(r.address),182    lat: r.lat == null ? null : Number(r.lat),183    lng: r.lng == null ? null : Number(r.lng),184    geoPrecision: r.geo_precision == null ? null : String(r.geo_precision),185    externalIds: (r.external_ids as Record<string, string | number>) ?? null,186  }));187}188189interface Resolution {190  id: string | null;191  how: "key" | "external_id" | "merge" | "pending" | "create" | "campus_link";192  match?: { candidate: FacilityCandidate; match: MatchScore } | null;193}194195/** true when a facility name designates a campus / park / multi-building site rather than one building. */196export function isCampusName(name: string | null | undefined, campusName?: string | null): boolean {197  return /\b(campus|park|cluster|hub|gigafactory|complex|estate|mega ?site)\b/i.test(`${name ?? ""} ${campusName ?? ""}`);198}199200/** Record scope from the source's statement or the name: campus designation → campus, building code (DC12, Hall 3, Building B) → building, else facility. */201export function inferRecordScope(nf: { name: string; recordScope?: string | null; campusName?: string | null }): "building" | "facility" | "campus" {202  if (nf.recordScope === "building" || nf.recordScope === "facility" || nf.recordScope === "campus") return nf.recordScope;203  if (isCampusName(nf.name)) return "campus";204  if (/\b(building|bldg|hall|data hall|phase)\s*[A-Z0-9]{1,3}\b/i.test(nf.name) || (nf.campusName && nf.campusName !== nf.name)) return "building";205  return "facility";206}207208async function resolveFacility(tx: Tx, ctx: IngestContext, nf: NormalizedFacility, operatorId: string | null, operatorName: string | null, countryIso2: string | null): Promise<Resolution> {209  const byKey = await facilityIdForKey(tx, ctx, nf.key);210  if (byKey) return { id: byKey, how: "key" };211  const byExt = await byExternalIds(tx, nf.externalIds);212  if (byExt) return { id: await followMerged(tx, byExt), how: "external_id" };213  const candidates = await loadCandidates(tx, nf, operatorId, countryIso2);214  if (!candidates.length) return { id: null, how: "create" };215  const geo = validGeo(nf.geo) ? nf.geo : null;216  const best = bestMatch(217    { name: nf.name, aliases: nf.aliases, operatorId, operatorName, countryIso2, city: nf.city ?? null, address: nf.address ?? null, lat: geo?.lat ?? null, lng: geo?.lng ?? null, geoPrecision: geo?.precision ?? null },218    candidates,219  );220  if (!best) return { id: null, how: "create" };221  const d: MatchDecision = decide(best.match.score);222  if (d === "merge") return { id: best.candidate.id, how: "merge", match: best };223  // campus vs one of its buildings: not a duplicate but a containment — create the record and link it to its parent224  if (best.match.reasons.includes("rule:campus-vs-building")) return { id: null, how: "campus_link", match: best };225  if (d === "pending") return { id: null, how: "pending", match: best };226  return { id: null, how: "create", match: best.match.score >= 0.3 ? best : null };227}228229function obs(value: unknown, p: Provenance, ctx: IngestContext): FieldObservation {230  return { value, sourceKind: ctx.run.sourceKind, confidence: p.confidence, isEstimate: !!p.isEstimate, observedAt: p.lastObserved || ctx.now, sourceId: p.sourceId || ctx.run.sourceId, url: p.url || ctx.doc?.url || null };231}232233function storedObs(prov: CurrentProvenance[], field: string, current: unknown): FieldObservation | null {234  if (current == null) return null;235  const b = backingObservation(prov, field, current);236  return { value: current, sourceKind: b?.sourceKind ?? null, confidence: b?.confidence ?? null, isEstimate: b?.isEstimate ?? false, observedAt: b?.lastObserved ?? null, sourceId: b?.sourceId ?? null, url: b?.url ?? null };237}238239/** Ingest one NormalizedFacility. */240export async function ingestFacility(tx: Tx, ctx: IngestContext, nf: NormalizedFacility): Promise<void> {241  const name = cleanText(nf.name);242  if (!name) throw new Error(`facility ${nf.key}: name is required`);243  if (!nf.key) throw new Error("facility key is required");244  const P = nf.provenance;245  const pf = (field: string) => provenanceFor(P, field, nf.facts);246  const url = P.url || ctx.doc?.url || "";247  if (!url) throw new Error(`facility ${nf.key}: provenance url is required`);248249  // --- related entities -----------------------------------------------------------------------------------------250  // operator: the source's value, else a conservative canonical-brand prefix of the name ("Equinix AM3" → Equinix),251  // recorded with method "inferred:name-prefix" and moderate confidence so the inference stays visible in provenance252  const inferred = nf.operatorName ? null : inferOperatorFromName(name);253  const operatorName = nf.operatorName ?? inferred?.name ?? null;254  const operatorProv: Provenance = inferred ? { ...pf("operatorName"), method: inferred.method, confidence: "moderate" } : pf("operatorName");255  if (inferred) ctx.stats.inferredOperators = (ctx.stats.inferredOperators ?? 0) + 1;256  const operator = operatorName ? await resolveOperator(tx, ctx, { name: operatorName, key: nf.operatorKey ?? null }) : null;257  const owner = nf.ownerName && normalizeName(nf.ownerName) !== normalizeName(operatorName ?? "") ? await resolveOperator(tx, ctx, { name: nf.ownerName }) : null;258  let geo = validGeo(nf.geo) ? nf.geo : null;259  let country = await safeCountry(tx, ctx, nf.countryIso2);260  if (!country && geo) country = await safeCountry(tx, ctx, countryFromPoint(geo.lat, geo.lng)); // offline Natural Earth polygons (OSM features rarely carry addr:country)261  // coordinates that fall in another country than the one the source states are wrong on one side — keep the262  // stated country, drop the point (a swapped lat/lng or a geocoder miss must never place a facility abroad)263  if (geo && nf.countryIso2 && geo.precision !== "metro") {264    const pointCountry = countryFromPoint(geo.lat, geo.lng);265    if (pointCountry && pointCountry !== nf.countryIso2.toUpperCase() && !BORDER_TOLERANT.has(`${nf.countryIso2.toUpperCase()}:${pointCountry}`)) {266      ctx.stats.geoCountryMismatch = (ctx.stats.geoCountryMismatch ?? 0) + 1;267      geo = null;268    }269  }270  const metro = await assignMetro(tx, { lat: geo?.lat, lng: geo?.lng, city: nf.city, countryIso2: country });271  country = country ?? metro.countryIso2;272273  // --- resolution -------------------------------------------------------------------------------------------------274  const res = await resolveFacility(tx, ctx, nf, operator?.id ?? null, operator?.name ?? operatorName, country);275  const existing = res.id ? await loadFacility(tx, res.id) : null;276  const prov = existing ? await loadCurrentProvenance(tx, "facility", String(existing.id)) : [];277  const isNew = !existing;278  const id = existing ? String(existing.id) : newId("facility");279  const campus = nf.campusName ? await resolveCampus(tx, ctx, { name: nf.campusName, key: nf.campusKey ?? null, operatorId: operator?.id ?? null, countryIso2: country, city: nf.city ?? null, lat: geo?.lat ?? null, lng: geo?.lng ?? null }) : null;280281  // --- field merge -----------------------------------------------------------------------------------------------282  const before: Row = existing ? { ...existing } : {};283  const next: Row = existing ? { ...existing } : { id, geoPrecision: "unknown", status: "unknown", facilityType: "unknown", confidence: "unverified", completeness: 0, sourceCount: 0, isAi: false, isHyperscale: false, mwIsEstimate: false, externalIds: {}, certifications: [], aliases: [], recordScope: "facility", aiEvidence: "unknown" };284  const observed: ObservedField[] = [];285  // containment: a building that matched its campus record (or names its campus record by key) points at the parent286  if (res.how === "campus_link" && res.match) {287    const candIsCampus = isCampusName(res.match.candidate.name) && !isCampusName(name);288    if (candIsCampus) { next.parentFacilityId = res.match.candidate.id; next.recordScope = "building"; }289    else if (isCampusName(name) && !isCampusName(res.match.candidate.name) && !ctx.run.dryRun) {290      // the incoming record is the campus: the stored building becomes its child291      await tx.execute(sql`update facilities set parent_facility_id = ${id}, record_scope = 'building', updated_at = now() where id = ${res.match.candidate.id} and parent_facility_id is null`);292    }293    ctx.stats.campusLinks = (ctx.stats.campusLinks ?? 0) + 1;294  }295  if (nf.parentFacilityKey && !next.parentFacilityId) { const pid = await facilityIdForKey(tx, ctx, nf.parentFacilityKey); if (pid && pid !== id) { next.parentFacilityId = pid; next.recordScope = "building"; } }296  if (nf.developerName) { const dev = await resolveOperator(tx, ctx, { name: nf.developerName }); if (dev) { next.developerId = dev.id; observed.push({ field: "developerName", value: dev.name, provenance: pf("developerName") }); } }297  if (nf.landownerName) { const lo = await resolveOperator(tx, ctx, { name: nf.landownerName }); if (lo) { next.landownerId = lo.id; observed.push({ field: "landownerName", value: lo.name, provenance: pf("landownerName") }); } }298  const incoming: Record<string, unknown> = {299    name,300    city: cleanText(nf.city),301    regionName: cleanText(nf.regionName),302    address: cleanText(nf.address),303    postalCode: cleanText(nf.postalCode),304    status: nf.status && nf.status !== "unknown" ? nf.status : null,305    facilityType: nf.facilityType && nf.facilityType !== "unknown" ? nf.facilityType : null,306    tier: cleanText(nf.tier),307    buildingSqm: nf.buildingSqm ?? null,308    siteAreaHa: nf.siteAreaHa ?? null,309    rackCount: nf.rackCount ?? null,310    pue: nf.pue ?? null,311    coolingType: cleanText(nf.coolingType),312    renewableClaim: cleanText(nf.renewableClaim),313    openedOn: nf.openedOn ?? null,314    constructionStartedOn: nf.constructionStartedOn ?? null,315    announcedOn: nf.announcedOn ?? null,316    website: normalizeWebsite(nf.website),317    description: cleanText(nf.description),318  };319  for (const field of SCALAR_FIELDS) {320    const v = incoming[field];321    if (v == null) continue;322    const p = pf(field);323    observed.push({ field, value: v, provenance: p });324    // an opening year before 1980 from a crowd-sourced / dataset source is the building's construction date325    // (OSM `start_date`), not the data center's: keep the observation in provenance, never in the column326    if (field === "openedOn" && !isPrimaryKind(ctx.run.sourceKind) && ctx.run.connectorId !== "wikidata" && Number(String(v).slice(0, 4)) < 1980) {327      ctx.stats.implausibleDropped = (ctx.stats.implausibleDropped ?? 0) + 1;328      continue;329    }330    const dec = shouldReplace(obs(v, p, ctx), storedObs(prov, field, existing?.[field]));331    if (dec.replace) next[field] = v;332  }333  // capacity figures go through the claim store: scope + semantics from the supporting sentence (structured operator334  // specs have none → the record's own scope), sanity engine, then the authority policy for the column335  const recordScope = inferRecordScope(nf);336  const campusDesignation = isCampusName(name, nf.campusName);337  const MW_PREDICATE: Record<(typeof MW_FIELDS)[number], CapacityPredicate> = { itCapacityMw: "it_capacity_mw", totalPowerMw: "current_power_mw", plannedPowerMw: "planned_power_mw" };338  for (const field of MW_FIELDS) {339    const v = nf[field];340    if (v == null || !Number.isFinite(v) || v <= 0) continue;341    const p = pf(field);342    const context = nf.claimContext?.[field] ?? null;343    const structured = !context; // a parser read a spec table / JSON field: the figure describes the record itself344    const sc = structured ? { scope: recordScope as ClaimScope, reason: "structured:record" } : classifyScope(context, recordScope as ClaimScope);345    const sem = structured ? { predicate: MW_PREDICATE[field], reason: "field" } : classifyCapacitySemantics(context);346    const predicate: CapacityPredicate = sem.predicate ?? MW_PREDICATE[field];347    const evidence = context ? findEvidence(context, v, "mw") ?? { text: context.slice(0, 600), start: 0, end: Math.min(600, context.length) } : null;348    const decision = await recordCapacityClaim(tx, ctx, "facility", id, { field, predicate, value: v, scope: sc.scope, scopeReason: sc.reason, evidence, context, previous: (existing?.[field] as number | null) ?? null, recordScope, campusDesignation, provenance: p, semanticsDefaulted: !structured && !sem.predicate });349    if (!decision.assign) { ctx.stats.unscopedClaims = (ctx.stats.unscopedClaims ?? 0) + 1; continue; }350    // utility / grid / ultimate figures live in their own columns, never in IT / total / planned351    const column = CAPACITY_COLUMN[predicate];352    const target = column && column !== "itCapacityMw" && column !== "totalPowerMw" && column !== "plannedPowerMw" ? column : field;353    observed.push({ field: target, value: v, provenance: { ...p, note: p.note ?? `${predicate} · ${sc.scope}` } });354    const dec = shouldReplaceMw(obs(v, p, ctx), storedObs(prov, target, existing?.[target]));355    if (dec.replace) { next[target] = v; if (target === field) { next.capacityScope = sc.scope; next.capacitySemantics = predicate; } }356  }357  next.recordScope = recordScope;358  // AI evidence: graded from the source text, never from one keyword359  {360    const ai = classifyAiEvidence(`${nf.aiEvidence === "confirmed" ? "AI campus" : ""} ${name} ${nf.description ?? ""} ${nf.facilityType === "ai" || nf.facilityType === "hpc" ? "AI data center" : ""}`);361    const level = nf.aiEvidence && nf.aiEvidence !== "unknown" ? nf.aiEvidence : ai.level;362    const rank: Record<string, number> = { unknown: 0, associated: 1, likely: 2, confirmed: 3 };363    if ((rank[level] ?? 0) > (rank[String(next.aiEvidence ?? "unknown")] ?? 0)) next.aiEvidence = level;364    if (level === "confirmed" || level === "likely") { next.isAi = true; if (nf.isAi !== true) observed.push({ field: "isAi", value: true, provenance: { ...P, method: `ai-evidence:${level}${ai.evidence ? `:${ai.evidence}` : ""}` } }); }365  }366  // operator / owner: authority policy like other fields (a news source never re-assigns an operator set by the operator itself);367  // an operator inferred from the name never replaces one stated by a source368  if (operator) {369    observed.push({ field: "operatorId", value: operator.id, provenance: operatorProv });370    observed.push({ field: "operatorName", value: operator.name, provenance: operatorProv });371    if (inferred ? !existing?.operatorId : shouldReplace(obs(operator.id, operatorProv, ctx), storedObs(prov, "operatorId", existing?.operatorId)).replace) next.operatorId = operator.id;372  }373  if (owner) {374    observed.push({ field: "ownerId", value: owner.id, provenance: pf("ownerName") });375    if (shouldReplace(obs(owner.id, pf("ownerName"), ctx), storedObs(prov, "ownerId", existing?.ownerId)).replace) next.ownerId = owner.id;376  }377  if (campus) next.campusId = next.campusId ?? campus.id;378  // geo: never degrade precision379  if (geo) {380    const p = pf("geo");381    observed.push({ field: "geo", value: { lat: geo.lat, lng: geo.lng, precision: geo.precision, source: geo.source }, provenance: p });382    const same = storedObs(prov, "geo", existing?.lat != null ? { lat: existing.lat, lng: existing.lng, precision: existing.geoPrecision, source: existing.geoSource } : null);383    const sameSource = !!same?.sourceId && same.sourceId === (p.sourceId || ctx.run.sourceId);384    if (shouldReplaceGeo(geo.precision, existing?.geoPrecision as string | null, validLatLng(existing?.lat, existing?.lng), sameSource)) {385      next.lat = geo.lat;386      next.lng = geo.lng;387      next.geoPrecision = geo.precision;388      next.geoSource = geo.source;389    }390  }391  if (validLatLng(next.lat, next.lng)) next.geohash = geohash(next.lat as number, next.lng as number, 7);392  // country / metro (metro follows the coordinates we keep)393  if (country) {394    observed.push({ field: "countryIso2", value: country, provenance: pf("countryIso2") });395    if (!next.countryIso2 || shouldReplace(obs(country, pf("countryIso2"), ctx), storedObs(prov, "countryIso2", existing?.countryIso2)).replace) next.countryIso2 = country;396  }397  {398    const m = validLatLng(next.lat, next.lng) ? await assignMetro(tx, { lat: next.lat as number, lng: next.lng as number, city: next.city as string | null, countryIso2: next.countryIso2 as string | null }) : metro;399    if (m.metroId) next.metroId = m.metroId;400    if (!next.countryIso2 && m.countryIso2) next.countryIso2 = m.countryIso2;401  }402  // flags accumulate (true sticks), hyperscaler operator implies hyperscale403  if (nf.isAi === true) {404    next.isAi = true;405    observed.push({ field: "isAi", value: true, provenance: pf("isAi") });406  }407  if (nf.isHyperscale === true || isHyperscalerName(operatorName)) {408    next.isHyperscale = true;409    if (nf.isHyperscale === true) observed.push({ field: "isHyperscale", value: true, provenance: pf("isHyperscale") });410  }411  if (nf.certifications?.length) {412    next.certifications = uniqStrings([...((next.certifications as string[]) ?? []), ...nf.certifications]);413    observed.push({ field: "certifications", value: nf.certifications, provenance: pf("certifications") });414  }415  if (nf.externalIds && Object.keys(nf.externalIds).length) {416    next.externalIds = { ...((next.externalIds as Record<string, string | number>) ?? {}), ...nf.externalIds };417    observed.push({ field: "externalIds", value: nf.externalIds, provenance: pf("externalIds") });418  }419  if (nf.carriers?.length) observed.push({ field: "carriers", value: nf.carriers, provenance: pf("carriers") });420  if (nf.cloudProviders?.length) observed.push({ field: "cloudProviders", value: nf.cloudProviders, provenance: pf("cloudProviders") });421  if (nf.ixps?.length) observed.push({ field: "ixps", value: nf.ixps, provenance: pf("ixps") });422  if (nf.campusName) observed.push({ field: "campusName", value: nf.campusName, provenance: pf("campusName") });423424  // --- persist ----------------------------------------------------------------------------------------------------425  const normalizedName = normalizeName(String(next.name));426  if (isNew) {427    const slug = await uniqueSlug(tx, "facilities", `${operator && !normalizeName(name).includes(normalizeName(operator.name)) ? `${operator.name} ` : ""}${name}`, (next.city as string | null) ?? (next.countryIso2 as string | null));428    next.slug = slug;429    next.confidence = res.how === "pending" ? "unverified" : "moderate";430    await tx.execute(sql`insert into facilities (id, slug, name, normalized_name, operator_id, owner_id, campus_id, metro_id, country_iso2, city, region_name, address, postal_code, lat, lng, geo_precision, geo_source, geohash,431        status, facility_type, tier, building_sqm, site_area_ha, it_capacity_mw, total_power_mw, planned_power_mw, mw_is_estimate, rack_count, pue, cooling_type, renewable_claim, opened_on, construction_started_on, announced_on,432        website, description, is_ai, is_hyperscale, certifications, confidence, completeness, external_ids, source_count, first_seen,433        parent_facility_id, record_scope, ai_evidence, utility_capacity_mw, grid_connection_mw, ultimate_campus_mw, capacity_scope, capacity_semantics, developer_id, landowner_id)434      values (${id}, ${slug}, ${next.name}, ${normalizedName}, ${next.operatorId ?? null}, ${next.ownerId ?? null}, ${next.campusId ?? null}, ${next.metroId ?? null}, ${next.countryIso2 ?? null}, ${next.city ?? null}, ${next.regionName ?? null},435        ${next.address ?? null}, ${next.postalCode ?? null}, ${next.lat ?? null}, ${next.lng ?? null}, ${next.geoPrecision ?? "unknown"}, ${next.geoSource ?? null}, ${next.geohash ?? null},436        ${next.status ?? "unknown"}, ${next.facilityType ?? "unknown"}, ${next.tier ?? null}, ${next.buildingSqm ?? null}, ${next.siteAreaHa ?? null}, ${next.itCapacityMw ?? null}, ${next.totalPowerMw ?? null}, ${next.plannedPowerMw ?? null}, false,437        ${next.rackCount ?? null}, ${next.pue ?? null}, ${next.coolingType ?? null}, ${next.renewableClaim ?? null}, ${next.openedOn ?? null}, ${next.constructionStartedOn ?? null}, ${next.announcedOn ?? null},438        ${next.website ?? null}, ${next.description ?? null}, ${!!next.isAi}, ${!!next.isHyperscale}, ${textArray((next.certifications as string[]) ?? [])}::text[], ${next.confidence}, 0, ${JSON.stringify(next.externalIds ?? {})}::jsonb, 0, ${ctx.now},439        ${next.parentFacilityId ?? null}, ${next.recordScope ?? "facility"}, ${next.aiEvidence ?? "unknown"}, ${next.utilityCapacityMw ?? null}, ${next.gridConnectionMw ?? null}, ${next.ultimateCampusMw ?? null}, ${next.capacityScope ?? null}, ${next.capacitySemantics ?? null}, ${next.developerId ?? null}, ${next.landownerId ?? null})`);440  } else {441    const sets = [];442    for (const [camel, col] of Object.entries(COLUMNS)) {443      if (["confidence", "completeness", "sourceCount", "lastVerified", "mwIsEstimate"].includes(camel)) continue;444      const a = before[camel];445      const b = next[camel];446      if (JSON.stringify(a ?? null) === JSON.stringify(b ?? null)) continue;447      if (camel === "certifications") sets.push(sql`${sql.identifier(col)} = ${textArray((b as string[]) ?? [])}::text[]`);448      else if (camel === "externalIds") sets.push(sql`${sql.identifier(col)} = ${JSON.stringify(b ?? {})}::jsonb`);449      else sets.push(sql`${sql.identifier(col)} = ${b as string | number | boolean | null}`);450    }451    if (String(before.name) !== String(next.name)) sets.push(sql`normalized_name = ${normalizedName}`);452    if (sets.length) {453      sets.push(sql`updated_at = now()`);454      await tx.execute(sql`update facilities set ${sql.join(sets, sql`, `)} where id = ${id}`);455      ctx.stats.updated++;456    } else ctx.stats.unchanged++;457    if (res.how === "merge") ctx.stats.merged++;458  }459  await upsertKey(tx, ctx, nf.key, "facility", id);460  addRef(ctx, "facility", id);461  bump(ctx, "facility");462463  // aliases: incoming aliases, the incoming name when it differs from the kept name, the previous name when renamed464  const aliasValues = uniqStrings([...(nf.aliases ?? []), name !== next.name ? name : null, existing && String(existing.name) !== String(next.name) ? String(existing.name) : null]).filter((a) => normalizeName(a) && a !== next.name);465  for (const a of aliasValues) await tx.execute(sql`insert into facility_aliases (facility_id, alias, normalized, source_id) values (${id}, ${a}, ${normalizeName(a)}, ${ctx.run.sourceId}) on conflict do nothing`);466467  // provenance for every observed field, then mark which observation backs each displayed value468  await writeProvenance(tx, ctx, "facility", id, observed, nf.key);469  await markWinners(tx, ctx, "facility", id, [...SCALAR_FIELDS, ...MW_FIELDS, "utilityCapacityMw", "gridConnectionMw", "ultimateCampusMw", "countryIso2"].map((f) => ({ field: f, value: next[f] })).concat([{ field: "operatorId", value: next.operatorId }, { field: "ownerId", value: next.ownerId }]));470471  // tenants / IXPs472  if (nf.carriers?.length) await attachTenants(tx, ctx, id, { names: nf.carriers, role: "carrier" });473  if (nf.cloudProviders?.length) await attachTenants(tx, ctx, id, { names: nf.cloudProviders, role: "cloud" });474  if (nf.carriers?.length || nf.cloudProviders?.length) await refreshTenantCounts(tx, id);475  if (nf.ixps?.length) await attachIxpsByName(tx, ctx, id, nf.ixps, (next.countryIso2 as string | null) ?? null);476477  // derived: mw_is_estimate, source_count, last_verified, confidence, completeness478  await refreshDerived(tx, id, res.how === "pending" ? "unverified" : null);479480  // --- reconciliation bookkeeping --------------------------------------------------------------------------------481  if (res.how === "merge" && res.match) {482    await tx.execute(sql`insert into entity_matches (id, connector_id, candidate_key, candidate, matched_facility_id, score, reasons, status, decided_by, decided_at)483      values (${newId("match")}, ${ctx.run.connectorId}, ${nf.key}, ${JSON.stringify(candidateJson(nf, id))}::jsonb, ${res.match.candidate.id}, ${res.match.match.score}, ${textArray(res.match.match.reasons)}::text[], 'auto_merged', 'system', ${ctx.now})`);484  } else if (res.how === "pending" && res.match) {485    await tx.execute(sql`insert into entity_matches (id, connector_id, candidate_key, candidate, matched_facility_id, score, reasons, status)486      values (${newId("match")}, ${ctx.run.connectorId}, ${nf.key}, ${JSON.stringify(candidateJson(nf, id))}::jsonb, ${res.match.candidate.id}, ${res.match.match.score}, ${textArray(res.match.match.reasons)}::text[], 'pending')`);487    ctx.stats.pendingMatches++;488  } else if (res.how === "create" && res.match) {489    await tx.execute(sql`insert into entity_matches (id, connector_id, candidate_key, candidate, matched_facility_id, score, reasons, status, decided_by, decided_at)490      values (${newId("match")}, ${ctx.run.connectorId}, ${nf.key}, ${JSON.stringify(candidateJson(nf, id))}::jsonb, ${res.match.candidate.id}, ${res.match.match.score}, ${textArray(res.match.match.reasons)}::text[], 'auto_created', 'system', ${ctx.now})`);491  } else if (res.how === "campus_link" && res.match) {492    await tx.execute(sql`insert into entity_matches (id, connector_id, candidate_key, candidate, matched_facility_id, score, reasons, status, decided_by, decided_at)493      values (${newId("match")}, ${ctx.run.connectorId}, ${nf.key}, ${JSON.stringify(candidateJson(nf, id))}::jsonb, ${res.match.candidate.id}, ${res.match.match.score}, ${textArray(res.match.match.reasons)}::text[], 'related_campus', 'system', ${ctx.now})`);494  }495  // duplicate candidates awaiting review are a quality flag too (they inflate every aggregate until decided)496  if (res.how === "pending" && res.match) {497    await writeQualityFlags(tx, ctx, "facility", id, [{ code: "possible_duplicate", severity: "warn", message: `possible duplicate of ${res.match.candidate.name} (score ${res.match.match.score})`, field: "name" }], { details: { matchedFacilityId: res.match.candidate.id, reasons: res.match.match.reasons } });498  }499500  // --- events ----------------------------------------------------------------------------------------------------501  const conf = (await tx.execute(sql`select confidence from facilities where id = ${id}`))[0]?.confidence as ConfidenceLevel | undefined;502  if (isNew) {503    ctx.stats.created++;504    await maybeOperatorExpansion(tx, ctx, { entityType: "facility", entityId: id, entityName: String(next.name), operatorId: (next.operatorId as string | null) ?? null, operatorName: operator?.name ?? null, countryIso2: (next.countryIso2 as string | null) ?? null, metroId: (next.metroId as string | null) ?? null, url, confidence: conf ?? "moderate", isAi: !!next.isAi });505    const mw = bestMw({ itCapacityMw: next.itCapacityMw as number | null, totalPowerMw: next.totalPowerMw as number | null, plannedPowerMw: next.plannedPowerMw as number | null });506    const pipeline = isPipelineStatus(next.status as string | null);507    const opName = operator?.name ?? null;508    const where = [next.city, next.countryIso2].filter(Boolean).join(", ");509    await recordEvent(tx, ctx, {510      entityType: "facility",511      entityId: id,512      eventType: "facility_discovered",513      title: `New facility indexed: ${next.name}${opName && !String(next.name).toLowerCase().includes(opName.toLowerCase()) ? ` (${opName})` : ""}${where ? ` — ${where}` : ""}`,514      summary: [pipeline ? `Status: ${String(next.status).replace(/_/g, " ")}.` : null, mw != null ? `${mw} MW${next.plannedPowerMw != null && next.itCapacityMw == null && next.totalPowerMw == null ? " planned" : ""}.` : null].filter(Boolean).join(" ") || null,515      newValue: { name: next.name, status: next.status, mw },516      significance: discoverySignificance(mw, pipeline),517      confidence: conf ?? "moderate",518      url,519      countryIso2: next.countryIso2 as string | null,520      operatorId: next.operatorId as string | null,521      metroId: next.metroId as string | null,522      reviewStatus: res.how === "pending" ? "pending" : "auto",523    });524  } else {525    const names = await operatorNames(tx, ctx, [before.operatorId as string | null, next.operatorId as string | null, before.ownerId as string | null, next.ownerId as string | null]);526    await emitDiffEvents(tx, ctx, {527      entityType: "facility",528      entityId: id,529      entityName: String(next.name),530      before,531      after: next,532      specs: TRACKED_FACILITY_FIELDS,533      url,534      confidence: conf ?? "moderate",535      countryIso2: next.countryIso2 as string | null,536      operatorId: next.operatorId as string | null,537      metroId: next.metroId as string | null,538      operatorNames: names,539    });540  }541}542543function candidateJson(nf: NormalizedFacility, createdFacilityId: string): Record<string, unknown> {544  return {545    createdFacilityId,546    key: nf.key,547    name: nf.name,548    operatorName: nf.operatorName ?? null,549    address: nf.address ?? null,550    city: nf.city ?? null,551    countryIso2: nf.countryIso2 ?? null,552    geo: nf.geo ?? null,553    status: nf.status ?? null,554    itCapacityMw: nf.itCapacityMw ?? null,555    totalPowerMw: nf.totalPowerMw ?? null,556    plannedPowerMw: nf.plannedPowerMw ?? null,557    externalIds: nf.externalIds ?? {},558    url: nf.provenance.url,559    sourceId: nf.provenance.sourceId,560  };561}562563/** Recompute mw_is_estimate, source_count, last_verified, confidence and completeness from provenance + columns. */564export async function refreshDerived(tx: Tx, id: string, forceConfidence: ConfidenceLevel | null = null): Promise<void> {565  const prov = await loadCurrentProvenance(tx, "facility", id);566  const row = await loadFacility(tx, id);567  if (!row) return;568  const mwRows = prov.filter((r) => (MW_FIELDS as readonly string[]).includes(r.field));569  const hasMw = row.itCapacityMw != null || row.totalPowerMw != null || row.plannedPowerMw != null;570  const mwIsEstimate = hasMw && mwRows.length > 0 && mwRows.every((r) => r.isEstimate);571  const { sourceIds, kinds, lastPrimaryObserved } = summarizeSources(prov);572  const onlyEstimates = prov.length > 0 && prov.every((r) => r.isEstimate);573  // a facility created from an ambiguous match stays "unverified" until the admin decides574  const pending = forceConfidence ? [] : await tx.execute(sql`select 1 from entity_matches where status = 'pending' and candidate->>'createdFacilityId' = ${id} limit 1`);575  const confidence = forceConfidence ?? (pending.length ? "unverified" : facilityConfidence({ sourceKinds: kinds, sourceCount: sourceIds.length, lastVerifiedIso: lastPrimaryObserved, onlyEstimates }));576  const counts = (await tx.execute(sql`select carriers_count, ixp_count, (select coalesce(max(priority), 0) from quality_flags q where q.entity_type = 'facility' and q.entity_id = ${id} and q.status = 'open') as review_priority from facilities where id = ${id}`))[0];577  const completeness = completenessScore({578    geoPrecision: row.geoPrecision as string,579    hasCoords: validLatLng(row.lat, row.lng),580    operatorId: row.operatorId as string | null,581    address: row.address as string | null,582    status: row.status as string,583    itCapacityMw: row.itCapacityMw as number | null,584    totalPowerMw: row.totalPowerMw as number | null,585    plannedPowerMw: row.plannedPowerMw as number | null,586    mwIsEstimate,587    facilityType: row.facilityType as string,588    openedOn: row.openedOn as string | null,589    website: row.website as string | null,590    description: row.description as string | null,591    carriersCount: (counts?.carriers_count as number | null) ?? null,592    ixpCount: (counts?.ixp_count as number | null) ?? null,593  });594  await tx.execute(sql`update facilities set mw_is_estimate = ${mwIsEstimate}, source_count = ${sourceIds.length}, last_verified = ${lastPrimaryObserved}, confidence = ${confidence}, completeness = ${completeness}, review_priority = ${Number(counts?.review_priority ?? 0)} where id = ${id}`);595}596597/** Admin approval of a pending match: fold `fromId` into `intoId` (keys, aliases, provenance, links, projects). */598export async function mergeFacilities(tx: Tx, fromId: string, intoId: string, decidedBy = "admin"): Promise<void> {599  if (fromId === intoId) return;600  const from = await loadFacility(tx, fromId);601  if (!from) throw new Error(`facility ${fromId} not found`);602  await tx.execute(sql`update entity_keys set entity_id = ${intoId} where entity_type = 'facility' and entity_id = ${fromId}`);603  await tx.execute(sql`insert into facility_aliases (facility_id, alias, normalized, source_id) select ${intoId}, alias, normalized, source_id from facility_aliases where facility_id = ${fromId} on conflict do nothing`);604  await tx.execute(sql`insert into facility_aliases (facility_id, alias, normalized, source_id) values (${intoId}, ${String(from.name)}, ${normalizeName(String(from.name))}, null) on conflict do nothing`);605  await tx.execute(sql`update provenance set entity_id = ${intoId} where entity_type = 'facility' and entity_id = ${fromId} and not exists (select 1 from provenance q where q.entity_type = 'facility' and q.entity_id = ${intoId} and q.field = provenance.field and q.source_id = provenance.source_id and q.url = provenance.url)`);606  await tx.execute(sql`delete from provenance where entity_type = 'facility' and entity_id = ${fromId}`);607  await tx.execute(sql`insert into facility_ixps (facility_id, ixp_id, source_id) select ${intoId}, ixp_id, source_id from facility_ixps where facility_id = ${fromId} on conflict do nothing`);608  await tx.execute(sql`insert into facility_tenants (facility_id, operator_id, role, asn, source_id) select ${intoId}, operator_id, role, asn, source_id from facility_tenants where facility_id = ${fromId} on conflict do nothing`);609  await tx.execute(sql`update projects set facility_id = ${intoId} where facility_id = ${fromId}`);610  await tx.execute(sql`update events set entity_id = ${intoId} where entity_type = 'facility' and entity_id = ${fromId}`);611  await tx.execute(sql`update facilities set external_ids = external_ids || (select external_ids from facilities where id = ${fromId}) where id = ${intoId}`);612  await tx.execute(sql`update facilities set merged_into = ${intoId}, updated_at = now() where id = ${fromId}`);613  await tx.execute(sql`update entity_matches set status = 'approved', decided_by = ${decidedBy}, decided_at = now() where status = 'pending' and matched_facility_id = ${intoId} and candidate->>'createdFacilityId' = ${fromId}`);614  await refreshTenantCounts(tx, intoId);615  await tx.execute(sql`update facilities f set ixp_count = (select count(*) from facility_ixps x where x.facility_id = f.id) where f.id = ${intoId}`);616  await refreshDerived(tx, intoId);617}618619620/** Extraction debugger: how would this normalized facility resolve right now (candidates, scores, decision)? Read-only. */621export async function previewFacilityResolution(tx: Tx, ctx: IngestContext, nf: NormalizedFacility): Promise<{ how: string; matchedId: string | null; candidates: Array<{ id: string; name: string; operatorName: string | null; city: string | null; score: number; reasons: string[]; distanceKm: number | null }> }> {622  const operator = nf.operatorName ? await resolveOperator(tx, ctx, { name: nf.operatorName, key: nf.operatorKey ?? null }) : null;623  const country = await safeCountry(tx, ctx, nf.countryIso2);624  const byKey = await facilityIdForKey(tx, ctx, nf.key);625  const byExt = byKey ? null : await byExternalIds(tx, nf.externalIds);626  const candidates = await loadCandidates(tx, nf, operator?.id ?? null, country);627  const geo = validGeo(nf.geo) ? nf.geo : null;628  const probe = { name: nf.name, aliases: nf.aliases, operatorId: operator?.id ?? null, operatorName: operator?.name ?? nf.operatorName ?? null, countryIso2: country, city: nf.city ?? null, address: nf.address ?? null, lat: geo?.lat ?? null, lng: geo?.lng ?? null, geoPrecision: geo?.precision ?? null };629  const scored = candidates.map((c) => { const m = scoreFacilityMatch(probe, c); return { id: c.id, name: c.name, operatorName: c.operatorName ?? null, city: c.city ?? null, score: m.score, reasons: m.reasons, distanceKm: m.distanceKm ?? null }; }).sort((a, b) => b.score - a.score).slice(0, 10);630  const best = scored[0];631  const how = byKey ? "key" : byExt ? "external_id" : !best ? "create" : best.reasons.includes("rule:campus-vs-building") ? "campus_link" : decide(best.score);632  return { how, matchedId: byKey ?? byExt ?? (best && (how === "merge") ? best.id : null), candidates: scored };633}634635/**636 * First record of an operator in a country (or a metro) → `operator_expansion` event ("Which companies are expanding637 * into new countries?"). Deterministic: counts other live facilities / projects of the operator in that country.638 */639export async function maybeOperatorExpansion(tx: Tx, ctx: IngestContext, i: { entityType: "facility" | "project"; entityId: string; entityName: string; operatorId: string | null; operatorName: string | null; countryIso2: string | null; metroId: string | null; url: string; confidence: ConfidenceLevel; isAi: boolean }): Promise<void> {640  if (!i.operatorId || !i.countryIso2) return;641  const others = (await tx.execute(sql`select (select count(*) from facilities f where f.operator_id = ${i.operatorId} and f.country_iso2 = ${i.countryIso2} and f.merged_into is null and f.id <> ${i.entityId})::int642      + (select count(*) from projects p where p.operator_id = ${i.operatorId} and p.country_iso2 = ${i.countryIso2} and p.merged_into is null and not p.hidden and p.id <> ${i.entityId})::int as n`))[0];643  if (Number(others?.n ?? 0) > 0) return;644  const country = (await tx.execute(sql`select name from countries where iso2 = ${i.countryIso2}`))[0]?.name;645  await recordEvent(tx, ctx, {646    entityType: i.entityType,647    entityId: i.entityId,648    eventType: "operator_expansion",649    title: `${i.operatorName ?? "Operator"} enters ${String(country ?? i.countryIso2)}: ${i.entityName}`.slice(0, 300),650    summary: `First ${i.entityType} of ${i.operatorName ?? "this operator"} indexed in ${String(country ?? i.countryIso2)}.`,651    newValue: { countryIso2: i.countryIso2, operatorId: i.operatorId, [i.entityType === "facility" ? "facilityId" : "projectId"]: i.entityId },652    significance: 70,653    confidence: i.confidence,654    url: i.url,655    countryIso2: i.countryIso2,656    operatorId: i.operatorId,657    metroId: i.metroId,658    projectId: i.entityType === "project" ? i.entityId : null,659    isAi: i.isAi,660    fingerprint: sha256(`operator_expansion|${i.operatorId}|${i.countryIso2}`),661  });662}663