spb/datacenterindex
Public
HTML 53.9%
TypeScript 44.5%
JavaScript 0.6%
SQL 0.5%
1/**2 * Ingest layer entry point — implements `IngestFn` from contract.ts.3 *4 * One Postgres transaction per batch; each entity runs inside its own savepoint so a bad record is rolled back5 * and counted as rejected without failing the batch. Dry runs execute everything and roll the batch back.6 * Entities are processed in dependency order: operators → campuses → facilities → IXPs → cloud regions →7 * projects → news → country statistics.8 */9import { getDb, sql } from "@dci/db";10import type { NormalizedEntity } from "@dci/core";11import type { IngestDocRef, IngestFn, IngestRun, IngestStats } from "./contract.js";12import { errorMessage, newContext, type IngestContext, type Tx } from "./common.js";13import { ingestCampus } from "./campuses.js";14import { ingestCloudRegion } from "./cloud-regions.js";15import { ingestCountryStat } from "./countries.js";16import { ingestFacility } from "./facilities.js";17import { ingestIxp } from "./ixps.js";18import { ingestNews } from "./news.js";19import { ingestOperator } from "./operators.js";20import { ingestProject } from "./projects.js";21import { flushObservations } from "./provenance.js";2223const ORDER: Record<NormalizedEntity["entityType"], number> = { operator: 0, campus: 1, facility: 2, ixp: 3, cloud_region: 4, project: 5, news_event: 6, country: 7 };2425class DryRunRollback extends Error {26 constructor() {27 super("dry-run rollback");28 }29}3031/** Make sure a `sources` row exists for this run (kind is needed by the field-merge authority policy). */32async function ensureSourceRow(tx: Tx, run: IngestRun): Promise<void> {33 await tx.execute(sql`insert into sources (id, connector_id, name, domain, kind, priority)34 values (${run.sourceId}, ${run.connectorId}, ${run.connectorId}, '', ${run.sourceKind}, ${run.sourcePriority || 3})35 on conflict (id) do nothing`);36}3738async function ingestOne(tx: Tx, ctx: IngestContext, e: NormalizedEntity): Promise<void> {39 switch (e.entityType) {40 case "operator":41 return ingestOperator(tx, ctx, e);42 case "campus":43 return ingestCampus(tx, ctx, e);44 case "facility":45 return ingestFacility(tx, ctx, e);46 case "ixp":47 return ingestIxp(tx, ctx, e);48 case "cloud_region":49 return ingestCloudRegion(tx, ctx, e);50 case "project":51 return ingestProject(tx, ctx, e);52 case "news_event":53 return ingestNews(tx, ctx, e);54 case "country":55 return ingestCountryStat(tx, ctx, e);56 default:57 throw new Error(`unknown entityType ${(e as { entityType: string }).entityType}`);58 }59}6061function snapshot(s: IngestStats): IngestStats {62 return { ...s, byType: { ...s.byType }, changes: [...s.changes], refs: [...s.refs] };63}6465export const ingestEntities: IngestFn = async (run: IngestRun, entities: NormalizedEntity[], doc: IngestDocRef | null): Promise<IngestStats> => {66 const ctx = newContext(run, doc);67 ctx.stats.received = entities.length;68 if (!entities.length) return ctx.stats;69 const ordered = [...entities].sort((a, b) => (ORDER[a.entityType] ?? 99) - (ORDER[b.entityType] ?? 99));70 const db = getDb();71 try {72 await db.transaction(async (tx) => {73 await ensureSourceRow(tx, run);74 for (const e of ordered) {75 const before = snapshot(ctx.stats);76 try {77 await tx.transaction(async (inner) => {78 await ingestOne(inner, ctx, e);79 });80 } catch (err) {81 // savepoint rolled back: restore counters, drop caches that may reference rolled-back rows82 ctx.stats = before;83 ctx.stats.rejected++;84 ctx.caches.operatorsByNorm.clear();85 ctx.caches.operatorNamesById.clear();86 ctx.caches.facilityIdByKey.clear();87 ctx.caches.operatorLexicon = null;88 console.warn(`[ingest] ${run.connectorId} ${e.entityType} ${(e as { key?: string }).key ?? "?"} rejected: ${errorMessage(err)}`);89 }90 }91 if (run.dryRun) throw new DryRunRollback();92 });93 } catch (err) {94 if (!(err instanceof DryRunRollback)) throw err;95 }96 await flushObservations(ctx);97 return ctx.stats;98};99100export type { IngestDocRef, IngestFn, IngestRun, IngestStats } from "./contract.js";101export { mergeFacilities, refreshDerived, previewFacilityResolution, isCampusName, inferRecordScope } from "./facilities.js";102export { hideProject, mergeProjects } from "./projects.js";103export { writeClaim, writeQualityFlags, recordCapacityClaim, recordInvestmentClaim, reviewPriority, bestCurrentClaim } from "./claims.js";104export { geocodeCity } from "./geocode.js";105export { resolveOperator, lookupOperatorId, inferOperatorKind, mergeOperators } from "./operators.js";106export { findCanonicalOperator, CANONICAL_OPERATORS, HYPERSCALERS } from "./canonical-operators.js";107export { assignMetro, loadMetros, nearestMetro, invalidateMetroCache } from "./metros.js";108export { scoreFacilityMatch, bestMatch, decide, completenessScore, facilityConfidence, shouldReplace, shouldReplaceMw, shouldReplaceGeo, AUTO_MERGE_THRESHOLD, PENDING_THRESHOLD } from "./match.js";109export { eventFingerprint, significanceFor, eventTitle, recordEvent, TRACKED_FACILITY_FIELDS, TRACKED_PROJECT_FIELDS } from "./events.js";110export { writeProvenance, loadCurrentProvenance } from "./provenance.js";111export { linkNews } from "./news.js";112