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.9 KB · 189 lines typescript
Raw Blame History
1/**2 * Shared plumbing for the ingest layer: transaction type, per-batch context, slug helpers, small utils.3 */4import { sql, type Db } from "@dci/db";5import { slugify, type ConfidenceLevel, type DetectedChange, type Provenance, type SourceKind } from "@dci/core";6import type { IngestDocRef, IngestRun, IngestStats } from "./contract.js";78/** Drizzle transaction handle (same query surface as Db). */9export type Tx = Parameters<Parameters<Db["transaction"]>[0]>[0];1011export const PRIMARY_SOURCE_KINDS: ReadonlySet<string> = new Set<SourceKind>(["operator", "government", "filing", "utility", "cloud_provider", "registry"]);1213/** Authority of a source kind when two sources disagree on a field (higher wins; ties → most recent). */14export const SOURCE_KIND_AUTHORITY: Record<string, number> = {15  operator: 4,16  government: 4,17  filing: 4,18  utility: 3.5,19  cloud_provider: 3.5,20  registry: 3,21  dataset: 2.5,22  community: 2,23  secondary: 1.5,24  news: 1,25};2627export const CONFIDENCE_AUTHORITY: Record<string, number> = { verified: 4, high: 3.5, moderate: 2.5, estimated: 1, unverified: 0.5 };2829export interface BatchCaches {30  operatorsByNorm: Map<string, string>; // normalizedName → operator id31  operatorNamesById: Map<string, string>;32  facilityIdByKey: Map<string, string>;33  countries: Set<string> | null;34  operatorLexicon: Array<{ id: string | null; name: string; patterns: RegExp[] }> | null;35}3637export interface IngestContext {38  run: IngestRun;39  doc: IngestDocRef | null;40  stats: IngestStats;41  caches: BatchCaches;42  now: string; // ISO43  day: string; // YYYY-MM-DD44  /** ClickHouse observation rows accumulated for the batch (flushed once, never fatal). */45  observations: Array<Record<string, unknown>>;46}4748export function newStats(): IngestStats {49  return { received: 0, created: 0, updated: 0, unchanged: 0, merged: 0, pendingMatches: 0, rejected: 0, events: 0, provenanceRows: 0, byType: {}, changes: [], refs: [] };50}5152export function newContext(run: IngestRun, doc: IngestDocRef | null): IngestContext {53  const now = new Date().toISOString();54  return {55    run,56    doc,57    stats: newStats(),58    caches: { operatorsByNorm: new Map(), operatorNamesById: new Map(), facilityIdByKey: new Map(), countries: null, operatorLexicon: null },59    now,60    day: now.slice(0, 10),61    observations: [],62  };63}6465export function addRef(ctx: IngestContext, type: string, id: string): void {66  if (!ctx.stats.refs.some((r) => r.type === type && r.id === id)) ctx.stats.refs.push({ type, id });67}6869export function addChange(ctx: IngestContext, c: DetectedChange): void {70  ctx.stats.changes.push(c);71}7273export function bump(ctx: IngestContext, type: string): void {74  ctx.stats.byType[type] = (ctx.stats.byType[type] ?? 0) + 1;75}7677/** Ensure a slug is unique within a table (appends a discriminator then a counter). */78export async function uniqueSlug(tx: Tx, table: string, base: string, discriminator?: string | null, excludeId?: string | null): Promise<string> {79  let root = slugify(base) || "item";80  const exists = async (s: string) => {81    const rows = await tx.execute(sql`select id from ${sql.identifier(table)} where slug = ${s} ${excludeId ? sql`and id <> ${excludeId}` : sql``} limit 1`);82    return rows.length > 0;83  };84  if (!(await exists(root))) return root;85  if (discriminator) {86    const withDisc = slugify(`${root}-${discriminator}`);87    if (withDisc !== root && !(await exists(withDisc))) return withDisc;88    root = withDisc;89  }90  for (let i = 2; i < 500; i++) {91    const s = `${root}-${i}`;92    if (!(await exists(s))) return s;93  }94  return `${root}-${Date.now().toString(36)}`;95}9697/** Load (once per batch) the set of ISO2 codes known to the countries table (FK target). */98export async function knownCountries(tx: Tx, ctx: IngestContext): Promise<Set<string>> {99  if (!ctx.caches.countries) {100    const rows = await tx.execute(sql`select iso2 from countries`);101    ctx.caches.countries = new Set(rows.map((r) => String(r.iso2)));102  }103  return ctx.caches.countries;104}105106/** Return the iso2 only when it exists in `countries` (otherwise the FK would fail). */107export async function safeCountry(tx: Tx, ctx: IngestContext, iso2: string | null | undefined): Promise<string | null> {108  if (!iso2) return null;109  const up = iso2.trim().toUpperCase();110  if (!/^[A-Z]{2}$/.test(up)) return null;111  const known = await knownCountries(tx, ctx);112  return known.has(up) ? up : null;113}114115export function websiteDomain(url: string | null | undefined): string | null {116  if (!url) return null;117  try {118    const u = new URL(/^https?:\/\//i.test(url) ? url : `https://${url}`);119    return u.hostname.toLowerCase().replace(/^www\./, "") || null;120  } catch {121    return null;122  }123}124125export function normalizeWebsite(url: string | null | undefined): string | null {126  if (!url) return null;127  try {128    const u = new URL(/^https?:\/\//i.test(url) ? url : `https://${url}`);129    u.hash = "";130    return u.toString().replace(/\/$/, "");131  } catch {132    return null;133  }134}135136export function uniqStrings(values: Array<string | null | undefined>): string[] {137  const seen = new Set<string>();138  const out: string[] = [];139  for (const v of values) {140    const t = v?.trim();141    if (!t) continue;142    const k = t.toLowerCase();143    if (seen.has(k)) continue;144    seen.add(k);145    out.push(t);146  }147  return out;148}149150export function numOrNull(v: unknown): number | null {151  if (v == null || v === "") return null;152  const n = typeof v === "number" ? v : Number(v);153  return Number.isFinite(n) ? n : null;154}155156export function isoOrNull(v: unknown): string | null {157  if (!v) return null;158  const d = new Date(String(v));159  return Number.isNaN(d.getTime()) ? null : d.toISOString();160}161162export function isPrimaryKind(kind: string | null | undefined): boolean {163  return !!kind && PRIMARY_SOURCE_KINDS.has(kind);164}165166/** Authority of an observation: source kind (or confidence fallback) minus estimate penalty. */167export function authority(kind: string | null | undefined, confidence: string | null | undefined, isEstimate: boolean | null | undefined): number {168  const base = kind && SOURCE_KIND_AUTHORITY[kind] != null ? SOURCE_KIND_AUTHORITY[kind]! : CONFIDENCE_AUTHORITY[confidence ?? "moderate"] ?? 2;169  return base - (isEstimate ? 1.5 : 0);170}171172/** Provenance for one field: the default entity provenance, overridden by a matching `facts[]` entry. */173export function provenanceFor(base: Provenance, field: string, facts: Array<{ field: string; provenance: Provenance }> | undefined): Provenance {174  const f = facts?.find((x) => x.field === field);175  return f ? { ...base, ...f.provenance } : base;176}177178export function conf(level: string | null | undefined): ConfidenceLevel {179  return (["verified", "high", "moderate", "estimated", "unverified"].includes(level ?? "") ? level : "moderate") as ConfidenceLevel;180}181182export function pct(n: number): number {183  return Math.round(n * 100) / 100;184}185186export function errorMessage(e: unknown): string {187  return e instanceof Error ? e.message : String(e);188}189