'use strict'; /** * snapshotSettlementService — settle `model_snapshots` ON A SCHEDULE. * * Spec: specs/read-integrity-harness.md §9 * * ── WHY THIS EXISTS ────────────────────────────────────────────────────── * `model_snapshots` is the retention table built for replay — it holds the * model's INPUTS (the frozen feature vector) alongside its prediction, which * `ledger_entries` does not. It was settled exactly once, by hand, via * `scripts/settle-model-snapshots.js`. Nothing ever settled it on a cron. * * Measured 2026-08-09: 104,834 rows were past-dated and settleable; 89,350 of * them had no outcome. Worse, of the 34,128 rows carrying the repaired * champion's marker — the only rows a forward re-audit may be measured on — * ZERO were settled. The accrual clock could never advance. * * ── OUTCOMES ONLY. NO CONTEXT RECONSTRUCTION. ──────────────────────────── * This joins each row to the realized stat from a box score, using the row's OWN * keys (game_date + player_key + stat + line + side). It does not fetch, infer, * or rebuild park, weather, defence, platoon, archetype or any other context. * A settle pass that reconstructed context would be re-deciding what the model * saw, which is precisely what the retention table exists to prevent. * * ── OUTCOME IS SIDE-ALIGNED ────────────────────────────────────────────── * `p_win` is expressed for the GRADED SIDE, so `outcome` must be too. A raw * `realized > line` indicator is the OVER perspective and would silently invert * the target on every under row, making calibration measure the wrong thing. * `actual_value` stores the raw realized stat; `outcome` stores whether the * graded side won. * * ── A ROW LOGGED AFTER FIRST PITCH IS NOT A PREDICTION ─────────────────── * The pipeline runs on UTC cron hours, so a 01:00-UTC cycle is 21:00 the * previous evening ET — same game date, three hours into the slate. Those rows * are refused rather than settled. Preserved verbatim from the hand-run script, * because dropping it would quietly admit post-hoc rows into every measurement. * * Every dependency is injectable, so the unit tests never touch a network. */ const { paginate } = require('../utils/safePaginate'); const { uniqueKeyFor } = require('../utils/tableKeys'); const { nameKey } = require('../utils/playerName'); const { knownNumber } = require('../utils/known'); const STATS = ['hits', 'total_bases', 'rbi', 'runs']; const PAGE = 1000; /** Eastern first pitch, conservatively. At or after this is in-game. */ const FIRST_PITCH_ET_HOUR = 19; /** * Dates fetched per run — bounded, because a cron that re-fetches 26 dates of * box scores every five hours is a quota problem pretending to be thoroughness. * * NEWEST FIRST, and that ordering is a correctness property, not a preference. * The first implementation drained oldest-first and DEADLOCKED: measured on * production, 308 rows across the eight oldest dates are structurally * unsettleable (184 logged after first pitch, 124 with no box-score line), so * the window re-processed the same dead dates on every run and `dates_remaining` * never moved. Newest-first also reaches the rows that matter — the repaired * champion's, which are the only ones a forward re-audit may be measured on. */ const MAX_DATES_PER_RUN = Number(process.env.SNAPSHOT_SETTLE_MAX_DATES || 8); /** Realized value per stat, from the box-score batting line. */ const FIELD = Object.freeze({ hits: (b) => knownNumber(b.hits), total_bases: (b) => knownNumber(b.totalBases), rbi: (b) => knownNumber(b.rbi), runs: (b) => knownNumber(b.runs), }); function isPreGame(capturedAt, gameDate) { if (!capturedAt || !gameDate) return false; const cap = new Date(capturedAt); if (Number.isNaN(cap.getTime())) return false; const et = new Date(cap.getTime() - 4 * 3600 * 1000); // EDT const etDate = et.toISOString().slice(0, 10); if (etDate < String(gameDate)) return true; // day before, fine if (etDate > String(gameDate)) return false; // day after, post-game return et.getUTCHours() < FIRST_PITCH_ET_HOUR; } function todayEt(now) { return new Intl.DateTimeFormat('en-CA', { timeZone: 'America/New_York', year: 'numeric', month: '2-digit', day: '2-digit', }).format(now || new Date()); } async function pool(items, fn, n = 6) { const out = []; let i = 0; await Promise.all(Array.from({ length: n }, async () => { while (i < items.length) { const idx = i; i += 1; try { out[idx] = await fn(items[idx]); } catch { out[idx] = null; } } })); return out.filter(Boolean); } /** * Box-score batting lines for a set of dates, keyed `${date}|${nameKey}`. * A doubleheader gives two lines; they are SUMMED, because the prop covers the * day rather than a game. */ async function battingLines(dates, deps) { const getJson = deps.getJson; const games = []; for (const d of dates) { try { const s = await getJson(`https://statsapi.mlb.com/api/v1/schedule?sportId=1&date=${d}`); for (const day of (s && s.dates) || []) { for (const g of day.games || []) { if (String(g.status && g.status.detailedState) === 'Final') { games.push({ pk: g.gamePk, date: g.officialDate || d }); } } } } catch { /* absent day — stays absent */ } } const lines = {}; const loaded = await pool(games, async (g) => { const box = await getJson(`https://statsapi.mlb.com/api/v1/game/${g.pk}/boxscore`); const out = []; for (const side of ['home', 'away']) { const t = box && box.teams && box.teams[side]; if (!t) continue; for (const id of t.batters || []) { const pl = t.players[`ID${id}`]; const b = pl && pl.stats && pl.stats.batting; if (!b || b.atBats == null) continue; // did not bat → absent, never zero out.push({ date: g.date, key: nameKey(pl.person && pl.person.fullName), hits: b.hits, totalBases: b.totalBases, rbi: b.rbi, runs: b.runs, }); } } return out; }, deps.concurrency || 6); for (const arr of loaded) { for (const r of arr) { const k = `${r.date}|${r.key}`; if (!lines[k]) lines[k] = { ...r }; else { lines[k].hits += r.hits; lines[k].totalBases += r.totalBases; lines[k].rbi += r.rbi; lines[k].runs += r.runs; } } } return lines; } /** * PURE — decide each row's outcome from the box-score index. * Separated so the whole decision rule is unit-testable with no I/O. */ function decide(snaps, lines) { const counts = { candidates: snaps.length, settled: 0, unresolvable: 0, orphaned: 0, post_hoc_logged: 0 }; const updates = []; const seen = new Set(); for (const s of snaps) { if (seen.has(s.id)) { const e = new Error(`INTEGRITY: duplicate snapshot id ${s.id}`); e.code = 'DUPLICATE_ROW'; throw e; } seen.add(s.id); if (!isPreGame(s.captured_at, s.game_date)) { counts.post_hoc_logged += 1; counts.unresolvable += 1; continue; } const line = knownNumber(s.line); if (line === null || !s.side) { counts.unresolvable += 1; continue; } const b = lines[`${s.game_date}|${s.player_key}`]; if (!b) { counts.orphaned += 1; continue; } const realized = FIELD[s.stat] ? FIELD[s.stat](b) : null; if (realized === null) { counts.unresolvable += 1; continue; } const over = realized > line; const won = String(s.side).toLowerCase() === 'under' ? !over : over; updates.push({ id: s.id, outcome: won ? 'hit' : 'miss', actual_value: realized }); counts.settled += 1; } // CONSERVATION — hard fail. Every candidate lands in exactly one bucket. const acc = counts.settled + counts.unresolvable + counts.orphaned; if (acc !== counts.candidates) { const e = new Error(`INTEGRITY: conservation violated ${acc} != ${counts.candidates}`); e.code = 'CONSERVATION'; throw e; } return { counts, updates }; } /** * Settle one pass. * * @param {object} deps * - sb supabase service client (required to do anything) * - getJson async (url) => json * - now () => Date * - write default true; false = dry run * - maxDates dates fetched this run (oldest first) * @returns {object} { skipped?, counts, dates, written, repaired_champion_settled } */ async function settleSnapshots(deps = {}) { const sb = deps.sb; if (!sb) return { skipped: 'supabase not configured', counts: null, written: 0 }; const now = (deps.now || (() => new Date()))(); const cutoff = todayEt(now); const write = deps.write !== false; // Unsettled, PAST-DATED rows only — a game that has not finished cannot be // settled, and asking would produce an orphan rather than an absence. const snaps = await paginate( () => sb.from('model_snapshots') .select('id, game_date, captured_at, stat, player_key, line, side, model_version') .eq('sport', 'mlb').in('stat', STATS).is('outcome', null) .lt('game_date', cutoff), { key: uniqueKeyFor('model_snapshots'), pageSize: PAGE, label: 'snapshotSettlement' }, ); if (!snaps.length) return { counts: { candidates: 0, settled: 0, unresolvable: 0, orphaned: 0, post_hoc_logged: 0 }, dates: [], written: 0, repaired_champion_settled: 0 }; // CHOOSE DATES FROM ROWS THAT COULD ACTUALLY SETTLE. // // `isPreGame` is pure, so a row logged after first pitch is known-unsettleable // WITHOUT fetching anything. Measured on production, 19,074 such rows sit in // the recent dates; letting them pick the window meant the same dead dates // were re-fetched on every run and the backlog never converged. Filtering // first is what makes the drain terminate. const settleable = snaps.filter((r) => isPreGame(r.captured_at, r.game_date)); // NEWEST FIRST: the eligible rows — the repaired champion's — are the newest, // and an older date whose remainder cannot settle must not block them. const allDates = [...new Set(settleable.map((r) => r.game_date))].sort().reverse(); const perWindow = deps.maxDates || MAX_DATES_PER_RUN; const maxWindows = deps.maxWindows || Number(process.env.SNAPSHOT_SETTLE_MAX_WINDOWS || 4); // ADVANCE PAST A DRAINED WINDOW. The newest dates keep rows that can never // settle (a player with no box-score line never gets one), so a fixed window // would sit on them forever while older settleable dates were never reached. // Move to the next window when this one produces nothing, bounded so a run // still costs a predictable number of box-score fetches. let dates = []; let counts = null; let updates = []; let windows = 0; for (let w = 0; w < maxWindows; w += 1) { const slice = allDates.slice(w * perWindow, (w + 1) * perWindow); if (!slice.length) break; windows = w + 1; /* eslint-disable no-await-in-loop */ const lines = await battingLines(slice, { getJson: deps.getJson, concurrency: deps.concurrency }); const res = decide(snaps.filter((r) => slice.includes(r.game_date)), lines); /* eslint-enable no-await-in-loop */ dates = slice; counts = res.counts; updates = res.updates; if (updates.length) break; // progress — stop here } if (!counts) return { counts: { candidates: 0, settled: 0, unresolvable: 0, orphaned: 0, post_hoc_logged: 0 }, dates: [], written: 0, repaired_champion_settled: 0 }; const inWindow = snaps.filter((r) => dates.includes(r.game_date)); let written = 0; let repairedSettled = 0; const byId = new Map(inWindow.map((r) => [r.id, r])); if (write && updates.length) { const settledAt = now.toISOString(); for (let i = 0; i < updates.length; i += 500) { const batch = updates.slice(i, i + 500); /* eslint-disable no-await-in-loop */ const results = await Promise.all(batch.map((u) => sb.from('model_snapshots') .update({ outcome: u.outcome, actual_value: u.actual_value, settled_at: settledAt, settlement_source: 'statsapi_boxscore', }) // IDEMPOTENT: only an unsettled row is written, so a re-run can never // overwrite an outcome that is already on the record. .eq('id', u.id).is('outcome', null))); /* eslint-enable no-await-in-loop */ results.forEach((r, k) => { if (r.error) return; written += 1; const row = byId.get(batch[k].id); if (row && row.model_version === require('../config/modelVersion').REPAIRED_CHAMPION_VERSION) { repairedSettled += 1; } }); } } return { counts, dates, dates_remaining: Math.max(0, allDates.length - (windows * perWindow)), windows_scanned: windows, // Rows that can never settle, counted rather than hidden: a row logged after // first pitch is not a prediction and will never become one. permanently_unsettleable: snaps.length - settleable.length, written, repaired_champion_settled: repairedSettled, mode: write ? 'write' : 'dry-run', }; } module.exports = { settleSnapshots, decide, isPreGame, battingLines, STATS, FIELD, MAX_DATES_PER_RUN, FIRST_PITCH_ET_HOUR, };