SPB Git forge
38commits 1branches 0releases
338.7 MBsize
maindefault branch
3 h agolast push
HTML 53.9% TypeScript 44.5% JavaScript 0.6% SQL 0.5%
23.1 KB · 190 lines typescript
Raw Blame History
1/**2 * Data-quality maintenance:3 *  - `snapshotAndCheck()` — daily JSON snapshots (global totals, per-operator / per-country totals, project stages,4 *    facility statuses) into `entity_snapshots` + automated regression checks against the previous snapshot5 *    (facilities drop > 5 %, known MW jumps > 20 %, one operator gains > 10 GW, one connector creates > 500 projects,6 *    one source changes hundreds of locations) → `system_alerts`.7 *  - `qualitySweep()` — deterministic flags over the live database (largest values, scope, duplicates, orphans,8 *    stale entities, project false-positive candidates…) into `quality_flags`, plus review priorities.9 *  - `dataGaps()` — counts for the admin data-gaps page (facilities without operator / coordinates / capacity…).10 */11import { getDb, sql, type Db } from "@dci/db";12import { sha256 } from "@dci/core";13import { reviewPriority } from "./ingest/claims.js";1415const day = () => new Date().toISOString().slice(0, 10);1617export interface SnapshotResult { day: string; snapshots: number; alerts: string[] }1819async function alert(db: Db, component: string, level: "warn" | "error", message: string, details: Record<string, unknown>): Promise<void> {20  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)`);21}2223export async function snapshotAndCheck(db: Db = getDb(), today = day()): Promise<SnapshotResult> {24  const alerts: string[] = [];25  const g = (await db.execute(sql`26    select27      (select count(*) from facilities where merged_into is null)::int as facilities,28      (select count(*) from operators)::int as operators,29      (select count(*) from projects where merged_into is null and not hidden)::int as projects,30      (select count(distinct country_iso2) from facilities where merged_into is null and country_iso2 is not null)::int as countries,31      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,32      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,33      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,34      (select count(*) from events where detected_at > now() - interval '24 hours')::int as events_24h,35      (select count(*) from claims)::int as claims,36      (select count(*) from quality_flags where status = 'open')::int as open_flags`))[0]!;37  const totals = Object.fromEntries(Object.entries(g).map(([k, v]) => [k, Number(v)]));38  const statuses = await db.execute(sql`select status, count(*)::int as n from facilities where merged_into is null group by 1`);39  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`);40  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,41      (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_mw42    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)`);43  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`);44  const rankings = await db.execute(sql`select key, rows from rankings where is_current`);45  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`);46  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`);4748  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<string, number> | undefined;49  const prevOps = new Map<string, Record<string, unknown>>();50  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<string, unknown>);5152  // ── regression checks53  if (previous) {54    const pf = Number(previous.facilities ?? 0), nf = totals.facilities!;55    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 }); }56    const pm = Number(previous.known_mw ?? 0), nm = totals.known_mw!;57    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 }); }58    const pp = Number(previous.projects ?? 0), np = totals.projects!;59    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 }); }60    const pc = Number(previous.countries ?? 0), nc = totals.countries!;61    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 }); }62  }63  for (const o of ops) {64    const prev = prevOps.get(String(o.id));65    const gain = Number(o.known_mw) + Number(o.project_mw) - (prev ? Number(prev.known_mw ?? 0) + Number(prev.project_mw ?? 0) : 0);66    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 }); }67  }68  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) }); }69  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) }); }7071  // ── snapshots72  let n = 0;73  await db.transaction(async (tx) => {74    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++; };75    await put("global_totals", "global", totals);76    await put("facility_status", "global", Object.fromEntries(statuses.map((r) => [String(r.status), Number(r.n)])));77    await put("project_stage", "global", Object.fromEntries(stages.map((r) => [String(r.status), { count: Number(r.n), mw: Number(r.mw) }])));78    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) });79    for (const c of countries) await put("country_totals", String(c.iso2), { facilities: Number(c.facilities), known_mw: Number(c.known_mw) });80    for (const r of rankings) await put("ranking", String(r.key), { rows: (r.rows as unknown[]).slice(0, 100) });81  });82  return { day: today, snapshots: n, alerts };83}8485export interface SweepResult { flagged: number; resolved: number; byCode: Record<string, number> }8687/** Batch quality sweep over the live database. Idempotent (dedupe keys). */88export async function qualitySweep(db: Db = getDb()): Promise<SweepResult> {89  const byCode: Record<string, number> = {};90  let flagged = 0, resolved = 0;91  const flag = async (entityType: string, entityId: string, code: string, severity: "warn" | "critical", message: string, field: string | null, details: Record<string, unknown>, priority: number) => {92    const dedupeKey = `${entityType}|${entityId}|${code}|${field ?? ""}`;93    await db.execute(sql`insert into quality_flags (id, entity_type, entity_id, code, severity, field, message, details, priority, status, dedupe_key)94      values (${`flg_${sha256(dedupeKey).slice(0, 20)}`}, ${entityType}, ${entityId}, ${code}, ${severity}, ${field}, ${message.slice(0, 1000)}, ${JSON.stringify(details)}::jsonb, ${priority}, 'open', ${dedupeKey})95      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`);96    byCode[code] = (byCode[code] ?? 0) + 1; flagged++;97  };98  const clear = async (code: string, keepIds: string[], entityType: string) => {99    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`);100    resolved += r.length;101  };102103  // suspicious MW on facilities (single site > 1 GW without campus designation, building > 500 MW, IT > total)104  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)'`);105  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 }));106  await clear("mw_single_site_gt_1000", bigFac.map((r) => String(r.id)), "facility");107  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`);108  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" }));109  await clear("mw_it_gt_total", itGtTotal.map((r) => String(r.id)), "facility");110  // tiny buildings with huge MW (the 0.779 → 779 class of unit errors)111  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`);112  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" }));113  await clear("mw_density_implausible", dense.map((r) => String(r.id)), "facility");114115  // projects: false-positive candidates, headline names, huge figures, missing location, unknown scope116  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 (117      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)'118      or (project_class is not null and project_class not in ('NEW_BUILD','EXPANSION','CONSTRUCTION_START','PERMIT','LAND_ACQUISITION','GRID_CONNECTION')))`);119  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 }));120  await clear("project_false_positive_candidate", fp.map((r) => String(r.id)), "project");121  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)`);122  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" }));123  await clear("project_title_like_name", headline.map((r) => String(r.id)), "project");124  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)'`);125  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 }));126  await clear("mw_project_gt_2000", bigPrj.map((r) => String(r.id)), "project");127  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`);128  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" }));129  await clear("inv_single_site_gt_50b", bigInv.map((r) => String(r.id)), "project");130  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`);131  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);132  await clear("project_no_location", noLoc.map((r) => String(r.id)), "project");133  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'))`);134  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" }));135  await clear("project_unknown_scope", unknownScope.map((r) => String(r.id)), "project");136137  // facilities: location / country mismatch, missing country, orphan operators, stale138  const noCountry = await db.execute(sql`select id, name, slug from facilities where merged_into is null and country_iso2 is null`);139  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);140  await clear("facility_no_country", noCountry.map((r) => String(r.id)), "facility");141  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'`);142  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);143  await clear("facility_stale_pipeline", stale.map((r) => String(r.id)), "facility");144  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'`);145  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);146  await clear("operator_orphan", orphanOps.map((r) => String(r.id)), "operator");147  const badSlug = await db.execute(sql`select id, name, slug from operators where slug ~ '^(item|op|operator)(-\\d+)?$' or slug = ''`);148  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);149  await clear("operator_bad_slug", badSlug.map((r) => String(r.id)), "operator");150151  // review priority roll-up152  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`);153  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`);154  return { flagged, resolved, byCode };155}156157export interface DataGaps { [key: string]: number }158159/** Counts for /admin/data-gaps — what we do not know, ranked for enrichment. */160export async function dataGaps(db: Db = getDb()): Promise<DataGaps> {161  const r = (await db.execute(sql`162    select163      (select count(*) from facilities where merged_into is null and operator_id is null)::int as facilities_without_operator,164      (select count(*) from facilities where merged_into is null and lat is null)::int as facilities_without_coordinates,165      (select count(*) from facilities where merged_into is null and geo_precision in ('city','metro','approximate') )::int as facilities_imprecise_coordinates,166      (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,167      (select count(*) from facilities where merged_into is null and status = 'unknown')::int as facilities_unknown_status,168      (select count(*) from facilities where merged_into is null and facility_type = 'unknown')::int as facilities_unknown_type,169      (select count(*) from facilities where merged_into is null and opened_on is null)::int as facilities_without_opening_date,170      (select count(*) from facilities where merged_into is null and source_count <= 1)::int as facilities_single_source,171      (select count(*) from facilities where merged_into is null and coalesce(last_verified, first_seen) < now() - interval '180 days')::int as facilities_stale_sources,172      (select count(*) from facilities where merged_into is null and country_iso2 is null)::int as facilities_without_country,173      (select count(*) from projects where merged_into is null and not hidden and lat is null)::int as projects_without_coordinates,174      (select count(*) from projects where merged_into is null and not hidden and country_iso2 is null)::int as projects_without_country,175      (select count(*) from projects where merged_into is null and not hidden and operator_id is null)::int as projects_without_operator,176      (select count(*) from projects where merged_into is null and not hidden and planned_mw is null)::int as projects_without_capacity,177      (select count(*) from projects where merged_into is null and not hidden and expected_opening is null)::int as projects_without_expected_opening,178      (select count(*) from projects where merged_into is null and not hidden and source_url is null)::int as projects_without_source,179      (select count(*) from projects where merged_into is null and not hidden and facility_id is null)::int as projects_without_facility_link,180      (select count(*) from operators where hq_country_iso2 is null)::int as operators_without_hq,181      (select count(*) from operators where website is null)::int as operators_without_website,182      (select count(*) from entity_matches where status = 'pending')::int as pending_matches,183      (select count(*) from quality_flags where status = 'open')::int as open_quality_flags,184      (select count(*) from claims where status = 'unscoped')::int as unscoped_claims,185      (select count(*) from claims where status = 'review')::int as claims_in_review,186      (select count(*) from cloud_regions where lat is null)::int as cloud_regions_without_coordinates,187      (select count(*) from ixps where metro_id is null)::int as ixps_without_metro`))[0]!;188  return Object.fromEntries(Object.entries(r).map(([k, v]) => [k, Number(v)]));189}190