/** * Ingest smoke test against the local databases (Postgres + ClickHouse). SYNTHETIC DATA ONLY — every record is * prefixed "Smoke Test", uses connector id "smoke-test" and is deleted at the end. * * set -a; source .env; set +a; node node_modules/tsx/dist/cli.mjs scripts/ingest-smoke.ts * * Run 1: 6 entities (3 facilities — one is a duplicate that must auto-merge —, 1 project, 1 news event, 1 cloud region). * Run 2: same entities with one changed MW → 0 created, exactly 1 event (capacity_changed). * Then computeRankings / computeDailyMetrics / refreshStats, cleanup, recompute so the database is left clean. */ import { closeDb, getDb, sql } from "@dci/db"; import { ensureClickHouse, getClickHouse } from "@dci/db/clickhouse"; import type { NormalizedCloudRegion, NormalizedEntity, NormalizedFacility, NormalizedNewsEvent, NormalizedProject, Provenance } from "@dci/core"; import { ingestEntities, type IngestRun } from "../apps/worker/src/ingest/index.js"; import { computeRankings, refreshStats } from "../apps/worker/src/rankings.js"; import { computeDailyMetrics } from "../apps/worker/src/metrics.js"; import { seedReferenceData } from "../apps/worker/src/seeds/index.js"; const CONNECTOR = "smoke-test"; const SOURCE = "src_smoke-test"; const run: IngestRun = { connectorId: CONNECTOR, sourceId: SOURCE, runId: `run_smoke_${Date.now().toString(36)}`, sourceKind: "operator", sourcePriority: 2, dryRun: false }; const doc = { documentId: "doc_smoke_test", url: "https://smoke-test.invalid/facilities", pageType: "facility_index", fetchedAt: new Date().toISOString() }; function prov(url: string, extra: Partial = {}): Provenance { const now = new Date().toISOString(); return { sourceId: SOURCE, connectorId: CONNECTOR, url, firstObserved: now, lastObserved: now, retrievedAt: now, confidence: "high", method: "smoke-test", extractorVersion: "smoke", ...extra }; } function entities(mwA: number): NormalizedEntity[] { const facA: NormalizedFacility = { entityType: "facility", key: "smoke-test:fac:1", name: "Smoke Test DC12", operatorName: "Smoke Test Operator", address: "1 Smoke Test Court, Ashburn, VA 20147", city: "Ashburn", regionName: "Virginia", countryIso2: "US", geo: { lat: 39.0401, lng: -77.4899, precision: "exact", source: "smoke-test" }, status: "operational", facilityType: "colocation", itCapacityMw: mwA, openedOn: "2016-05", website: "https://smoke-test.invalid/dc12", carriers: ["Smoke Test Carrier Networks"], ixps: ["SMOKE-IX"], externalIds: { smoke_test: "fac-1" }, provenance: prov("https://smoke-test.invalid/dc12"), }; // duplicate of A from a "second page": must auto-merge (same operator, code DC12, 30 m apart) const facB: NormalizedFacility = { entityType: "facility", key: "smoke-test:fac:2", name: "Smoke Test DC12 Data Center", operatorName: "Smoke Test Operator", city: "Ashburn", countryIso2: "US", geo: { lat: 39.0403, lng: -77.4901, precision: "exact", source: "smoke-test" }, status: "operational", totalPowerMw: 60, provenance: prov("https://smoke-test.invalid/dc12-alt"), }; const facC: NormalizedFacility = { entityType: "facility", key: "smoke-test:fac:3", name: "Smoke Test Montréal Campus", operatorName: "Smoke Test Operator", city: "Montréal", regionName: "Quebec", countryIso2: "CA", geo: { lat: 45.5017, lng: -73.5673, precision: "city", source: "smoke-test" }, status: "under_construction", facilityType: "hyperscale", plannedPowerMw: 120, isAi: true, facts: [{ field: "plannedPowerMw", value: 120, provenance: prov("https://smoke-test.invalid/montreal", { isEstimate: true, method: "smoke-estimate" }) }], provenance: prov("https://smoke-test.invalid/montreal"), }; const project: NormalizedProject = { entityType: "project", key: "smoke-test:prj:1", name: "Smoke Test Hyperscale Campus Phase 2", operatorName: "Smoke Test Operator", facilityKey: "smoke-test:fac:3", city: "Montréal", countryIso2: "CA", status: "announced", announcedOn: "2026-08", expectedOpening: "2028-Q2", plannedMw: 200, investmentUsd: 1.2e9, sourceUrl: "https://smoke-test.invalid/news/1", timeline: [{ date: "2026-08-15", type: "project_announced", description: "Phase 2 announced (smoke test)", url: "https://smoke-test.invalid/news/1" }], provenance: prov("https://smoke-test.invalid/projects/phase-2"), }; const news: NormalizedNewsEvent = { entityType: "news_event", key: "smoke-test:news:1", title: "Smoke Test Operator announces 200 MW campus in Montréal, Canada", url: "https://smoke-test.invalid/news/1", publishedAt: "2026-08-15T12:00:00Z", summary: "Synthetic news item used by the ingest smoke test.", pageType: "press_release", eventType: "project_announced", mentions: { operators: ["Smoke Test Operator"], countriesIso2: ["CA"], mw: [200] }, provenance: prov("https://smoke-test.invalid/news/1"), }; const region: NormalizedCloudRegion = { entityType: "cloud_region", key: "smoke-test:region:smoke-east-1", providerName: "Smoke Test Cloud", code: "smoke-east-1", name: "Smoke East (Ashburn)", city: "Ashburn", countryIso2: "US", geo: { lat: 39.04, lng: -77.49, precision: "city", source: "smoke-test" }, availabilityZones: 3, launchedOn: "2020", status: "operational", provenance: prov("https://smoke-test.invalid/regions"), }; return [facA, facB, facC, project, news, region]; } function assert(cond: unknown, msg: string): void { if (!cond) throw new Error(`ASSERTION FAILED: ${msg}`); } async function count(q: ReturnType): Promise { const r = await getDb().execute(q); return Number(r[0]?.n ?? 0); } async function cleanup(): Promise { const db = getDb(); const keyed = async (type: string) => (await db.execute(sql`select entity_id from entity_keys where connector_id = ${CONNECTOR} and entity_type = ${type}`)).map((r) => String(r.entity_id)); const facIds = await keyed("facility"); const prjIds = await keyed("project"); const crIds = await keyed("cloud_region"); const ixIds = await keyed("ixp"); await db.execute(sql`delete from events where source_id = ${SOURCE}`); await db.execute(sql`delete from provenance where connector_id = ${CONNECTOR}`); await db.execute(sql`delete from news_items where connector_id = ${CONNECTOR}`); await db.execute(sql`delete from entity_matches where connector_id = ${CONNECTOR}`); await db.execute(sql`delete from project_timeline where source_id = ${SOURCE}`); if (prjIds.length) await db.execute(sql`delete from projects where id in ${prjIds}`); if (crIds.length) await db.execute(sql`delete from cloud_regions where id in ${crIds}`); if (facIds.length) await db.execute(sql`delete from facilities where id in ${facIds}`); // aliases/tenants/ixps cascade await db.execute(sql`delete from ixps where name = 'SMOKE-IX' ${ixIds.length ? sql`or id in ${ixIds}` : sql``}`); await db.execute(sql`delete from entity_keys where connector_id = ${CONNECTOR}`); await db.execute(sql`delete from facility_tenants where operator_id in (select id from operators where name like 'Smoke Test%')`); await db.execute(sql`delete from operators where name like 'Smoke Test%'`); await db.execute(sql`delete from sources where id = ${SOURCE}`); try { await getClickHouse().command({ query: `ALTER TABLE observations DELETE WHERE connector_id = 'smoke-test'` }); } catch { /* ClickHouse optional */ } } async function main() { const db = getDb(); try { await ensureClickHouse(); } catch (e) { console.warn(`clickhouse unavailable: ${(e as Error).message}`); } const seeded = await seedReferenceData(db); console.log(`seeds: ${seeded.countries} countries, ${seeded.metros} metros`); await cleanup(); // in case a previous run aborted try { // ---- run 1 const s1 = await ingestEntities(run, entities(36), doc); console.log("run 1", JSON.stringify({ ...s1, changes: s1.changes.length, refs: s1.refs.length })); assert(s1.rejected === 0, `run 1 rejected=${s1.rejected}`); const facCount = await count(sql`select count(*)::int as n from facilities f where f.id in (select entity_id from entity_keys where connector_id = ${CONNECTOR} and entity_type = 'facility') and f.merged_into is null`); assert(facCount === 2, `expected 2 facilities after merge, got ${facCount}`); const merged = await count(sql`select count(*)::int as n from entity_matches where connector_id = ${CONNECTOR} and status = 'auto_merged'`); assert(merged === 1, `expected 1 auto_merged match, got ${merged}`); const keys = await db.execute(sql`select key, entity_id from entity_keys where connector_id = ${CONNECTOR} and key in ('smoke-test:fac:1','smoke-test:fac:2')`); assert(keys.length === 2 && keys[0]!.entity_id === keys[1]!.entity_id, "fac:1 and fac:2 must point to the same facility"); const facA = (await db.execute(sql`select f.* from facilities f join entity_keys k on k.entity_id = f.id where k.key = 'smoke-test:fac:1'`))[0]!; assert(Number(facA.it_capacity_mw) === 36 && Number(facA.total_power_mw) === 60, "merged facility keeps IT 36 MW and total 60 MW"); assert(facA.metro_id != null, "Ashburn facility assigned to a metro"); assert(Number(facA.carriers_count) === 1 && Number(facA.ixp_count) === 1, `tenants/ixps counted (${facA.carriers_count}/${facA.ixp_count})`); assert(Number(facA.completeness) >= 70, `completeness ${facA.completeness} ≥ 70`); const facC = (await db.execute(sql`select f.* from facilities f join entity_keys k on k.entity_id = f.id where k.key = 'smoke-test:fac:3'`))[0]!; assert(facC.mw_is_estimate === true, "estimate-only MW flagged mw_is_estimate"); assert(facC.country_iso2 === "CA" && facC.metro_id != null, "Montréal facility → CA + metro"); const discovered = await count(sql`select count(*)::int as n from events where source_id = ${SOURCE} and event_type = 'facility_discovered'`); assert(discovered === 2, `expected 2 facility_discovered events, got ${discovered}`); const news = await count(sql`select count(*)::int as n from news_items where connector_id = ${CONNECTOR} and 'CA' = any(country_iso2s) and mw = 200`); assert(news === 1, "news item linked to CA with 200 MW"); const newsEvt = await count(sql`select count(*)::int as n from events where source_id = ${SOURCE} and entity_type = 'news_event'`); assert(newsEvt === 1, "news item with eventType produced one event"); const prj = (await db.execute(sql`select p.*, (select count(*) from project_timeline t where t.project_id = p.id) as tl from projects p join entity_keys k on k.entity_id = p.id where k.key = 'smoke-test:prj:1'`))[0]!; assert(prj.facility_id === facC.id && Number(prj.tl) === 1, "project linked to facility with 1 timeline row"); const provRows = await count(sql`select count(*)::int as n from provenance where connector_id = ${CONNECTOR} and is_current`); assert(provRows >= 25, `provenance rows ${provRows} ≥ 25`); // ---- run 2 (idempotent except one MW change) const s2 = await ingestEntities(run, entities(48), doc); console.log("run 2", JSON.stringify({ ...s2, changes: s2.changes, refs: s2.refs.length })); assert(s2.created === 0, `run 2 created=${s2.created} (expected 0)`); assert(s2.rejected === 0, `run 2 rejected=${s2.rejected}`); assert(s2.events === 1, `run 2 events=${s2.events} (expected exactly 1 capacity_changed)`); assert(s2.changes.length === 1 && s2.changes[0]!.field === "itCapacityMw" && s2.changes[0]!.newValue === 48, "the only change is itCapacityMw → 48"); const cap = (await db.execute(sql`select significance, title from events where source_id = ${SOURCE} and event_type = 'capacity_changed' and new_value = '48'::jsonb`))[0]!; assert(Number(cap.significance) === 80, `+33 % MW change → significance 80 (got ${cap.significance})`); const facCount2 = await count(sql`select count(*)::int as n from facilities f where f.id in (select entity_id from entity_keys where connector_id = ${CONNECTOR} and entity_type = 'facility') and f.merged_into is null`); assert(facCount2 === 2, "still 2 facilities"); const provRows2 = await count(sql`select count(*)::int as n from provenance where connector_id = ${CONNECTOR} and is_current`); assert(provRows2 === provRows, `provenance row count unchanged on re-ingest (${provRows} → ${provRows2})`); // ---- run 3: pure repeat → nothing at all const s3 = await ingestEntities(run, entities(48), doc); assert(s3.created === 0 && s3.events === 0 && s3.changes.length === 0, `run 3 must be a no-op (created=${s3.created} events=${s3.events})`); // ---- dry run must not persist const before = await count(sql`select count(*)::int as n from facilities where name like 'Smoke Test%'`); await ingestEntities({ ...run, dryRun: true }, [{ ...(entities(48)[0] as NormalizedFacility), key: "smoke-test:fac:dry", name: "Smoke Test Dry Run", geo: { lat: 47.6, lng: -122.3, precision: "exact", source: "smoke-test" }, externalIds: {} }], doc); const after = await count(sql`select count(*)::int as n from facilities where name like 'Smoke Test%'`); assert(before === after, "dry run persisted nothing"); // ---- rankings / metrics / stats const r = await computeRankings(db); console.log("rankings", r); assert(r.rankings >= 20 && r.rows > 0, "rankings computed"); const top = (await db.execute(sql`select rows from rankings where key = 'countries_facilities' and is_current`))[0]!; assert(Array.isArray(top.rows) && (top.rows as Array<{ id: string }>).some((x) => x.id === "US"), "US appears in countries_facilities"); const m = await computeDailyMetrics(db); console.log("metrics", m); assert(m.rows > 0, "daily metrics written"); const st = await refreshStats(db); console.log("stats", st); const opStats = (await db.execute(sql`select stats from operators where name = 'Smoke Test Operator'`))[0]!; assert(Number((opStats.stats as Record).facilityCount) === 2, "operator stats facilityCount = 2"); console.log("\nSMOKE TEST PASSED"); } finally { await cleanup(); // leave aggregates consistent with the clean database await computeRankings(db); await computeDailyMetrics(db); await refreshStats(db); const left = await count(sql`select (select count(*) from facilities where name like 'Smoke Test%') + (select count(*) from operators where name like 'Smoke Test%') + (select count(*) from events where source_id = ${SOURCE}) + (select count(*) from provenance where connector_id = ${CONNECTOR}) as n`); console.log(`cleanup: ${left} smoke rows left`); await closeDb(); } } main().catch((e) => { console.error(e); process.exit(1); });