/** * Data-quality maintenance: * - `snapshotAndCheck()` — daily JSON snapshots (global totals, per-operator / per-country totals, project stages, * facility statuses) into `entity_snapshots` + automated regression checks against the previous snapshot * (facilities drop > 5 %, known MW jumps > 20 %, one operator gains > 10 GW, one connector creates > 500 projects, * one source changes hundreds of locations) → `system_alerts`. * - `qualitySweep()` — deterministic flags over the live database (largest values, scope, duplicates, orphans, * stale entities, project false-positive candidates…) into `quality_flags`, plus review priorities. * - `dataGaps()` — counts for the admin data-gaps page (facilities without operator / coordinates / capacity…). */ import { getDb, sql, type Db } from "@dci/db"; import { sha256 } from "@dci/core"; import { reviewPriority } from "./ingest/claims.js"; const day = () => new Date().toISOString().slice(0, 10); export interface SnapshotResult { day: string; snapshots: number; alerts: string[] } async function alert(db: Db, component: string, level: "warn" | "error", message: string, details: Record): Promise { await db.execute(sql`insert into system_alerts (id, level, component, message, details) values (${`alr_${Date.now().toString(36)}_${Math.random().toString(36).slice(2, 6)}`}, ${level}, ${component}, ${message.slice(0, 1000)}, ${JSON.stringify(details)}::jsonb)`); } export async function snapshotAndCheck(db: Db = getDb(), today = day()): Promise { const alerts: string[] = []; const g = (await db.execute(sql` select (select count(*) from facilities where merged_into is null)::int as facilities, (select count(*) from operators)::int as operators, (select count(*) from projects where merged_into is null and not hidden)::int as projects, (select count(distinct country_iso2) from facilities where merged_into is null and country_iso2 is not null)::int as countries, coalesce((select sum(coalesce(it_capacity_mw, total_power_mw)) from facilities where merged_into is null and status in ('operational','partially_operational','expansion') and not exists (select 1 from facilities c where c.parent_facility_id = facilities.id and c.merged_into is null and coalesce(c.it_capacity_mw, c.total_power_mw) is not null)), 0)::double precision as known_mw, coalesce((select sum(coalesce(planned_power_mw, it_capacity_mw, total_power_mw)) from facilities where merged_into is null and status = 'under_construction'), 0)::double precision as construction_mw, coalesce((select sum(planned_mw) from projects where merged_into is null and not hidden and status in ('rumored','proposed','announced','permitting','approved','under_construction','delayed')), 0)::double precision as project_mw, (select count(*) from events where detected_at > now() - interval '24 hours')::int as events_24h, (select count(*) from claims)::int as claims, (select count(*) from quality_flags where status = 'open')::int as open_flags`))[0]!; const totals = Object.fromEntries(Object.entries(g).map(([k, v]) => [k, Number(v)])); const statuses = await db.execute(sql`select status, count(*)::int as n from facilities where merged_into is null group by 1`); const stages = await db.execute(sql`select status, count(*)::int as n, coalesce(sum(planned_mw), 0)::double precision as mw from projects where merged_into is null and not hidden group by 1`); const ops = await db.execute(sql`select o.id, o.slug, count(f.id)::int as facilities, coalesce(sum(coalesce(f.it_capacity_mw, f.total_power_mw)), 0)::double precision as known_mw, (select coalesce(sum(p.planned_mw), 0) from projects p where p.operator_id = o.id and p.merged_into is null and not p.hidden)::double precision as project_mw from operators o left join facilities f on f.operator_id = o.id and f.merged_into is null group by o.id, o.slug having count(f.id) > 0 or exists (select 1 from projects p where p.operator_id = o.id and not p.hidden)`); const countries = await db.execute(sql`select country_iso2 as iso2, count(*)::int as facilities, coalesce(sum(coalesce(it_capacity_mw, total_power_mw)), 0)::double precision as known_mw from facilities where merged_into is null and country_iso2 is not null group by 1`); const rankings = await db.execute(sql`select key, rows from rankings where is_current`); const connectors = await db.execute(sql`select connector_id, coalesce(sum((stats->>'created')::int), 0)::int as created, count(*)::int as runs from connector_runs where started_at > now() - interval '24 hours' and status <> 'running' and task <> 'full' group by 1`); const locChanges = await db.execute(sql`select p.connector_id, count(*)::int as n from provenance p where p.field = 'geo' and p.last_observed > now() - interval '24 hours' and p.first_observed < p.last_observed - interval '1 hour' group by 1`); const previous = (await db.execute(sql`select payload from entity_snapshots where kind = 'global_totals' and key = 'global' and day < ${today}::date order by day desc limit 1`))[0]?.payload as Record | undefined; const prevOps = new Map>(); for (const r of await db.execute(sql`select key, payload from entity_snapshots where kind = 'operator_totals' and day = (select max(day) from entity_snapshots where kind = 'operator_totals' and day < ${today}::date)`)) prevOps.set(String(r.key), r.payload as Record); // ── regression checks if (previous) { const pf = Number(previous.facilities ?? 0), nf = totals.facilities!; if (pf > 100 && nf < pf * 0.95) { const m = `facilities dropped ${pf} → ${nf} (${(((nf - pf) / pf) * 100).toFixed(1)} %)`; alerts.push(m); await alert(db, "regression:facilities", "error", m, { previous: pf, now: nf }); } const pm = Number(previous.known_mw ?? 0), nm = totals.known_mw!; if (pm > 500 && nm > pm * 1.2) { const m = `known MW jumped ${Math.round(pm)} → ${Math.round(nm)} (+${(((nm - pm) / pm) * 100).toFixed(1)} %)`; alerts.push(m); await alert(db, "regression:known_mw", "warn", m, { previous: pm, now: nm }); } const pp = Number(previous.projects ?? 0), np = totals.projects!; if (pp > 50 && np > pp * 1.5) { const m = `projects jumped ${pp} → ${np}`; alerts.push(m); await alert(db, "regression:projects", "warn", m, { previous: pp, now: np }); } const pc = Number(previous.countries ?? 0), nc = totals.countries!; if (pc > 20 && nc < pc - 3) { const m = `countries dropped ${pc} → ${nc}`; alerts.push(m); await alert(db, "regression:countries", "warn", m, { previous: pc, now: nc }); } } for (const o of ops) { const prev = prevOps.get(String(o.id)); const gain = Number(o.known_mw) + Number(o.project_mw) - (prev ? Number(prev.known_mw ?? 0) + Number(prev.project_mw ?? 0) : 0); if (prev && gain > 10_000) { const m = `operator ${String(o.slug)} gained ${Math.round(gain)} MW in one day`; alerts.push(m); await alert(db, "regression:operator_mw", "error", m, { operator: o.slug, gainMw: gain }); } } for (const c of connectors) if (Number(c.created) > 500) { const m = `connector ${String(c.connector_id)} created ${Number(c.created)} records in 24 h`; alerts.push(m); await alert(db, `regression:connector:${String(c.connector_id)}`, "warn", m, { created: Number(c.created), runs: Number(c.runs) }); } for (const l of locChanges) if (Number(l.n) > 200) { const m = `connector ${String(l.connector_id)} changed ${Number(l.n)} locations in 24 h`; alerts.push(m); await alert(db, `regression:locations:${String(l.connector_id)}`, "warn", m, { changed: Number(l.n) }); } // ── snapshots let n = 0; await db.transaction(async (tx) => { const put = async (kind: string, key: string, payload: unknown) => { await tx.execute(sql`insert into entity_snapshots (day, kind, key, payload) values (${today}::date, ${kind}, ${key}, ${JSON.stringify(payload)}::jsonb) on conflict (day, kind, key) do update set payload = excluded.payload, created_at = now()`); n++; }; await put("global_totals", "global", totals); await put("facility_status", "global", Object.fromEntries(statuses.map((r) => [String(r.status), Number(r.n)]))); await put("project_stage", "global", Object.fromEntries(stages.map((r) => [String(r.status), { count: Number(r.n), mw: Number(r.mw) }]))); for (const o of ops) await put("operator_totals", String(o.id), { slug: o.slug, facilities: Number(o.facilities), known_mw: Number(o.known_mw), project_mw: Number(o.project_mw) }); for (const c of countries) await put("country_totals", String(c.iso2), { facilities: Number(c.facilities), known_mw: Number(c.known_mw) }); for (const r of rankings) await put("ranking", String(r.key), { rows: (r.rows as unknown[]).slice(0, 100) }); }); return { day: today, snapshots: n, alerts }; } export interface SweepResult { flagged: number; resolved: number; byCode: Record } /** Batch quality sweep over the live database. Idempotent (dedupe keys). */ export async function qualitySweep(db: Db = getDb()): Promise { const byCode: Record = {}; let flagged = 0, resolved = 0; const flag = async (entityType: string, entityId: string, code: string, severity: "warn" | "critical", message: string, field: string | null, details: Record, priority: number) => { const dedupeKey = `${entityType}|${entityId}|${code}|${field ?? ""}`; await db.execute(sql`insert into quality_flags (id, entity_type, entity_id, code, severity, field, message, details, priority, status, dedupe_key) values (${`flg_${sha256(dedupeKey).slice(0, 20)}`}, ${entityType}, ${entityId}, ${code}, ${severity}, ${field}, ${message.slice(0, 1000)}, ${JSON.stringify(details)}::jsonb, ${priority}, 'open', ${dedupeKey}) on conflict (dedupe_key) do update set message = excluded.message, details = excluded.details, priority = excluded.priority, updated_at = now(), status = case when quality_flags.status = 'dismissed' then 'dismissed' when quality_flags.status = 'resolved' and quality_flags.resolved_by <> 'system' then 'resolved' else 'open' end`); byCode[code] = (byCode[code] ?? 0) + 1; flagged++; }; const clear = async (code: string, keepIds: string[], entityType: string) => { const r = await db.execute(sql`update quality_flags set status = 'resolved', resolution = 'auto: condition cleared', resolved_by = 'system', resolved_at = now(), updated_at = now() where code = ${code} and entity_type = ${entityType} and status = 'open' ${keepIds.length ? sql`and entity_id not in ${keepIds}` : sql``} returning id`); resolved += r.length; }; // suspicious MW on facilities (single site > 1 GW without campus designation, building > 500 MW, IT > total) const bigFac = await db.execute(sql`select id, name, slug, coalesce(it_capacity_mw, total_power_mw) as mw, record_scope, confidence from facilities where merged_into is null and coalesce(it_capacity_mw, total_power_mw) > 1000 and record_scope <> 'campus' and name !~* '(campus|park|complex|hub|cluster|mega ?site)'`); for (const r of bigFac) await flag("facility", String(r.id), "mw_single_site_gt_1000", "critical", `${String(r.name)}: ${Number(r.mw)} MW on a single ${String(r.record_scope)} record without a campus designation`, "itCapacityMw", { slug: r.slug, mw: Number(r.mw) }, reviewPriority({ mw: Number(r.mw), confidence: String(r.confidence), severity: "critical", homepageVisible: true })); await clear("mw_single_site_gt_1000", bigFac.map((r) => String(r.id)), "facility"); const itGtTotal = await db.execute(sql`select id, name, slug, it_capacity_mw, total_power_mw from facilities where merged_into is null and it_capacity_mw is not null and total_power_mw is not null and it_capacity_mw > total_power_mw * 1.05`); for (const r of itGtTotal) await flag("facility", String(r.id), "mw_it_gt_total", "warn", `${String(r.name)}: IT capacity ${Number(r.it_capacity_mw)} MW exceeds total power ${Number(r.total_power_mw)} MW`, "itCapacityMw", { slug: r.slug }, reviewPriority({ mw: Number(r.it_capacity_mw), severity: "warn" })); await clear("mw_it_gt_total", itGtTotal.map((r) => String(r.id)), "facility"); // tiny buildings with huge MW (the 0.779 → 779 class of unit errors) const dense = await db.execute(sql`select id, name, slug, it_capacity_mw, building_sqm from facilities where merged_into is null and it_capacity_mw is not null and building_sqm is not null and building_sqm > 0 and it_capacity_mw / building_sqm > 0.05`); for (const r of dense) await flag("facility", String(r.id), "mw_density_implausible", "critical", `${String(r.name)}: ${Number(r.it_capacity_mw)} MW in ${Number(r.building_sqm)} m² (${(Number(r.it_capacity_mw) * 1000 / Number(r.building_sqm)).toFixed(0)} kW/m²) — unit error?`, "itCapacityMw", { slug: r.slug }, reviewPriority({ mw: Number(r.it_capacity_mw), severity: "critical" })); await clear("mw_density_implausible", dense.map((r) => String(r.id)), "facility"); // projects: false-positive candidates, headline names, huge figures, missing location, unknown scope const fp = await db.execute(sql`select id, name, slug, planned_mw, investment_usd, project_class, evidence_level from projects where merged_into is null and not hidden and ( name ~* '\\m(appoint|names? [A-Z]\\w+ [A-Z]\\w+ as|joins|welcomes|promot|retire|award|shortlist|finalist|webinar|podcast|interview|market to (surpass|reach|hit|grow)|market size|cagr|report|survey|headquarters|head office|academy|workforce|partner page|homepage|bubble|comparing|why |how |what |ppa\\M|power purchase|sustainab|net.zero|carbon|financ|refinanc|raises \\$|series [a-e]\\M|bond|loan|credit facility|earnings|results|revenue|acquires [A-Z]\\w+ (group|holdings|inc|ltd|llc)|merger|takeover)' or (project_class is not null and project_class not in ('NEW_BUILD','EXPANSION','CONSTRUCTION_START','PERMIT','LAND_ACQUISITION','GRID_CONNECTION')))`); for (const r of fp) await flag("project", String(r.id), "project_false_positive_candidate", "critical", `${String(r.name)}: headline / class (${String(r.project_class ?? "n/a")}) does not describe physical development`, "name", { slug: r.slug, plannedMw: r.planned_mw, investmentUsd: r.investment_usd }, reviewPriority({ mw: Number(r.planned_mw ?? 0), investmentUsd: Number(r.investment_usd ?? 0), severity: "critical", homepageVisible: true })); await clear("project_false_positive_candidate", fp.map((r) => String(r.id)), "project"); const headline = await db.execute(sql`select id, name, slug, planned_mw from projects where merged_into is null and not hidden and (name ~ '\\s[|]\\s|\\s[–—]\\s| - ' or length(name) > 90)`); for (const r of headline) await flag("project", String(r.id), "project_title_like_name", "warn", `${String(r.name)}: project name looks like an article headline`, "name", { slug: r.slug }, reviewPriority({ mw: Number(r.planned_mw ?? 0), severity: "warn" })); await clear("project_title_like_name", headline.map((r) => String(r.id)), "project"); const bigPrj = await db.execute(sql`select id, name, slug, planned_mw, capacity_scope from projects where merged_into is null and not hidden and planned_mw >= 2000 and name !~* '(campus|park|complex|hub|cluster|gigafactory)'`); for (const r of bigPrj) await flag("project", String(r.id), "mw_project_gt_2000", "critical", `${String(r.name)}: ${Number(r.planned_mw)} MW planned (scope ${String(r.capacity_scope ?? "unknown")}) — verify it is one site`, "plannedMw", { slug: r.slug }, reviewPriority({ mw: Number(r.planned_mw), severity: "critical", homepageVisible: true })); await clear("mw_project_gt_2000", bigPrj.map((r) => String(r.id)), "project"); const bigInv = await db.execute(sql`select id, name, slug, investment_usd from projects where merged_into is null and not hidden and investment_usd > 50e9`); for (const r of bigInv) await flag("project", String(r.id), "inv_single_site_gt_50b", "critical", `${String(r.name)}: $${(Number(r.investment_usd) / 1e9).toFixed(1)}B on one project — verify scope`, "investmentUsd", { slug: r.slug }, reviewPriority({ investmentUsd: Number(r.investment_usd), severity: "critical" })); await clear("inv_single_site_gt_50b", bigInv.map((r) => String(r.id)), "project"); const noLoc = await db.execute(sql`select id, name, slug from projects where merged_into is null and not hidden and country_iso2 is null and lat is null`); for (const r of noLoc) await flag("project", String(r.id), "project_no_location", "warn", `${String(r.name)}: no country and no coordinates`, "countryIso2", { slug: r.slug }, 15); await clear("project_no_location", noLoc.map((r) => String(r.id)), "project"); const unknownScope = await db.execute(sql`select id, name, slug, planned_mw from projects where merged_into is null and not hidden and planned_mw is not null and (capacity_scope is null or capacity_scope not in ('building','facility','campus'))`); for (const r of unknownScope) await flag("project", String(r.id), "project_unknown_scope", "warn", `${String(r.name)}: ${Number(r.planned_mw)} MW with no site-scoped evidence`, "plannedMw", { slug: r.slug }, reviewPriority({ mw: Number(r.planned_mw), severity: "warn" })); await clear("project_unknown_scope", unknownScope.map((r) => String(r.id)), "project"); // facilities: location / country mismatch, missing country, orphan operators, stale const noCountry = await db.execute(sql`select id, name, slug from facilities where merged_into is null and country_iso2 is null`); for (const r of noCountry) await flag("facility", String(r.id), "facility_no_country", "warn", `${String(r.name)}: no country`, "countryIso2", { slug: r.slug }, 8); await clear("facility_no_country", noCountry.map((r) => String(r.id)), "facility"); const stale = await db.execute(sql`select id, name, slug, last_verified from facilities where merged_into is null and status in ('announced','under_construction','permitting','approved') and coalesce(last_verified, first_seen) < now() - interval '400 days'`); for (const r of stale) await flag("facility", String(r.id), "facility_stale_pipeline", "warn", `${String(r.name)}: pipeline status not verified for over 400 days`, "status", { slug: r.slug, lastVerified: r.last_verified }, 12); await clear("facility_stale_pipeline", stale.map((r) => String(r.id)), "facility"); const orphanOps = await db.execute(sql`select o.id, o.name, o.slug from operators o where not exists (select 1 from facilities f where f.operator_id = o.id or f.owner_id = o.id) and not exists (select 1 from projects p where p.operator_id = o.id) and not exists (select 1 from cloud_regions c where c.provider_id = o.id) and not exists (select 1 from facility_tenants t where t.operator_id = o.id) and o.created_at < now() - interval '7 days'`); for (const r of orphanOps) await flag("operator", String(r.id), "operator_orphan", "warn", `${String(r.name)}: operator with no facility, project, region or tenancy`, null, { slug: r.slug }, 5); await clear("operator_orphan", orphanOps.map((r) => String(r.id)), "operator"); const badSlug = await db.execute(sql`select id, name, slug from operators where slug ~ '^(item|op|operator)(-\\d+)?$' or slug = ''`); for (const r of badSlug) await flag("operator", String(r.id), "operator_bad_slug", "warn", `${String(r.name)}: slug "${String(r.slug)}" is not derived from the name (non-Latin script?)`, "slug", {}, 6); await clear("operator_bad_slug", badSlug.map((r) => String(r.id)), "operator"); // review priority roll-up await db.execute(sql`update facilities f set review_priority = coalesce((select max(priority) from quality_flags q where q.entity_type = 'facility' and q.entity_id = f.id and q.status = 'open'), 0) where merged_into is null`); await db.execute(sql`update projects p set review_priority = coalesce((select max(priority) from quality_flags q where q.entity_type = 'project' and q.entity_id = p.id and q.status = 'open'), 0) where merged_into is null`); return { flagged, resolved, byCode }; } export interface DataGaps { [key: string]: number } /** Counts for /admin/data-gaps — what we do not know, ranked for enrichment. */ export async function dataGaps(db: Db = getDb()): Promise { const r = (await db.execute(sql` select (select count(*) from facilities where merged_into is null and operator_id is null)::int as facilities_without_operator, (select count(*) from facilities where merged_into is null and lat is null)::int as facilities_without_coordinates, (select count(*) from facilities where merged_into is null and geo_precision in ('city','metro','approximate') )::int as facilities_imprecise_coordinates, (select count(*) from facilities where merged_into is null and coalesce(it_capacity_mw, total_power_mw, planned_power_mw) is null and parent_facility_id is null)::int as facilities_without_capacity, (select count(*) from facilities where merged_into is null and status = 'unknown')::int as facilities_unknown_status, (select count(*) from facilities where merged_into is null and facility_type = 'unknown')::int as facilities_unknown_type, (select count(*) from facilities where merged_into is null and opened_on is null)::int as facilities_without_opening_date, (select count(*) from facilities where merged_into is null and source_count <= 1)::int as facilities_single_source, (select count(*) from facilities where merged_into is null and coalesce(last_verified, first_seen) < now() - interval '180 days')::int as facilities_stale_sources, (select count(*) from facilities where merged_into is null and country_iso2 is null)::int as facilities_without_country, (select count(*) from projects where merged_into is null and not hidden and lat is null)::int as projects_without_coordinates, (select count(*) from projects where merged_into is null and not hidden and country_iso2 is null)::int as projects_without_country, (select count(*) from projects where merged_into is null and not hidden and operator_id is null)::int as projects_without_operator, (select count(*) from projects where merged_into is null and not hidden and planned_mw is null)::int as projects_without_capacity, (select count(*) from projects where merged_into is null and not hidden and expected_opening is null)::int as projects_without_expected_opening, (select count(*) from projects where merged_into is null and not hidden and source_url is null)::int as projects_without_source, (select count(*) from projects where merged_into is null and not hidden and facility_id is null)::int as projects_without_facility_link, (select count(*) from operators where hq_country_iso2 is null)::int as operators_without_hq, (select count(*) from operators where website is null)::int as operators_without_website, (select count(*) from entity_matches where status = 'pending')::int as pending_matches, (select count(*) from quality_flags where status = 'open')::int as open_quality_flags, (select count(*) from claims where status = 'unscoped')::int as unscoped_claims, (select count(*) from claims where status = 'review')::int as claims_in_review, (select count(*) from cloud_regions where lat is null)::int as cloud_regions_without_coordinates, (select count(*) from ixps where metro_id is null)::int as ixps_without_metro`))[0]!; return Object.fromEntries(Object.entries(r).map(([k, v]) => [k, Number(v)])); }