File size: 9,080 Bytes
ee888e1 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 | // @ts-check
// Source freshness registry. Mirrors the table in
// docs/internal/pro-regional-intelligence-appendix-scoring.md "Source Freshness Registry".
//
// Each entry maps a Redis key (or key prefix) to its expected max-age.
// The snapshot writer marks inputs as stale or missing based on this table
// and feeds those flags into SnapshotMeta.snapshot_confidence.
import { CII_RISK_SCORE_CACHE_KEYS } from '../_cii-risk-cache-keys.mjs';
/**
* @typedef {object} SourceFreshnessSpec
* @property {string} key - Redis key (literal, no template variables)
* @property {number} maxAgeMin - Maximum acceptable age in minutes
* @property {string[]} feedsAxes - Which balance axes / sections this input drives
* @property {string=} metaKey - Optional companion seed-meta key carrying
* {fetchedAt, recordCount}. Used when the primary payload has no
* top-level timestamp field. classifyInputs() will prefer the meta key's
* timestamp when both are available so a stalled seeder can be detected
* even if the data payload is still served from a previous write.
*/
/**
* Only keys that compute modules actually consume via sources['...'].
* Keys must be added here in lockstep with new compute consumers, never
* speculatively. Drift between this list and the consumers is an alerting
* blind spot (a missing key drags down snapshot_confidence and a present
* key with no consumer wastes a Redis read).
*
* @type {SourceFreshnessSpec[]}
*/
export const FRESHNESS_REGISTRY = [
{ key: CII_RISK_SCORE_CACHE_KEYS.stale, maxAgeMin: 30, feedsAxes: ['domestic_fragility', 'coercive_pressure'], metaKey: 'seed-meta:intelligence:risk-scores' },
{ key: 'forecast:predictions:v2', maxAgeMin: 180, feedsAxes: ['scenarios', 'actors'] },
{ key: 'supply_chain:chokepoints:v4', maxAgeMin: 30, feedsAxes: ['maritime_access', 'corridors'] },
{ key: 'supply_chain:transit-summaries:v1', maxAgeMin: 30, feedsAxes: ['maritime_access'], metaKey: 'seed-meta:supply_chain:transit-summaries' },
{ key: 'intelligence:cross-source-signals:v1', maxAgeMin: 45, feedsAxes: ['coercive_pressure', 'evidence'], metaKey: 'seed-meta:intelligence:cross-source-signals' },
{ key: 'relay:oref:history:v1', maxAgeMin: 15, feedsAxes: ['coercive_pressure', 'triggers'], metaKey: 'seed-meta:relay:oref:history' },
{ key: 'economic:macro-signals:v1', maxAgeMin: 60, feedsAxes: ['capital_stress'] },
{ key: 'economic:national-debt:v1', maxAgeMin: 86400, feedsAxes: ['capital_stress'], metaKey: 'seed-meta:economic:national-debt' }, // monthly seed (30d cron), 60d window absorbs one missed run — mirrors api/health.js nationalDebt. metaKey is the primary freshness source (payload's seededAt is also recognized by extractTimestamp as a fallback).
{ key: 'economic:stress-index:v1', maxAgeMin: 120, feedsAxes: ['capital_stress'] },
{ key: 'energy:mix:v1:_all', maxAgeMin: 50400, feedsAxes: ['energy_vulnerability'], metaKey: 'seed-meta:economic:owid-energy-mix' },
{ key: 'economic:eu-gas-storage:v1', maxAgeMin: 2880, feedsAxes: ['energy_vulnerability'], metaKey: 'seed-meta:economic:eu-gas-storage' }, // runSeed writes seed-meta with numeric fetchedAt; metaKey takes priority over payload fields, defense-in-depth against a future refactor that strips fetchedAt/seededAt from the payload (regression risk that produced #3728).
{ key: 'economic:spr:v1', maxAgeMin: 10080, feedsAxes: ['energy_buffer'] },
// Mobility v1 (Phase 2 PR2) — feed the MobilityState block via mobility.mjs.
// maxAgeMin matches each seeder's cron interval + safety buffer.
// The aviation/gpsjam payloads have no top-level timestamp field, so they
// rely on companion seed-meta:* keys (written by the seeders via
// writeFreshnessMetadata / upstashSet) for stale detection. Without these
// metaKey hints, classifyInputs would fall back to "undated = fresh" and
// miss stalled seeders entirely.
{ key: 'aviation:delays:faa:v1', maxAgeMin: 60, feedsAxes: ['mobility'], metaKey: 'seed-meta:aviation:faa' },
{ key: 'aviation:delays:intl:v3', maxAgeMin: 90, feedsAxes: ['mobility'], metaKey: 'seed-meta:aviation:intl' },
{ key: 'aviation:notam:closures:v2', maxAgeMin: 120, feedsAxes: ['mobility'], metaKey: 'seed-meta:aviation:notam' },
{ key: 'intelligence:gpsjam:v2', maxAgeMin: 1440, feedsAxes: ['mobility', 'airspace'], metaKey: 'seed-meta:intelligence:gpsjam' }, // gpsjam.org is a DAILY source (restored from Wingbits, PR #4987); 1440min matches api/health.js gpsjam.maxStaleMin so snapshots don't mark it stale hours after a healthy daily seed.
// military:flights:v1 already carries top-level fetchedAt, no metaKey needed.
{ key: 'military:flights:v1', maxAgeMin: 30, feedsAxes: ['mobility', 'reroute_intensity'] },
];
/** Every metaKey referenced by FRESHNESS_REGISTRY, for pre-fetching. */
export const ALL_META_KEYS = FRESHNESS_REGISTRY
.map((s) => s.metaKey)
.filter((k) => typeof k === 'string' && k.length > 0);
export const ALL_INPUT_KEYS = FRESHNESS_REGISTRY.map((s) => s.key);
/**
* Classify each input as fresh, stale, or missing.
*
* Timestamp resolution order per input:
* 1. If the spec has a `metaKey`, use metaPayloads[metaKey].fetchedAt.
* This is the canonical signal for sources whose data payload lacks
* a top-level timestamp (FAA alerts, AviationStack, NOTAM, GPS jam).
* 2. Otherwise, pull a timestamp from the primary payload via
* extractTimestamp (fetchedAt, generatedAt, timestamp, updatedAt,
* lastUpdate, seededAt).
* 3. If neither yields a timestamp, classify as stale. A present-but-undated
* payload cannot be proven fresh — defaulting to fresh let stalled
* seeders silently inflate snapshot_confidence (#3728). Forcing stale
* pressures upstream to emit a timestamp.
*
* @param {Record<string, unknown>} payloads - Map of key -> raw value (or null)
* @param {Record<string, unknown>} [metaPayloads] - Map of metaKey -> raw value (or null)
* @returns {{ fresh: string[]; stale: string[]; missing: string[] }}
*/
export function classifyInputs(payloads, metaPayloads = {}) {
const fresh = [];
const stale = [];
const missing = [];
const now = Date.now();
for (const spec of FRESHNESS_REGISTRY) {
const payload = payloads[spec.key];
if (payload === null || payload === undefined) {
missing.push(spec.key);
continue;
}
// Prefer the companion seed-meta:*.fetchedAt when the spec declares one.
// This is the only way to detect a stalled seeder for payloads that
// don't carry a top-level timestamp of their own.
let ts = null;
if (spec.metaKey) {
const meta = metaPayloads[spec.metaKey];
ts = extractTimestamp(meta);
}
if (ts === null) ts = extractTimestamp(payload);
if (ts === null) {
// Present but undated — classify as stale. We cannot prove freshness,
// and the previous "default to fresh" behavior fabricated
// snapshot_confidence whenever a seeder stalled with no timestamp.
stale.push(spec.key);
continue;
}
const ageMin = (now - ts) / 60_000;
if (ageMin > spec.maxAgeMin) {
stale.push(spec.key);
} else {
fresh.push(spec.key);
}
}
return { fresh, stale, missing };
}
/**
* Resolve the effective timestamp for an input, preferring the metaKey
* (when declared) over the payload's own timestamp. Returns null if neither
* source carries a parseable timestamp.
*
* Exported for snapshot-meta.mjs, which needs the per-input timestamp to
* derive valid_until from the registry's maxAgeMin.
*
* @param {SourceFreshnessSpec} spec
* @param {unknown} payload
* @param {Record<string, unknown>} metaPayloads
* @returns {number | null}
*/
export function resolveInputTimestamp(spec, payload, metaPayloads) {
let ts = null;
if (spec.metaKey) {
const meta = metaPayloads[spec.metaKey];
ts = extractTimestamp(meta);
}
if (ts === null) ts = extractTimestamp(payload);
return ts;
}
/** Pull a timestamp out of common payload shapes; null if none found. */
function extractTimestamp(payload) {
if (typeof payload !== 'object' || payload === null) return null;
const obj = payload;
// `seededAt` is the convention used by runSeed-based seeders that wrap
// data in a { ...data, seededAt: ISOString } shape (seed-national-debt,
// seed-iea-oil-stocks, seed-eurostat-country-data, etc.). Without it here,
// those seeds got classified "present but undated" → always fresh,
// silently masking stalled crons.
for (const field of ['fetchedAt', 'generatedAt', 'timestamp', 'updatedAt', 'lastUpdate', 'seededAt']) {
if (typeof obj[field] === 'number') return obj[field];
if (typeof obj[field] === 'string') {
const parsed = Date.parse(obj[field]);
if (Number.isFinite(parsed)) return parsed;
}
}
return null;
}
|