diff --git a/scripts/teeth-artifact-governance.js b/scripts/teeth-artifact-governance.js index 3155f5c..2bad416 100644 --- a/scripts/teeth-artifact-governance.js +++ b/scripts/teeth-artifact-governance.js @@ -209,6 +209,46 @@ logic(26, 'PerformanceDistribution becomes servable', () => { return { caught: !/chain\.chainAcross\(/.test(snap), detail: 'no chainAcross call' }; }); +// ── FORWARD MONITOR (this tranche) ─────────────────────────────────────── +logic(27, 'forward monitor claimed active with no production caller', () => { + const sched = codeOf(fs.readFileSync(path.join(ROOT, 'src/snapshotScheduler.js'), 'utf8')); + const defined = sched.includes('const calibrationMonitorTick = async () =>'); + const invoked = sched.includes('await calibrationMonitorTick();'); + const exported = sched.includes('calibrationMonitorTick };'); + // the previous tranche had the CONTRACT and none of these three + return { caught: defined && invoked && exported, + detail: `defined=${defined} invoked_on_tick=${invoked} exported=${exported}` }; +}); + +inject(28, 'forward monitor refits the active artifact', + 'src/services/model/forwardMonitor.js', + `const { knownNumber } = require('../../utils/known');`, + `const { knownNumber } = require('../../utils/known'); +const _cal = require('./calibration'); +const _refit = () => _cal.fitIsotonic([], {});`, + 'tests/unit/forwardMonitor.test.js'); + +inject(29, 'low forward N reported HEALTHY', + 'src/services/model/forwardMonitor.js', + ` if (n < MIN_FORWARD_ROWS) {`, + ` if (false) {`, + 'tests/unit/forwardMonitor.test.js'); + +inject(30, 'the monitor scores evidence the fit already saw', + 'src/services/model/forwardMonitor.js', + ` if (artifact.training_cutoff && String(r.date) <= String(artifact.training_cutoff)) continue;`, + ` if (false) continue;`, + 'tests/unit/forwardMonitor.test.js'); + +logic(31, 'live serving turns on before all gates pass', () => { + const snap = codeOf(fs.readFileSync(path.join(ROOT, 'src/services/snapshotService.js'), 'utf8')); + const deployedEmpty = /CALIBRATION_DEPLOYED\s*=\s*Object\.freeze\(\[\s*\]\)/.test(snap); + const stage = registry.PROMOTED['mlb:hits'].stage === registry.STAGE.APPROVED_FOR_SHADOW; + const notLive = A.approved_for_live === false; + return { caught: deployedEmpty && stage && notLive, + detail: `CALIBRATION_DEPLOYED empty=${deployedEmpty}; stage=${registry.PROMOTED['mlb:hits'].stage}; approved_for_live=${A.approved_for_live}` }; +}); + const landed = results.filter((r) => r.landed).length; console.log(JSON.stringify({ teeth_landed: `${landed}/${results.length}`, results }, null, 2)); process.exit(landed === results.length ? 0 : 1); diff --git a/src/services/model/currentEraSource.js b/src/services/model/currentEraSource.js index 58d9f55..9869962 100644 --- a/src/services/model/currentEraSource.js +++ b/src/services/model/currentEraSource.js @@ -37,7 +37,7 @@ const CANONICAL_FIELDS = Object.freeze(['id', 'p_win', 'outcome', 'game_date', ' * `before` is STRICT — a Read may never train on its own outcome. * The quarantine exclusion mirrors the legacy path so the two remain comparable. */ -async function loadRows(sb, { sport, stat, modelVersion, before } = {}) { +async function loadRows(sb, { sport, stat, modelVersion, before, after = null } = {}) { if (!sb) throw new Error('currentEraSource.loadRows: no client'); if (!sport || !stat || !modelVersion || !before) { throw new Error('currentEraSource.loadRows: sport, stat, modelVersion and before are all required'); @@ -48,7 +48,11 @@ async function loadRows(sb, { sport, stat, modelVersion, before } = {}) { .eq('sport', sport).is('user_id', null).eq('stat', stat) .eq('model_version', modelVersion) // THE RESTRICTION .in('outcome', ['hit', 'miss']).not('p_win', 'is', null) - .lt('game_date', before), // strictly before + .lt('game_date', before) // strictly before + // FORWARD WINDOW (optional): observations the artifact was NOT fitted on. + // Strictly after, so the training cutoff date itself can never be scored + // as forward evidence. + .gt('game_date', after || '0001-01-01'), { key: 'id', pageSize: 1000, label: `currentEraSource(${sport}/${stat}/${modelVersion})` }, ); return rows diff --git a/src/services/model/forwardMonitor.js b/src/services/model/forwardMonitor.js new file mode 100644 index 0000000..313c533 --- /dev/null +++ b/src/services/model/forwardMonitor.js @@ -0,0 +1,203 @@ +'use strict'; + +/** + * forwardMonitor — EVALUATE THE FROZEN ARTIFACT. NEVER CHANGE IT. + * + * ── WHY THIS EXISTS ────────────────────────────────────────────────────── + * The previous tranche declared that "future settled outcomes are forward + * evaluation evidence" and then shipped no evaluator. Traced 2026-09-03: the + * only production reference to the artifact machinery outside its own directory + * is the shadow builder, and `calibrationRegistry.reverify` has ZERO production + * callers. The monitoring contract existed in a comment. + * + * That is the same shape as `statModel.js` / `correlateValidator.js`, which were + * cited for months and never existed. A contract with no caller is a plan. + * + * ── WHAT IT MAY AND MAY NOT DO ─────────────────────────────────────────── + * It scores the frozen curve against outcomes that settled STRICTLY AFTER the + * artifact's `training_cutoff` — evidence the fit never saw. It does not fit, + * refit, promote, or demote: this module imports no fitter, and a test asserts + * that. A future candidate is a separate, deliberate act. + */ + +const { knownNumber } = require('../../utils/known'); + +/** + * HEALTH. Mirrors `lineageCoverage`'s vocabulary deliberately — one operational + * language for "measured / not measured / bad", rather than a second dialect. + */ +const HEALTH = Object.freeze({ + HEALTHY: 'HEALTHY', + INSUFFICIENT_SAMPLE: 'INSUFFICIENT_SAMPLE', + DRIFT_WARNING: 'DRIFT_WARNING', + INVALID: 'INVALID', + AUDIT_UNAVAILABLE: 'AUDIT_UNAVAILABLE', +}); + +/** + * The tolerance the artifact was certified under — `fitPolicy`'s band tolerance + * and `calibration.certifyBands`' MAX_BIN_ERROR are both 0.05. The monitor uses + * the same number so it cannot be stricter or laxer than the thing it watches. + */ +const TOLERANCE = 0.05; + +/** + * MINIMUM USEFUL FORWARD N — derived, not chosen. + * + * To decide whether the observed calibration error exceeds TOLERANCE, the + * estimate needs a standard error comfortably inside it. At the worst-case + * binomial variance (p = 0.5) and a two-standard-error criterion: + * + * SE = sqrt(0.25 / n) <= TOLERANCE / 2 => n >= 0.25 / (0.025^2) = 400 + * + * Below this the monitor cannot distinguish drift from noise, so it says + * INSUFFICIENT_SAMPLE. It never says HEALTHY on evidence it could not read. + */ +const MIN_FORWARD_ROWS = 400; +const MIN_BAND_ROWS = 100; + +const BANDS = Object.freeze([[0.50, 0.60], [0.60, 0.70], [0.70, 0.80]]); +const r5 = (v) => (v == null || !Number.isFinite(v) ? null : Math.round(v * 100000) / 100000); +const r3 = (v) => (v == null || !Number.isFinite(v) ? null : Math.round(v * 1000) / 1000); + +function wilson(k, n, z = 1.96) { + if (!n) return null; + const p = k / n, d = 1 + (z * z) / n; + const c = (p + (z * z) / (2 * n)) / d; + const h = (z * Math.sqrt((p * (1 - p)) / n + (z * z) / (4 * n * n))) / d; + return [Math.max(0, c - h), Math.min(1, c + h)]; +} + +/** + * Score the frozen artifact on forward outcomes. + * + * @param {object} artifact the promoted, frozen artifact + * @param {Array<{p, won, date, model_version}>} rows settled AFTER training_cutoff + * @param {function} applyCurve (artifact, p) -> served probability or null + */ +function evaluate(artifact, rows, applyCurve) { + if (!artifact) { + return { health: HEALTH.AUDIT_UNAVAILABLE, reason: 'no promoted artifact', n: 0, healthy: null }; + } + if (artifact.servable !== true) { + return { health: HEALTH.INVALID, reason: 'active artifact is not servable', + artifact_id: artifact.artifact_id, n: 0, healthy: false }; + } + if (!Array.isArray(rows)) { + return { health: HEALTH.AUDIT_UNAVAILABLE, reason: 'forward evidence could not be read', + artifact_id: artifact.artifact_id, n: 0, healthy: null }; + } + + // Wrong-era rows are not this artifact's evidence. Scoring them would measure + // a different forecaster, which is the defect this whole line of work removed. + const foreign = rows.filter((r) => r.model_version && r.model_version !== artifact.model_version).length; + if (foreign > 0) { + return { health: HEALTH.INVALID, reason: `forward evidence contains ${foreign} wrong-era rows`, + artifact_id: artifact.artifact_id, n: rows.length, healthy: false }; + } + + const scored = []; + for (const r of rows) { + const raw = knownNumber(r.p); + const won = knownNumber(r.won); + if (raw === null || won === null) continue; + // Strictly after the cutoff — evidence the fit never saw. + if (artifact.training_cutoff && String(r.date) <= String(artifact.training_cutoff)) continue; + const served = applyCurve(artifact, raw); + if (served === null) continue; // outside certified support + scored.push({ raw, served, won }); + } + + const n = scored.length; + const base = { + artifact_id: artifact.artifact_id, + model_version: artifact.model_version, + procedure_version: artifact.procedure_version, + knot_digest: artifact.knot_digest, + training_cutoff: artifact.training_cutoff, + n, + min_required: MIN_FORWARD_ROWS, + tolerance: TOLERANCE, + }; + + if (n < MIN_FORWARD_ROWS) { + // NOT healthy, NOT drift. "We cannot tell yet" is its own answer. + return { ...base, health: HEALTH.INSUFFICIENT_SAMPLE, healthy: null, + reason: `${n} forward rows in support; ${MIN_FORWARD_ROWS} needed to resolve a ${TOLERANCE} error` }; + } + + const E = 1e-12; + let brier = 0, ll = 0, sumServed = 0, hits = 0; + for (const s of scored) { + brier += (s.served - s.won) ** 2; + const p = Math.min(1 - E, Math.max(E, s.served)); + ll += -(s.won * Math.log(p) + (1 - s.won) * Math.log(1 - p)); + sumServed += s.served; hits += s.won; + } + brier /= n; ll /= n; + const observed = hits / n; + const predicted = sumServed / n; + const error = predicted - observed; + + const bands = BANDS.map(([lo, hi]) => { + const sel = scored.filter((s) => s.raw >= lo && s.raw < hi); + const k = sel.reduce((a, s) => a + s.won, 0); + const mean = sel.length ? sel.reduce((a, s) => a + s.served, 0) / sel.length : null; + const ci = wilson(k, sel.length); + return { band: `${lo.toFixed(2)}-${hi.toFixed(2)}`, n: sel.length, + mean_served: r3(mean), observed: sel.length ? r3(k / sel.length) : null, + observed_ci95: ci ? [r3(ci[0]), r3(ci[1])] : null, + error: sel.length ? r3(mean - k / sel.length) : null, + resolvable: sel.length >= MIN_BAND_ROWS }; + }); + + // ECE over the served value, same shape as the certification used. + const binsN = 10, acc = Array.from({ length: binsN }, () => ({ n: 0, sp: 0, sy: 0 })); + for (const s of scored) { + const b = Math.min(binsN - 1, Math.floor(s.served * binsN)); + acc[b].n++; acc[b].sp += s.served; acc[b].sy += s.won; + } + let ece = 0; + for (const b of acc) if (b.n) ece += (b.n / n) * Math.abs(b.sp / b.n - b.sy / b.n); + + const driftBands = bands.filter((b) => b.resolvable && b.error != null && Math.abs(b.error) > TOLERANCE); + const drift = Math.abs(error) > TOLERANCE || driftBands.length > 0; + + return { ...base, + health: drift ? HEALTH.DRIFT_WARNING : HEALTH.HEALTHY, + healthy: !drift, + brier: r5(brier), logloss: r5(ll), ece: r5(ece), + predicted: r3(predicted), observed: r3(observed), error: r3(error), + observed_ci95: (() => { const ci = wilson(hits, n); return ci ? [r3(ci[0]), r3(ci[1])] : null; })(), + bands, + drift_bands: driftBands.map((b) => b.band), + reason: drift + ? `calibration error ${r3(error)} / bands ${driftBands.map((b) => b.band).join(',') || 'none'} exceed tolerance ${TOLERANCE}` + : null, + }; +} + +/** Throttle, mirroring lineageCoverage.coverageDue. */ +function monitorDue(lastAtMs, nowMs, everyMs = 6 * 3600 * 1000) { + if (!Number.isFinite(nowMs)) return false; + if (!Number.isFinite(lastAtMs)) return true; + return nowMs - lastAtMs >= everyMs; +} + +/** One alert per state transition, not one per tick. */ +function monitorAlarm(prevKey, result) { + const key = `${result.health}:${result.artifact_id || 'none'}`; + if (key === prevKey) return { alert: false, key }; + if (result.health === HEALTH.HEALTHY || result.health === HEALTH.INSUFFICIENT_SAMPLE) { + return { alert: false, key }; + } + const messages = { + [HEALTH.DRIFT_WARNING]: `Calibration DRIFT on ${result.artifact_id}: ${result.reason}. The artifact is frozen and unchanged; a new candidate needs certifying.`, + [HEALTH.INVALID]: `Calibration artifact INVALID: ${result.reason}. Nothing certified is being served.`, + [HEALTH.AUDIT_UNAVAILABLE]: `Calibration forward monitor COULD NOT RUN (${result.reason}). This is not a health result.`, + }; + return { alert: true, key, priority: result.health === HEALTH.DRIFT_WARNING ? 'default' : 'high', + message: messages[result.health] || `Calibration monitor ${result.health}.` }; +} + +module.exports = { HEALTH, TOLERANCE, MIN_FORWARD_ROWS, MIN_BAND_ROWS, evaluate, monitorDue, monitorAlarm }; diff --git a/src/snapshotScheduler.js b/src/snapshotScheduler.js index 5ae132b..b172671 100644 --- a/src/snapshotScheduler.js +++ b/src/snapshotScheduler.js @@ -194,10 +194,60 @@ function startSnapshotScheduler(opts = {}) { } catch { /* the monitor must never break the scheduler */ } }; + // ── AUTOMATIC FORWARD CALIBRATION MONITOR ─────────────────────────────── + // + // The artifact-governance tranche declared that future settled outcomes are + // forward evaluation evidence and then shipped no evaluator: traced + // 2026-09-03, `calibrationRegistry.reverify` had ZERO production callers and + // the only reference to the artifact machinery outside its own directory was + // the shadow builder. A monitoring contract with no callsite is a plan. + // + // This is the callsite. It SCORES the frozen artifact on outcomes that + // settled strictly after its training_cutoff and never touches it — the + // module imports no fitter and a test asserts that. Every failure path is + // swallowed: a monitor that can take down the pipeline it watches is worse + // than no monitor. + const forwardMonitor = opts.forwardMonitor || require('./services/model/forwardMonitor'); + const artifactRegistry = opts.artifactRegistry || require('./services/model/artifactRegistry'); + const eraSource = opts.currentEraSource || require('./services/model/currentEraSource'); + const MONITOR_EVERY_MS = Number.parseInt(process.env.CALIBRATION_MONITOR_EVERY_MS || '', 10) + || 6 * 60 * 60 * 1000; + let lastMonitorAtMs = null; + let lastMonitorAlertKey = null; + const calibrationMonitorTick = async () => { + try { + const nowMs = now().getTime(); + if (!forwardMonitor.monitorDue(lastMonitorAtMs, nowMs, MONITOR_EVERY_MS)) return; + lastMonitorAtMs = nowMs; + const artifact = artifactRegistry.load('mlb', 'hits'); + if (!artifact) return; // nothing promoted: nothing to watch + const sb = opts.supabase || require('./utils/supabase').getSupabaseServiceClient(); + let rows = null; + if (sb) { + rows = await eraSource.loadRows(sb, { + sport: artifact.sport, stat: artifact.stat, modelVersion: artifact.model_version, + before: now().toISOString().slice(0, 10), + after: artifact.training_cutoff, // evidence the fit never saw + }); + } + const result = forwardMonitor.evaluate(artifact, rows, artifactRegistry.applyCurve); + console.log(`[calibration-monitor] ${artifact.artifact_id} ${result.health}` + + ` n=${result.n}/${result.min_required}` + + `${result.error != null ? ` err=${result.error}` : ''}` + + `${result.reason ? ` reason=${result.reason}` : ''}`); + const alarm = forwardMonitor.monitorAlarm(lastMonitorAlertKey, result); + lastMonitorAlertKey = alarm.key; + if (alarm.alert) { + await notify(alarm.message, { title: 'VYNDR calibration', priority: alarm.priority, tags: ['chart_with_downwards_trend'] }); + } + } catch { /* the monitor must never break the scheduler */ } + }; + const tick = async () => { await checkOverdue(); await pulseTick(); await coverageTick(); + await calibrationMonitorTick(); const d = now(); if (d.getUTCMinutes() !== 0) return; const h = d.getUTCHours(); @@ -504,7 +554,7 @@ function startSnapshotScheduler(opts = {}) { console.log(`[settlement] armed — ${settleTags} outcomes + ledger settle pass runs FIRST at each snapshot slot (${HOURS_UTC.join(',')} UTC), idempotent re-runs`); // Session 8 — same verifiability rule: every watchdog states itself at boot. console.log(`[opsWatch] armed — settle alarms (throw + morning zero-settle), failure pager (${failureTracker.threshold} consecutive), quota daily check, pulse ${PULSE_HOUR_UTC}:00 UTC`); - return { interval, tick, refreshTick, pulseTick, coverageTick }; + return { interval, tick, refreshTick, pulseTick, coverageTick, calibrationMonitorTick }; } module.exports = { startSnapshotScheduler, HOURS_UTC, mostRecentExpectedSlot, isSnapshotOverdue }; diff --git a/tests/unit/currentEraSource.test.js b/tests/unit/currentEraSource.test.js index e90f9d1..2669910 100644 --- a/tests/unit/currentEraSource.test.js +++ b/tests/unit/currentEraSource.test.js @@ -18,14 +18,16 @@ function client(rows) { in: (c, v) => { f[`in_${c}`] = v; return q; }, not: () => q, lt: (c, v) => { f[`lt_${c}`] = v; return q; }, + gt: (c, v) => { f[`gt_${c}`] = v; return q; }, order: () => q, range: async () => { let out = rows; for (const [k, v] of Object.entries(f)) { - if (k.startsWith('in_') || k.startsWith('lt_')) continue; + if (k.startsWith('in_') || k.startsWith('lt_') || k.startsWith('gt_')) continue; out = out.filter((r) => r[k] === v); } if (f.lt_game_date) out = out.filter((r) => r.game_date < f.lt_game_date); + if (f.gt_game_date) out = out.filter((r) => r.game_date > f.gt_game_date); return { data: out, error: null, count: out.length }; }, }; @@ -55,6 +57,20 @@ describe('the era restriction is in the QUERY', () => { expect(src.eraAudit([{ model_version: ERA }], ERA).wrong_era_rows).toBe(0); }); + it('the optional forward window is STRICTLY after — the cutoff date is never scored', async () => { + const { client: c, filters } = client(ROWS); + const out = await src.loadRows(c, { sport: 'mlb', stat: 'hits', modelVersion: ERA, + before: '2026-09-03', after: '2026-08-12' }); + expect(filters.gt_game_date).toBe('2026-08-12'); + expect(out.map((r) => r.id)).toEqual(['2']); // 08-12 itself excluded + }); + + it('omitting the forward window loads the whole eligible history', async () => { + const { client: c } = client(ROWS); + const out = await src.loadRows(c, { sport: 'mlb', stat: 'hits', modelVersion: ERA, before: '2026-09-03' }); + expect(out.map((r) => r.id).sort()).toEqual(['1', '2']); + }); + it('refuses to run without an explicit era and horizon', async () => { const { client: c } = client(ROWS); await expect(src.loadRows(c, { sport: 'mlb', stat: 'hits', before: '2026-09-03' })).rejects.toThrow(); diff --git a/tests/unit/forwardMonitor.test.js b/tests/unit/forwardMonitor.test.js new file mode 100644 index 0000000..de7081e --- /dev/null +++ b/tests/unit/forwardMonitor.test.js @@ -0,0 +1,160 @@ +'use strict'; + +/** + * A MONITORING CONTRACT WITH NO CALLSITE IS A PLAN. + * + * The previous tranche declared forward evaluation and shipped none. These + * tests hold both halves: the evaluator behaves, AND production actually calls it. + */ +const fm = require('../../src/services/model/forwardMonitor'); +const registry = require('../../src/services/model/artifactRegistry'); + +const A = registry.load('mlb', 'hits'); +const AFTER = '2026-09-02'; // strictly after training_cutoff 2026-09-01 + +/** Rows whose outcomes track the frozen curve, optionally shifted to force drift. */ +function rows(n, shift = 0, over = {}) { + return Array.from({ length: n }, (_, i) => { + const raw = Math.round((0.50 + (i % 29) / 100) * 1000) / 1000; + const served = registry.applyCurve(A, raw) ?? 0.6; + return { p: raw, won: ((i * 2654435761) % 1000) / 1000 < served + shift ? 1 : 0, + date: AFTER, model_version: A.model_version, ...over }; + }); +} + +describe('the monitor evaluates and refuses to guess', () => { + it('says HEALTHY only on enough evidence', () => { + const r = fm.evaluate(A, rows(1200), registry.applyCurve); + expect(r.health).toBe(fm.HEALTH.HEALTHY); + expect(r.healthy).toBe(true); + expect(r.n).toBeGreaterThanOrEqual(fm.MIN_FORWARD_ROWS); + }); + + it('LOW N IS NOT HEALTHY AND NOT DRIFT — it is its own answer', () => { + const r = fm.evaluate(A, rows(50), registry.applyCurve); + expect(r.health).toBe(fm.HEALTH.INSUFFICIENT_SAMPLE); + expect(r.healthy).toBeNull(); // never false, never true + expect(r.reason).toContain(String(fm.MIN_FORWARD_ROWS)); + }); + + it('the sample floor is DERIVED from the certified tolerance, not chosen', () => { + // SE = sqrt(0.25/n) <= TOLERANCE/2 => n >= 0.25 / (TOLERANCE/2)^2 + expect(fm.MIN_FORWARD_ROWS).toBe(Math.ceil(0.25 / ((fm.TOLERANCE / 2) ** 2))); + expect(fm.TOLERANCE).toBe(0.05); // same number certifyBands used + }); + + it('flags drift when the frozen curve stops matching outcomes', () => { + const r = fm.evaluate(A, rows(1200, 0.20), registry.applyCurve); + expect(r.health).toBe(fm.HEALTH.DRIFT_WARNING); + expect(r.healthy).toBe(false); + expect(r.drift_bands.length).toBeGreaterThan(0); + }); + + it('refuses an unservable artifact and wrong-era evidence', () => { + expect(fm.evaluate({ ...A, servable: false }, rows(1200), registry.applyCurve).health) + .toBe(fm.HEALTH.INVALID); + expect(fm.evaluate(A, rows(1200, 0, { model_version: 'engine1@2026-07-20' }), registry.applyCurve).health) + .toBe(fm.HEALTH.INVALID); + }); + + it('an unreadable read is AUDIT_UNAVAILABLE, never a health verdict', () => { + const r = fm.evaluate(A, null, registry.applyCurve); + expect(r.health).toBe(fm.HEALTH.AUDIT_UNAVAILABLE); + expect(r.healthy).toBeNull(); + expect(fm.evaluate(null, rows(1200), registry.applyCurve).health).toBe(fm.HEALTH.AUDIT_UNAVAILABLE); + }); + + it('scores ONLY evidence the fit never saw', () => { + const atCutoff = rows(1200).map((r) => ({ ...r, date: A.training_cutoff })); + expect(fm.evaluate(A, atCutoff, registry.applyCurve).n).toBe(0); + const before = rows(1200).map((r) => ({ ...r, date: '2026-08-01' })); + expect(fm.evaluate(A, before, registry.applyCurve).n).toBe(0); + }); + + it('scores ONLY rows inside certified support', () => { + const outside = rows(1200).map((r, i) => ({ ...r, p: i % 2 ? 0.9 : 0.4 })); + expect(fm.evaluate(A, outside, registry.applyCurve).n).toBe(0); + }); +}); + +describe('the monitor cannot change what it watches', () => { + it('imports no fitter and performs no write', () => { + const src = require('fs').readFileSync( + require('path').join(__dirname, '../../src/services/model/forwardMonitor.js'), 'utf8'); + const code = src.replace(/\/\*[\s\S]*?\*\//g, '').replace(/^\s*\/\/.*$/gm, ''); + // token-level, not substring: 'promote' also appears inside the reason + // string 'no promoted artifact', which is a description, not a call. + for (const forbidden of ['fitIsotonic', 'fitPlatt', 'writeFileSync', '.upsert(', '.update(', + 'promote(', 'artifactRegistry', 'PROMOTED']) { + expect(code).not.toContain(forbidden); + } + }); + + it('evaluating does not mutate the artifact', () => { + const before = JSON.stringify(A); + fm.evaluate(A, rows(1200, 0.2), registry.applyCurve); + expect(JSON.stringify(registry.load('mlb', 'hits'))).toBe(before); + }); +}); + +describe('alarm discipline', () => { + it('does not alert on HEALTHY or INSUFFICIENT_SAMPLE', () => { + for (const h of [fm.HEALTH.HEALTHY, fm.HEALTH.INSUFFICIENT_SAMPLE]) { + expect(fm.monitorAlarm(null, { health: h, artifact_id: 'x' }).alert).toBe(false); + } + }); + + it('alerts once per state transition, not once per tick', () => { + const first = fm.monitorAlarm(null, { health: fm.HEALTH.DRIFT_WARNING, artifact_id: 'x', reason: 'r' }); + expect(first.alert).toBe(true); + expect(fm.monitorAlarm(first.key, { health: fm.HEALTH.DRIFT_WARNING, artifact_id: 'x', reason: 'r' }).alert).toBe(false); + }); + + it('says plainly that a drift alert has NOT changed the artifact', () => { + const a = fm.monitorAlarm(null, { health: fm.HEALTH.DRIFT_WARNING, artifact_id: 'x', reason: 'r' }); + expect(a.message).toMatch(/frozen and unchanged/); + }); +}); + +describe('THE CALLSITE — production actually runs it', () => { + it('the scheduler exposes and invokes the monitor on its tick', async () => { + const src = require('fs').readFileSync( + require('path').join(__dirname, '../../src/snapshotScheduler.js'), 'utf8'); + expect(src).toContain('await calibrationMonitorTick();'); + expect(src).toContain('calibrationMonitorTick };'); + }); + + it('the tick calls evaluate with the promoted artifact and forward-only rows', async () => { + const sched = require('../../src/snapshotScheduler'); + const calls = []; + const fake = { + monitorDue: () => true, + evaluate: (artifact, rws) => { calls.push({ artifact, rws }); return { health: 'HEALTHY', healthy: true, n: 999, min_required: 400, artifact_id: artifact.artifact_id }; }, + monitorAlarm: () => ({ alert: false, key: 'k' }), + HEALTH: fm.HEALTH, + }; + const loaded = []; + const prev = process.env.SNAPSHOT_CRON; + process.env.SNAPSHOT_CRON = '1'; // the scheduler is inert unless armed + const s = sched.startSnapshotScheduler({ + forwardMonitor: fake, + artifactRegistry: registry, + currentEraSource: { loadRows: async (sb, args) => { loaded.push(args); return rows(500); } }, + supabase: {}, + now: () => new Date('2026-09-03T05:00:00Z'), + runAllSnapshots: async () => ({}), + }); + try { + expect(s).not.toBeNull(); // armed, or the assertions below are vacuous + await s.calibrationMonitorTick(); + expect(calls).toHaveLength(1); + expect(calls[0].artifact.artifact_id).toBe(A.artifact_id); + // forward-only: bounded strictly after the artifact's training cutoff + expect(loaded[0].after).toBe(A.training_cutoff); + expect(loaded[0].modelVersion).toBe(A.model_version); + } finally { + if (s && s.interval) clearInterval(s.interval); + if (prev === undefined) delete process.env.SNAPSHOT_CRON; else process.env.SNAPSHOT_CRON = prev; + } + }); +});