/** * Daily metric snapshots → `daily_metrics` (Postgres, upsert on day/metric/dim) mirrored to ClickHouse `entity_daily`. * Dimensions: global, country:, operator: (top 100 by facility count), metro:. */ import { getDb, sql, type Db } from "@dci/db"; import { chInsert } from "@dci/db/clickhouse"; import { aggregateCountries, aggregateMetros, aggregateOperators, type Aggregate } from "./rankings.js"; // stats refresh lives with the aggregates in rankings.ts; re-exported here for the runtime (cli/main import it from metrics). export { refreshStats } from "./rankings.js"; export interface MetricsResult { day: string; rows: number; dims: number; } type MetricRow = { day: string; metric: string; dim: string; value: number }; function facilityMetrics(dim: string, a: Aggregate, day: string): MetricRow[] { return [ { day, metric: "facilities_total", dim, value: a.facilities }, { day, metric: "operational", dim, value: a.operational }, { day, metric: "under_construction", dim, value: a.construction }, { day, metric: "planned", dim, value: a.planned }, { day, metric: "known_mw", dim, value: a.knownMw }, { day, metric: "construction_mw", dim, value: a.constructionMw }, { day, metric: "planned_mw", dim, value: a.plannedMw }, { day, metric: "mw_coverage", dim, value: a.coverage }, ]; } async function eventCounts(db: Db, column: "country_iso2" | "operator_id" | "metro_id"): Promise> { const rows = await db.execute(sql`select ${sql.identifier(column)} as dim, count(*)::int as n from events where detected_at >= now() - interval '24 hours' and ${sql.identifier(column)} is not null group by 1`); return new Map(rows.map((r) => [String(r.dim), Number(r.n)])); } export async function computeDailyMetrics(db: Db = getDb(), day = new Date().toISOString().slice(0, 10)): Promise { const [countries, metros, operators, evCountry, evOperator, evMetro] = await Promise.all([aggregateCountries(db), aggregateMetros(db), aggregateOperators(db, 100), eventCounts(db, "country_iso2"), eventCounts(db, "operator_id"), eventCounts(db, "metro_id")]); const g = (await db.execute(sql` select (select count(*) from facilities where merged_into is null)::int as facilities, (select count(*) from facilities where merged_into is null and status in ('operational','partially_operational','expansion'))::int as operational, (select count(*) from facilities where merged_into is null and status = 'under_construction')::int as construction, (select count(*) from facilities where merged_into is null and status in ('rumored','proposed','announced','permitting','approved','delayed'))::int as planned, coalesce((select sum(coalesce(it_capacity_mw, total_power_mw)) from facilities where merged_into is null and status in ('operational','partially_operational','expansion')), 0) 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) as construction_mw, coalesce((select sum(coalesce(planned_power_mw, it_capacity_mw, total_power_mw)) from facilities where merged_into is null and status in ('rumored','proposed','announced','permitting','approved','delayed')), 0) as planned_mw, (select count(*) from facilities where merged_into is null and coalesce(it_capacity_mw, total_power_mw, planned_power_mw) is not null)::int as with_mw, (select count(*) from events where detected_at >= now() - interval '24 hours')::int as events_24h, (select count(*) from events where detected_at >= now() - interval '7 days')::int as events_7d, (select count(*) from sources)::int as sources, (select count(*) from documents)::int as documents, (select count(*) from operators)::int as operators, (select count(distinct country_iso2) from facilities where merged_into is null and country_iso2 is not null)::int as countries, (select count(*) from metros)::int as metros, (select count(*) from cloud_regions where status <> 'retired')::int as cloud_regions, (select count(*) from ixps)::int as ixps, (select count(*) from projects)::int as projects, (select count(*) from news_items)::int as news_items, (select count(*) from provenance where is_current)::int as provenance_rows, (select count(*) from entity_matches where status = 'pending')::int as pending_matches`))[0]!; const n = (k: string) => Number(g[k] ?? 0); const rows: MetricRow[] = []; const push = (metric: string, dim: string, value: number) => rows.push({ day, metric, dim, value: Number.isFinite(value) ? value : 0 }); for (const k of ["facilities_total:facilities", "operational:operational", "under_construction:construction", "planned:planned", "known_mw:known_mw", "construction_mw:construction_mw", "planned_mw:planned_mw", "events_24h:events_24h", "events_7d:events_7d", "sources:sources", "documents:documents", "operators:operators", "countries:countries", "metros:metros", "cloud_regions:cloud_regions", "ixps:ixps", "projects:projects", "news_items:news_items", "provenance_rows:provenance_rows", "pending_matches:pending_matches"]) { const [metric, col] = k.split(":") as [string, string]; push(metric, "global", n(col)); } push("mw_coverage", "global", n("facilities") ? Math.round((n("with_mw") / n("facilities")) * 1000) / 1000 : 0); for (const a of countries) { if (!a.facilities && !a.cloudRegions && !a.projects) continue; const dim = `country:${a.id}`; rows.push(...facilityMetrics(dim, a, day)); push("events_24h", dim, evCountry.get(a.id) ?? 0); push("cloud_regions", dim, a.cloudRegions); push("projects", dim, a.projects); } for (const a of operators) { if (!a.facilities) continue; const dim = `operator:${a.id}`; rows.push(...facilityMetrics(dim, a, day)); push("events_24h", dim, evOperator.get(a.id) ?? 0); push("countries", dim, a.countries); } for (const a of metros) { if (!a.facilities) continue; const dim = `metro:${a.id}`; rows.push(...facilityMetrics(dim, a, day)); push("events_24h", dim, evMetro.get(a.id) ?? 0); push("operators", dim, a.operators); } await db.transaction(async (tx) => { for (let i = 0; i < rows.length; i += 500) { const chunk = rows.slice(i, i + 500); const values = sql.join(chunk.map((r) => sql`(${r.day}::date, ${r.metric}, ${r.dim}, ${r.value})`), sql`, `); await tx.execute(sql`insert into daily_metrics (day, metric, dim, value) values ${values} on conflict (day, metric, dim) do update set value = excluded.value`); } }); try { await chInsert("entity_daily", rows.map((r) => ({ day: r.day, metric: r.metric, dim: r.dim, value: r.value }))); } catch (e) { console.warn(`[metrics] clickhouse mirror skipped: ${(e as Error).message}`); } return { day, rows: rows.length, dims: new Set(rows.map((r) => r.dim)).size }; }