The monitoring contract gets a callsite
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, behind a flag that is off. The previous tranche declared that future settled outcomes are forward evaluation evidence and shipped no evaluator — the contract lived in a comment. That is the same shape as statModel.js and correlateValidator.js, cited for months and never present. `forwardMonitor` scores the FROZEN artifact on outcomes settled strictly after its training_cutoff, and `snapshotScheduler.calibrationMonitorTick` runs it on the existing per-minute cadence, throttled 6h, every failure swallowed. A test asserts the tick is defined, invoked AND exported, and drives it with a fake to prove it passes the promoted artifact and a forward-only window. It cannot change what it watches: the module imports no fitter and no registry, and a test greps the stripped source for fitIsotonic, fitPlatt, writeFileSync, upsert, update, promote( , artifactRegistry and PROMOTED. LOW N IS ITS OWN ANSWER. Below the floor it reports INSUFFICIENT_SAMPLE with `healthy: null` — never false, never true. The floor is DERIVED, not chosen: resolving an error of the certified tolerance at two standard errors needs n >= 0.25/(0.05/2)^2 = 400, and a test recomputes it from TOLERANCE so the two cannot drift apart. TOLERANCE is 0.05, the same number certifyBands used, so the monitor can be neither stricter nor laxer than the thing it watches. Wrong-era forward rows are INVALID, not scored — scoring them would measure a different forecaster, which is the defect this line of work removed. An unreadable read is AUDIT_UNAVAILABLE, which is not a health verdict. A drift alert says in its own text that the artifact is frozen and unchanged, because the alert is not a demotion. `currentEraSource.loadRows` gains an optional `after` bound, strictly greater so the cutoff date itself can never be scored as forward evidence. SHADOW REMAINS BLOCKED: the production variable is still absent (configuration_source "default") after a restart at 05:09:02Z, so no shadow cohort was obtained and no replay was substituted for one. Artifact unchanged: mlb-hits-isotonic@2026-09-03, knots 5ae940ea163b7da2. Live OFF. Suite 405/405, 5,654 passed. Teeth 31/31 + 10/10 + 23/23. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01CQJeAG8vcDoL5zkiaJyVb8
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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 };
|
||||
@@ -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 };
|
||||
|
||||
Reference in New Issue
Block a user