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%
15.5 KB · 255 lines typescript
Raw Blame History
1/**2 * Change events: field diffs old→new on tracked fields → `events` rows deduplicated by fingerprint3 * (entityType|entityId|eventType|JSON(newValue)|day). Titles are templated so the feed reads naturally.4 */5import { sql } from "@dci/db";6import { formatPartialDate, newId, sha256, type ConfidenceLevel, type DetectedChange, type EventType } from "@dci/core";7import { addChange, type IngestContext, type Tx } from "./common.js";89export interface TrackedFieldSpec {10  field: string;11  eventType: EventType;12  label: string;13  kind: "status" | "mw" | "planned_mw" | "date" | "operator" | "owner" | "type" | "name" | "text";14}1516export const TRACKED_FACILITY_FIELDS: TrackedFieldSpec[] = [17  { field: "status", eventType: "status_changed", label: "Status", kind: "status" },18  { field: "itCapacityMw", eventType: "capacity_changed", label: "IT capacity", kind: "mw" },19  { field: "totalPowerMw", eventType: "capacity_changed", label: "Total power", kind: "mw" },20  { field: "plannedPowerMw", eventType: "planned_capacity_changed", label: "Planned capacity", kind: "planned_mw" },21  { field: "openedOn", eventType: "opening_date_changed", label: "Opening date", kind: "date" },22  { field: "announcedOn", eventType: "facility_updated", label: "Announcement date", kind: "date" },23  { field: "constructionStartedOn", eventType: "facility_updated", label: "Construction start", kind: "date" },24  { field: "operatorId", eventType: "operator_changed", label: "Operator", kind: "operator" },25  { field: "ownerId", eventType: "owner_changed", label: "Owner", kind: "owner" },26  { field: "facilityType", eventType: "facility_updated", label: "Facility type", kind: "type" },27  { field: "name", eventType: "facility_updated", label: "Name", kind: "name" },28];2930export const TRACKED_PROJECT_FIELDS: TrackedFieldSpec[] = [31  { field: "status", eventType: "project_status_changed", label: "Status", kind: "status" },32  { field: "plannedMw", eventType: "planned_capacity_changed", label: "Planned capacity", kind: "planned_mw" },33  { field: "expectedOpening", eventType: "opening_date_changed", label: "Expected opening", kind: "date" },34  { field: "announcedOn", eventType: "facility_updated", label: "Announcement date", kind: "date" },35  { field: "operatorId", eventType: "operator_changed", label: "Operator", kind: "operator" },36  { field: "investmentUsd", eventType: "investment_announced", label: "Investment", kind: "text" },37  { field: "name", eventType: "facility_updated", label: "Name", kind: "name" },38];3940/** Significance 0–100: status 90, MW ≥ 20 % change 80 else 50, dates 60, operator 85, others 20. */41export function significanceFor(spec: TrackedFieldSpec, oldValue: unknown, newValue: unknown): number {42  switch (spec.kind) {43    case "status":44      return 90;45    case "operator":46      return 85;47    case "owner":48      return 70;49    case "mw":50    case "planned_mw": {51      const o = typeof oldValue === "number" ? oldValue : null;52      const n = typeof newValue === "number" ? newValue : null;53      if (o == null || n == null || o === 0) return n != null && o == null ? 50 : 50;54      return Math.abs(n - o) / Math.abs(o) >= 0.2 ? 80 : 50;55    }56    case "date":57      return 60;58    default:59      return 20;60  }61}6263export function eventFingerprint(entityType: string, entityId: string | null, eventType: string, newValue: unknown, day: string): string {64  return sha256(`${entityType}|${entityId ?? ""}|${eventType}|${JSON.stringify(newValue ?? null)}|${day}`);65}6667const fmtMw = (v: unknown) => (typeof v === "number" ? `${Number.isInteger(v) ? v : v.toFixed(1)} MW` : "unknown");68const fmtStatus = (v: unknown) => (typeof v === "string" ? v.replace(/_/g, " ") : "unknown");69const 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");7071export function eventTitle(entityName: string, spec: TrackedFieldSpec, oldValue: unknown, newValue: unknown, names?: { old?: string | null; new?: string | null }): { title: string; summary: string } {72  switch (spec.kind) {73    case "status":74      return { title: `${entityName}: status changed to ${fmtStatus(newValue)}`, summary: `Status changed from ${fmtStatus(oldValue)} to ${fmtStatus(newValue)}.` };75    case "mw":76      return { title: `${entityName}: ${spec.label.toLowerCase()} changed to ${fmtMw(newValue)}`, summary: `${spec.label} changed from ${fmtMw(oldValue)} to ${fmtMw(newValue)}.` };77    case "planned_mw":78      return { title: `${entityName}: planned capacity changed to ${fmtMw(newValue)}`, summary: `Planned capacity changed from ${fmtMw(oldValue)} to ${fmtMw(newValue)}.` };79    case "date":80      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)}.` };81    case "operator":82    case "owner":83      return { title: `${entityName}: ${spec.label.toLowerCase()} changed to ${names?.new ?? "unknown"}`, summary: `${spec.label} changed from ${names?.old ?? "unknown"} to ${names?.new ?? "unknown"}.` };84    case "name":85      return { title: `${String(oldValue)} renamed ${String(newValue)}`, summary: `Name changed from "${String(oldValue)}" to "${String(newValue)}".` };86    case "type":87      return { title: `${entityName}: facility type is now ${fmtStatus(newValue)}`, summary: `Facility type changed from ${fmtStatus(oldValue)} to ${fmtStatus(newValue)}.` };88    default:89      if (spec.field === "investmentUsd") return { title: `${entityName}: investment now ${fmtMoney(newValue)}`, summary: `Announced investment changed from ${fmtMoney(oldValue)} to ${fmtMoney(newValue)}.` };90      return { title: `${entityName}: ${spec.label.toLowerCase()} updated`, summary: `${spec.label} changed from ${JSON.stringify(oldValue)} to ${JSON.stringify(newValue)}.` };91  }92}9394export interface EventInput {95  entityType: string;96  entityId: string | null;97  eventType: EventType;98  title: string;99  summary?: string | null;100  oldValue?: unknown;101  newValue?: unknown;102  significance: number;103  confidence?: ConfidenceLevel;104  effectiveDate?: string | null;105  url: string;106  countryIso2?: string | null;107  operatorId?: string | null;108  metroId?: string | null;109  projectId?: string | null;110  reviewStatus?: "auto" | "pending";111  /** override the fingerprint day/value (default: newValue + ctx.day) */112  fingerprint?: string;113  isAi?: boolean;114  /** documents describing the same announcement share a cluster id */115  clusterId?: string | null;116}117118/** Insert an event unless the same fingerprint already exists today. Returns the id when inserted. */119export async function recordEvent(tx: Tx, ctx: IngestContext, e: EventInput): Promise<string | null> {120  const fp = e.fingerprint ?? eventFingerprint(e.entityType, e.entityId, e.eventType, e.newValue, ctx.day);121  const id = newId("event");122  const rows = await tx.execute(sql`123    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,124      significance, confidence, review_status, country_iso2, operator_id, metro_id, project_id, fingerprint, run_id, source_kind, is_ai, cluster_id)125    values (${id}, ${e.entityType}, ${e.entityId}, ${e.eventType}, ${ctx.now}, ${e.effectiveDate ?? null},126      ${e.oldValue === undefined ? null : JSON.stringify(e.oldValue)}::jsonb, ${e.newValue === undefined ? null : JSON.stringify(e.newValue)}::jsonb,127      ${ctx.run.sourceId}, ${ctx.doc?.documentId ?? null}, ${e.url}, ${e.title.slice(0, 300)}, ${e.summary ?? null},128      ${Math.max(0, Math.min(100, Math.round(e.significance)))}, ${e.confidence ?? "moderate"}, ${e.reviewStatus ?? "auto"},129      ${e.countryIso2 ?? null}, ${e.operatorId ?? null}, ${e.metroId ?? null}, ${e.projectId ?? null}, ${fp}, ${ctx.run.runId}, ${ctx.run.sourceKind}, ${!!e.isAi}, ${e.clusterId ?? null})130    on conflict (fingerprint) do nothing131    returning id`);132  if (!rows.length) return null;133  ctx.stats.events++;134  return id;135}136137export interface DiffEventsInput {138  entityType: "facility" | "project";139  entityId: string;140  entityName: string;141  before: Record<string, unknown>;142  after: Record<string, unknown>;143  specs: TrackedFieldSpec[];144  url: string;145  confidence: ConfidenceLevel;146  countryIso2?: string | null;147  operatorId?: string | null;148  metroId?: string | null;149  projectId?: string | null;150  /** operator id → display name (for operator/owner changes) */151  operatorNames?: Map<string, string>;152}153154function same(a: unknown, b: unknown): boolean {155  if (typeof a === "number" && typeof b === "number") return Math.abs(a - b) < 1e-9;156  return JSON.stringify(a ?? null) === JSON.stringify(b ?? null);157}158159/** Emit one event per tracked field whose value changed (old non-null → different new). Returns the number of events written. */160export async function emitDiffEvents(tx: Tx, ctx: IngestContext, i: DiffEventsInput): Promise<number> {161  let n = 0;162  for (const spec of i.specs) {163    const oldValue = i.before[spec.field] ?? null;164    const newValue = i.after[spec.field] ?? null;165    if (newValue == null || same(oldValue, newValue)) continue;166    // a value appearing for the first time is an enrichment, not a change — except status/dates/MW which the feed cares about167    if (oldValue == null && (spec.kind === "name" || spec.kind === "type" || spec.kind === "text")) continue;168    if (oldValue == null && spec.kind === "operator") continue;169    if (oldValue == null && spec.kind === "owner") continue;170    if (oldValue == null && spec.kind === "date" && spec.field !== "openedOn" && spec.field !== "expectedOpening") continue;171    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;172    const { title, summary } = eventTitle(i.entityName, spec, oldValue, newValue, names);173    const significance = oldValue == null ? Math.min(50, significanceFor(spec, oldValue, newValue)) : significanceFor(spec, oldValue, newValue);174    const change: DetectedChange = { field: spec.field, oldValue, newValue, significance, eventType: spec.eventType, description: summary };175    addChange(ctx, change);176    const id = await recordEvent(tx, ctx, {177      entityType: i.entityType,178      entityId: i.entityId,179      eventType: spec.eventType,180      title,181      summary,182      oldValue,183      newValue,184      significance,185      confidence: i.confidence,186      effectiveDate: spec.kind === "date" ? (newValue as string) : null,187      url: i.url,188      countryIso2: i.countryIso2 ?? null,189      operatorId: i.operatorId ?? null,190      metroId: i.metroId ?? null,191      projectId: i.projectId ?? null,192    });193    if (id) n++;194  }195  return n;196}197198/** Significance of a discovery: 65 when MW ≥ 50 or the facility is in the pipeline, else 40. */199export function discoverySignificance(mw: number | null, pipeline: boolean): number {200  return (mw != null && mw >= 50) || pipeline ? 65 : 40;201}202203/**204 * Event clustering (docs/CLAIMS.md): one underlying announcement covered by an operator PR and several outlets must205 * be ONE event with several evidence documents. Deterministic signals only: same event family, same operator (or206 * same country + same city when no operator), publication within ±3 days, MW overlap (±10 %) when both sides have one,207 * or ≥ 0.6 headline token overlap. Returns the cluster id to store on the new event (existing cluster → the primary208 * event's evidence_count is bumped and, when the incoming source is primary and the existing one is not, the cluster209 * primary is re-pointed).210 */211const EVENT_FAMILY: Record<string, string> = {212  project_announced: "announcement", expansion_announced: "announcement", investment_announced: "announcement", phase_announced: "announcement",213  construction_started: "construction", planning_filed: "planning", planning_approved: "planning", land_acquired: "planning",214  facility_opened: "opening", cloud_region_announced: "cloud", cloud_region_launched: "cloud",215  acquisition: "corporate", partnership: "corporate", customer_agreement: "corporate", executive_change: "corporate", operator_expansion: "corporate",216  power_agreement: "power", grid_connection: "power", grid_constraint: "power", utility_event: "power",217  project_delayed: "setback", project_cancelled: "setback", closure: "setback", incident: "setback",218};219const 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"]);220export function headlineTokens(s: string): Set<string> { 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))); }221export 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); }222223export 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 }224225export async function findEventCluster(tx: Tx, ctx: IngestContext, p: ClusterProbe): Promise<{ clusterId: string; joined: boolean }> {226  const family = EVENT_FAMILY[p.eventType] ?? p.eventType;227  const types = Object.entries(EVENT_FAMILY).filter(([, f]) => f === family).map(([t]) => t);228  if (!types.includes(p.eventType)) types.push(p.eventType);229  if (!p.operatorId && !p.metroId && !p.countryIso2) return { clusterId: newId("cluster"), joined: false };230  const rows = await tx.execute(sql`231    select id, cluster_id, title, new_value, source_kind, evidence_count, operator_id, metro_id, country_iso2 from events232    where event_type in ${types} and detected_at > now() - interval '45 days'233      and coalesce(effective_date, detected_at::date::text) between (${p.day}::date - 3)::text and (${p.day}::date + 3)::text234      and (${p.operatorId ?? null}::text is not null and operator_id = ${p.operatorId ?? null}235           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})))236    order by detected_at asc limit 50`);237  for (const r of rows) {238    const nv = (r.new_value ?? {}) as Record<string, unknown>;239    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;240    const mwOk = p.mw != null && mwOther != null ? Math.abs(p.mw - mwOther) <= 0.1 * Math.max(p.mw, mwOther) : null;241    if (mwOk === false) continue;242    const overlap = headlineOverlap(p.title, String(r.title));243    const sameCity = !!p.city && new RegExp(`\\b${p.city.replace(/[.*+?^${}()|[\]\\]/g, "\\$&")}\\b`, "i").test(String(r.title));244    if (mwOk === true || overlap >= 0.6 || (sameCity && overlap >= 0.3)) {245      const clusterId = r.cluster_id ? String(r.cluster_id) : newId("cluster");246      if (!ctx.run.dryRun) {247        if (!r.cluster_id) await tx.execute(sql`update events set cluster_id = ${clusterId} where id = ${r.id}`);248        await tx.execute(sql`update events set evidence_count = evidence_count + 1 where cluster_id = ${clusterId}`);249      }250      return { clusterId, joined: true };251    }252  }253  return { clusterId: newId("cluster"), joined: false };254}255