/** * Per-field provenance: one row per (entity, field, source, url), refreshed on every observation. * Older values from the same source for the same field are marked is_current=false. A ClickHouse * `observations` row is queued per field (flushed at batch end, never fatal). */ import { sql } from "@dci/db"; import { chInsert } from "@dci/db/clickhouse"; import { stableId, type Provenance } from "@dci/core"; import type { IngestContext, Tx } from "./common.js"; export interface CurrentProvenance { field: string; sourceId: string; value: unknown; confidence: string; isEstimate: boolean; lastObserved: string; sourceKind: string | null; url: string; } /** Current provenance rows for one entity, with the source kind joined from `sources` when available. */ export async function loadCurrentProvenance(tx: Tx, entityType: string, entityId: string): Promise { const rows = await tx.execute(sql` select p.field, p.source_id, p.value, p.confidence, p.is_estimate, p.last_observed, p.url, s.kind as source_kind 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 = true`); return rows.map((r) => ({ field: String(r.field), sourceId: String(r.source_id), value: r.value, confidence: String(r.confidence), isEstimate: Boolean(r.is_estimate), lastObserved: String(r.last_observed), url: String(r.url), sourceKind: r.source_kind == null ? null : String(r.source_kind), })); } /** The stored observation that currently backs a column value (same value; else the most authoritative row). */ export function backingObservation(rows: CurrentProvenance[], field: string, currentValue: unknown): CurrentProvenance | null { const forField = rows.filter((r) => r.field === field); if (!forField.length) return null; const exact = forField.find((r) => JSON.stringify(r.value) === JSON.stringify(currentValue)); if (exact) return exact; return forField.sort((a, b) => Date.parse(b.lastObserved) - Date.parse(a.lastObserved))[0] ?? null; } export interface ObservedField { field: string; value: unknown; provenance: Provenance; entityKey?: string; /** claim scope of the observation (building / facility / campus / …) when known */ scope?: string | null; } export function provenanceId(entityType: string, entityId: string, field: string, sourceId: string, url: string): string { return stableId("provenance", `${entityType}|${entityId}|${field}|${sourceId}|${url}`); } /** Upsert provenance rows for the given fields; returns the number of rows written. */ export async function writeProvenance(tx: Tx, ctx: IngestContext, entityType: string, entityId: string, fields: ObservedField[], entityKey?: string): Promise { let n = 0; for (const f of fields) { if (f.value === undefined || f.value === null || f.value === "") continue; const p = f.provenance; const url = p.url || ctx.doc?.url || ""; if (!url) continue; const sourceId = p.sourceId || ctx.run.sourceId; const id = provenanceId(entityType, entityId, f.field, sourceId, url); const valueJson = JSON.stringify(f.value); const observed = p.lastObserved || ctx.now; await tx.execute(sql` insert into provenance (id, entity_type, entity_id, field, value, source_id, connector_id, document_id, url, first_observed, last_observed, retrieved_at, confidence, is_estimate, method, extractor_version, is_current, note, run_id, scope) values (${id}, ${entityType}, ${entityId}, ${f.field}, ${valueJson}::jsonb, ${sourceId}, ${p.connectorId || ctx.run.connectorId}, ${p.documentId ?? ctx.doc?.documentId ?? null}, ${url}, ${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}) on conflict (entity_type, entity_id, field, source_id, url) do update set value = excluded.value, last_observed = excluded.last_observed, retrieved_at = excluded.retrieved_at, confidence = excluded.confidence, is_estimate = excluded.is_estimate, method = excluded.method, extractor_version = excluded.extractor_version, document_id = coalesce(excluded.document_id, provenance.document_id), is_current = true, note = excluded.note, run_id = excluded.run_id, scope = coalesce(excluded.scope, provenance.scope)`); // the same source reporting the field from another URL earlier → superseded 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`); n++; ctx.observations.push({ ts: ctx.now.replace("T", " ").replace("Z", ""), connector_id: ctx.run.connectorId, source_id: sourceId, document_id: ctx.doc?.documentId ?? "", entity_type: entityType, entity_key: entityKey ?? f.entityKey ?? "", entity_id: entityId, field: f.field, value: typeof f.value === "string" ? f.value.slice(0, 2000) : valueJson.slice(0, 2000), value_num: typeof f.value === "number" && Number.isFinite(f.value) ? f.value : null, confidence: p.confidence ?? "moderate", is_estimate: p.isEstimate ? 1 : 0, method: p.method ?? "", extractor_version: p.extractorVersion ?? "", run_id: ctx.run.runId, }); } ctx.stats.provenanceRows += n; return n; } /** Flush queued ClickHouse observations. Never throws. */ export async function flushObservations(ctx: IngestContext): Promise { if (!ctx.observations.length || ctx.run.dryRun) { ctx.observations = []; return; } const rows = ctx.observations; ctx.observations = []; try { await chInsert("observations", rows); } catch (e) { console.warn(`[ingest] clickhouse observations skipped: ${(e as Error).message}`); } } /** Distinct source ids / kinds for an entity (for source_count, confidence and last_verified). */ export function summarizeSources(rows: CurrentProvenance[]): { sourceIds: string[]; kinds: string[]; lastPrimaryObserved: string | null } { const ids = new Set(); const kinds = new Set(); let lastPrimary: string | null = null; for (const r of rows) { ids.add(r.sourceId); if (r.sourceKind) kinds.add(r.sourceKind); if (r.sourceKind && ["operator", "government", "filing", "utility", "cloud_provider", "registry"].includes(r.sourceKind)) { if (!lastPrimary || Date.parse(r.lastObserved) > Date.parse(lastPrimary)) lastPrimary = r.lastObserved; } } return { sourceIds: [...ids], kinds: [...kinds], lastPrimaryObserved: lastPrimary }; }