spb/datacenterindex
Public
HTML 53.9%
TypeScript 44.5%
JavaScript 0.6%
SQL 0.5%
1/**2 * Per-field provenance: one row per (entity, field, source, url), refreshed on every observation.3 * Older values from the same source for the same field are marked is_current=false. A ClickHouse4 * `observations` row is queued per field (flushed at batch end, never fatal).5 */6import { sql } from "@dci/db";7import { chInsert } from "@dci/db/clickhouse";8import { stableId, type Provenance } from "@dci/core";9import type { IngestContext, Tx } from "./common.js";1011export interface CurrentProvenance {12 field: string;13 sourceId: string;14 value: unknown;15 confidence: string;16 isEstimate: boolean;17 lastObserved: string;18 sourceKind: string | null;19 url: string;20}2122/** Current provenance rows for one entity, with the source kind joined from `sources` when available. */23export async function loadCurrentProvenance(tx: Tx, entityType: string, entityId: string): Promise<CurrentProvenance[]> {24 const rows = await tx.execute(sql`25 select p.field, p.source_id, p.value, p.confidence, p.is_estimate, p.last_observed, p.url, s.kind as source_kind26 from provenance p left join sources s on s.id = p.source_id27 where p.entity_type = ${entityType} and p.entity_id = ${entityId} and p.is_current = true`);28 return rows.map((r) => ({29 field: String(r.field),30 sourceId: String(r.source_id),31 value: r.value,32 confidence: String(r.confidence),33 isEstimate: Boolean(r.is_estimate),34 lastObserved: String(r.last_observed),35 url: String(r.url),36 sourceKind: r.source_kind == null ? null : String(r.source_kind),37 }));38}3940/** The stored observation that currently backs a column value (same value; else the most authoritative row). */41export function backingObservation(rows: CurrentProvenance[], field: string, currentValue: unknown): CurrentProvenance | null {42 const forField = rows.filter((r) => r.field === field);43 if (!forField.length) return null;44 const exact = forField.find((r) => JSON.stringify(r.value) === JSON.stringify(currentValue));45 if (exact) return exact;46 return forField.sort((a, b) => Date.parse(b.lastObserved) - Date.parse(a.lastObserved))[0] ?? null;47}4849export interface ObservedField {50 field: string;51 value: unknown;52 provenance: Provenance;53 entityKey?: string;54 /** claim scope of the observation (building / facility / campus / …) when known */55 scope?: string | null;56}5758export function provenanceId(entityType: string, entityId: string, field: string, sourceId: string, url: string): string {59 return stableId("provenance", `${entityType}|${entityId}|${field}|${sourceId}|${url}`);60}6162/** Upsert provenance rows for the given fields; returns the number of rows written. */63export async function writeProvenance(tx: Tx, ctx: IngestContext, entityType: string, entityId: string, fields: ObservedField[], entityKey?: string): Promise<number> {64 let n = 0;65 for (const f of fields) {66 if (f.value === undefined || f.value === null || f.value === "") continue;67 const p = f.provenance;68 const url = p.url || ctx.doc?.url || "";69 if (!url) continue;70 const sourceId = p.sourceId || ctx.run.sourceId;71 const id = provenanceId(entityType, entityId, f.field, sourceId, url);72 const valueJson = JSON.stringify(f.value);73 const observed = p.lastObserved || ctx.now;74 await tx.execute(sql`75 insert into provenance (id, entity_type, entity_id, field, value, source_id, connector_id, document_id, url, first_observed, last_observed, retrieved_at,76 confidence, is_estimate, method, extractor_version, is_current, note, run_id, scope)77 values (${id}, ${entityType}, ${entityId}, ${f.field}, ${valueJson}::jsonb, ${sourceId}, ${p.connectorId || ctx.run.connectorId}, ${p.documentId ?? ctx.doc?.documentId ?? null}, ${url},78 ${p.firstObserved || observed}, ${observed}, ${p.retrievedAt || ctx.now}, ${p.confidence ?? "moderate"}, ${!!p.isEstimate}, ${p.method ?? null}, ${p.extractorVersion ?? null}, true, ${p.note ?? null}, ${ctx.run.runId}, ${f.scope ?? null})79 on conflict (entity_type, entity_id, field, source_id, url) do update set80 value = excluded.value, last_observed = excluded.last_observed, retrieved_at = excluded.retrieved_at, confidence = excluded.confidence,81 is_estimate = excluded.is_estimate, method = excluded.method, extractor_version = excluded.extractor_version, document_id = coalesce(excluded.document_id, provenance.document_id),82 is_current = true, note = excluded.note, run_id = excluded.run_id, scope = coalesce(excluded.scope, provenance.scope)`);83 // the same source reporting the field from another URL earlier → superseded84 await tx.execute(sql`update provenance set is_current = false where entity_type = ${entityType} and entity_id = ${entityId} and field = ${f.field} and source_id = ${sourceId} and id <> ${id} and is_current = true`);85 n++;86 ctx.observations.push({87 ts: ctx.now.replace("T", " ").replace("Z", ""),88 connector_id: ctx.run.connectorId,89 source_id: sourceId,90 document_id: ctx.doc?.documentId ?? "",91 entity_type: entityType,92 entity_key: entityKey ?? f.entityKey ?? "",93 entity_id: entityId,94 field: f.field,95 value: typeof f.value === "string" ? f.value.slice(0, 2000) : valueJson.slice(0, 2000),96 value_num: typeof f.value === "number" && Number.isFinite(f.value) ? f.value : null,97 confidence: p.confidence ?? "moderate",98 is_estimate: p.isEstimate ? 1 : 0,99 method: p.method ?? "",100 extractor_version: p.extractorVersion ?? "",101 run_id: ctx.run.runId,102 });103 }104 ctx.stats.provenanceRows += n;105 return n;106}107108/** Flush queued ClickHouse observations. Never throws. */109export async function flushObservations(ctx: IngestContext): Promise<void> {110 if (!ctx.observations.length || ctx.run.dryRun) {111 ctx.observations = [];112 return;113 }114 const rows = ctx.observations;115 ctx.observations = [];116 try {117 await chInsert("observations", rows);118 } catch (e) {119 console.warn(`[ingest] clickhouse observations skipped: ${(e as Error).message}`);120 }121}122123/** Distinct source ids / kinds for an entity (for source_count, confidence and last_verified). */124export function summarizeSources(rows: CurrentProvenance[]): { sourceIds: string[]; kinds: string[]; lastPrimaryObserved: string | null } {125 const ids = new Set<string>();126 const kinds = new Set<string>();127 let lastPrimary: string | null = null;128 for (const r of rows) {129 ids.add(r.sourceId);130 if (r.sourceKind) kinds.add(r.sourceKind);131 if (r.sourceKind && ["operator", "government", "filing", "utility", "cloud_provider", "registry"].includes(r.sourceKind)) {132 if (!lastPrimary || Date.parse(r.lastObserved) > Date.parse(lastPrimary)) lastPrimary = r.lastObserved;133 }134 }135 return { sourceIds: [...ids], kinds: [...kinds], lastPrimaryObserved: lastPrimary };136}137