/** * sourceHistory: a merged, dated trail of where an entity's data came from — provenance observations, * detected events and page versions of the documents that reference the entity. Shared by facilities, * operators and projects. */ import type { EventDTO, ProvenanceDTO, SourceHistoryItem, SourceRef } from "@dci/core"; import { pg } from "./sql.js"; import { int, iso, reqStr, type Row } from "./rows.js"; import { asSourceKind, sourceRef } from "./dto.js"; const ANNOUNCE_EVENTS = new Set(["project_announced", "phase_announced", "expansion_announced", "cloud_region_announced", "investment_announced", "power_agreement", "planning_filed"]); const CHANGE_EVENTS = new Set(["status_changed", "capacity_changed", "planned_capacity_changed", "opening_date_changed", "operator_changed", "owner_changed", "project_status_changed", "planning_approved", "construction_started", "facility_opened", "acquisition", "closure", "cloud_region_launched", "cloud_region_updated"]); const DETECT_EVENTS = new Set(["facility_discovered", "page_changed", "incident", "news"]); function fmtValue(v: unknown): string { if (v == null) return "—"; if (typeof v === "object") { const s = JSON.stringify(v); return s.length > 80 ? s.slice(0, 77) + "…" : s; } const s = String(v); return s.length > 80 ? s.slice(0, 77) + "…" : s; } export interface DocVersionRow { fetchedAt: string; url: string; sourceId: string; sourceName: string; sourceKind: string; significance: number; changes: number } /** Document versions for documents whose entity_refs contain (type,id). */ export async function documentVersionsFor(entityType: string, entityId: string, limit = 60): Promise { const sql = pg(); const refJson = JSON.stringify([{ type: entityType, id: entityId }]); const rows = await sql` select v.fetched_at, d.url, d.source_id, s.name as source_name, s.kind as source_kind, v.significance, jsonb_array_length(coalesce(v.detected_changes, '[]'::jsonb)) as changes from documents d join document_versions v on v.document_id = d.id left join sources s on s.id = d.source_id where d.entity_refs @> ${refJson}::jsonb order by v.fetched_at desc limit ${limit}`; return rows.map((r) => ({ fetchedAt: reqStr(iso(r.fetched_at)), url: reqStr(r.url), sourceId: reqStr(r.source_id), sourceName: reqStr(r.source_name, reqStr(r.source_id)), sourceKind: reqStr(r.source_kind, "secondary"), significance: int(r.significance), changes: int(r.changes) })); } export function buildSourceHistory(provenance: ProvenanceDTO[], events: EventDTO[], versions: DocVersionRow[], max = 150): SourceHistoryItem[] { const items: SourceHistoryItem[] = []; const seen = new Set(); const push = (it: SourceHistoryItem) => { const k = `${it.date.slice(0, 10)}|${it.sourceId}|${it.kind}|${it.description}`; if (seen.has(k)) return; seen.add(k); items.push(it); }; for (const p of provenance) { push({ date: p.firstObserved, sourceId: p.sourceId, sourceName: p.sourceName, sourceKind: p.sourceKind, url: p.url, kind: "observed", description: `${p.field} observed: ${fmtValue(p.value)}${p.isEstimate ? " (estimate)" : ""}` }); if (p.lastObserved.slice(0, 10) !== p.firstObserved.slice(0, 10)) push({ date: p.lastObserved, sourceId: p.sourceId, sourceName: p.sourceName, sourceKind: p.sourceKind, url: p.url, kind: "verified", description: `${p.field} re-confirmed (${fmtValue(p.value)})` }); } for (const e of events) { const kind: SourceHistoryItem["kind"] = ANNOUNCE_EVENTS.has(e.eventType) ? "announced" : CHANGE_EVENTS.has(e.eventType) ? "changed" : DETECT_EVENTS.has(e.eventType) ? "detected" : "updated"; const desc = e.oldValue != null && e.newValue != null && kind === "changed" ? `${e.title} (${fmtValue(e.oldValue)} → ${fmtValue(e.newValue)})` : e.title; push({ date: e.detectedAt, sourceId: e.sourceId, sourceName: e.sourceName, sourceKind: e.sourceKind, url: e.url, kind, description: desc }); } for (const v of versions) { push({ date: v.fetchedAt, sourceId: v.sourceId, sourceName: v.sourceName, sourceKind: asSourceKind(v.sourceKind), url: v.url, kind: "changed", description: v.changes > 0 ? `Source page changed — ${v.changes} field change${v.changes > 1 ? "s" : ""} detected (significance ${v.significance})` : `Source page changed (significance ${v.significance})` }); } items.sort((a, b) => (a.date < b.date ? 1 : a.date > b.date ? -1 : 0)); return items.slice(0, max); } /** SourceRef[] for a set of source ids (envelope `sources`). */ export async function sourceRefsFor(ids: Iterable): Promise { const list = [...new Set([...ids].filter(Boolean))]; if (!list.length) return []; const sql = pg(); const rows = await sql`select id, name, domain, kind, url, license, attribution, redistribution, attribution_required from sources where id = any(${list}) order by name`; return rows.map(sourceRef); } export function sourceIdsOf(...lists: Array>): Set { const s = new Set(); for (const l of lists) for (const x of l) if (x.sourceId) s.add(x.sourceId); return s; } /** Distinct sources behind a set of entities (current provenance rows), capped — for LIST envelopes. */ export async function sourcesForEntities(entityType: string, ids: string[], cap = 30): Promise { const list = [...new Set(ids.filter(Boolean))]; if (!list.length) return []; const sql = pg(); const rows = await sql` select s.id, s.name, s.domain, s.kind, s.url, s.license, s.attribution, s.redistribution, s.attribution_required, count(*) as n from provenance p join sources s on s.id = p.source_id where p.entity_type = ${entityType} and p.entity_id = any(${list}) and p.is_current group by s.id order by n desc, s.name limit ${cap}`; return rows.map(sourceRef); } /** Distinct sources of a set of source ids (events / map points), capped. */ export async function sourcesForIds(ids: Iterable, cap = 30): Promise { const list = [...new Set([...ids].filter(Boolean))].slice(0, 500); if (!list.length) return []; const sql = pg(); const rows = await sql`select id, name, domain, kind, url, license, attribution, redistribution, attribution_required from sources where id = any(${list}) order by name limit ${cap}`; return rows.map(sourceRef); }