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;
}