11277a1b99
Transport truth says every intended write succeeded. Materialization truth says
every identity that should exist actually exists. The retention writer could
only report the first, and the gap is not theoretical.
THE CONFLICT IDENTITY, traced to the real index:
model_snapshots_cycle_prop_uniq UNIQUE (snapshot_id, player_key, stat, line, side)
written by `upsert(..., { onConflict: same five columns, ignoreDuplicates: true })`.
Measured in production: no key column is ever NULL (0 of 328,262 rows), so
NULLS DISTINCT never applies and the identity is plain column equality. `line`
is an unconstrained numeric, so identity normalises it — 0.5 in and "0.50" out
must not read as two identities for one stored row.
PRIOR-CYCLE COLLISION IS IMPOSSIBLE. `snapshot_id` is in the identity and is a
fresh UUID per cycle, so no row can be suppressed by an earlier cycle.
Append-only chronology across cycles is safe, and `captured_at` is not in the
identity, so a cohort cannot be split by timing.
INTRA-CYCLE COLLISION IS REAL, AND WE CAUSED IT. `canonical_event_id` is NOT in
the identity. A doubleheader — same hitter, same stat, same line, two genuinely
different games — is ONE identity. Demonstrated through the real collector: 4
outbound rows, 2 distinct identities, 2 rows discarded by ignoreDuplicates with
no error, `written` counting all 4 and the terminal status reading COMPLETE.
Before event-aware dedupe the second game was dropped before grading, so the
collision could not arise; that fix moved the loss downstream into retention.
The conflict identity is NOT changed here — that is a separate decision with its
own before/after. This makes the loss visible instead of silent.
EXPECTED vs ACTUAL. `expectedMaterialization(rows)` derives the identity set
from the FINAL outbound payload using the exact database identity — never from
`attempted`, which counts rows sent, not identities that can exist.
`reconcileMaterialization` compares SETS, not counts: two sets of equal size can
still differ, and a cohort that swapped one identity for another passes every
count test ever written. A collision passes set equality by construction (the
discarded row was never in the expected set) while real rows were lost, so
collision_count > 0 fails the cohort on its own.
A cohort is evidence-complete only when transport is COMPLETE, missing = 0,
extra = 0, and collisions = 0.
OBSERVABILITY stayed minimal. `last_retention` was already PER SPORT (a Map
keyed by sport), so no fix was needed there and the route is UNCHANGED — the new
fields ride the existing entry: outbound_rows, expected_materialized_count,
outbound_collision_count, expected_identity_digest. Counts and a digest only,
never the identities, which carry player names. The expected set is the one
materialization fact unrecoverable from the database afterwards, which is why it
is the only thing recorded at runtime.
A collision leaves transport COMPLETE, so the existing failure alert could never
see it. It now has its own high-severity alert naming the counts, the cycle and
the build, and says the cohort is not evidence-complete.
Seven teeth, each injection verified present, against a GREEN baseline of 63:
1 attempted===written as evidence completeness -> 1 fail
2 COUNT(*) equality instead of set equality -> 1 fail
3 snapshot_id dropped from expected identity -> 4 fail
4 unexpected collision allowed to qualify -> 1 fail
5 single global last_retention slot -> 2 fail
6 partial chunk failure treated as usable -> 3 fail
7 collision loses its announcement -> 1 fail
Restored byte-identically (retention 742f116473d97f49, snapshot 81129facbabeb280).
Three brittle assertions repaired, with the reason recorded: two windowed on a
byte count that a neighbouring block outgrew — a test failing because of its
neighbour, not its subject — now windowed to syntactic landmarks; and one
counted TERMINAL.COMPLETE occurrences, which a legitimate comparison
incremented. It now asserts one DECISION and one READ.
persist() and createCollector are BYTE-IDENTICAL. onConflict and
ignoreDuplicates appear in the diff only as prose. Model, event, ledger,
calibration, chain, lineage config, and the status route: UNCHANGED. Zero
lineage/publication files, zero cacheSet changes, zero web paths. Lineage OFF.
Schema contract unchanged: release 64, prod 67, prod-only 3 (debt, not
authorized), missing in prod 0.
384 suites / 5,144 tests pass. web tsc exit 0.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01CQJeAG8vcDoL5zkiaJyVb8
987 lines
42 KiB
JavaScript
987 lines
42 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');
|
|
|
|
/** 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;
|
|
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: nameKey(player),
|
|
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,
|
|
});
|
|
}
|
|
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',
|
|
]);
|
|
|
|
/** 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 found = [];
|
|
const CH = 100; // never send an unbounded id list — it becomes a URL
|
|
for (let i = 0; i < naturalKeys.length; i += CH) {
|
|
const { data, error } = await supabase
|
|
.from('model_snapshots')
|
|
.select(['id', 'read_id', 'read_natural_key', 'game_id', 'claim_digest',
|
|
'revision_ordinal', 'lineage_action', 'supersedes_id', 'captured_at',
|
|
...LINEAGE_FETCH_CLAIM].join(', '))
|
|
.in('read_natural_key', naturalKeys.slice(i, i + CH));
|
|
if (error) throw new Error(error.message);
|
|
if (Array.isArray(data)) found.push(...data);
|
|
}
|
|
return found;
|
|
});
|
|
|
|
const existing = await fetchExisting(keys);
|
|
if (existing === null) return out; // no database configured — leave NULL
|
|
|
|
const byKey = new Map();
|
|
for (const e of existing) {
|
|
// 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);
|
|
}
|
|
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 supabase
|
|
.from('model_snapshots')
|
|
.upsert(chunk, { onConflict: 'snapshot_id,player_key,stat,line,side' });
|
|
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 supabase
|
|
.from('model_snapshots')
|
|
.upsert(chunk, { onConflict: 'snapshot_id,player_key,stat,line,side' });
|
|
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 { error } = await supabase
|
|
.from('model_snapshots')
|
|
.upsert(chunk, {
|
|
onConflict: 'snapshot_id,player_key,stat,line,side',
|
|
ignoreDuplicates: true,
|
|
});
|
|
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.
|
|
*/
|
|
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();
|
|
}
|
|
|
|
/* ------------------------------------------------------------------ *
|
|
* MATERIALIZATION IDENTITY
|
|
*
|
|
* TRANSPORT COMPLETE and MATERIALIZATION COMPLETE are different facts.
|
|
*
|
|
* The write is `upsert(..., { onConflict: 'snapshot_id,player_key,stat,line,side',
|
|
* ignoreDuplicates: true })`, and the production index behind it is
|
|
* model_snapshots_cycle_prop_uniq UNIQUE (snapshot_id, player_key, stat, line, side)
|
|
* 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 of model_snapshots_cycle_prop_uniq, in index order. */
|
|
const CONFLICT_IDENTITY = Object.freeze(['snapshot_id', 'player_key', 'stat', 'line', 'side']);
|
|
|
|
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,
|
|
MATERIALIZATION,
|
|
CONFLICT_IDENTITY,
|
|
rowIdentity,
|
|
identityDigest,
|
|
expectedMaterialization,
|
|
reconcileMaterialization,
|
|
classifyPersist,
|
|
isRetentionFailure,
|
|
recordTerminal,
|
|
lastRetention,
|
|
resetTerminal,
|
|
REPAIRED_CHAMPION_VERSION,
|
|
codeSha,
|
|
rowsFromSides,
|
|
createCollector,
|
|
mergeEnrichment,
|
|
mergeChainShadow,
|
|
persist,
|
|
attachLineage,
|
|
commitPublication,
|
|
recoverFromFork,
|
|
isSupersedesConflict,
|
|
FORK_RETRY_LIMIT,
|
|
LINEAGE_KEYS,
|
|
LINEAGE_FETCH_CLAIM,
|
|
newSnapshotId,
|
|
__internals: { numOrNull, intOrNull, boolOrNull },
|
|
};
|