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%
6.9 KB · 107 lines typescript
Raw Blame History
1/**2 * Daily metric snapshots → `daily_metrics` (Postgres, upsert on day/metric/dim) mirrored to ClickHouse `entity_daily`.3 * Dimensions: global, country:<iso2>, operator:<id> (top 100 by facility count), metro:<id>.4 */5import { getDb, sql, type Db } from "@dci/db";6import { chInsert } from "@dci/db/clickhouse";7import { aggregateCountries, aggregateMetros, aggregateOperators, type Aggregate } from "./rankings.js";89// stats refresh lives with the aggregates in rankings.ts; re-exported here for the runtime (cli/main import it from metrics).10export { refreshStats } from "./rankings.js";1112export interface MetricsResult {13  day: string;14  rows: number;15  dims: number;16}1718type MetricRow = { day: string; metric: string; dim: string; value: number };1920function facilityMetrics(dim: string, a: Aggregate, day: string): MetricRow[] {21  return [22    { day, metric: "facilities_total", dim, value: a.facilities },23    { day, metric: "operational", dim, value: a.operational },24    { day, metric: "under_construction", dim, value: a.construction },25    { day, metric: "planned", dim, value: a.planned },26    { day, metric: "known_mw", dim, value: a.knownMw },27    { day, metric: "construction_mw", dim, value: a.constructionMw },28    { day, metric: "planned_mw", dim, value: a.plannedMw },29    { day, metric: "mw_coverage", dim, value: a.coverage },30  ];31}3233async function eventCounts(db: Db, column: "country_iso2" | "operator_id" | "metro_id"): Promise<Map<string, number>> {34  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`);35  return new Map(rows.map((r) => [String(r.dim), Number(r.n)]));36}3738export async function computeDailyMetrics(db: Db = getDb(), day = new Date().toISOString().slice(0, 10)): Promise<MetricsResult> {39  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")]);40  const g = (await db.execute(sql`41    select42      (select count(*) from facilities where merged_into is null)::int as facilities,43      (select count(*) from facilities where merged_into is null and status in ('operational','partially_operational','expansion'))::int as operational,44      (select count(*) from facilities where merged_into is null and status = 'under_construction')::int as construction,45      (select count(*) from facilities where merged_into is null and status in ('rumored','proposed','announced','permitting','approved','delayed'))::int as planned,46      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,47      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,48      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,49      (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,50      (select count(*) from events where detected_at >= now() - interval '24 hours')::int as events_24h,51      (select count(*) from events where detected_at >= now() - interval '7 days')::int as events_7d,52      (select count(*) from sources)::int as sources,53      (select count(*) from documents)::int as documents,54      (select count(*) from operators)::int as operators,55      (select count(distinct country_iso2) from facilities where merged_into is null and country_iso2 is not null)::int as countries,56      (select count(*) from metros)::int as metros,57      (select count(*) from cloud_regions where status <> 'retired')::int as cloud_regions,58      (select count(*) from ixps)::int as ixps,59      (select count(*) from projects)::int as projects,60      (select count(*) from news_items)::int as news_items,61      (select count(*) from provenance where is_current)::int as provenance_rows,62      (select count(*) from entity_matches where status = 'pending')::int as pending_matches`))[0]!;63  const n = (k: string) => Number(g[k] ?? 0);64  const rows: MetricRow[] = [];65  const push = (metric: string, dim: string, value: number) => rows.push({ day, metric, dim, value: Number.isFinite(value) ? value : 0 });66  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"]) {67    const [metric, col] = k.split(":") as [string, string];68    push(metric, "global", n(col));69  }70  push("mw_coverage", "global", n("facilities") ? Math.round((n("with_mw") / n("facilities")) * 1000) / 1000 : 0);71  for (const a of countries) {72    if (!a.facilities && !a.cloudRegions && !a.projects) continue;73    const dim = `country:${a.id}`;74    rows.push(...facilityMetrics(dim, a, day));75    push("events_24h", dim, evCountry.get(a.id) ?? 0);76    push("cloud_regions", dim, a.cloudRegions);77    push("projects", dim, a.projects);78  }79  for (const a of operators) {80    if (!a.facilities) continue;81    const dim = `operator:${a.id}`;82    rows.push(...facilityMetrics(dim, a, day));83    push("events_24h", dim, evOperator.get(a.id) ?? 0);84    push("countries", dim, a.countries);85  }86  for (const a of metros) {87    if (!a.facilities) continue;88    const dim = `metro:${a.id}`;89    rows.push(...facilityMetrics(dim, a, day));90    push("events_24h", dim, evMetro.get(a.id) ?? 0);91    push("operators", dim, a.operators);92  }93  await db.transaction(async (tx) => {94    for (let i = 0; i < rows.length; i += 500) {95      const chunk = rows.slice(i, i + 500);96      const values = sql.join(chunk.map((r) => sql`(${r.day}::date, ${r.metric}, ${r.dim}, ${r.value})`), sql`, `);97      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`);98    }99  });100  try {101    await chInsert("entity_daily", rows.map((r) => ({ day: r.day, metric: r.metric, dim: r.dim, value: r.value })));102  } catch (e) {103    console.warn(`[metrics] clickhouse mirror skipped: ${(e as Error).message}`);104  }105  return { day, rows: rows.length, dims: new Set(rows.map((r) => r.dim)).size };106}107