/** * Shared plumbing for the ingest layer: transaction type, per-batch context, slug helpers, small utils. */ import { sql, type Db } from "@dci/db"; import { slugify, type ConfidenceLevel, type DetectedChange, type Provenance, type SourceKind } from "@dci/core"; import type { IngestDocRef, IngestRun, IngestStats } from "./contract.js"; /** Drizzle transaction handle (same query surface as Db). */ export type Tx = Parameters[0]>[0]; export const PRIMARY_SOURCE_KINDS: ReadonlySet = new Set(["operator", "government", "filing", "utility", "cloud_provider", "registry"]); /** Authority of a source kind when two sources disagree on a field (higher wins; ties → most recent). */ export const SOURCE_KIND_AUTHORITY: Record = { operator: 4, government: 4, filing: 4, utility: 3.5, cloud_provider: 3.5, registry: 3, dataset: 2.5, community: 2, secondary: 1.5, news: 1, }; export const CONFIDENCE_AUTHORITY: Record = { verified: 4, high: 3.5, moderate: 2.5, estimated: 1, unverified: 0.5 }; export interface BatchCaches { operatorsByNorm: Map; // normalizedName → operator id operatorNamesById: Map; facilityIdByKey: Map; countries: Set | null; operatorLexicon: Array<{ id: string | null; name: string; patterns: RegExp[] }> | null; } export interface IngestContext { run: IngestRun; doc: IngestDocRef | null; stats: IngestStats; caches: BatchCaches; now: string; // ISO day: string; // YYYY-MM-DD /** ClickHouse observation rows accumulated for the batch (flushed once, never fatal). */ observations: Array>; } export function newStats(): IngestStats { return { received: 0, created: 0, updated: 0, unchanged: 0, merged: 0, pendingMatches: 0, rejected: 0, events: 0, provenanceRows: 0, byType: {}, changes: [], refs: [] }; } export function newContext(run: IngestRun, doc: IngestDocRef | null): IngestContext { const now = new Date().toISOString(); return { run, doc, stats: newStats(), caches: { operatorsByNorm: new Map(), operatorNamesById: new Map(), facilityIdByKey: new Map(), countries: null, operatorLexicon: null }, now, day: now.slice(0, 10), observations: [], }; } export function addRef(ctx: IngestContext, type: string, id: string): void { if (!ctx.stats.refs.some((r) => r.type === type && r.id === id)) ctx.stats.refs.push({ type, id }); } export function addChange(ctx: IngestContext, c: DetectedChange): void { ctx.stats.changes.push(c); } export function bump(ctx: IngestContext, type: string): void { ctx.stats.byType[type] = (ctx.stats.byType[type] ?? 0) + 1; } /** Ensure a slug is unique within a table (appends a discriminator then a counter). */ export async function uniqueSlug(tx: Tx, table: string, base: string, discriminator?: string | null, excludeId?: string | null): Promise { let root = slugify(base) || "item"; const exists = async (s: string) => { const rows = await tx.execute(sql`select id from ${sql.identifier(table)} where slug = ${s} ${excludeId ? sql`and id <> ${excludeId}` : sql``} limit 1`); return rows.length > 0; }; if (!(await exists(root))) return root; if (discriminator) { const withDisc = slugify(`${root}-${discriminator}`); if (withDisc !== root && !(await exists(withDisc))) return withDisc; root = withDisc; } for (let i = 2; i < 500; i++) { const s = `${root}-${i}`; if (!(await exists(s))) return s; } return `${root}-${Date.now().toString(36)}`; } /** Load (once per batch) the set of ISO2 codes known to the countries table (FK target). */ export async function knownCountries(tx: Tx, ctx: IngestContext): Promise> { if (!ctx.caches.countries) { const rows = await tx.execute(sql`select iso2 from countries`); ctx.caches.countries = new Set(rows.map((r) => String(r.iso2))); } return ctx.caches.countries; } /** Return the iso2 only when it exists in `countries` (otherwise the FK would fail). */ export async function safeCountry(tx: Tx, ctx: IngestContext, iso2: string | null | undefined): Promise { if (!iso2) return null; const up = iso2.trim().toUpperCase(); if (!/^[A-Z]{2}$/.test(up)) return null; const known = await knownCountries(tx, ctx); return known.has(up) ? up : null; } export function websiteDomain(url: string | null | undefined): string | null { if (!url) return null; try { const u = new URL(/^https?:\/\//i.test(url) ? url : `https://${url}`); return u.hostname.toLowerCase().replace(/^www\./, "") || null; } catch { return null; } } export function normalizeWebsite(url: string | null | undefined): string | null { if (!url) return null; try { const u = new URL(/^https?:\/\//i.test(url) ? url : `https://${url}`); u.hash = ""; return u.toString().replace(/\/$/, ""); } catch { return null; } } export function uniqStrings(values: Array): string[] { const seen = new Set(); const out: string[] = []; for (const v of values) { const t = v?.trim(); if (!t) continue; const k = t.toLowerCase(); if (seen.has(k)) continue; seen.add(k); out.push(t); } return out; } export function numOrNull(v: unknown): number | null { if (v == null || v === "") return null; const n = typeof v === "number" ? v : Number(v); return Number.isFinite(n) ? n : null; } export function isoOrNull(v: unknown): string | null { if (!v) return null; const d = new Date(String(v)); return Number.isNaN(d.getTime()) ? null : d.toISOString(); } export function isPrimaryKind(kind: string | null | undefined): boolean { return !!kind && PRIMARY_SOURCE_KINDS.has(kind); } /** Authority of an observation: source kind (or confidence fallback) minus estimate penalty. */ export function authority(kind: string | null | undefined, confidence: string | null | undefined, isEstimate: boolean | null | undefined): number { const base = kind && SOURCE_KIND_AUTHORITY[kind] != null ? SOURCE_KIND_AUTHORITY[kind]! : CONFIDENCE_AUTHORITY[confidence ?? "moderate"] ?? 2; return base - (isEstimate ? 1.5 : 0); } /** Provenance for one field: the default entity provenance, overridden by a matching `facts[]` entry. */ export function provenanceFor(base: Provenance, field: string, facts: Array<{ field: string; provenance: Provenance }> | undefined): Provenance { const f = facts?.find((x) => x.field === field); return f ? { ...base, ...f.provenance } : base; } export function conf(level: string | null | undefined): ConfidenceLevel { return (["verified", "high", "moderate", "estimated", "unverified"].includes(level ?? "") ? level : "moderate") as ConfidenceLevel; } export function pct(n: number): number { return Math.round(n * 100) / 100; } export function errorMessage(e: unknown): string { return e instanceof Error ? e.message : String(e); }