/** * Claim-first read helpers shared by detail payloads and the /history, /claims, /provenance endpoints: * claims with their winner flag, dated history points (observed provenance + claims + change events) and the * DataQualitySummary (sources, completeness, claims by status, open quality flags, pending duplicate matches). */ import type { ClaimDTO, DataQualitySummary, EntityHistory, EventDTO, HistoryPoint, ProvenanceDTO } from "@dci/core"; import { CAPACITY_COLUMN } from "@dci/core"; import { pg, claimCols, eventCols, eventJoins } from "./sql.js"; import { int, iso, num, reqStr, str, type Row } from "./rows.js"; import { asSeverity, claimDto, eventDto, provenanceDto } from "./dto.js"; const MW_FIELDS = new Set(["itCapacityMw", "totalPowerMw", "plannedPowerMw", "utilityCapacityMw", "gridConnectionMw", "ultimateCampusMw", "plannedMw"]); const MW_EVENTS = new Set(["capacity_changed", "planned_capacity_changed"]); /** Claims about a subject, winner flag derived from the displayed column value (same predicate column, same value). */ export async function claimsFor(subjectType: string, subjectId: string, opts: { status?: string[]; predicate?: string; limit?: number } = {}): Promise { const sql = pg(); const rows = await sql` select ${claimCols(sql)} from claims k left join sources s on s.id = k.source_id where k.subject_type = ${subjectType} and k.subject_id = ${subjectId} ${opts.status?.length ? sql`and k.status = any(${opts.status})` : sql``} ${opts.predicate ? sql`and k.predicate = ${opts.predicate}` : sql``} order by k.status = 'current' desc, k.authority_tier asc, k.last_observed desc limit ${opts.limit ?? 500}`; const winners = await winnerValues(subjectType, subjectId); return rows.map((r) => { const pred = reqStr(r.predicate); const col = subjectType === "project" ? (pred === "planned_power_mw" || pred === "phase_mw" || pred === "it_capacity_mw" ? "plannedMw" : /usd$/.test(pred) ? "investmentUsd" : null) : (CAPACITY_COLUMN as Record)[pred] ?? null; const v = num(r.value); const isWinner = str(r.status) === "current" && col != null && v != null && winners.get(col) != null && Math.abs(winners.get(col)! - v) < 1e-9; return claimDto(r, isWinner); }); } async function winnerValues(subjectType: string, subjectId: string): Promise> { const sql = pg(); const out = new Map(); if (subjectType === "facility" || subjectType === "campus") { const r = (await sql`select it_capacity_mw, total_power_mw, planned_power_mw, utility_capacity_mw, grid_connection_mw, ultimate_campus_mw from facilities where id = ${subjectId}`)[0]; if (r) { out.set("itCapacityMw", num(r.it_capacity_mw)); out.set("totalPowerMw", num(r.total_power_mw)); out.set("plannedPowerMw", num(r.planned_power_mw)); out.set("utilityCapacityMw", num(r.utility_capacity_mw)); out.set("gridConnectionMw", num(r.grid_connection_mw)); out.set("ultimateCampusMw", num(r.ultimate_campus_mw)); } } else if (subjectType === "project") { const r = (await sql`select planned_mw, investment_usd from projects where id = ${subjectId}`)[0]; if (r) { out.set("plannedMw", num(r.planned_mw)); out.set("investmentUsd", num(r.investment_usd)); } } return out; } /** All provenance rows (current and superseded) for an entity — the /provenance endpoint. */ export async function provenanceAll(entityType: string, entityId: string, limit = 1000): Promise { const sql = pg(); const rows = await sql` select p.field, p.value, p.source_id, s.name as source_name, s.kind as source_kind, p.url, p.first_observed, p.last_observed, p.retrieved_at, p.confidence, p.is_estimate, p.method, p.is_winner, p.scope, p.run_id, p.document_id, p.is_current from provenance p left join sources s on s.id = p.source_id where p.entity_type = ${entityType} and p.entity_id = ${entityId} order by p.is_current desc, p.field, p.last_observed desc limit ${limit}`; return rows.map(provenanceDto); } function pointFromProvenance(p: ProvenanceDTO): HistoryPoint { return { date: p.firstObserved, field: p.field, predicate: null, value: p.value, sourceId: p.sourceId, sourceName: p.sourceName, sourceKind: p.sourceKind, url: p.url, kind: "observed" }; } function pointFromClaim(c: ClaimDTO, field: string): HistoryPoint { return { date: c.publishedAt && /^\d{4}/.test(c.publishedAt) ? c.publishedAt : c.firstObserved, field, predicate: c.predicate, value: c.value ?? c.valueText, sourceId: c.sourceId, sourceName: c.sourceName, sourceKind: c.sourceKind, url: c.url, claimId: c.id, kind: "claim" }; } function pointFromEvent(e: EventDTO, field: string): HistoryPoint { return { date: e.effectiveDate && /^\d{4}-\d{2}-\d{2}/.test(e.effectiveDate) ? e.effectiveDate : e.detectedAt, field, predicate: null, value: e.newValue, oldValue: e.oldValue, sourceId: e.sourceId, sourceName: e.sourceName, sourceKind: e.sourceKind, url: e.url, eventId: e.id, kind: "changed" }; } function fieldOfEvent(e: EventDTO): string { switch (e.eventType) { case "planned_capacity_changed": return "plannedPowerMw"; case "capacity_changed": return "itCapacityMw"; case "status_changed": case "project_status_changed": return "status"; case "opening_date_changed": return "openedOn"; case "operator_changed": return "operatorId"; case "owner_changed": return "ownerId"; default: return e.eventType; } } const byDate = (a: HistoryPoint, b: HistoryPoint) => (a.date < b.date ? -1 : a.date > b.date ? 1 : 0); /** MW history for a facility / project: winner + non-winner observations, claims and capacity_changed events, dated. */ export function capacityHistory(provenance: ProvenanceDTO[], claims: ClaimDTO[], events: EventDTO[]): HistoryPoint[] { const out: HistoryPoint[] = []; for (const p of provenance) if (MW_FIELDS.has(p.field)) out.push(pointFromProvenance(p)); for (const c of claims) if (/mw$/.test(c.predicate) && c.status !== "rejected") out.push(pointFromClaim(c, (CAPACITY_COLUMN as Record)[c.predicate] ?? c.predicate)); for (const e of events) if (MW_EVENTS.has(e.eventType)) out.push(pointFromEvent(e, fieldOfEvent(e))); return out.sort(byDate); } /** Full field history (EntityHistory): every field's dated observations + claims, and the change list. */ export function entityHistory(entityType: EntityHistory["entityType"], entityId: string, provenance: ProvenanceDTO[], claims: ClaimDTO[], events: EventDTO[]): EntityHistory { const fields: Record = {}; const push = (f: string, p: HistoryPoint) => { (fields[f] ??= []).push(p); }; for (const p of provenance) push(p.field, pointFromProvenance(p)); for (const c of claims) if (c.status !== "rejected") push((CAPACITY_COLUMN as Record)[c.predicate] ?? (/usd$/.test(c.predicate) ? "investmentUsd" : c.predicate), pointFromClaim(c, (CAPACITY_COLUMN as Record)[c.predicate] ?? c.predicate)); const changes: HistoryPoint[] = []; for (const e of events) { if (e.oldValue == null && e.newValue == null) continue; if (!/changed|opened|started|approved|filed|cancelled|delayed|acquisition|closure/.test(e.eventType)) continue; const f = fieldOfEvent(e); const hp = pointFromEvent(e, f); changes.push(hp); push(f, hp); } for (const k of Object.keys(fields)) fields[k]!.sort(byDate); return { entityType, entityId, fields, changes: changes.sort(byDate) }; } /** Events for a subject (entity_type + id, plus project_id for projects), newest first. */ export async function eventsForSubject(subjectType: string, subjectId: string, limit = 200): Promise { const sql = pg(); const rows = await sql`select ${eventCols(sql)} from events e ${eventJoins(sql)} where ((e.entity_type = ${subjectType} and e.entity_id = ${subjectId}) ${subjectType === "project" ? sql`or e.project_id = ${subjectId}` : sql``}) and e.review_status <> 'rejected' order by e.detected_at desc limit ${limit}`; return rows.map((r) => eventDto(r, null)); } /** DataQualitySummary from provenance / claims / quality_flags / entity_matches. */ export async function dataQualityFor(entityType: string, entityId: string, completeness: number, lastVerified: string | null): Promise { const sql = pg(); const [prov, cl, flags, dup] = await Promise.all([ sql`select count(distinct p.source_id)::int as sources, count(distinct p.source_id) filter (where s.kind in ('operator','government','filing','utility','cloud_provider','registry'))::int as primary_sources, count(distinct p.field)::int as fields from provenance p left join sources s on s.id = p.source_id where p.entity_type = ${entityType} and p.entity_id = ${entityId} and p.is_current`, sql`select count(*)::int as total, count(*) filter (where status = 'current')::int as current, count(*) filter (where status = 'unscoped')::int as unscoped, count(*) filter (where status = 'review')::int as review from claims where subject_type = ${entityType} and subject_id = ${entityId}`, sql`select code, severity, message, field from quality_flags where entity_type = ${entityType} and entity_id = ${entityId} and status = 'open' order by priority desc, created_at desc limit 50`, entityType === "facility" ? sql`select 1 from entity_matches where status = 'pending' and (matched_facility_id = ${entityId} or candidate->>'createdFacilityId' = ${entityId}) limit 1` : Promise.resolve([] as Row[]), ]); const p = prov[0] ?? {}; const c = cl[0] ?? {}; return { sourceCount: int(p.sources), primarySourceCount: int(p.primary_sources), lastVerified: lastVerified ? iso(lastVerified) : null, completeness, fieldsWithProvenance: int(p.fields), claimsTotal: int(c.total), claimsCurrent: int(c.current), claimsUnscoped: int(c.unscoped), claimsInReview: int(c.review), openFlags: flags.map((f) => ({ code: reqStr(f.code), severity: asSeverity(f.severity), message: reqStr(f.message), field: str(f.field) })), pendingDuplicate: dup.length > 0, }; }