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%
10.0 KB · 147 lines typescript
Raw Blame History
1/**2 * Claim-first read helpers shared by detail payloads and the /history, /claims, /provenance endpoints:3 * claims with their winner flag, dated history points (observed provenance + claims + change events) and the4 * DataQualitySummary (sources, completeness, claims by status, open quality flags, pending duplicate matches).5 */6import type { ClaimDTO, DataQualitySummary, EntityHistory, EventDTO, HistoryPoint, ProvenanceDTO } from "@dci/core";7import { CAPACITY_COLUMN } from "@dci/core";8import { pg, claimCols, eventCols, eventJoins } from "./sql.js";9import { int, iso, num, reqStr, str, type Row } from "./rows.js";10import { asSeverity, claimDto, eventDto, provenanceDto } from "./dto.js";1112const MW_FIELDS = new Set(["itCapacityMw", "totalPowerMw", "plannedPowerMw", "utilityCapacityMw", "gridConnectionMw", "ultimateCampusMw", "plannedMw"]);13const MW_EVENTS = new Set(["capacity_changed", "planned_capacity_changed"]);1415/** Claims about a subject, winner flag derived from the displayed column value (same predicate column, same value). */16export async function claimsFor(subjectType: string, subjectId: string, opts: { status?: string[]; predicate?: string; limit?: number } = {}): Promise<ClaimDTO[]> {17  const sql = pg();18  const rows = await sql<Row[]>`19    select ${claimCols(sql)} from claims k left join sources s on s.id = k.source_id20    where k.subject_type = ${subjectType} and k.subject_id = ${subjectId}21      ${opts.status?.length ? sql`and k.status = any(${opts.status})` : sql``}22      ${opts.predicate ? sql`and k.predicate = ${opts.predicate}` : sql``}23    order by k.status = 'current' desc, k.authority_tier asc, k.last_observed desc limit ${opts.limit ?? 500}`;24  const winners = await winnerValues(subjectType, subjectId);25  return rows.map((r) => {26    const pred = reqStr(r.predicate);27    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<string, string | null>)[pred] ?? null;28    const v = num(r.value);29    const isWinner = str(r.status) === "current" && col != null && v != null && winners.get(col) != null && Math.abs(winners.get(col)! - v) < 1e-9;30    return claimDto(r, isWinner);31  });32}3334async function winnerValues(subjectType: string, subjectId: string): Promise<Map<string, number | null>> {35  const sql = pg();36  const out = new Map<string, number | null>();37  if (subjectType === "facility" || subjectType === "campus") {38    const r = (await sql<Row[]>`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];39    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)); }40  } else if (subjectType === "project") {41    const r = (await sql<Row[]>`select planned_mw, investment_usd from projects where id = ${subjectId}`)[0];42    if (r) { out.set("plannedMw", num(r.planned_mw)); out.set("investmentUsd", num(r.investment_usd)); }43  }44  return out;45}4647/** All provenance rows (current and superseded) for an entity — the /provenance endpoint. */48export async function provenanceAll(entityType: string, entityId: string, limit = 1000): Promise<ProvenanceDTO[]> {49  const sql = pg();50  const rows = await sql<Row[]>`51    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_current52    from provenance p left join sources s on s.id = p.source_id53    where p.entity_type = ${entityType} and p.entity_id = ${entityId}54    order by p.is_current desc, p.field, p.last_observed desc limit ${limit}`;55  return rows.map(provenanceDto);56}5758function pointFromProvenance(p: ProvenanceDTO): HistoryPoint {59  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" };60}61function pointFromClaim(c: ClaimDTO, field: string): HistoryPoint {62  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" };63}64function pointFromEvent(e: EventDTO, field: string): HistoryPoint {65  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" };66}6768function fieldOfEvent(e: EventDTO): string {69  switch (e.eventType) {70    case "planned_capacity_changed": return "plannedPowerMw";71    case "capacity_changed": return "itCapacityMw";72    case "status_changed": case "project_status_changed": return "status";73    case "opening_date_changed": return "openedOn";74    case "operator_changed": return "operatorId";75    case "owner_changed": return "ownerId";76    default: return e.eventType;77  }78}7980const byDate = (a: HistoryPoint, b: HistoryPoint) => (a.date < b.date ? -1 : a.date > b.date ? 1 : 0);8182/** MW history for a facility / project: winner + non-winner observations, claims and capacity_changed events, dated. */83export function capacityHistory(provenance: ProvenanceDTO[], claims: ClaimDTO[], events: EventDTO[]): HistoryPoint[] {84  const out: HistoryPoint[] = [];85  for (const p of provenance) if (MW_FIELDS.has(p.field)) out.push(pointFromProvenance(p));86  for (const c of claims) if (/mw$/.test(c.predicate) && c.status !== "rejected") out.push(pointFromClaim(c, (CAPACITY_COLUMN as Record<string, string | null>)[c.predicate] ?? c.predicate));87  for (const e of events) if (MW_EVENTS.has(e.eventType)) out.push(pointFromEvent(e, fieldOfEvent(e)));88  return out.sort(byDate);89}9091/** Full field history (EntityHistory): every field's dated observations + claims, and the change list. */92export function entityHistory(entityType: EntityHistory["entityType"], entityId: string, provenance: ProvenanceDTO[], claims: ClaimDTO[], events: EventDTO[]): EntityHistory {93  const fields: Record<string, HistoryPoint[]> = {};94  const push = (f: string, p: HistoryPoint) => { (fields[f] ??= []).push(p); };95  for (const p of provenance) push(p.field, pointFromProvenance(p));96  for (const c of claims) if (c.status !== "rejected") push((CAPACITY_COLUMN as Record<string, string | null>)[c.predicate] ?? (/usd$/.test(c.predicate) ? "investmentUsd" : c.predicate), pointFromClaim(c, (CAPACITY_COLUMN as Record<string, string | null>)[c.predicate] ?? c.predicate));97  const changes: HistoryPoint[] = [];98  for (const e of events) {99    if (e.oldValue == null && e.newValue == null) continue;100    if (!/changed|opened|started|approved|filed|cancelled|delayed|acquisition|closure/.test(e.eventType)) continue;101    const f = fieldOfEvent(e);102    const hp = pointFromEvent(e, f);103    changes.push(hp);104    push(f, hp);105  }106  for (const k of Object.keys(fields)) fields[k]!.sort(byDate);107  return { entityType, entityId, fields, changes: changes.sort(byDate) };108}109110/** Events for a subject (entity_type + id, plus project_id for projects), newest first. */111export async function eventsForSubject(subjectType: string, subjectId: string, limit = 200): Promise<EventDTO[]> {112  const sql = pg();113  const rows = await sql<Row[]>`select ${eventCols(sql)} from events e ${eventJoins(sql)}114    where ((e.entity_type = ${subjectType} and e.entity_id = ${subjectId}) ${subjectType === "project" ? sql`or e.project_id = ${subjectId}` : sql``}) and e.review_status <> 'rejected'115    order by e.detected_at desc limit ${limit}`;116  return rows.map((r) => eventDto(r, null));117}118119/** DataQualitySummary from provenance / claims / quality_flags / entity_matches. */120export async function dataQualityFor(entityType: string, entityId: string, completeness: number, lastVerified: string | null): Promise<DataQualitySummary> {121  const sql = pg();122  const [prov, cl, flags, dup] = await Promise.all([123    sql<Row[]>`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 fields124      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`,125    sql<Row[]>`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}`,126    sql<Row[]>`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`,127    entityType === "facility"128      ? sql<Row[]>`select 1 from entity_matches where status = 'pending' and (matched_facility_id = ${entityId} or candidate->>'createdFacilityId' = ${entityId}) limit 1`129      : Promise.resolve([] as Row[]),130  ]);131  const p = prov[0] ?? {};132  const c = cl[0] ?? {};133  return {134    sourceCount: int(p.sources),135    primarySourceCount: int(p.primary_sources),136    lastVerified: lastVerified ? iso(lastVerified) : null,137    completeness,138    fieldsWithProvenance: int(p.fields),139    claimsTotal: int(c.total),140    claimsCurrent: int(c.current),141    claimsUnscoped: int(c.unscoped),142    claimsInReview: int(c.review),143    openFlags: flags.map((f) => ({ code: reqStr(f.code), severity: asSeverity(f.severity), message: reqStr(f.message), field: str(f.field) })),144    pendingDuplicate: dup.length > 0,145  };146}147