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%
6.7 KB · 137 lines typescript
Raw Blame History
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