/** * Ingest layer entry point — implements `IngestFn` from contract.ts. * * One Postgres transaction per batch; each entity runs inside its own savepoint so a bad record is rolled back * and counted as rejected without failing the batch. Dry runs execute everything and roll the batch back. * Entities are processed in dependency order: operators → campuses → facilities → IXPs → cloud regions → * projects → news → country statistics. */ import { getDb, sql } from "@dci/db"; import type { NormalizedEntity } from "@dci/core"; import type { IngestDocRef, IngestFn, IngestRun, IngestStats } from "./contract.js"; import { errorMessage, newContext, type IngestContext, type Tx } from "./common.js"; import { ingestCampus } from "./campuses.js"; import { ingestCloudRegion } from "./cloud-regions.js"; import { ingestCountryStat } from "./countries.js"; import { ingestFacility } from "./facilities.js"; import { ingestIxp } from "./ixps.js"; import { ingestNews } from "./news.js"; import { ingestOperator } from "./operators.js"; import { ingestProject } from "./projects.js"; import { flushObservations } from "./provenance.js"; const ORDER: Record = { operator: 0, campus: 1, facility: 2, ixp: 3, cloud_region: 4, project: 5, news_event: 6, country: 7 }; class DryRunRollback extends Error { constructor() { super("dry-run rollback"); } } /** Make sure a `sources` row exists for this run (kind is needed by the field-merge authority policy). */ async function ensureSourceRow(tx: Tx, run: IngestRun): Promise { await tx.execute(sql`insert into sources (id, connector_id, name, domain, kind, priority) values (${run.sourceId}, ${run.connectorId}, ${run.connectorId}, '', ${run.sourceKind}, ${run.sourcePriority || 3}) on conflict (id) do nothing`); } async function ingestOne(tx: Tx, ctx: IngestContext, e: NormalizedEntity): Promise { switch (e.entityType) { case "operator": return ingestOperator(tx, ctx, e); case "campus": return ingestCampus(tx, ctx, e); case "facility": return ingestFacility(tx, ctx, e); case "ixp": return ingestIxp(tx, ctx, e); case "cloud_region": return ingestCloudRegion(tx, ctx, e); case "project": return ingestProject(tx, ctx, e); case "news_event": return ingestNews(tx, ctx, e); case "country": return ingestCountryStat(tx, ctx, e); default: throw new Error(`unknown entityType ${(e as { entityType: string }).entityType}`); } } function snapshot(s: IngestStats): IngestStats { return { ...s, byType: { ...s.byType }, changes: [...s.changes], refs: [...s.refs] }; } export const ingestEntities: IngestFn = async (run: IngestRun, entities: NormalizedEntity[], doc: IngestDocRef | null): Promise => { const ctx = newContext(run, doc); ctx.stats.received = entities.length; if (!entities.length) return ctx.stats; const ordered = [...entities].sort((a, b) => (ORDER[a.entityType] ?? 99) - (ORDER[b.entityType] ?? 99)); const db = getDb(); try { await db.transaction(async (tx) => { await ensureSourceRow(tx, run); for (const e of ordered) { const before = snapshot(ctx.stats); try { await tx.transaction(async (inner) => { await ingestOne(inner, ctx, e); }); } catch (err) { // savepoint rolled back: restore counters, drop caches that may reference rolled-back rows ctx.stats = before; ctx.stats.rejected++; ctx.caches.operatorsByNorm.clear(); ctx.caches.operatorNamesById.clear(); ctx.caches.facilityIdByKey.clear(); ctx.caches.operatorLexicon = null; console.warn(`[ingest] ${run.connectorId} ${e.entityType} ${(e as { key?: string }).key ?? "?"} rejected: ${errorMessage(err)}`); } } if (run.dryRun) throw new DryRunRollback(); }); } catch (err) { if (!(err instanceof DryRunRollback)) throw err; } await flushObservations(ctx); return ctx.stats; }; export type { IngestDocRef, IngestFn, IngestRun, IngestStats } from "./contract.js"; export { mergeFacilities, refreshDerived, previewFacilityResolution, isCampusName, inferRecordScope } from "./facilities.js"; export { hideProject, mergeProjects } from "./projects.js"; export { writeClaim, writeQualityFlags, recordCapacityClaim, recordInvestmentClaim, reviewPriority, bestCurrentClaim } from "./claims.js"; export { geocodeCity } from "./geocode.js"; export { resolveOperator, lookupOperatorId, inferOperatorKind, mergeOperators } from "./operators.js"; export { findCanonicalOperator, CANONICAL_OPERATORS, HYPERSCALERS } from "./canonical-operators.js"; export { assignMetro, loadMetros, nearestMetro, invalidateMetroCache } from "./metros.js"; export { scoreFacilityMatch, bestMatch, decide, completenessScore, facilityConfidence, shouldReplace, shouldReplaceMw, shouldReplaceGeo, AUTO_MERGE_THRESHOLD, PENDING_THRESHOLD } from "./match.js"; export { eventFingerprint, significanceFor, eventTitle, recordEvent, TRACKED_FACILITY_FIELDS, TRACKED_PROJECT_FIELDS } from "./events.js"; export { writeProvenance, loadCurrentProvenance } from "./provenance.js"; export { linkNews } from "./news.js";