981b26010a
Shadow converged at 18:09:07Z and the 19:00Z slot produced cohort 0353c551.
It immediately falsified something no test had asked: 513 PUBLISHED NON-HITS
rows came back CERTIFIED_CALIBRATED with a served probability drawn from the
mlb-hits curve — total_bases 246, runs 149, rbi 125, walks 109, outs 28,
strikeouts 24, hits_allowed 20, earned_runs 17.
Two bugs, one on top of the other. `mergeProbabilityContract` passed only
{model_version, p_win}, dropping the row's identity; and the service's resolve
then stamped the {sport, stat} it had been BUILT with onto every read. So all
3,000 rows in the batch resolved as mlb hits.
The governance tests could not see it. They asked "does build() refuse another
stat?" — it does, and always did — and then exercised the merge with
hits-only rows. Production sends one mixed batch. The regression test now drives
the REAL collector with hits, total_bases, rbi, runs, walks, strikeouts and
home_runs at the same p_win and requires hits certified and every other stat
neither certified nor numeric.
Fixed in three layers, because one would have been the same single point that
just failed:
1. the service no longer substitutes its own identity — the row's decides,
and a read naming no stat resolves to no contract, which is UNSUPPORTED;
2. the merge carries the row's sport and stat;
3. probabilityContract refuses an artifact whose own sport/stat disagree with
the contract it is being used under, independent of plumbing.
NO USER IMPACT. Shadow only: every block carries servable:false, live serving is
OFF, CALIBRATION_DEPLOYED is [], and the anonymous payload showed zero
calibration fields before and after. But this is exactly the defect that would
have served a hits calibration curve for strikeouts on the day live was enabled,
and only a real cohort surfaced it.
Two teeth were themselves wrong. Both runners checked "retention identity
changes" by grepping the diff for `stat:`, which fired on `stat: r.stat` — a
line that READS identity to hand it to a reader, not one that changes what
identifies a row. A guard that cannot tell those apart blocks the fix for the
defect it exists to protect against. Both are now behavioural: build a row
through the real collector and compare the identity tuple.
Artifact unchanged: mlb-hits-isotonic@2026-09-03, knot 5ae940ea163b7da2.
Suite 405/405, 5,659 passed. Teeth 34/34 + 10/10 + 23/23.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01CQJeAG8vcDoL5zkiaJyVb8
1285 lines
57 KiB
JavaScript
1285 lines
57 KiB
JavaScript
'use strict';
|
|
|
|
/**
|
|
* MODEL SNAPSHOT RETENTION (Session 64, Phase 2 — priority zero).
|
|
*
|
|
* As of 2026-07-20 the only store of model history was `ledger_entries`: 640
|
|
* rows, 6 game days, and NO MODEL INPUTS. Every other warehouse table in the
|
|
* schema was empty. That meant we could score the grades we emitted but could
|
|
* not replay a different model against the same conditions — the only question
|
|
* a backtest exists to answer, and the gate the metrics-engine north star
|
|
* requires ("does this metric predict better than without it?").
|
|
*
|
|
* This writes `model_snapshots`: one APPEND-ONLY row per graded prop per
|
|
* snapshot cycle, carrying the market values, the model output, the refusal
|
|
* (if any), and the FEATURE VECTOR that produced it.
|
|
*
|
|
* CONTRACT: best-effort, exactly like the ledger write. A retention failure
|
|
* must NEVER break a snapshot. Every path returns a summary; none throw.
|
|
*
|
|
* Two deliberate choices:
|
|
* - REFUSALS ARE STORED. The ledger drops them, so a gate that refuses props
|
|
* which would have won is invisible. That is lost edge we cannot measure.
|
|
* - EVERY ROW IS STAMPED with model_version + code_sha. A backtest that mixes
|
|
* model eras is worthless, and `ledger_entries` is already permanently
|
|
* contaminated across the 2026-07-19 fix boundary with no way to separate.
|
|
*/
|
|
|
|
const crypto = require('crypto');
|
|
const { normalizeName, nameKey } = require('../utils/playerName');
|
|
const { semanticPlayerKey } = require('./model/participantIdentity');
|
|
|
|
/** ET calendar date of an ISO timestamp (shares gameBinder's rule). */
|
|
function etDateOf(iso) {
|
|
if (!iso) return null;
|
|
const t = new Date(iso);
|
|
if (Number.isNaN(t.getTime())) return null;
|
|
return new Intl.DateTimeFormat('en-CA', {
|
|
timeZone: 'America/New_York', year: 'numeric', month: '2-digit', day: '2-digit',
|
|
}).format(t);
|
|
}
|
|
|
|
/**
|
|
* Bump when the grading model changes in a way that makes rows non-comparable.
|
|
* This is the marker `ledger_entries` never had.
|
|
*/
|
|
/**
|
|
* CHAMPION VERSION — the eligibility marker for every forward re-audit.
|
|
*
|
|
* Bumped when the forecaster itself changes, so a settled row is
|
|
* self-identifying: rows tagged `engine1@2026-08-07-fullwindow` were produced by
|
|
* the REPAIRED champion (full season log, recency weight 0.20); anything earlier
|
|
* came from the retired ten-game forecaster.
|
|
*
|
|
* This is what makes the re-audit rule mechanical rather than a promise.
|
|
* Calibration may only be re-fit, and factor verdicts may only be re-audited, on
|
|
* rows carrying the current marker — never on reconstructions of a retired
|
|
* forecast, and never on a mixture of the two, which is the trap that would
|
|
* otherwise be invisible once both generations sit in the same table.
|
|
*/
|
|
// FIX A3 — ONE SOURCE (src/config/modelVersion.js). Re-exported so existing
|
|
// importers keep working, but never re-declared: the twin default in
|
|
// ledgerService is exactly how these two drifted apart.
|
|
const { MODEL_VERSION, REPAIRED_CHAMPION_VERSION } = require('../config/modelVersion');
|
|
|
|
function codeSha() {
|
|
return process.env.SOURCE_COMMIT || process.env.GIT_SHA || process.env.COOLIFY_GIT_COMMIT_SHA || null;
|
|
}
|
|
|
|
/** Strict numeric parse — `Number(null) === 0` is this codebase's classic
|
|
* fabrication bug, so absent must stay absent. */
|
|
function numOrNull(v) {
|
|
if (v == null || v === '') return null;
|
|
const n = Number(v);
|
|
return Number.isFinite(n) ? n : null;
|
|
}
|
|
|
|
function intOrNull(v) {
|
|
const n = numOrNull(typeof v === 'string' ? v.replace('+', '') : v);
|
|
return n == null ? null : Math.round(n);
|
|
}
|
|
|
|
function boolOrNull(v) {
|
|
return typeof v === 'boolean' ? v : null;
|
|
}
|
|
|
|
/**
|
|
* Build retention rows from ONE prop's graded sides (both over and under,
|
|
* graded or refused). `base` is the prop the grader was called with.
|
|
*/
|
|
function rowsFromSides(base, sides, ctx = {}) {
|
|
const rows = [];
|
|
// `published` is declared on every row and defaults FALSE. A capture is not a
|
|
// publication until the winner signal says so; defaulting true would assert
|
|
// that 65% of these rows were shown to a user.
|
|
|
|
const list = Array.isArray(sides) ? sides : [];
|
|
for (const s of list) {
|
|
if (!s) continue;
|
|
const player = s.player || base.player;
|
|
if (!player) continue;
|
|
// SEMANTIC PLAYER IDENTITY. When MLB has proven who this participant is,
|
|
// the identity comes from the league's own name for them, so every book's
|
|
// spelling converges. When it has not, the raw spelling is used exactly as
|
|
// before — we never guess a human into existence. The raw provider string
|
|
// is untouched here and survives verbatim in closing_captures.
|
|
const identity = semanticPlayerKey({
|
|
player,
|
|
mlb_person_id: s.mlb_person_id ?? base.mlb_person_id ?? null,
|
|
canonical_player_name: s.canonical_player_name ?? base.canonical_player_name ?? null,
|
|
});
|
|
const stat = s.stat_type || base.stat_type;
|
|
const line = numOrNull(s.line != null ? s.line : base.line);
|
|
const side = String(s.direction || '').toLowerCase();
|
|
if (!stat || line == null || (side !== 'over' && side !== 'under')) continue;
|
|
|
|
const refused = !!(s.insufficient_data || !s.grade);
|
|
|
|
rows.push({
|
|
snapshot_id: ctx.snapshotId,
|
|
captured_at: ctx.capturedAt,
|
|
cycle_hour_utc: ctx.cycleHourUtc ?? null,
|
|
model_version: MODEL_VERSION,
|
|
code_sha: codeSha(),
|
|
|
|
sport: ctx.sport,
|
|
game_id: ctx.gameIdFor ? ctx.gameIdFor(base, s) : (base.game_id || `${ctx.sport}:${ctx.gameDate}`),
|
|
// Session 64 (Order 1.5) — the GAME's date, from the bound game_time,
|
|
// never the snapshot clock. ctx.gameDate is only a last resort for props
|
|
// the binder could not tie to a real game.
|
|
game_date: etDateOf(base && base.game_time) || ctx.gameDate,
|
|
player_key: identity.key,
|
|
// RAW SOURCE NAME, deliberately. The identity above converges on the
|
|
// league's record; the NAME on the row stays the provider's, so the row
|
|
// still carries the representation it was published under. Identity and
|
|
// provenance are different jobs and this row does both.
|
|
player_name: normalizeName(player).display || player,
|
|
team: s.team || base.team || null,
|
|
opponent: s.opponent || base.opponent || null,
|
|
stat,
|
|
line,
|
|
side,
|
|
|
|
book: s.book || base.book || null,
|
|
book_odds: intOrNull(s.book_odds),
|
|
over_odds: intOrNull(base.over_odds),
|
|
under_odds: intOrNull(base.under_odds),
|
|
fair_odds: intOrNull(s.fair_odds),
|
|
fair_prob: numOrNull(s.fair_prob),
|
|
overround: numOrNull(s.overround),
|
|
devig_method: s.devig_method || null,
|
|
|
|
grade: s.grade || null,
|
|
grade_11: s._grade_11 || null,
|
|
confidence: numOrNull(s.confidence),
|
|
confidence_basis: s.confidence_basis || null,
|
|
p_win: numOrNull(s.p_win),
|
|
ev_pct: numOrNull(s.ev_pct),
|
|
projection: numOrNull(s.projection),
|
|
edge_pct: numOrNull(s.edge_pct),
|
|
takeable: boolOrNull(s.takeable),
|
|
value: boolOrNull(s.value),
|
|
archetype: s.archetype || null,
|
|
|
|
refused,
|
|
// Set TRUE only by the collector's publication signal (see createCollector).
|
|
published: false,
|
|
// CANONICAL EVENT IDENTITY. Declared on every row; null when the sport
|
|
// has no resolver or the event could not be resolved. `game_id` below
|
|
// stays as the LEGACY DERIVED label — useful for diagnostics and
|
|
// compatibility, never authoritative for identity.
|
|
canonical_event_id: (s.canonical_event_id ?? base.canonical_event_id) ?? null,
|
|
event_identity_source: (s.event_identity_source ?? base.event_identity_source) ?? null,
|
|
event_identity_method: (s.event_identity_method ?? base.event_identity_method) ?? null,
|
|
event_identity_version: (s.event_identity_version ?? base.event_identity_version) ?? null,
|
|
event_occurrence: (s.event_occurrence ?? base.event_occurrence) ?? null,
|
|
refusal_reason: refused
|
|
? (s.suppressed_reason || (s.insufficient_data ? 'insufficient_data' : 'no_grade'))
|
|
: null,
|
|
|
|
// The counterfactual enabler. Absent on pre-feature refusals (the juice
|
|
// gate runs before features are computed) — honestly null, never faked.
|
|
features: s._features && Object.keys(s._features).length ? s._features : null,
|
|
|
|
// FIX A5 — the RAW inputs the hits factors read, frozen at grade time so a
|
|
// re-audit reads evidence instead of reconstructing context. Inputs only:
|
|
// the multiplier is recomputable from them and is deliberately not stored.
|
|
// Null on every non-hits row and on any row with no factor context.
|
|
factor_inputs: s.factor_inputs || null,
|
|
|
|
// The SHADOW CHAIN read. Declared here (always present, usually null) so
|
|
// every row in a batch carries the same keys — PostgREST builds a bulk
|
|
// insert from the FIRST row's shape, so a column that appears only on some
|
|
// rows is silently dropped for the whole batch. Filled by
|
|
// `mergeChainShadow` after enrichment; never read by anything served.
|
|
chain_shadow: null,
|
|
|
|
// The PROBABILITY CONTRACT shadow — what the certified serving contract
|
|
// WOULD serve, beside what was actually served. Declared always, for the
|
|
// same first-row-shape reason as chain_shadow. Filled by
|
|
// `mergeProbabilityContract`; SHADOW ONLY, read by nothing served.
|
|
probability_contract: null,
|
|
});
|
|
}
|
|
return rows;
|
|
}
|
|
|
|
/** A collector to hand to gradeAndCacheSlate's `onGraded` hook. */
|
|
function createCollector(ctx) {
|
|
const rows = [];
|
|
// Index by (player|stat|line|side) so the publication signal can find the row
|
|
// it already collected without re-deriving the winner rule.
|
|
const index = new Map();
|
|
const keyOf = (r) => [r.player_key, r.stat, r.line, r.side].join('|');
|
|
return {
|
|
rows,
|
|
onGraded(base, sides) {
|
|
try {
|
|
const made = rowsFromSides(base, sides, ctx);
|
|
for (const r of made) index.set(keyOf(r), r);
|
|
rows.push(...made);
|
|
} catch { /* collection never affects grading */ }
|
|
},
|
|
/**
|
|
* PUBLICATION SIGNAL — fired only for the side that actually reaches the
|
|
* slate.
|
|
*
|
|
* ── WHY THIS EXISTS ──────────────────────────────────────────────────
|
|
* `onGraded` fires with BOTH sides, graded AND refused, BEFORE any
|
|
* filtering. Only the higher-confidence graded side becomes the served
|
|
* Read. MEASURED on mlb 2026-08-26: of 13,012 collected rows, 3,859 were
|
|
* refusals and only 4,569 matched a prop that reached the ledger — 8,443
|
|
* (64.9%) describe a state the user was never shown.
|
|
*
|
|
* Nothing on the row said which was which, so treating every capture as a
|
|
* published claim would have built "publication history" for model activity
|
|
* that was never published. This marks the difference at the one place that
|
|
* knows it: the moment the winner is chosen.
|
|
*/
|
|
onPublished(base, winner) {
|
|
try {
|
|
if (!winner) return;
|
|
const made = rowsFromSides(base, [winner], ctx);
|
|
for (const r of made) {
|
|
const hit = index.get(keyOf(r));
|
|
// `published` alone. A companion `published_side` was written here
|
|
// and is NOT a model_snapshots column — it made PostgREST reject the
|
|
// ENTIRE retention batch with a 400, silently, for every sport. The
|
|
// side is already on the row (`side`), so the flag was redundant as
|
|
// well as invalid.
|
|
if (hit) { hit.published = true; }
|
|
}
|
|
} catch { /* signalling never affects grading */ }
|
|
},
|
|
};
|
|
}
|
|
|
|
|
|
/**
|
|
* DUAL-WRITE: attach append-only Read lineage to rows about to be persisted.
|
|
*
|
|
* ── WHY THIS IS BEST-EFFORT AND WHY IT MUST DECLARE EVERY KEY ────────────
|
|
* Retention is already best-effort relative to the ledger; lineage sits one
|
|
* layer further out, so a lineage failure must never cost a retention row and
|
|
* never reach the grade. Every failure path below leaves the row persistable
|
|
* with NULL lineage, which reads as "unknown" — the honest value.
|
|
*
|
|
* Lineage keys are declared on EVERY row even when unresolved. PostgREST builds
|
|
* a bulk insert from the FIRST row's shape, so a key present on only some rows
|
|
* is dropped for the whole batch — the same defect that silently killed
|
|
* retention for three days when `chain_shadow` was missing from the schema.
|
|
*
|
|
* NOT AUTHORITATIVE. Nothing reads these columns to serve a product surface;
|
|
* this records truth in parallel so it can be compared against the live path
|
|
* before any authority changes.
|
|
*/
|
|
const LINEAGE_KEYS = Object.freeze([
|
|
'read_id', 'read_natural_key', 'claim_digest', 'revision_ordinal',
|
|
'supersedes_id', 'recaptures_id', 'lineage_action', 'lineage_state',
|
|
'lineage_version', 'change_type', 'claim_schema_version',
|
|
'digest_algorithm_version',
|
|
]);
|
|
|
|
/**
|
|
* ── WHAT A COMPLETED LINEAGE ACTION IS ───────────────────────────────────
|
|
* A row carrying `read_natural_key` is NOT lineage history. The key is stamped
|
|
* on every candidate BEFORE the family lookup runs, so a failed resolution
|
|
* leaves it behind on a row that never became an action.
|
|
*
|
|
* A completed action is a row that answers all four questions the chronology
|
|
* asks of it: WHICH Read (`read_id`), WHAT it did (`lineage_action`), WHAT it
|
|
* claimed (`claim_digest` + the two version stamps that make the digest
|
|
* interpretable), and WHERE it sits (`revision_ordinal`). Miss any one and the
|
|
* row cannot serve as history, a head, a parent or a recapture predecessor.
|
|
*
|
|
* This is stated as what an action IS — not as a rule shaped to exclude one
|
|
* known batch of failed rows.
|
|
*/
|
|
const VALID_LINEAGE_ACTION_FIELDS = Object.freeze([
|
|
'read_id', 'read_natural_key', 'lineage_action', 'claim_digest',
|
|
'revision_ordinal', 'lineage_state', 'lineage_version',
|
|
'claim_schema_version', 'digest_algorithm_version',
|
|
]);
|
|
|
|
/** Does this physical row qualify as a completed lineage action? */
|
|
function isValidLineageAction(row) {
|
|
if (!row) return false;
|
|
for (const f of VALID_LINEAGE_ACTION_FIELDS) {
|
|
const v = row[f];
|
|
if (v === null || v === undefined || v === '') return false;
|
|
}
|
|
return true;
|
|
}
|
|
|
|
/**
|
|
* ── WHY THE FAMILY IS FETCHED BY (sport, game_date) ──────────────────────
|
|
* `readLineage.readNaturalKey` builds the key as
|
|
* sport | game_date | player_key | stat | side | line [| #event]
|
|
* so SPORT AND GAME_DATE ARE COMPONENTS OF THE KEY ITSELF. Two rows sharing a
|
|
* natural key necessarily share both. Restricting the lookup to the (sport,
|
|
* game_date) pairs present in the requested keys is therefore LOSSLESS BY
|
|
* CONSTRUCTION, not an optimisation that trades recall for speed. A test
|
|
* asserts the derivation against the key builder so the two cannot drift.
|
|
*
|
|
* That bound is what makes the read scale with THE SLATE rather than with
|
|
* history: one date's valid actions, however many years of chronology sit
|
|
* behind it.
|
|
*
|
|
* WHY NOT KEEP THE IN-LIST. Measured 2026-08-28: `read_natural_key` has NO
|
|
* pg_stats row at all (the last autoanalyze predates the column ever being
|
|
* populated), so the planner falls back to a default per-value selectivity.
|
|
* At 100 keys it estimated 172,409 rows and chose a sequential scan of 344,818
|
|
* — 8.5s, then 57014. The plan for this shape is
|
|
* `Index Scan using model_snapshots_lineage_family_idx, cost 0.28..1.92`,
|
|
* and it does not depend on an estimate being good.
|
|
*/
|
|
function familyScopesFrom(naturalKeys) {
|
|
const scopes = new Map();
|
|
for (const k of naturalKeys || []) {
|
|
const parts = String(k).split('|');
|
|
const sport = parts[0];
|
|
const gameDate = parts[1];
|
|
if (!sport || !gameDate) continue;
|
|
scopes.set(`${sport}|${gameDate}`, { sport, game_date: gameDate });
|
|
}
|
|
return [...scopes.values()];
|
|
}
|
|
|
|
/** Claim fields the resolver needs from EXISTING rows to classify a change. */
|
|
const LINEAGE_FETCH_CLAIM = Object.freeze([
|
|
'canonical_event_id',
|
|
'line', 'side', 'book', 'over_odds', 'under_odds',
|
|
'p_win', 'grade', 'confidence', 'projection', 'edge_pct', 'ev_pct',
|
|
'takeable', 'value', 'fair_odds', 'fair_prob', 'confidence_basis',
|
|
'refused', 'refusal_reason',
|
|
]);
|
|
|
|
function blankLineage() {
|
|
const o = {};
|
|
for (const k of LINEAGE_KEYS) o[k] = null;
|
|
return o;
|
|
}
|
|
|
|
/**
|
|
* ── DUAL-WRITE FAILURE CONTRACT ──────────────────────────────────────────
|
|
* Lineage lives on the SAME ROW as the retention record and is attached before
|
|
* the single upsert that writes it. That is not a convenience; it is what makes
|
|
* two of the four failure quadrants structurally impossible:
|
|
*
|
|
* legacy OK + lineage OK one insert carries both.
|
|
* legacy OK + lineage FAILS the row still persists with NULL lineage. The
|
|
* claim is retained; only its position in the
|
|
* chronology is unknown, and `out.lineage.error`
|
|
* plus the not_published/refused counters make
|
|
* the gap measurable rather than silent.
|
|
* legacy FAILS + lineage OK IMPOSSIBLE. There is no separate lineage
|
|
* write, so a failed retention insert cannot
|
|
* leave behind a lineage record claiming a Read
|
|
* was published. This is the quadrant that would
|
|
* manufacture false publication history, and the
|
|
* write ordering removes it rather than guarding
|
|
* against it.
|
|
* legacy FAILS + lineage n/a nothing is written at all.
|
|
*
|
|
* Ordering matters and is deliberate: resolve first, then one write. A separate
|
|
* lineage table written after the fact would reintroduce the impossible
|
|
* quadrant as a real one.
|
|
*/
|
|
async function attachLineage(rows, deps = {}) {
|
|
const out = {
|
|
attempted: Array.isArray(rows) ? rows.length : 0,
|
|
origins: 0, revisions: 0, recaptures: 0, forks: 0, refused: 0, unresolved: 0,
|
|
// A capture that was never the served Read is NOT a lineage failure. It is
|
|
// counted separately so the parity gap stays legible: lumping it into
|
|
// `refused` would make a healthy slate look broken.
|
|
not_published: 0,
|
|
change_types: {},
|
|
error: null,
|
|
};
|
|
const list = Array.isArray(rows) ? rows : [];
|
|
// Declare the keys first, unconditionally. Even a total failure leaves a
|
|
// persistable batch shape.
|
|
for (const r of list) Object.assign(r, blankLineage());
|
|
if (list.length === 0) return out;
|
|
|
|
try {
|
|
const lineage = deps.lineage || require('./read/readLineage');
|
|
const mintReadId = deps.mintReadId || (() => require('crypto').randomUUID());
|
|
|
|
const keyed = [];
|
|
for (const r of list) {
|
|
const key = lineage.readNaturalKey(r);
|
|
if (!key) { out.unresolved += 1; continue; }
|
|
r.read_natural_key = key;
|
|
keyed.push(r);
|
|
}
|
|
if (keyed.length === 0) return out;
|
|
|
|
// ONE bounded read of the existing families for these natural keys.
|
|
const keys = [...new Set(keyed.map((r) => r.read_natural_key))];
|
|
const fetchExisting = deps.fetchExisting || (async (naturalKeys) => {
|
|
const getClient = deps.getClient || require('../utils/supabase').getSupabaseServiceClient;
|
|
const supabase = getClient();
|
|
if (!supabase) return null;
|
|
const { paginate } = require('../utils/safePaginate');
|
|
const columns = ['id', 'read_id', 'read_natural_key', 'game_id', 'claim_digest',
|
|
'revision_ordinal', 'lineage_action', 'supersedes_id', 'captured_at',
|
|
'lineage_state', 'lineage_version', 'claim_schema_version',
|
|
'digest_algorithm_version', ...LINEAGE_FETCH_CLAIM].join(', ');
|
|
const wanted = new Set(naturalKeys);
|
|
const found = [];
|
|
out.lookup = { scopes: 0, rows_scanned: 0, valid_actions: 0, invalid_excluded: 0 };
|
|
for (const scope of familyScopesFrom(naturalKeys)) {
|
|
// ONE index-backed range per (sport, game_date), walked with the
|
|
// repository's keyset paginator — never a hand-rolled second walk, and
|
|
// it THROWS rather than treating a failed page as end-of-data.
|
|
// eslint-disable-next-line no-await-in-loop
|
|
const rows = await paginate(() => supabase
|
|
.from('model_snapshots')
|
|
.select(columns)
|
|
.eq('sport', scope.sport)
|
|
.eq('game_date', scope.game_date)
|
|
.not('lineage_action', 'is', null), { key: 'id', label: 'attachLineage.fetchExisting' });
|
|
out.lookup.scopes += 1;
|
|
out.lookup.rows_scanned += rows.length;
|
|
for (const r of rows) {
|
|
if (!wanted.has(r.read_natural_key)) continue;
|
|
// The predicate is applied in code as well as in the query. The query
|
|
// narrows what travels; this decides what COUNTS, and a row that half
|
|
// resolved must never be mistaken for history by either.
|
|
if (!isValidLineageAction(r)) { out.lookup.invalid_excluded += 1; continue; }
|
|
out.lookup.valid_actions += 1;
|
|
found.push(r);
|
|
}
|
|
}
|
|
return found;
|
|
});
|
|
|
|
const existing = await fetchExisting(keys);
|
|
if (existing === null) return out; // no database configured — leave NULL
|
|
// Second line of defence: an injected or future fetcher must not be able to
|
|
// smuggle a half-resolved row into the chronology either.
|
|
const validExisting = existing.filter(isValidLineageAction);
|
|
out.invalid_rows_excluded = existing.length - validExisting.length;
|
|
|
|
const byKey = new Map();
|
|
for (const e of validExisting) {
|
|
// The prior claim rides alongside so `classifyChange` can say WHY the
|
|
// published state advanced, not merely that the digest differs.
|
|
const claim = {};
|
|
for (const f of LINEAGE_FETCH_CLAIM) claim[f] = e[f];
|
|
const entry = { ...e, claim };
|
|
if (!byKey.has(e.read_natural_key)) byKey.set(e.read_natural_key, []);
|
|
byKey.get(e.read_natural_key).push(entry);
|
|
}
|
|
|
|
for (const r of keyed) {
|
|
const res = lineage.resolveLineage({
|
|
candidate: r,
|
|
existing: byKey.get(r.read_natural_key) || [],
|
|
mintReadId,
|
|
});
|
|
if (!res.ok) {
|
|
if (res.refused === 'not_published') out.not_published += 1;
|
|
else if (res.action === lineage.LINEAGE_ACTION.FORK_DETECTED) out.forks += 1;
|
|
else out.refused += 1;
|
|
// A refusal leaves NULL lineage on a row that is still persisted. The
|
|
// claim is retained; only its position in the chronology is unknown.
|
|
continue;
|
|
}
|
|
r.read_id = res.read_id;
|
|
r.claim_digest = res.claim_digest;
|
|
r.revision_ordinal = res.revision_ordinal;
|
|
r.supersedes_id = res.supersedes_id ?? null;
|
|
r.recaptures_id = res.recaptures_id ?? null;
|
|
r.lineage_action = res.action;
|
|
r.lineage_state = res.lineage_state;
|
|
r.lineage_version = res.lineage_version;
|
|
r.change_type = res.change_type || null;
|
|
r.claim_schema_version = res.claim_schema_version || null;
|
|
r.digest_algorithm_version = res.digest_algorithm_version || null;
|
|
if (res.change_type) {
|
|
out.change_types[res.change_type] = (out.change_types[res.change_type] || 0) + 1;
|
|
}
|
|
|
|
if (res.action === lineage.LINEAGE_ACTION.ORIGIN) out.origins += 1;
|
|
else if (res.action === lineage.LINEAGE_ACTION.REVISION) out.revisions += 1;
|
|
else if (res.action === lineage.LINEAGE_ACTION.RECAPTURE) out.recaptures += 1;
|
|
|
|
// Newly minted read_ids must be visible to later rows in the SAME batch,
|
|
// or two rows of one Read would each mint an id and fork it immediately.
|
|
const fam = byKey.get(r.read_natural_key) || [];
|
|
const claim = {};
|
|
for (const f of LINEAGE_FETCH_CLAIM) claim[f] = r[f];
|
|
fam.push({
|
|
id: null, read_id: r.read_id, read_natural_key: r.read_natural_key,
|
|
game_id: r.game_id, claim_digest: r.claim_digest,
|
|
revision_ordinal: r.revision_ordinal, lineage_action: r.lineage_action,
|
|
supersedes_id: r.supersedes_id, captured_at: r.captured_at, claim,
|
|
});
|
|
byKey.set(r.read_natural_key, fam);
|
|
}
|
|
} catch (e) {
|
|
out.error = e && e.message ? e.message : String(e);
|
|
}
|
|
|
|
// ── FAILURE ATOMICITY ───────────────────────────────────────────────────
|
|
// `read_natural_key` is stamped on every candidate BEFORE the family lookup,
|
|
// so a lookup that throws leaves it behind on a row that never became an
|
|
// action. That is the shape the 2026-08-28 scheduled canary wrote 1,119 times.
|
|
//
|
|
// It is LINEAGE FAMILY IDENTITY, not publication provenance, so a failed
|
|
// attempt has no claim to it. `publication_id` / `published_at` are NOT
|
|
// touched here: the slate really was published, and erasing that to make the
|
|
// failure look tidier would delete a true fact to hide a false one.
|
|
//
|
|
// The result: a failed row carries no lineage-specific state at all, so it
|
|
// cannot be mistaken for history by a future resolver, a query, or a reader.
|
|
let cleared = 0;
|
|
for (const r of list) {
|
|
if (isValidLineageAction(r)) continue;
|
|
if (LINEAGE_KEYS.some((k) => r[k] !== null && r[k] !== undefined)) cleared += 1;
|
|
Object.assign(r, blankLineage());
|
|
}
|
|
out.incomplete_cleared = cleared;
|
|
return out;
|
|
}
|
|
|
|
/** The unique index that prevents a forked chain, named so the guard is legible. */
|
|
const SUPERSEDES_CONSTRAINT = 'model_snapshots_supersedes_unique';
|
|
const FORK_RETRY_LIMIT = 2;
|
|
|
|
function isSupersedesConflict(error) {
|
|
if (!error) return false;
|
|
const m = `${error.message || ''} ${error.details || ''}`;
|
|
return m.includes(SUPERSEDES_CONSTRAINT) || (error && error.code === '23505' && m.includes('supersedes'));
|
|
}
|
|
|
|
/**
|
|
* Recover a batch that lost a supersedes race.
|
|
*
|
|
* Bounded on purpose: two attempts, then an explicit failure. An unbounded loop
|
|
* against a writer that keeps winning would spin forever, and the correct answer
|
|
* after a couple of misses is to report the gap rather than keep trying.
|
|
*
|
|
* Sequence per attempt: re-resolve the batch against the CURRENT head, then
|
|
* write. Re-resolution is what makes this safe -- if the rival already published
|
|
* an identical claim, `attachLineage` returns RECAPTURE and the row lands
|
|
* idempotently; if our claim is still materially different it appends after the
|
|
* new head instead of trying to supersede a parent that is no longer the head.
|
|
*/
|
|
async function recoverFromFork(chunk, deps = {}) {
|
|
const out = { recovered: false, written: 0, attempts: 0, error: null };
|
|
const supabase = deps.supabase;
|
|
for (let attempt = 1; attempt <= FORK_RETRY_LIMIT; attempt += 1) {
|
|
out.attempts = attempt;
|
|
try {
|
|
// Clear stale lineage so re-resolution sees the row as a fresh candidate.
|
|
for (const r of chunk) Object.assign(r, blankLineage());
|
|
// eslint-disable-next-line no-await-in-loop
|
|
await attachLineage(chunk, deps);
|
|
// eslint-disable-next-line no-await-in-loop
|
|
const { error } = await upsertSnapshotChunk(supabase, chunk);
|
|
if (!error) { out.recovered = true; out.written = chunk.length; return out; }
|
|
if (!isSupersedesConflict(error)) { out.error = error.message; return out; }
|
|
out.error = error.message;
|
|
} catch (e) {
|
|
out.error = e && e.message ? e.message : String(e);
|
|
return out;
|
|
}
|
|
}
|
|
return out;
|
|
}
|
|
|
|
/**
|
|
* COMMIT PUBLICATION - the only place a claim becomes "published".
|
|
*
|
|
* -- WHY THIS IS SEPARATE FROM THE CAPTURE WRITE -------------------------
|
|
* The authoritative served slate is ONE atomic Redis write,
|
|
* `cacheSet('snapshot:{sport}:latest', snapshot)`. Retention persists BEFORE
|
|
* that write (measured: persistRetention at snapshotService:1137, the slate
|
|
* write at :1152), and the winner is selected earlier still. So neither capture
|
|
* nor winner-selection can mean "published":
|
|
*
|
|
* winner selected -> intended for the slate
|
|
* retention persisted -> the capture is on the record
|
|
* REDIS WRITE SUCCEEDS -> the claim was actually made available <-- HERE
|
|
*
|
|
* Setting the marker any earlier lets lineage assert a publication that a failed
|
|
* Redis write never made. That is the one thing lineage must never do.
|
|
*
|
|
* -- PUBLICATION IS SLATE-LEVEL, BECAUSE THE WRITE IS --------------------
|
|
* The whole envelope commits or none of it does, so every served row shares one
|
|
* `publication_id` and one `published_at`. Modelling per-row publication would
|
|
* claim a granularity the transport does not have.
|
|
*
|
|
* -- FAILURE DIRECTION IS DELIBERATE ------------------------------------
|
|
* If this update fails after a successful Redis write, rows stay
|
|
* `published: false` with NULL lineage. The record then UNDERSTATES what was
|
|
* served, which is recoverable and observable. The opposite error - claiming a
|
|
* publication that did not happen - is not.
|
|
*/
|
|
async function commitPublication(spec, deps = {}) {
|
|
const s = spec || {};
|
|
const out = {
|
|
publication_id: s.publicationId || null,
|
|
published_at: s.publishedAt || null,
|
|
candidates: 0, published: 0, skipped: false, error: null, lineage: null,
|
|
};
|
|
const served = Array.isArray(s.rows) ? s.rows.filter((r) => r && r.published === true) : [];
|
|
out.candidates = served.length;
|
|
if (served.length === 0) return out;
|
|
|
|
try {
|
|
const getClient = deps.getClient || require('../utils/supabase').getSupabaseServiceClient;
|
|
const supabase = getClient();
|
|
if (!supabase) { out.skipped = true; return out; }
|
|
|
|
for (const r of served) {
|
|
r.published_at = out.published_at;
|
|
r.publication_id = out.publication_id;
|
|
}
|
|
// Lineage resolves ONLY over confirmed-published rows.
|
|
if (deps.lineage !== false) {
|
|
out.lineage = await attachLineage(served, { ...deps, getClient });
|
|
}
|
|
|
|
const CHUNK = 250;
|
|
for (let i = 0; i < served.length; i += CHUNK) {
|
|
const chunk = served.slice(i, i + CHUNK);
|
|
// eslint-disable-next-line no-await-in-loop
|
|
const { error } = await upsertSnapshotChunk(supabase, chunk);
|
|
if (!error) { out.published += chunk.length; continue; }
|
|
|
|
// ── CONCURRENCY RECOVERY ────────────────────────────────────────────
|
|
// `model_snapshots_supersedes_unique` rejects a second row superseding the
|
|
// same parent. That constraint is the authority against a forked history,
|
|
// and hitting it means another writer advanced the chain between our
|
|
// resolve and our write. The loser must not silently drop a materially new
|
|
// published state, and must not loop.
|
|
if (!isSupersedesConflict(error)) { out.error = error.message; break; }
|
|
// eslint-disable-next-line no-await-in-loop
|
|
const rec = await recoverFromFork(chunk, { ...deps, getClient, supabase });
|
|
out.conflicts = (out.conflicts || 0) + 1;
|
|
if (rec.recovered) {
|
|
out.recovered = (out.recovered || 0) + 1;
|
|
out.published += rec.written;
|
|
} else {
|
|
out.recovery_failed = (out.recovery_failed || 0) + 1;
|
|
out.error = rec.error || 'fork recovery exhausted';
|
|
break;
|
|
}
|
|
}
|
|
} catch (e) {
|
|
out.error = e && e.message ? e.message : String(e);
|
|
}
|
|
|
|
// ── EXACT PARITY-GAP IDENTITY ──────────────────────────────────────────
|
|
// "lineage failed" is not an answer. When the product published successfully
|
|
// and the record did not, the system must be able to say WHICH EXACT SERVED
|
|
// CLAIM is missing, without re-deriving it from player/stat/line/date.
|
|
//
|
|
// Everything needed to recover it deterministically is captured here from the
|
|
// in-memory rows that were actually served.
|
|
if (out.error || (out.lineage && out.lineage.error) || out.published < out.candidates) {
|
|
out.parity_gap = Object.freeze({
|
|
sport: served[0] ? served[0].sport : null,
|
|
publication_id: out.publication_id,
|
|
published_at: out.published_at,
|
|
attempted_at: (deps.now || (() => new Date().toISOString()))(),
|
|
reason: out.error || (out.lineage && out.lineage.error) || 'partial_commit',
|
|
missing_count: out.candidates - out.published,
|
|
// The exact rows, by the identity the retention table itself keys on.
|
|
missing: Object.freeze(served.slice(out.published).map((r) => Object.freeze({
|
|
snapshot_id: r.snapshot_id || null,
|
|
player_key: r.player_key || null,
|
|
stat: r.stat || null,
|
|
line: r.line === undefined ? null : r.line,
|
|
side: r.side || null,
|
|
canonical_event_id: r.canonical_event_id || null,
|
|
read_id: r.read_id || null,
|
|
read_natural_key: r.read_natural_key || null,
|
|
claim_digest: r.claim_digest || null,
|
|
captured_at: r.captured_at || null,
|
|
}))),
|
|
});
|
|
}
|
|
return out;
|
|
}
|
|
|
|
/**
|
|
* Persist rows. Chunked, upsert-on-conflict-ignore so a retried cycle can never
|
|
* duplicate. No-ops (successfully) without Supabase env, so tests and local dev
|
|
* never touch a database.
|
|
*/
|
|
async function persist(rows, deps = {}) {
|
|
const out = { attempted: Array.isArray(rows) ? rows.length : 0, written: 0, skipped: false, error: null };
|
|
if (!out.attempted) return out;
|
|
try {
|
|
const getClient = deps.getClient || require('../utils/supabase').getSupabaseServiceClient;
|
|
const supabase = getClient();
|
|
if (!supabase) { out.skipped = true; return out; }
|
|
// NOTE: lineage is NOT resolved here any more.
|
|
//
|
|
// Capture happens BEFORE the authoritative Redis slate write
|
|
// (`snapshot:{sport}:latest`), so resolving lineage at this point would
|
|
// record a published claim for a slate that may never be served. Publication
|
|
// is committed by `commitPublication()` after that write succeeds.
|
|
//
|
|
// Rows are persisted with `published: false` and NULL lineage, which is the
|
|
// truthful state of a capture that has not yet been served.
|
|
const CHUNK = 250;
|
|
for (let i = 0; i < rows.length; i += CHUNK) {
|
|
const chunk = rows.slice(i, i + CHUNK);
|
|
const res = await upsertSnapshotChunk(supabase, chunk, { ignoreDuplicates: true });
|
|
if (res.fellBack) out.legacy_conflict_fallback = true;
|
|
const { error } = res;
|
|
if (error) { out.error = error.message; break; }
|
|
out.written += chunk.length;
|
|
}
|
|
} catch (e) {
|
|
out.error = e && e.message ? e.message : String(e);
|
|
}
|
|
return out;
|
|
}
|
|
|
|
/**
|
|
* Fill archetype / team / opponent onto collected rows from the ENRICHED grades.
|
|
*
|
|
* Retention collects at GRADE time, which is the only moment the feature vector
|
|
* exists — but archetype/team/opponent are attached later, during snapshot
|
|
* enrichment. Capturing at grade time alone left all three permanently null,
|
|
* which specifically blocks the archetype-baselined metrics work.
|
|
*
|
|
* CONTRACT: this ONLY fills those three fields. It must never touch `features`
|
|
* or any model output — grade-time values are the record, and enrichment must
|
|
* not rewrite history. Unmatched rows (e.g. refusals, which never reach the
|
|
* enriched slate) pass through untouched with the fields left null: honestly
|
|
* absent, not guessed.
|
|
*/
|
|
function mergeEnrichment(rows, enrichedGrades) {
|
|
if (!Array.isArray(rows) || !rows.length) return rows || [];
|
|
const byPlayer = new Map();
|
|
for (const g of enrichedGrades || []) {
|
|
const raw = g && (g.player || g.player_name);
|
|
if (!raw) continue;
|
|
const k = nameKey(raw);
|
|
// First enriched grade per player wins; archetype/team are player-level.
|
|
if (!byPlayer.has(k)) {
|
|
byPlayer.set(k, {
|
|
archetype: g.archetype ?? null,
|
|
team: g.team ?? null,
|
|
opponent: g.opponent ?? null,
|
|
});
|
|
}
|
|
}
|
|
return rows.map((r) => {
|
|
const e = byPlayer.get(r.player_key);
|
|
if (!e) return r;
|
|
return {
|
|
...r,
|
|
archetype: r.archetype ?? e.archetype ?? null,
|
|
team: r.team ?? e.team ?? null,
|
|
opponent: r.opponent ?? e.opponent ?? null,
|
|
};
|
|
});
|
|
}
|
|
|
|
/**
|
|
* Attach the SHADOW CHAIN read to collected rows — the (chain_p, counter_p,
|
|
* outcome) triple's first two thirds.
|
|
*
|
|
* Like `mergeEnrichment` this fills ONE field and touches nothing else. It runs
|
|
* later than grade time for the same structural reason the archetype does: the
|
|
* chain needs the statcast rows and park factors the enrichment pass loaded, and
|
|
* pulling that forward into the grader would put per-prop I/O on the serving
|
|
* path.
|
|
*
|
|
* THE SIDE ALIGNMENT IS THE LOAD-BEARING PART. The chain computes P(over the
|
|
* line); `p_win` on the row is expressed for the graded SIDE. Storing the raw
|
|
* over-probability against an under row's `p_win` would invert every comparison
|
|
* made from it afterwards, silently — so `alignToSide` is given the row's own
|
|
* side and the row's own counter probability, and both halves of the triple end
|
|
* up pointing the same way. `outcome` is written by the ordinary settle pass and
|
|
* is already side-aligned, which completes it.
|
|
*
|
|
* A row the chain could not read is left NULL. It is not a zero probability and
|
|
* not an average hitter — the chain refusing to read someone is a fact worth
|
|
* keeping, and a fabricated third of a triple would poison the adjudication this
|
|
* column exists to enable.
|
|
*/
|
|
/**
|
|
* Attach the SHADOW CHAIN read to collected rows — the (chain_p, counter_p,
|
|
* outcome) triple's first two thirds.
|
|
*
|
|
* Like `mergeEnrichment` this fills ONE field and touches nothing else. It runs
|
|
* later than grade time for the same structural reason the archetype does: the
|
|
* chain needs the statcast rows and park factors the enrichment pass loaded, and
|
|
* pulling that forward into the grader would put per-prop I/O on the serving
|
|
* path.
|
|
*
|
|
* THE SIDE ALIGNMENT IS THE LOAD-BEARING PART. The chain computes P(over the
|
|
* line); `p_win` on the row is expressed for the graded SIDE. Storing the raw
|
|
* over-probability against an under row's `p_win` would invert every comparison
|
|
* made from it afterwards, silently — so `alignToSide` is given the row's own
|
|
* side and the row's own counter probability, and both halves of the triple end
|
|
* up pointing the same way. `outcome` is written by the ordinary settle pass and
|
|
* is already side-aligned, which completes it.
|
|
*
|
|
* A row the chain could not read is left NULL. It is not a zero probability and
|
|
* not an average hitter — the chain refusing to read someone is a fact worth
|
|
* keeping, and a fabricated third of a triple would poison the adjudication this
|
|
* column exists to enable.
|
|
*/
|
|
/**
|
|
* PROBABILITY CONTRACT SHADOW.
|
|
*
|
|
* Records, per row: the raw model belief, what the certified contract would
|
|
* serve, the state, the estimator identity, and what EV/Kelly/VALUE would be
|
|
* under the actionability law. It writes to its own column and mutates nothing
|
|
* a user sees — `p_win`, `confidence`, `grade`, `ev_pct`, `value` and
|
|
* `takeable` on the row are untouched.
|
|
*
|
|
* The raw probability is ALWAYS carried, including on refusals: it is model
|
|
* evidence, and a row that records only the refusal cannot be re-adjudicated.
|
|
*/
|
|
function mergeProbabilityContract(rows, contract) {
|
|
if (!Array.isArray(rows) || !rows.length) return rows || [];
|
|
if (!contract || typeof contract.resolve !== 'function') return rows;
|
|
const pc = require('./model/probabilityContract');
|
|
|
|
return rows.map((r) => {
|
|
let res;
|
|
try {
|
|
// The ROW's sport and stat travel with it. Omitting them let the
|
|
// resolver substitute its own and calibrate every stat as hits.
|
|
res = contract.resolve({
|
|
sport: r.sport, stat: r.stat,
|
|
model_version: r.model_version, p_win: numOrNull(r.p_win),
|
|
});
|
|
} catch { return r; }
|
|
if (!res) return r;
|
|
let derived;
|
|
// The SIDE'S price. A retention row carries book_odds plus both sides'
|
|
// prices; using the wrong side would price the opposite bet.
|
|
const sideOdds = r.book_odds != null ? r.book_odds
|
|
: (String(r.side) === 'under' ? r.under_odds : r.over_odds);
|
|
try { derived = pc.derivedClaims(res, sideOdds); } catch { derived = null; }
|
|
return {
|
|
...r,
|
|
probability_contract: {
|
|
raw_model_probability: res.raw_model_probability,
|
|
served_probability: res.served_probability,
|
|
probability_state: res.probability_state,
|
|
estimator_type: res.estimator_type,
|
|
estimator_version: res.estimator_version,
|
|
certification_version: res.certification_version,
|
|
model_version: res.model_version,
|
|
reason: res.reason,
|
|
// WHICH FROZEN ARTIFACT PRODUCED THIS — its IDENTITY, not its body.
|
|
// The curve is committed in the repository and addressable by
|
|
// `artifact_id`, so embedding it on every row would store the same ~900
|
|
// bytes thousands of times per snapshot to say something the id already
|
|
// says. `source_digest` pins the exact observation set it was fitted
|
|
// on, so the row remains reconstructable.
|
|
artifact: res.artifact ? {
|
|
artifact_id: res.artifact.artifact_id,
|
|
procedure_version: res.artifact.procedure_version,
|
|
model_version: res.artifact.model_version,
|
|
fit_as_of: res.artifact.fit_as_of,
|
|
training_cutoff: res.artifact.training_cutoff,
|
|
fit_n: res.artifact.fit_n,
|
|
source_digest: res.artifact.source_digest,
|
|
knot_digest: res.artifact.knot_digest,
|
|
served_curve_digest: res.artifact.served_curve_digest,
|
|
certified_bands: res.artifact.certified_bands,
|
|
stage: res.artifact.stage,
|
|
servable: res.artifact.servable,
|
|
} : null,
|
|
derived: derived ? {
|
|
available: derived.available,
|
|
ev_pct: derived.ev_pct,
|
|
kelly_pct: derived.kelly ? derived.kelly.pct : null,
|
|
value: derived.value,
|
|
} : null,
|
|
// SHADOW. Nothing here has been served to anyone.
|
|
servable: false,
|
|
},
|
|
};
|
|
});
|
|
}
|
|
|
|
function mergeChainShadow(rows, shadow) {
|
|
if (!Array.isArray(rows) || !rows.length) return rows || [];
|
|
const byKey = shadow && shadow.byKey;
|
|
if (!byKey || typeof byKey.get !== 'function') return rows;
|
|
const cs = require('./model/chainShadow');
|
|
|
|
return rows.map((r) => {
|
|
const block = byKey.get(cs.shadowKey(r.player_key, r.stat, r.line));
|
|
if (!block) return r;
|
|
const aligned = cs.alignToSide(block, r.side, r.p_win);
|
|
return aligned ? { ...r, chain_shadow: aligned } : r;
|
|
});
|
|
}
|
|
|
|
function newSnapshotId() {
|
|
return crypto.randomUUID();
|
|
}
|
|
|
|
/* ------------------------------------------------------------------ *
|
|
* RETENTION CONFLICT TARGET — event-aware, mixed-fleet safe
|
|
*
|
|
* SEMANTIC IDENTITY. Two outbound rows are the SAME retention proposition
|
|
* within one cycle when they share: the cycle, the EVENT, the participant, the
|
|
* stat, the line and the side. Provider/book is deliberately NOT part of it —
|
|
* collapsing books is `dedupeProps`'s actual job, and the price anchor is
|
|
* chosen later.
|
|
*
|
|
* The EVENT component uses the strongest truthful label available and never
|
|
* fabricates one: `canonical_event_id` where a sport has a resolver (MLB, where
|
|
* `admitForGrading` makes it non-null for every row that can reach retention),
|
|
* and `game_id` otherwise (NOT NULL in the schema, and the only event label
|
|
* sports without a resolver possess). Both columns are in the identity, so the
|
|
* weaker label still discriminates where the stronger one is absent.
|
|
*
|
|
* NULLS NOT DISTINCT is load-bearing. `canonical_event_id` is NULL for every
|
|
* non-MLB row, and under PostgreSQL's DEFAULT null semantics two NULLs are
|
|
* DISTINCT — measured: the same NBA proposition inserted twice produced TWO
|
|
* rows, i.e. every retry would duplicate for ever. With NULLS NOT DISTINCT the
|
|
* same test produces one row and retry idempotency holds.
|
|
*
|
|
* MIXED-FLEET SAFETY. A rollout serves old and new containers at once
|
|
* (measured: 11 of 12 probes new, 1 old). The two writers need different
|
|
* indexes, and no schema state satisfies both:
|
|
*
|
|
* old index present -> old writer works; new writer ERRORS 23505 on a
|
|
* doubleheader (the legacy index rejects a row the
|
|
* event-aware identity considers distinct)
|
|
* old index dropped -> new writer works; old writer ERRORS 42P10
|
|
*
|
|
* A bare `ON CONFLICT DO NOTHING` would have bridged this, but PostgREST does
|
|
* NOT emit one: `ignoreDuplicates` without `onConflict` was measured raising a
|
|
* real duplicate-key error, so that bridge does not exist through this client.
|
|
*
|
|
* So the writer bridges it instead. It targets the event-aware identity and, on
|
|
* exactly the two errors that mean "the schema is not in the state I expect",
|
|
* retries the SAME chunk on the legacy target. A failed chunk rolls back
|
|
* atomically (measured: 0 rows), so the retry cannot double-write. The result
|
|
* is a writer that is correct in every schema state:
|
|
*
|
|
* legacy only -> 42P10 -> legacy target -> legacy semantics, no outage
|
|
* both present -> 23505 -> legacy target -> legacy semantics, no outage
|
|
* new only -> primary succeeds -> doubleheaders kept
|
|
*
|
|
* The fallback is a BRIDGE, not a resting place: while it fires, doubleheader
|
|
* rows are still lost, and `outbound_collision_count` still reports it.
|
|
* ------------------------------------------------------------------ */
|
|
|
|
const LEGACY_CONFLICT_INDEX = 'model_snapshots_cycle_prop_uniq';
|
|
const LEGACY_CONFLICT = 'snapshot_id,player_key,stat,line,side';
|
|
const RETENTION_CONFLICT = 'snapshot_id,game_id,canonical_event_id,player_key,stat,line,side';
|
|
|
|
/**
|
|
* Is this error "the schema is not in the state the event-aware target
|
|
* expects"? Deliberately narrow: a supersedes conflict is also a 23505, and
|
|
* swallowing THAT would destroy the forked-history guard, so the legacy index
|
|
* must be named.
|
|
*/
|
|
function isLegacyConflictBlock(error) {
|
|
if (!error) return false;
|
|
const code = String(error.code || '');
|
|
const msg = String(error.message || '');
|
|
if (code === '42P10' || /no unique or exclusion constraint matching/i.test(msg)) return true;
|
|
return (code === '23505' || /duplicate key value/i.test(msg)) && msg.includes(LEGACY_CONFLICT_INDEX);
|
|
}
|
|
|
|
/**
|
|
* One chunk write. Returns the error (never throws) plus which target actually
|
|
* carried it, so a caller can report that the bridge is still in use.
|
|
*/
|
|
async function upsertSnapshotChunk(supabase, chunk, opts = {}) {
|
|
const first = await supabase.from('model_snapshots')
|
|
.upsert(chunk, { onConflict: RETENTION_CONFLICT, ...opts });
|
|
if (!first.error) return { error: null, target: RETENTION_CONFLICT, fellBack: false };
|
|
if (!isLegacyConflictBlock(first.error)) {
|
|
return { error: first.error, target: RETENTION_CONFLICT, fellBack: false };
|
|
}
|
|
const second = await supabase.from('model_snapshots')
|
|
.upsert(chunk, { onConflict: LEGACY_CONFLICT, ...opts });
|
|
return { error: second.error || null, target: LEGACY_CONFLICT, fellBack: true };
|
|
}
|
|
|
|
/* ------------------------------------------------------------------ *
|
|
* MATERIALIZATION IDENTITY
|
|
*
|
|
* TRANSPORT COMPLETE and MATERIALIZATION COMPLETE are different facts.
|
|
*
|
|
* The write is `upsert(..., { onConflict: RETENTION_CONFLICT, ignoreDuplicates:
|
|
* true })`. Until migration 049 retires it, the LEGACY index
|
|
* model_snapshots_cycle_prop_uniq UNIQUE (snapshot_id, player_key, stat, line, side)
|
|
* is still enforced and the writer falls back to it (see the bridge above)
|
|
* so two outbound rows sharing that tuple collapse to ONE stored row and the
|
|
* loser is discarded WITHOUT AN ERROR. `written` counts rows in committed
|
|
* chunks, so transport reports both as written and the status reads COMPLETE.
|
|
*
|
|
* `snapshot_id` IS in the identity and is a fresh UUID per cycle, so a row can
|
|
* never be suppressed by a PREVIOUS cycle — append-only chronology across
|
|
* cycles is safe. `canonical_event_id` is NOT in the identity, so two DISTINCT
|
|
* events inside ONE cycle (a doubleheader: same hitter, same stat, same line,
|
|
* both games) DO collide.
|
|
*
|
|
* No duplicate is intentional here: `rowsFromSides` emits one row per (base,
|
|
* side) and bases are already deduped by event::player::stat::line. Any
|
|
* duplicate identity in an outbound payload is therefore an UNEXPECTED
|
|
* COLLISION, and a cohort carrying one does not qualify as evidence.
|
|
*
|
|
* The conflict identity is NOT changed here. This only makes the loss visible.
|
|
* ------------------------------------------------------------------ */
|
|
|
|
/**
|
|
* The exact columns the database conflict target names, in index order —
|
|
* DERIVED from RETENTION_CONFLICT rather than restated, so the identity the
|
|
* materialization check expects can never drift from the identity the database
|
|
* enforces. Restating it is how the two silently disagreed before.
|
|
*/
|
|
const CONFLICT_IDENTITY = Object.freeze(RETENTION_CONFLICT.split(','));
|
|
|
|
const IDENTITY_SEP = '\u001f';
|
|
const IDENTITY_NULL = '\u0000NULL';
|
|
|
|
/**
|
|
* The database identity of one row. `line` is an unconstrained `numeric`, so it
|
|
* is normalised through Number() — otherwise 0.5 on the way in and "0.50" on
|
|
* the way out would read as two identities for one stored row.
|
|
*
|
|
* Every key column is NOT NULL in production (measured: 0 nulls across 328,262
|
|
* rows), so PostgreSQL's NULLS DISTINCT behaviour never applies. A null still
|
|
* gets its own marker rather than collapsing to '' — an empty string is a real
|
|
* value and must not be confused with an absent one.
|
|
*/
|
|
function rowIdentity(row) {
|
|
return CONFLICT_IDENTITY.map((c) => {
|
|
const v = row ? row[c] : undefined;
|
|
if (v === null || v === undefined) return IDENTITY_NULL;
|
|
return c === 'line' ? String(Number(v)) : String(v);
|
|
}).join(IDENTITY_SEP);
|
|
}
|
|
|
|
function identityDigest(identities) {
|
|
return crypto.createHash('sha256').update([...identities].sort().join('\n')).digest('hex').slice(0, 32);
|
|
}
|
|
|
|
/**
|
|
* What SHOULD materialize, derived from the FINAL outbound payload using the
|
|
* exact database identity — never from `attempted`, which counts rows sent, not
|
|
* identities that can exist.
|
|
*/
|
|
function expectedMaterialization(rows) {
|
|
const list = Array.isArray(rows) ? rows : [];
|
|
const seen = new Set();
|
|
let collisions = 0;
|
|
for (const r of list) {
|
|
const id = rowIdentity(r);
|
|
if (seen.has(id)) collisions += 1;
|
|
else seen.add(id);
|
|
}
|
|
return {
|
|
outbound_rows: list.length,
|
|
expected_identities: seen.size,
|
|
collision_count: collisions,
|
|
digest: identityDigest(seen),
|
|
identities: seen,
|
|
};
|
|
}
|
|
|
|
const MATERIALIZATION = Object.freeze({
|
|
COMPLETE: 'MATERIALIZATION_COMPLETE',
|
|
MISSING: 'MATERIALIZATION_MISSING',
|
|
EXTRA: 'MATERIALIZATION_EXTRA',
|
|
COLLISION: 'MATERIALIZATION_UNEXPECTED_COLLISION',
|
|
TRANSPORT_FAILED: 'MATERIALIZATION_UNPROVEN_TRANSPORT_FAILED',
|
|
NOT_RECONCILED: 'MATERIALIZATION_NOT_RECONCILED',
|
|
});
|
|
|
|
/**
|
|
* Exact SET comparison, not a count comparison. Two sets of equal size can
|
|
* still differ, and a cohort that swapped one identity for another would pass
|
|
* every count test ever written.
|
|
*
|
|
* A cohort qualifies ONLY when transport completed, the sets are equal, AND the
|
|
* outbound payload held no colliding identity — a collision passes set equality
|
|
* by construction (the discarded row was never in the expected set) while real
|
|
* rows were lost.
|
|
*/
|
|
function reconcileMaterialization({ expected, actualIdentities, transportStatus }) {
|
|
if (!expected || typeof expected.expected_identities !== 'number') {
|
|
throw new TypeError('reconcileMaterialization requires an expectedMaterialization() result');
|
|
}
|
|
const actual = actualIdentities instanceof Set ? actualIdentities : new Set(actualIdentities || []);
|
|
const exp = expected.identities instanceof Set ? expected.identities : new Set();
|
|
const missing = [...exp].filter((i) => !actual.has(i));
|
|
const extra = [...actual].filter((i) => !exp.has(i));
|
|
const out = {
|
|
expected_materialized_count: expected.expected_identities,
|
|
actual_materialized_count: actual.size,
|
|
outbound_rows: expected.outbound_rows,
|
|
missing_identity_count: missing.length,
|
|
extra_identity_count: extra.length,
|
|
collision_count: expected.collision_count,
|
|
status: MATERIALIZATION.NOT_RECONCILED,
|
|
};
|
|
if (transportStatus && transportStatus !== TERMINAL.COMPLETE) {
|
|
out.status = MATERIALIZATION.TRANSPORT_FAILED;
|
|
return out;
|
|
}
|
|
if (missing.length) out.status = MATERIALIZATION.MISSING;
|
|
else if (extra.length) out.status = MATERIALIZATION.EXTRA;
|
|
else if (expected.collision_count > 0) out.status = MATERIALIZATION.COLLISION;
|
|
else out.status = MATERIALIZATION.COMPLETE;
|
|
return out;
|
|
}
|
|
|
|
/* ------------------------------------------------------------------ *
|
|
* TERMINAL RETENTION STATE
|
|
*
|
|
* `persist()` writes in CHUNKS and STOPS ON THE FIRST FAILED CHUNK. Chunks
|
|
* committed before the failure are already durable, so a failed cycle can
|
|
* leave real, valid-looking rows behind. Row presence under a snapshot_id is
|
|
* therefore NOT completion evidence, and neither is a uniform captured_at.
|
|
*
|
|
* The invariant: written < attempted with attempted > 0 is a FAILED cycle.
|
|
* A partial cycle is never degraded success.
|
|
*
|
|
* `written` counts rows in successfully COMMITTED CHUNKS — not rows inserted.
|
|
* The upsert uses ignoreDuplicates, so a re-run legitimately inserts far fewer
|
|
* database rows than it writes. Comparing `written` to count(*) for a
|
|
* snapshot_id will disagree by design; that is not a partial write.
|
|
* ------------------------------------------------------------------ */
|
|
|
|
const TERMINAL = Object.freeze({
|
|
/** attempted === 0 — a refusal-only slate is still a legitimate cycle. */
|
|
NOTHING_TO_PERSIST: 'NOTHING_TO_PERSIST',
|
|
/** No database configured (dev/test). Not a failure. */
|
|
SKIPPED_NO_DATABASE: 'SKIPPED_NO_DATABASE',
|
|
/** attempted > 0, written === attempted, no unresolved error. */
|
|
COMPLETE: 'COMPLETE',
|
|
/** attempted > 0, written === 0 — the first chunk failed. */
|
|
FAILED_ZERO_WRITE: 'FAILED_ZERO_WRITE',
|
|
/** attempted > 0, 0 < written < attempted — a later chunk failed. */
|
|
FAILED_PARTIAL: 'FAILED_PARTIAL',
|
|
/**
|
|
* Counts look complete but an error is unresolved. Unreachable through
|
|
* today's loop (it breaks before crediting a failed chunk), and kept
|
|
* because the alternative is reporting COMPLETE with an error in hand.
|
|
*/
|
|
FAILED_UNRESOLVED_ERROR: 'FAILED_UNRESOLVED_ERROR',
|
|
});
|
|
|
|
const FAILURE_STATUSES = Object.freeze([
|
|
TERMINAL.FAILED_ZERO_WRITE, TERMINAL.FAILED_PARTIAL, TERMINAL.FAILED_UNRESOLVED_ERROR,
|
|
]);
|
|
|
|
function isRetentionFailure(status) { return FAILURE_STATUSES.includes(status); }
|
|
|
|
/**
|
|
* Classify the EXACT object `persist()` returned. It never recomputes
|
|
* attempted or written — a second calculation could disagree with the writer,
|
|
* and then the status would describe something that did not happen.
|
|
*/
|
|
function classifyPersist(result) {
|
|
if (!result || typeof result.attempted !== 'number' || typeof result.written !== 'number') {
|
|
throw new TypeError('classifyPersist requires the exact persist() result');
|
|
}
|
|
const { attempted, written, skipped, error } = result;
|
|
if (attempted === 0) return TERMINAL.NOTHING_TO_PERSIST;
|
|
if (skipped) return TERMINAL.SKIPPED_NO_DATABASE;
|
|
if (written === 0) return TERMINAL.FAILED_ZERO_WRITE;
|
|
if (written < attempted) return TERMINAL.FAILED_PARTIAL;
|
|
// written >= attempted from here.
|
|
if (error) return TERMINAL.FAILED_UNRESOLVED_ERROR;
|
|
return written === attempted ? TERMINAL.COMPLETE : TERMINAL.FAILED_UNRESOLVED_ERROR;
|
|
}
|
|
|
|
/**
|
|
* Latest terminal retention result per sport, for the internal status probe.
|
|
* In-memory and per-process on purpose: this is an observability surface, not
|
|
* a record. The record is model_snapshots.
|
|
*/
|
|
const lastTerminal = new Map();
|
|
|
|
function recordTerminal({ sport, snapshotId, result, completedAt, rows }) {
|
|
const status = classifyPersist(result);
|
|
// Derived from the FINAL outbound payload, which does not survive the call.
|
|
// This is the one materialization fact unrecoverable from the database
|
|
// afterwards, which is why it is the only one recorded at runtime.
|
|
const expected = rows ? expectedMaterialization(rows) : null;
|
|
const entry = Object.freeze({
|
|
sport: sport || null,
|
|
snapshot_id: snapshotId || null,
|
|
attempted: result.attempted,
|
|
written: result.written,
|
|
status,
|
|
completed_at: completedAt || new Date().toISOString(),
|
|
code_sha: codeSha(),
|
|
// Message only — never row payloads.
|
|
error_summary: result.error ? String(result.error).slice(0, 300) : null,
|
|
// Counts and a digest only — never the identities, which carry player names.
|
|
outbound_rows: expected ? expected.outbound_rows : null,
|
|
expected_materialized_count: expected ? expected.expected_identities : null,
|
|
outbound_collision_count: expected ? expected.collision_count : null,
|
|
expected_identity_digest: expected ? expected.digest : null,
|
|
});
|
|
if (sport) lastTerminal.set(sport, entry);
|
|
return entry;
|
|
}
|
|
|
|
function lastRetention() {
|
|
const out = {};
|
|
for (const [sport, entry] of lastTerminal) out[sport] = entry;
|
|
return out;
|
|
}
|
|
|
|
function resetTerminal() { lastTerminal.clear(); }
|
|
|
|
module.exports = {
|
|
MODEL_VERSION,
|
|
TERMINAL,
|
|
RETENTION_CONFLICT,
|
|
LEGACY_CONFLICT,
|
|
LEGACY_CONFLICT_INDEX,
|
|
isLegacyConflictBlock,
|
|
upsertSnapshotChunk,
|
|
MATERIALIZATION,
|
|
CONFLICT_IDENTITY,
|
|
rowIdentity,
|
|
identityDigest,
|
|
expectedMaterialization,
|
|
reconcileMaterialization,
|
|
classifyPersist,
|
|
isRetentionFailure,
|
|
recordTerminal,
|
|
lastRetention,
|
|
resetTerminal,
|
|
REPAIRED_CHAMPION_VERSION,
|
|
codeSha,
|
|
rowsFromSides,
|
|
createCollector,
|
|
mergeEnrichment,
|
|
mergeChainShadow,
|
|
mergeProbabilityContract,
|
|
persist,
|
|
attachLineage,
|
|
commitPublication,
|
|
isValidLineageAction,
|
|
VALID_LINEAGE_ACTION_FIELDS,
|
|
familyScopesFrom,
|
|
recoverFromFork,
|
|
isSupersedesConflict,
|
|
FORK_RETRY_LIMIT,
|
|
LINEAGE_KEYS,
|
|
LINEAGE_FETCH_CLAIM,
|
|
newSnapshotId,
|
|
__internals: { numOrNull, intOrNull, boolOrNull },
|
|
};
|