/** Country statistics (population, GDP, electricity, renewable share) → countries columns + stats jsonb + provenance. */ import { sql } from "@dci/db"; import type { NormalizedCountryStat } from "@dci/core"; import { addRef, bump, knownCountries, type IngestContext, type Tx } from "./common.js"; import { writeProvenance } from "./provenance.js"; export async function ingestCountryStat(tx: Tx, ctx: IngestContext, c: NormalizedCountryStat): Promise { const iso2 = (c.iso2 ?? c.key).trim().toUpperCase(); const known = await knownCountries(tx, ctx); if (!known.has(iso2)) throw new Error(`country ${iso2}: unknown ISO2 (seed countries first)`); const year = c.year ?? null; const existing = (await tx.execute(sql`select population, gdp_usd, electricity_twh, renewable_share, stats_year, stats from countries where iso2 = ${iso2}`))[0]!; const existingYear = existing.stats_year == null ? null : Number(existing.stats_year); // newer (or same-year) statistics win; an older vintage only fills blanks const newer = year == null || existingYear == null || year >= existingYear; const pick = (incoming: number | null | undefined, current: unknown) => (incoming == null ? (current as number | null) : newer || current == null ? incoming : (current as number)); const population = pick(c.population, existing.population == null ? null : Number(existing.population)); const gdp = pick(c.gdpUsd, existing.gdp_usd); const twh = pick(c.electricityTwh, existing.electricity_twh); const ren = pick(c.renewableShare, existing.renewable_share); const stats = { ...((existing.stats as Record) ?? {}) }; const indicators = { ...((stats.indicators as Record) ?? {}) }; for (const [k, v] of Object.entries({ population: c.population, gdpUsd: c.gdpUsd, electricityTwh: c.electricityTwh, renewableShare: c.renewableShare })) { if (v == null) continue; indicators[k] = { value: v, year, sourceId: c.provenance.sourceId, url: c.provenance.url, observedAt: ctx.now }; } stats.indicators = indicators; const changed = population !== (existing.population == null ? null : Number(existing.population)) || gdp !== existing.gdp_usd || twh !== existing.electricity_twh || ren !== existing.renewable_share; await tx.execute(sql`update countries set population = ${population}, gdp_usd = ${gdp}, electricity_twh = ${twh}, renewable_share = ${ren}, stats_year = ${newer && year != null ? year : existingYear}, stats = ${JSON.stringify(stats)}::jsonb, updated_at = now() where iso2 = ${iso2}`); if (changed) ctx.stats.updated++; else ctx.stats.unchanged++; bump(ctx, "country"); addRef(ctx, "country", iso2); await writeProvenance( tx, ctx, "country", iso2, [ { field: "population", value: c.population ?? null }, { field: "gdpUsd", value: c.gdpUsd ?? null }, { field: "electricityTwh", value: c.electricityTwh ?? null }, { field: "renewableShare", value: c.renewableShare ?? null }, { field: "statsYear", value: year }, ].map((f) => ({ ...f, provenance: c.provenance })), c.key, ); }