/** * Facility reconciliation + persistence (docs/RECONCILIATION.md): * 1. entity_keys exact key → same facility (update) * 2. shared external ids (peeringdb_fac, osm, wikidata, …) → same facility * 3. candidates within the match radius of the coordinates, or same normalized name in the same country, or same * operator in the same country (facility-code matching) → weighted score * ≥ 0.92 auto-merge · 0.60–0.92 create-as-new + pending entity_matches row (nothing is lost, admin merges) · < 0.60 create * Field merge follows the authority policy in match.ts; coordinates never lose precision. */ import { sql, textArray } from "@dci/db"; import { cleanText, geohash, newId, normalizeName, sha256, validLatLng, type ConfidenceLevel, type NormalizedFacility, type Provenance } from "@dci/core"; import { addRef, bump, isPrimaryKind, provenanceFor, safeCountry, uniqueSlug, uniqStrings, normalizeWebsite, type IngestContext, type Tx } from "./common.js"; import { isHyperscalerName } from "./canonical-operators.js"; import { inferOperatorFromName } from "./operator-inference.js"; /** Country pairs whose city-level points legitimately straddle a border / enclave (stated:point-in). */ export const BORDER_TOLERANT: ReadonlySet = 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"]); import { resolveCampus } from "./campuses.js"; import { discoverySignificance, emitDiffEvents, recordEvent, TRACKED_FACILITY_FIELDS } from "./events.js"; import { bboxAround, validGeo } from "./geo.js"; import { countryFromPoint } from "./country-lookup.js"; import { facilityIdForKey, upsertKey } from "./keys.js"; import { attachIxpsByName } from "./ixps.js"; import { bestMatch, bestMw, completenessScore, decide, facilityConfidence, isPipelineStatus, scoreFacilityMatch, shouldReplace, shouldReplaceGeo, shouldReplaceMw, type FacilityCandidate, type FieldObservation, type MatchDecision, type MatchScore } from "./match.js"; import { CAPACITY_COLUMN, classifyAiEvidence, classifyCapacitySemantics, classifyScope, findEvidence, type CapacityPredicate, type ClaimScope } from "@dci/core"; import { markWinners, recordCapacityClaim, writeQualityFlags } from "./claims.js"; import { assignMetro } from "./metros.js"; import { operatorNames, resolveOperator } from "./operators.js"; import { backingObservation, loadCurrentProvenance, summarizeSources, writeProvenance, type CurrentProvenance, type ObservedField } from "./provenance.js"; import { attachTenants, refreshTenantCounts } from "./tenants.js"; /** Stored facility columns we merge into (camelCase keys ↔ snake_case columns). */ const COLUMNS: Record = { name: "name", operatorId: "operator_id", ownerId: "owner_id", campusId: "campus_id", metroId: "metro_id", countryIso2: "country_iso2", city: "city", regionName: "region_name", address: "address", postalCode: "postal_code", lat: "lat", lng: "lng", geoPrecision: "geo_precision", geoSource: "geo_source", geohash: "geohash", status: "status", facilityType: "facility_type", tier: "tier", buildingSqm: "building_sqm", siteAreaHa: "site_area_ha", itCapacityMw: "it_capacity_mw", totalPowerMw: "total_power_mw", plannedPowerMw: "planned_power_mw", mwIsEstimate: "mw_is_estimate", rackCount: "rack_count", pue: "pue", coolingType: "cooling_type", renewableClaim: "renewable_claim", openedOn: "opened_on", constructionStartedOn: "construction_started_on", announcedOn: "announced_on", website: "website", description: "description", isAi: "is_ai", isHyperscale: "is_hyperscale", certifications: "certifications", confidence: "confidence", completeness: "completeness", externalIds: "external_ids", sourceCount: "source_count", lastVerified: "last_verified", parentFacilityId: "parent_facility_id", recordScope: "record_scope", aiEvidence: "ai_evidence", utilityCapacityMw: "utility_capacity_mw", gridConnectionMw: "grid_connection_mw", ultimateCampusMw: "ultimate_campus_mw", capacityScope: "capacity_scope", capacitySemantics: "capacity_semantics", developerId: "developer_id", landownerId: "landowner_id", }; const SCALAR_FIELDS = ["name", "city", "regionName", "address", "postalCode", "status", "facilityType", "tier", "buildingSqm", "siteAreaHa", "rackCount", "pue", "coolingType", "renewableClaim", "openedOn", "constructionStartedOn", "announcedOn", "website", "description"] as const; const MW_FIELDS = ["itCapacityMw", "totalPowerMw", "plannedPowerMw"] as const; type Row = Record; async function loadFacility(tx: Tx, id: string): Promise { 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}`); if (!rows[0]) return null; const r = rows[0]; // camelCase view of the row const out: Row = {}; for (const [camel, col] of Object.entries(COLUMNS)) out[camel] = r[col] ?? null; out.id = r.id; out.slug = r.slug; out.mergedInto = r.merged_into ?? null; out.aliases = Array.isArray(r.alias_list) ? r.alias_list : []; out.normalizedName = r.normalized_name; return out; } async function followMerged(tx: Tx, id: string): Promise { let cur = id; for (let i = 0; i < 5; i++) { const r = await tx.execute(sql`select merged_into from facilities where id = ${cur}`); const m = r[0]?.merged_into; if (!m) break; cur = String(m); } return cur; } /** * External-id namespaces that identify ONE facility. Anything else stored in `externalIds` (an operator's Wikidata * id, a state name, an investment figure, a shared `ref` tag…) is informative but must never link records — a * shared `operator_wikidata` once folded 121 AWS OpenStreetMap features into a single facility. */ export const IDENTIFYING_EXTERNAL_ID_KEYS: ReadonlySet = new Set([ "osm", "wikidata", "wikipedia_en", "peeringdb_fac", "peeringdb", "pdb_fac", "dcmap", "geonames", "facebook_page", "meta_info_sheet", "meta_location", "google_detail_page", "google_location", "equinix_ibx", "digitalrealty_node", "digitalrealty_site_code", "ntt_slug", "cyrusone_slug", "stack_slug", "qts_slug", "vantage_slug", "coresite_code", "switch_slug", "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", ]); /** * ALLOWLIST ONLY. `state_code`, `region_code`, `postal_code`, `source_url`, `osm_ref` (a shared `ref` tag), `stack_campus` * (one campus, many buildings) or `investment_currency` all end in an id-looking suffix and would fold unrelated * facilities into one row (the `*_code` heuristic once folded 121 AWS OpenStreetMap features). A new connector that * needs its key to link records adds it here explicitly. */ export function isIdentifyingExternalId(key: string): boolean { return IDENTIFYING_EXTERNAL_ID_KEYS.has(key); } 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 === "" || !isIdentifyingExternalId(k)) continue; 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`); if (rows[0]) return String(rows[0].id); } return null; } async function loadCandidates(tx: Tx, nf: NormalizedFacility, operatorId: string | null, countryIso2: string | null): Promise { const norm = normalizeName(nf.name); const aliasNorms = uniqStrings([norm, ...(nf.aliases ?? []).map((a) => normalizeName(a))]).filter(Boolean); const geo = validGeo(nf.geo) ? nf.geo : null; const box = geo ? bboxAround(geo.lat, geo.lng, 40) : null; const conds = []; if (box) conds.push(sql`(f.lat between ${box.minLat} and ${box.maxLat} and f.lng between ${box.minLng} and ${box.maxLng})`); if (aliasNorms.length) { 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})`); conds.push(sql`exists (select 1 from facility_aliases a where a.facility_id = f.id and a.normalized in ${aliasNorms})`); } if (operatorId && countryIso2) conds.push(sql`(f.operator_id = ${operatorId} and f.country_iso2 = ${countryIso2})`); if (!conds.length) return []; // deterministic and relevance-ordered: same operator first, then same normalized name, then nearest — a dense metro // (Ashburn, Dallas, Singapore) has far more than 400 rows in a 40 km box and the true duplicate must be on the first page const orderGeo = geo ? sql`, ((f.lat - ${geo.lat}) * (f.lat - ${geo.lat}) + (f.lng - ${geo.lng}) * (f.lng - ${geo.lng})) asc nulls last` : sql``; const rows = await tx.execute(sql` 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, (select coalesce(array_agg(alias), '{}'::text[]) from facility_aliases a where a.facility_id = f.id) as aliases 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 `)}) 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 asc limit 400`); return rows.map((r) => ({ id: String(r.id), name: String(r.name), normalizedName: String(r.normalized_name), aliases: Array.isArray(r.aliases) ? (r.aliases as string[]) : [], operatorId: r.operator_id == null ? null : String(r.operator_id), operatorName: r.operator_name == null ? null : String(r.operator_name), countryIso2: r.country_iso2 == null ? null : String(r.country_iso2), city: r.city == null ? null : String(r.city), address: r.address == null ? null : String(r.address), lat: r.lat == null ? null : Number(r.lat), lng: r.lng == null ? null : Number(r.lng), geoPrecision: r.geo_precision == null ? null : String(r.geo_precision), externalIds: (r.external_ids as Record) ?? null, })); } interface Resolution { id: string | null; how: "key" | "external_id" | "merge" | "pending" | "create" | "campus_link"; match?: { candidate: FacilityCandidate; match: MatchScore } | null; } /** true when a facility name designates a campus / park / multi-building site rather than one building. */ export function isCampusName(name: string | null | undefined, campusName?: string | null): boolean { return /\b(campus|park|cluster|hub|gigafactory|complex|estate|mega ?site)\b/i.test(`${name ?? ""} ${campusName ?? ""}`); } /** Record scope from the source's statement or the name: campus designation → campus, building code (DC12, Hall 3, Building B) → building, else facility. */ export function inferRecordScope(nf: { name: string; recordScope?: string | null; campusName?: string | null }): "building" | "facility" | "campus" { if (nf.recordScope === "building" || nf.recordScope === "facility" || nf.recordScope === "campus") return nf.recordScope; if (isCampusName(nf.name)) return "campus"; 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"; return "facility"; } async function resolveFacility(tx: Tx, ctx: IngestContext, nf: NormalizedFacility, operatorId: string | null, operatorName: string | null, countryIso2: string | null): Promise { const byKey = await facilityIdForKey(tx, ctx, nf.key); if (byKey) return { id: byKey, how: "key" }; const byExt = await byExternalIds(tx, nf.externalIds); if (byExt) return { id: await followMerged(tx, byExt), how: "external_id" }; const candidates = await loadCandidates(tx, nf, operatorId, countryIso2); if (!candidates.length) return { id: null, how: "create" }; const geo = validGeo(nf.geo) ? nf.geo : null; const best = bestMatch( { 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 }, candidates, ); if (!best) return { id: null, how: "create" }; const d: MatchDecision = decide(best.match.score); if (d === "merge") return { id: best.candidate.id, how: "merge", match: best }; // campus vs one of its buildings: not a duplicate but a containment — create the record and link it to its parent if (best.match.reasons.includes("rule:campus-vs-building")) return { id: null, how: "campus_link", match: best }; if (d === "pending") return { id: null, how: "pending", match: best }; return { id: null, how: "create", match: best.match.score >= 0.3 ? best : null }; } function obs(value: unknown, p: Provenance, ctx: IngestContext): FieldObservation { 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 }; } function storedObs(prov: CurrentProvenance[], field: string, current: unknown): FieldObservation | null { if (current == null) return null; const b = backingObservation(prov, field, current); 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 }; } /** Ingest one NormalizedFacility. */ export async function ingestFacility(tx: Tx, ctx: IngestContext, nf: NormalizedFacility): Promise { const name = cleanText(nf.name); if (!name) throw new Error(`facility ${nf.key}: name is required`); if (!nf.key) throw new Error("facility key is required"); const P = nf.provenance; const pf = (field: string) => provenanceFor(P, field, nf.facts); const url = P.url || ctx.doc?.url || ""; if (!url) throw new Error(`facility ${nf.key}: provenance url is required`); // --- related entities ----------------------------------------------------------------------------------------- // operator: the source's value, else a conservative canonical-brand prefix of the name ("Equinix AM3" → Equinix), // recorded with method "inferred:name-prefix" and moderate confidence so the inference stays visible in provenance const inferred = nf.operatorName ? null : inferOperatorFromName(name); const operatorName = nf.operatorName ?? inferred?.name ?? null; const operatorProv: Provenance = inferred ? { ...pf("operatorName"), method: inferred.method, confidence: "moderate" } : pf("operatorName"); if (inferred) ctx.stats.inferredOperators = (ctx.stats.inferredOperators ?? 0) + 1; const operator = operatorName ? await resolveOperator(tx, ctx, { name: operatorName, key: nf.operatorKey ?? null }) : null; const owner = nf.ownerName && normalizeName(nf.ownerName) !== normalizeName(operatorName ?? "") ? await resolveOperator(tx, ctx, { name: nf.ownerName }) : null; let geo = validGeo(nf.geo) ? nf.geo : null; let country = await safeCountry(tx, ctx, nf.countryIso2); if (!country && geo) country = await safeCountry(tx, ctx, countryFromPoint(geo.lat, geo.lng)); // offline Natural Earth polygons (OSM features rarely carry addr:country) // coordinates that fall in another country than the one the source states are wrong on one side — keep the // stated country, drop the point (a swapped lat/lng or a geocoder miss must never place a facility abroad) if (geo && nf.countryIso2 && geo.precision !== "metro") { const pointCountry = countryFromPoint(geo.lat, geo.lng); if (pointCountry && pointCountry !== nf.countryIso2.toUpperCase() && !BORDER_TOLERANT.has(`${nf.countryIso2.toUpperCase()}:${pointCountry}`)) { ctx.stats.geoCountryMismatch = (ctx.stats.geoCountryMismatch ?? 0) + 1; geo = null; } } const metro = await assignMetro(tx, { lat: geo?.lat, lng: geo?.lng, city: nf.city, countryIso2: country }); country = country ?? metro.countryIso2; // --- resolution ------------------------------------------------------------------------------------------------- const res = await resolveFacility(tx, ctx, nf, operator?.id ?? null, operator?.name ?? operatorName, country); const existing = res.id ? await loadFacility(tx, res.id) : null; const prov = existing ? await loadCurrentProvenance(tx, "facility", String(existing.id)) : []; const isNew = !existing; const id = existing ? String(existing.id) : newId("facility"); 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; // --- field merge ----------------------------------------------------------------------------------------------- const before: Row = existing ? { ...existing } : {}; 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" }; const observed: ObservedField[] = []; // containment: a building that matched its campus record (or names its campus record by key) points at the parent if (res.how === "campus_link" && res.match) { const candIsCampus = isCampusName(res.match.candidate.name) && !isCampusName(name); if (candIsCampus) { next.parentFacilityId = res.match.candidate.id; next.recordScope = "building"; } else if (isCampusName(name) && !isCampusName(res.match.candidate.name) && !ctx.run.dryRun) { // the incoming record is the campus: the stored building becomes its child 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`); } ctx.stats.campusLinks = (ctx.stats.campusLinks ?? 0) + 1; } if (nf.parentFacilityKey && !next.parentFacilityId) { const pid = await facilityIdForKey(tx, ctx, nf.parentFacilityKey); if (pid && pid !== id) { next.parentFacilityId = pid; next.recordScope = "building"; } } 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") }); } } 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") }); } } const incoming: Record = { name, city: cleanText(nf.city), regionName: cleanText(nf.regionName), address: cleanText(nf.address), postalCode: cleanText(nf.postalCode), status: nf.status && nf.status !== "unknown" ? nf.status : null, facilityType: nf.facilityType && nf.facilityType !== "unknown" ? nf.facilityType : null, tier: cleanText(nf.tier), buildingSqm: nf.buildingSqm ?? null, siteAreaHa: nf.siteAreaHa ?? null, rackCount: nf.rackCount ?? null, pue: nf.pue ?? null, coolingType: cleanText(nf.coolingType), renewableClaim: cleanText(nf.renewableClaim), openedOn: nf.openedOn ?? null, constructionStartedOn: nf.constructionStartedOn ?? null, announcedOn: nf.announcedOn ?? null, website: normalizeWebsite(nf.website), description: cleanText(nf.description), }; for (const field of SCALAR_FIELDS) { const v = incoming[field]; if (v == null) continue; const p = pf(field); observed.push({ field, value: v, provenance: p }); // an opening year before 1980 from a crowd-sourced / dataset source is the building's construction date // (OSM `start_date`), not the data center's: keep the observation in provenance, never in the column if (field === "openedOn" && !isPrimaryKind(ctx.run.sourceKind) && ctx.run.connectorId !== "wikidata" && Number(String(v).slice(0, 4)) < 1980) { ctx.stats.implausibleDropped = (ctx.stats.implausibleDropped ?? 0) + 1; continue; } const dec = shouldReplace(obs(v, p, ctx), storedObs(prov, field, existing?.[field])); if (dec.replace) next[field] = v; } // capacity figures go through the claim store: scope + semantics from the supporting sentence (structured operator // specs have none → the record's own scope), sanity engine, then the authority policy for the column const recordScope = inferRecordScope(nf); const campusDesignation = isCampusName(name, nf.campusName); const MW_PREDICATE: Record<(typeof MW_FIELDS)[number], CapacityPredicate> = { itCapacityMw: "it_capacity_mw", totalPowerMw: "current_power_mw", plannedPowerMw: "planned_power_mw" }; for (const field of MW_FIELDS) { const v = nf[field]; if (v == null || !Number.isFinite(v) || v <= 0) continue; const p = pf(field); const context = nf.claimContext?.[field] ?? null; const structured = !context; // a parser read a spec table / JSON field: the figure describes the record itself const sc = structured ? { scope: recordScope as ClaimScope, reason: "structured:record" } : classifyScope(context, recordScope as ClaimScope); const sem = structured ? { predicate: MW_PREDICATE[field], reason: "field" } : classifyCapacitySemantics(context); const predicate: CapacityPredicate = sem.predicate ?? MW_PREDICATE[field]; const evidence = context ? findEvidence(context, v, "mw") ?? { text: context.slice(0, 600), start: 0, end: Math.min(600, context.length) } : null; 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 }); if (!decision.assign) { ctx.stats.unscopedClaims = (ctx.stats.unscopedClaims ?? 0) + 1; continue; } // utility / grid / ultimate figures live in their own columns, never in IT / total / planned const column = CAPACITY_COLUMN[predicate]; const target = column && column !== "itCapacityMw" && column !== "totalPowerMw" && column !== "plannedPowerMw" ? column : field; observed.push({ field: target, value: v, provenance: { ...p, note: p.note ?? `${predicate} · ${sc.scope}` } }); const dec = shouldReplaceMw(obs(v, p, ctx), storedObs(prov, target, existing?.[target])); if (dec.replace) { next[target] = v; if (target === field) { next.capacityScope = sc.scope; next.capacitySemantics = predicate; } } } next.recordScope = recordScope; // AI evidence: graded from the source text, never from one keyword { const ai = classifyAiEvidence(`${nf.aiEvidence === "confirmed" ? "AI campus" : ""} ${name} ${nf.description ?? ""} ${nf.facilityType === "ai" || nf.facilityType === "hpc" ? "AI data center" : ""}`); const level = nf.aiEvidence && nf.aiEvidence !== "unknown" ? nf.aiEvidence : ai.level; const rank: Record = { unknown: 0, associated: 1, likely: 2, confirmed: 3 }; if ((rank[level] ?? 0) > (rank[String(next.aiEvidence ?? "unknown")] ?? 0)) next.aiEvidence = level; 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}` : ""}` } }); } } // operator / owner: authority policy like other fields (a news source never re-assigns an operator set by the operator itself); // an operator inferred from the name never replaces one stated by a source if (operator) { observed.push({ field: "operatorId", value: operator.id, provenance: operatorProv }); observed.push({ field: "operatorName", value: operator.name, provenance: operatorProv }); if (inferred ? !existing?.operatorId : shouldReplace(obs(operator.id, operatorProv, ctx), storedObs(prov, "operatorId", existing?.operatorId)).replace) next.operatorId = operator.id; } if (owner) { observed.push({ field: "ownerId", value: owner.id, provenance: pf("ownerName") }); if (shouldReplace(obs(owner.id, pf("ownerName"), ctx), storedObs(prov, "ownerId", existing?.ownerId)).replace) next.ownerId = owner.id; } if (campus) next.campusId = next.campusId ?? campus.id; // geo: never degrade precision if (geo) { const p = pf("geo"); observed.push({ field: "geo", value: { lat: geo.lat, lng: geo.lng, precision: geo.precision, source: geo.source }, provenance: p }); const same = storedObs(prov, "geo", existing?.lat != null ? { lat: existing.lat, lng: existing.lng, precision: existing.geoPrecision, source: existing.geoSource } : null); const sameSource = !!same?.sourceId && same.sourceId === (p.sourceId || ctx.run.sourceId); if (shouldReplaceGeo(geo.precision, existing?.geoPrecision as string | null, validLatLng(existing?.lat, existing?.lng), sameSource)) { next.lat = geo.lat; next.lng = geo.lng; next.geoPrecision = geo.precision; next.geoSource = geo.source; } } if (validLatLng(next.lat, next.lng)) next.geohash = geohash(next.lat as number, next.lng as number, 7); // country / metro (metro follows the coordinates we keep) if (country) { observed.push({ field: "countryIso2", value: country, provenance: pf("countryIso2") }); if (!next.countryIso2 || shouldReplace(obs(country, pf("countryIso2"), ctx), storedObs(prov, "countryIso2", existing?.countryIso2)).replace) next.countryIso2 = country; } { 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; if (m.metroId) next.metroId = m.metroId; if (!next.countryIso2 && m.countryIso2) next.countryIso2 = m.countryIso2; } // flags accumulate (true sticks), hyperscaler operator implies hyperscale if (nf.isAi === true) { next.isAi = true; observed.push({ field: "isAi", value: true, provenance: pf("isAi") }); } if (nf.isHyperscale === true || isHyperscalerName(operatorName)) { next.isHyperscale = true; if (nf.isHyperscale === true) observed.push({ field: "isHyperscale", value: true, provenance: pf("isHyperscale") }); } if (nf.certifications?.length) { next.certifications = uniqStrings([...((next.certifications as string[]) ?? []), ...nf.certifications]); observed.push({ field: "certifications", value: nf.certifications, provenance: pf("certifications") }); } if (nf.externalIds && Object.keys(nf.externalIds).length) { next.externalIds = { ...((next.externalIds as Record) ?? {}), ...nf.externalIds }; observed.push({ field: "externalIds", value: nf.externalIds, provenance: pf("externalIds") }); } if (nf.carriers?.length) observed.push({ field: "carriers", value: nf.carriers, provenance: pf("carriers") }); if (nf.cloudProviders?.length) observed.push({ field: "cloudProviders", value: nf.cloudProviders, provenance: pf("cloudProviders") }); if (nf.ixps?.length) observed.push({ field: "ixps", value: nf.ixps, provenance: pf("ixps") }); if (nf.campusName) observed.push({ field: "campusName", value: nf.campusName, provenance: pf("campusName") }); // --- persist ---------------------------------------------------------------------------------------------------- const normalizedName = normalizeName(String(next.name)); if (isNew) { 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)); next.slug = slug; next.confidence = res.how === "pending" ? "unverified" : "moderate"; 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, 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, website, description, is_ai, is_hyperscale, certifications, confidence, completeness, external_ids, source_count, first_seen, parent_facility_id, record_scope, ai_evidence, utility_capacity_mw, grid_connection_mw, ultimate_campus_mw, capacity_scope, capacity_semantics, developer_id, landowner_id) 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}, ${next.address ?? null}, ${next.postalCode ?? null}, ${next.lat ?? null}, ${next.lng ?? null}, ${next.geoPrecision ?? "unknown"}, ${next.geoSource ?? null}, ${next.geohash ?? null}, ${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, ${next.rackCount ?? null}, ${next.pue ?? null}, ${next.coolingType ?? null}, ${next.renewableClaim ?? null}, ${next.openedOn ?? null}, ${next.constructionStartedOn ?? null}, ${next.announcedOn ?? null}, ${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}, ${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})`); } else { const sets = []; for (const [camel, col] of Object.entries(COLUMNS)) { if (["confidence", "completeness", "sourceCount", "lastVerified", "mwIsEstimate"].includes(camel)) continue; const a = before[camel]; const b = next[camel]; if (JSON.stringify(a ?? null) === JSON.stringify(b ?? null)) continue; if (camel === "certifications") sets.push(sql`${sql.identifier(col)} = ${textArray((b as string[]) ?? [])}::text[]`); else if (camel === "externalIds") sets.push(sql`${sql.identifier(col)} = ${JSON.stringify(b ?? {})}::jsonb`); else sets.push(sql`${sql.identifier(col)} = ${b as string | number | boolean | null}`); } if (String(before.name) !== String(next.name)) sets.push(sql`normalized_name = ${normalizedName}`); if (sets.length) { sets.push(sql`updated_at = now()`); await tx.execute(sql`update facilities set ${sql.join(sets, sql`, `)} where id = ${id}`); ctx.stats.updated++; } else ctx.stats.unchanged++; if (res.how === "merge") ctx.stats.merged++; } await upsertKey(tx, ctx, nf.key, "facility", id); addRef(ctx, "facility", id); bump(ctx, "facility"); // aliases: incoming aliases, the incoming name when it differs from the kept name, the previous name when renamed 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); 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`); // provenance for every observed field, then mark which observation backs each displayed value await writeProvenance(tx, ctx, "facility", id, observed, nf.key); 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 }])); // tenants / IXPs if (nf.carriers?.length) await attachTenants(tx, ctx, id, { names: nf.carriers, role: "carrier" }); if (nf.cloudProviders?.length) await attachTenants(tx, ctx, id, { names: nf.cloudProviders, role: "cloud" }); if (nf.carriers?.length || nf.cloudProviders?.length) await refreshTenantCounts(tx, id); if (nf.ixps?.length) await attachIxpsByName(tx, ctx, id, nf.ixps, (next.countryIso2 as string | null) ?? null); // derived: mw_is_estimate, source_count, last_verified, confidence, completeness await refreshDerived(tx, id, res.how === "pending" ? "unverified" : null); // --- reconciliation bookkeeping -------------------------------------------------------------------------------- if (res.how === "merge" && res.match) { await tx.execute(sql`insert into entity_matches (id, connector_id, candidate_key, candidate, matched_facility_id, score, reasons, status, decided_by, decided_at) 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})`); } else if (res.how === "pending" && res.match) { await tx.execute(sql`insert into entity_matches (id, connector_id, candidate_key, candidate, matched_facility_id, score, reasons, status) 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')`); ctx.stats.pendingMatches++; } else if (res.how === "create" && res.match) { await tx.execute(sql`insert into entity_matches (id, connector_id, candidate_key, candidate, matched_facility_id, score, reasons, status, decided_by, decided_at) 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})`); } else if (res.how === "campus_link" && res.match) { await tx.execute(sql`insert into entity_matches (id, connector_id, candidate_key, candidate, matched_facility_id, score, reasons, status, decided_by, decided_at) 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})`); } // duplicate candidates awaiting review are a quality flag too (they inflate every aggregate until decided) if (res.how === "pending" && res.match) { 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 } }); } // --- events ---------------------------------------------------------------------------------------------------- const conf = (await tx.execute(sql`select confidence from facilities where id = ${id}`))[0]?.confidence as ConfidenceLevel | undefined; if (isNew) { ctx.stats.created++; 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 }); const mw = bestMw({ itCapacityMw: next.itCapacityMw as number | null, totalPowerMw: next.totalPowerMw as number | null, plannedPowerMw: next.plannedPowerMw as number | null }); const pipeline = isPipelineStatus(next.status as string | null); const opName = operator?.name ?? null; const where = [next.city, next.countryIso2].filter(Boolean).join(", "); await recordEvent(tx, ctx, { entityType: "facility", entityId: id, eventType: "facility_discovered", title: `New facility indexed: ${next.name}${opName && !String(next.name).toLowerCase().includes(opName.toLowerCase()) ? ` (${opName})` : ""}${where ? ` — ${where}` : ""}`, 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, newValue: { name: next.name, status: next.status, mw }, significance: discoverySignificance(mw, pipeline), confidence: conf ?? "moderate", url, countryIso2: next.countryIso2 as string | null, operatorId: next.operatorId as string | null, metroId: next.metroId as string | null, reviewStatus: res.how === "pending" ? "pending" : "auto", }); } else { 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]); await emitDiffEvents(tx, ctx, { entityType: "facility", entityId: id, entityName: String(next.name), before, after: next, specs: TRACKED_FACILITY_FIELDS, url, confidence: conf ?? "moderate", countryIso2: next.countryIso2 as string | null, operatorId: next.operatorId as string | null, metroId: next.metroId as string | null, operatorNames: names, }); } } function candidateJson(nf: NormalizedFacility, createdFacilityId: string): Record { return { createdFacilityId, key: nf.key, name: nf.name, operatorName: nf.operatorName ?? null, address: nf.address ?? null, city: nf.city ?? null, countryIso2: nf.countryIso2 ?? null, geo: nf.geo ?? null, status: nf.status ?? null, itCapacityMw: nf.itCapacityMw ?? null, totalPowerMw: nf.totalPowerMw ?? null, plannedPowerMw: nf.plannedPowerMw ?? null, externalIds: nf.externalIds ?? {}, url: nf.provenance.url, sourceId: nf.provenance.sourceId, }; } /** Recompute mw_is_estimate, source_count, last_verified, confidence and completeness from provenance + columns. */ export async function refreshDerived(tx: Tx, id: string, forceConfidence: ConfidenceLevel | null = null): Promise { const prov = await loadCurrentProvenance(tx, "facility", id); const row = await loadFacility(tx, id); if (!row) return; const mwRows = prov.filter((r) => (MW_FIELDS as readonly string[]).includes(r.field)); const hasMw = row.itCapacityMw != null || row.totalPowerMw != null || row.plannedPowerMw != null; const mwIsEstimate = hasMw && mwRows.length > 0 && mwRows.every((r) => r.isEstimate); const { sourceIds, kinds, lastPrimaryObserved } = summarizeSources(prov); const onlyEstimates = prov.length > 0 && prov.every((r) => r.isEstimate); // a facility created from an ambiguous match stays "unverified" until the admin decides const pending = forceConfidence ? [] : await tx.execute(sql`select 1 from entity_matches where status = 'pending' and candidate->>'createdFacilityId' = ${id} limit 1`); const confidence = forceConfidence ?? (pending.length ? "unverified" : facilityConfidence({ sourceKinds: kinds, sourceCount: sourceIds.length, lastVerifiedIso: lastPrimaryObserved, onlyEstimates })); 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]; const completeness = completenessScore({ geoPrecision: row.geoPrecision as string, hasCoords: validLatLng(row.lat, row.lng), operatorId: row.operatorId as string | null, address: row.address as string | null, status: row.status as string, itCapacityMw: row.itCapacityMw as number | null, totalPowerMw: row.totalPowerMw as number | null, plannedPowerMw: row.plannedPowerMw as number | null, mwIsEstimate, facilityType: row.facilityType as string, openedOn: row.openedOn as string | null, website: row.website as string | null, description: row.description as string | null, carriersCount: (counts?.carriers_count as number | null) ?? null, ixpCount: (counts?.ixp_count as number | null) ?? null, }); 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}`); } /** Admin approval of a pending match: fold `fromId` into `intoId` (keys, aliases, provenance, links, projects). */ export async function mergeFacilities(tx: Tx, fromId: string, intoId: string, decidedBy = "admin"): Promise { if (fromId === intoId) return; const from = await loadFacility(tx, fromId); if (!from) throw new Error(`facility ${fromId} not found`); await tx.execute(sql`update entity_keys set entity_id = ${intoId} where entity_type = 'facility' and entity_id = ${fromId}`); 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`); 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`); 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)`); await tx.execute(sql`delete from provenance where entity_type = 'facility' and entity_id = ${fromId}`); 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`); 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`); await tx.execute(sql`update projects set facility_id = ${intoId} where facility_id = ${fromId}`); await tx.execute(sql`update events set entity_id = ${intoId} where entity_type = 'facility' and entity_id = ${fromId}`); await tx.execute(sql`update facilities set external_ids = external_ids || (select external_ids from facilities where id = ${fromId}) where id = ${intoId}`); await tx.execute(sql`update facilities set merged_into = ${intoId}, updated_at = now() where id = ${fromId}`); 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}`); await refreshTenantCounts(tx, intoId); 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}`); await refreshDerived(tx, intoId); } /** Extraction debugger: how would this normalized facility resolve right now (candidates, scores, decision)? Read-only. */ export 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 }> }> { const operator = nf.operatorName ? await resolveOperator(tx, ctx, { name: nf.operatorName, key: nf.operatorKey ?? null }) : null; const country = await safeCountry(tx, ctx, nf.countryIso2); const byKey = await facilityIdForKey(tx, ctx, nf.key); const byExt = byKey ? null : await byExternalIds(tx, nf.externalIds); const candidates = await loadCandidates(tx, nf, operator?.id ?? null, country); const geo = validGeo(nf.geo) ? nf.geo : null; 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 }; 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); const best = scored[0]; const how = byKey ? "key" : byExt ? "external_id" : !best ? "create" : best.reasons.includes("rule:campus-vs-building") ? "campus_link" : decide(best.score); return { how, matchedId: byKey ?? byExt ?? (best && (how === "merge") ? best.id : null), candidates: scored }; } /** * First record of an operator in a country (or a metro) → `operator_expansion` event ("Which companies are expanding * into new countries?"). Deterministic: counts other live facilities / projects of the operator in that country. */ export 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 { if (!i.operatorId || !i.countryIso2) return; 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})::int + (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]; if (Number(others?.n ?? 0) > 0) return; const country = (await tx.execute(sql`select name from countries where iso2 = ${i.countryIso2}`))[0]?.name; await recordEvent(tx, ctx, { entityType: i.entityType, entityId: i.entityId, eventType: "operator_expansion", title: `${i.operatorName ?? "Operator"} enters ${String(country ?? i.countryIso2)}: ${i.entityName}`.slice(0, 300), summary: `First ${i.entityType} of ${i.operatorName ?? "this operator"} indexed in ${String(country ?? i.countryIso2)}.`, newValue: { countryIso2: i.countryIso2, operatorId: i.operatorId, [i.entityType === "facility" ? "facilityId" : "projectId"]: i.entityId }, significance: 70, confidence: i.confidence, url: i.url, countryIso2: i.countryIso2, operatorId: i.operatorId, metroId: i.metroId, projectId: i.entityType === "project" ? i.entityId : null, isAi: i.isAi, fingerprint: sha256(`operator_expansion|${i.operatorId}|${i.countryIso2}`), }); }