/** * Change events: field diffs old→new on tracked fields → `events` rows deduplicated by fingerprint * (entityType|entityId|eventType|JSON(newValue)|day). Titles are templated so the feed reads naturally. */ import { sql } from "@dci/db"; import { formatPartialDate, newId, sha256, type ConfidenceLevel, type DetectedChange, type EventType } from "@dci/core"; import { addChange, type IngestContext, type Tx } from "./common.js"; export interface TrackedFieldSpec { field: string; eventType: EventType; label: string; kind: "status" | "mw" | "planned_mw" | "date" | "operator" | "owner" | "type" | "name" | "text"; } export const TRACKED_FACILITY_FIELDS: TrackedFieldSpec[] = [ { field: "status", eventType: "status_changed", label: "Status", kind: "status" }, { field: "itCapacityMw", eventType: "capacity_changed", label: "IT capacity", kind: "mw" }, { field: "totalPowerMw", eventType: "capacity_changed", label: "Total power", kind: "mw" }, { field: "plannedPowerMw", eventType: "planned_capacity_changed", label: "Planned capacity", kind: "planned_mw" }, { field: "openedOn", eventType: "opening_date_changed", label: "Opening date", kind: "date" }, { field: "announcedOn", eventType: "facility_updated", label: "Announcement date", kind: "date" }, { field: "constructionStartedOn", eventType: "facility_updated", label: "Construction start", kind: "date" }, { field: "operatorId", eventType: "operator_changed", label: "Operator", kind: "operator" }, { field: "ownerId", eventType: "owner_changed", label: "Owner", kind: "owner" }, { field: "facilityType", eventType: "facility_updated", label: "Facility type", kind: "type" }, { field: "name", eventType: "facility_updated", label: "Name", kind: "name" }, ]; export const TRACKED_PROJECT_FIELDS: TrackedFieldSpec[] = [ { field: "status", eventType: "project_status_changed", label: "Status", kind: "status" }, { field: "plannedMw", eventType: "planned_capacity_changed", label: "Planned capacity", kind: "planned_mw" }, { field: "expectedOpening", eventType: "opening_date_changed", label: "Expected opening", kind: "date" }, { field: "announcedOn", eventType: "facility_updated", label: "Announcement date", kind: "date" }, { field: "operatorId", eventType: "operator_changed", label: "Operator", kind: "operator" }, { field: "investmentUsd", eventType: "investment_announced", label: "Investment", kind: "text" }, { field: "name", eventType: "facility_updated", label: "Name", kind: "name" }, ]; /** Significance 0–100: status 90, MW ≥ 20 % change 80 else 50, dates 60, operator 85, others 20. */ export function significanceFor(spec: TrackedFieldSpec, oldValue: unknown, newValue: unknown): number { switch (spec.kind) { case "status": return 90; case "operator": return 85; case "owner": return 70; case "mw": case "planned_mw": { const o = typeof oldValue === "number" ? oldValue : null; const n = typeof newValue === "number" ? newValue : null; if (o == null || n == null || o === 0) return n != null && o == null ? 50 : 50; return Math.abs(n - o) / Math.abs(o) >= 0.2 ? 80 : 50; } case "date": return 60; default: return 20; } } export function eventFingerprint(entityType: string, entityId: string | null, eventType: string, newValue: unknown, day: string): string { return sha256(`${entityType}|${entityId ?? ""}|${eventType}|${JSON.stringify(newValue ?? null)}|${day}`); } const fmtMw = (v: unknown) => (typeof v === "number" ? `${Number.isInteger(v) ? v : v.toFixed(1)} MW` : "unknown"); const fmtStatus = (v: unknown) => (typeof v === "string" ? v.replace(/_/g, " ") : "unknown"); const fmtMoney = (v: unknown) => (typeof v === "number" ? (v >= 1e9 ? `$${(v / 1e9).toFixed(1)} billion` : v >= 1e6 ? `$${Math.round(v / 1e6)} million` : `$${Math.round(v).toLocaleString("en-US")}`) : "unknown"); export function eventTitle(entityName: string, spec: TrackedFieldSpec, oldValue: unknown, newValue: unknown, names?: { old?: string | null; new?: string | null }): { title: string; summary: string } { switch (spec.kind) { case "status": return { title: `${entityName}: status changed to ${fmtStatus(newValue)}`, summary: `Status changed from ${fmtStatus(oldValue)} to ${fmtStatus(newValue)}.` }; case "mw": return { title: `${entityName}: ${spec.label.toLowerCase()} changed to ${fmtMw(newValue)}`, summary: `${spec.label} changed from ${fmtMw(oldValue)} to ${fmtMw(newValue)}.` }; case "planned_mw": return { title: `${entityName}: planned capacity changed to ${fmtMw(newValue)}`, summary: `Planned capacity changed from ${fmtMw(oldValue)} to ${fmtMw(newValue)}.` }; case "date": return { title: `${entityName}: ${spec.label.toLowerCase()} now ${formatPartialDate(newValue as string | null)}`, summary: `${spec.label} changed from ${formatPartialDate(oldValue as string | null)} to ${formatPartialDate(newValue as string | null)}.` }; case "operator": case "owner": return { title: `${entityName}: ${spec.label.toLowerCase()} changed to ${names?.new ?? "unknown"}`, summary: `${spec.label} changed from ${names?.old ?? "unknown"} to ${names?.new ?? "unknown"}.` }; case "name": return { title: `${String(oldValue)} renamed ${String(newValue)}`, summary: `Name changed from "${String(oldValue)}" to "${String(newValue)}".` }; case "type": return { title: `${entityName}: facility type is now ${fmtStatus(newValue)}`, summary: `Facility type changed from ${fmtStatus(oldValue)} to ${fmtStatus(newValue)}.` }; default: if (spec.field === "investmentUsd") return { title: `${entityName}: investment now ${fmtMoney(newValue)}`, summary: `Announced investment changed from ${fmtMoney(oldValue)} to ${fmtMoney(newValue)}.` }; return { title: `${entityName}: ${spec.label.toLowerCase()} updated`, summary: `${spec.label} changed from ${JSON.stringify(oldValue)} to ${JSON.stringify(newValue)}.` }; } } export interface EventInput { entityType: string; entityId: string | null; eventType: EventType; title: string; summary?: string | null; oldValue?: unknown; newValue?: unknown; significance: number; confidence?: ConfidenceLevel; effectiveDate?: string | null; url: string; countryIso2?: string | null; operatorId?: string | null; metroId?: string | null; projectId?: string | null; reviewStatus?: "auto" | "pending"; /** override the fingerprint day/value (default: newValue + ctx.day) */ fingerprint?: string; isAi?: boolean; /** documents describing the same announcement share a cluster id */ clusterId?: string | null; } /** Insert an event unless the same fingerprint already exists today. Returns the id when inserted. */ export async function recordEvent(tx: Tx, ctx: IngestContext, e: EventInput): Promise { const fp = e.fingerprint ?? eventFingerprint(e.entityType, e.entityId, e.eventType, e.newValue, ctx.day); const id = newId("event"); const rows = await tx.execute(sql` insert into events (id, entity_type, entity_id, event_type, detected_at, effective_date, old_value, new_value, source_id, document_id, url, title, summary, significance, confidence, review_status, country_iso2, operator_id, metro_id, project_id, fingerprint, run_id, source_kind, is_ai, cluster_id) values (${id}, ${e.entityType}, ${e.entityId}, ${e.eventType}, ${ctx.now}, ${e.effectiveDate ?? null}, ${e.oldValue === undefined ? null : JSON.stringify(e.oldValue)}::jsonb, ${e.newValue === undefined ? null : JSON.stringify(e.newValue)}::jsonb, ${ctx.run.sourceId}, ${ctx.doc?.documentId ?? null}, ${e.url}, ${e.title.slice(0, 300)}, ${e.summary ?? null}, ${Math.max(0, Math.min(100, Math.round(e.significance)))}, ${e.confidence ?? "moderate"}, ${e.reviewStatus ?? "auto"}, ${e.countryIso2 ?? null}, ${e.operatorId ?? null}, ${e.metroId ?? null}, ${e.projectId ?? null}, ${fp}, ${ctx.run.runId}, ${ctx.run.sourceKind}, ${!!e.isAi}, ${e.clusterId ?? null}) on conflict (fingerprint) do nothing returning id`); if (!rows.length) return null; ctx.stats.events++; return id; } export interface DiffEventsInput { entityType: "facility" | "project"; entityId: string; entityName: string; before: Record; after: Record; specs: TrackedFieldSpec[]; url: string; confidence: ConfidenceLevel; countryIso2?: string | null; operatorId?: string | null; metroId?: string | null; projectId?: string | null; /** operator id → display name (for operator/owner changes) */ operatorNames?: Map; } function same(a: unknown, b: unknown): boolean { if (typeof a === "number" && typeof b === "number") return Math.abs(a - b) < 1e-9; return JSON.stringify(a ?? null) === JSON.stringify(b ?? null); } /** Emit one event per tracked field whose value changed (old non-null → different new). Returns the number of events written. */ export async function emitDiffEvents(tx: Tx, ctx: IngestContext, i: DiffEventsInput): Promise { let n = 0; for (const spec of i.specs) { const oldValue = i.before[spec.field] ?? null; const newValue = i.after[spec.field] ?? null; if (newValue == null || same(oldValue, newValue)) continue; // a value appearing for the first time is an enrichment, not a change — except status/dates/MW which the feed cares about if (oldValue == null && (spec.kind === "name" || spec.kind === "type" || spec.kind === "text")) continue; if (oldValue == null && spec.kind === "operator") continue; if (oldValue == null && spec.kind === "owner") continue; if (oldValue == null && spec.kind === "date" && spec.field !== "openedOn" && spec.field !== "expectedOpening") continue; const names = spec.kind === "operator" || spec.kind === "owner" ? { old: oldValue ? i.operatorNames?.get(String(oldValue)) ?? null : null, new: i.operatorNames?.get(String(newValue)) ?? null } : undefined; const { title, summary } = eventTitle(i.entityName, spec, oldValue, newValue, names); const significance = oldValue == null ? Math.min(50, significanceFor(spec, oldValue, newValue)) : significanceFor(spec, oldValue, newValue); const change: DetectedChange = { field: spec.field, oldValue, newValue, significance, eventType: spec.eventType, description: summary }; addChange(ctx, change); const id = await recordEvent(tx, ctx, { entityType: i.entityType, entityId: i.entityId, eventType: spec.eventType, title, summary, oldValue, newValue, significance, confidence: i.confidence, effectiveDate: spec.kind === "date" ? (newValue as string) : null, url: i.url, countryIso2: i.countryIso2 ?? null, operatorId: i.operatorId ?? null, metroId: i.metroId ?? null, projectId: i.projectId ?? null, }); if (id) n++; } return n; } /** Significance of a discovery: 65 when MW ≥ 50 or the facility is in the pipeline, else 40. */ export function discoverySignificance(mw: number | null, pipeline: boolean): number { return (mw != null && mw >= 50) || pipeline ? 65 : 40; } /** * Event clustering (docs/CLAIMS.md): one underlying announcement covered by an operator PR and several outlets must * be ONE event with several evidence documents. Deterministic signals only: same event family, same operator (or * same country + same city when no operator), publication within ±3 days, MW overlap (±10 %) when both sides have one, * or ≥ 0.6 headline token overlap. Returns the cluster id to store on the new event (existing cluster → the primary * event's evidence_count is bumped and, when the incoming source is primary and the existing one is not, the cluster * primary is re-pointed). */ const EVENT_FAMILY: Record = { project_announced: "announcement", expansion_announced: "announcement", investment_announced: "announcement", phase_announced: "announcement", construction_started: "construction", planning_filed: "planning", planning_approved: "planning", land_acquired: "planning", facility_opened: "opening", cloud_region_announced: "cloud", cloud_region_launched: "cloud", acquisition: "corporate", partnership: "corporate", customer_agreement: "corporate", executive_change: "corporate", operator_expansion: "corporate", power_agreement: "power", grid_connection: "power", grid_constraint: "power", utility_event: "power", project_delayed: "setback", project_cancelled: "setback", closure: "setback", incident: "setback", }; const STOP = new Set(["the", "a", "an", "of", "in", "at", "to", "for", "and", "on", "with", "its", "new", "data", "center", "centre", "centers", "centres", "campus", "mw", "gw", "announces", "announced", "plans", "build", "project", "facility", "datacenter", "datacentre"]); export function headlineTokens(s: string): Set { return new Set(s.toLowerCase().normalize("NFKD").replace(/[^a-z0-9 ]/g, " ").split(/\s+/).filter((t) => t.length >= 3 && !STOP.has(t) && !/^\d+$/.test(t))); } export function headlineOverlap(a: string, b: string): number { const x = headlineTokens(a), y = headlineTokens(b); if (!x.size || !y.size) return 0; let n = 0; for (const t of x) if (y.has(t)) n++; return n / Math.min(x.size, y.size); } export interface ClusterProbe { eventType: EventType; title: string; operatorId?: string | null; countryIso2?: string | null; metroId?: string | null; city?: string | null; mw?: number | null; day: string; sourceKind: string } export async function findEventCluster(tx: Tx, ctx: IngestContext, p: ClusterProbe): Promise<{ clusterId: string; joined: boolean }> { const family = EVENT_FAMILY[p.eventType] ?? p.eventType; const types = Object.entries(EVENT_FAMILY).filter(([, f]) => f === family).map(([t]) => t); if (!types.includes(p.eventType)) types.push(p.eventType); if (!p.operatorId && !p.metroId && !p.countryIso2) return { clusterId: newId("cluster"), joined: false }; const rows = await tx.execute(sql` select id, cluster_id, title, new_value, source_kind, evidence_count, operator_id, metro_id, country_iso2 from events where event_type in ${types} and detected_at > now() - interval '45 days' and coalesce(effective_date, detected_at::date::text) between (${p.day}::date - 3)::text and (${p.day}::date + 3)::text and (${p.operatorId ?? null}::text is not null and operator_id = ${p.operatorId ?? null} or (${p.operatorId ?? null}::text is null and operator_id is null and ${p.countryIso2 ?? null}::text is not null and country_iso2 = ${p.countryIso2 ?? null} and (${p.metroId ?? null}::text is null or metro_id is null or metro_id = ${p.metroId ?? null}))) order by detected_at asc limit 50`); for (const r of rows) { const nv = (r.new_value ?? {}) as Record; const mwOther = typeof nv.mw === "number" ? nv.mw : typeof nv.plannedMw === "number" ? nv.plannedMw : Array.isArray(nv.mw) ? Number((nv.mw as number[])[0]) : null; const mwOk = p.mw != null && mwOther != null ? Math.abs(p.mw - mwOther) <= 0.1 * Math.max(p.mw, mwOther) : null; if (mwOk === false) continue; const overlap = headlineOverlap(p.title, String(r.title)); const sameCity = !!p.city && new RegExp(`\\b${p.city.replace(/[.*+?^${}()|[\]\\]/g, "\\$&")}\\b`, "i").test(String(r.title)); if (mwOk === true || overlap >= 0.6 || (sameCity && overlap >= 0.3)) { const clusterId = r.cluster_id ? String(r.cluster_id) : newId("cluster"); if (!ctx.run.dryRun) { if (!r.cluster_id) await tx.execute(sql`update events set cluster_id = ${clusterId} where id = ${r.id}`); await tx.execute(sql`update events set evidence_count = evidence_count + 1 where cluster_id = ${clusterId}`); } return { clusterId, joined: true }; } } return { clusterId: newId("cluster"), joined: false }; }