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