MLB canonical event identity, impossible-binding refusal, event-aware dedupe, publication commit

Release-isolated slice built from 41ba38e. Ships ONLY the event-integrity +
publication + lineage-canary closure; the 90-path development tree stays
undeployed.

- canonical MLB event identity from statsapi gamePk (mlb:gamepk:<pk>), with
  event_identity_source/method/version recorded. The id is canonical; the
  binding is derived and says so.
- IMPOSSIBLE-BINDING REFUSAL. Verified in prod 2026-08-26: Joe Mack (Marlins)
  was bound to Dodgers@Braves and Yandy Diaz (Rays) to Rangers@WhiteSox, both
  from one book in the 01:00/03:01 UTC cycles after their own games began. Root
  cause is source market data, not the binder. A prop whose player's team is not
  an event participant now refuses; unknown team preserves uncertainty.
- event-aware dedupe: books still collapse, events no longer do. An unresolved
  MLB event fractures rather than falling back to the collision-prone
  date+teams key.
- publication commit moved AFTER the authoritative Redis slate write, with
  exact parity-gap identity when the product publishes and the record does not.
- lineage dual-write behind LINEAGE_CANARY_SPORTS, DISABLED for this deploy.

Excluded deliberately: WNBA feed/chain, market ontology, PerformanceDistribution,
calibration certification, truth diagnostics, applyRevision Phase-1, and the
analyzeViaEngine1 confidence-rounding change (a served field).

Suite 380/5,040/0 from this worktree; web tsc exit 0; champion output identical
to production.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01CQJeAG8vcDoL5zkiaJyVb8
This commit is contained in:
Kev
2026-08-27 16:27:36 -04:00
parent 41ba38ed4e
commit 352016790a
16 changed files with 4217 additions and 6 deletions
+509
View File
@@ -0,0 +1,509 @@
'use strict';
/**
* gradeFreeze — is the served forecast genuinely untouched by the shadows?
*
* ── THE INVARIANT ────────────────────────────────────────────────────────
* Shadow layers (A5 `factor_inputs`, A6 `would_fire`, chain-v1 `chain_shadow`)
* exist so a challenger can be measured against the served number without
* becoming it. Every one of them is documented as "reaches nothing" — and every
* one of those statements is currently the generator's own account of its own
* code. A5 is the reason that is not good enough: three factors were built,
* correct, wired, described as live for two sessions, and fired on 0 of 596
* rows. Nobody was lying; nobody had measured.
*
* So this module re-derives the invariant from what is actually persisted and
* from what the real code actually emits, by four independent arms:
*
* 1. DIFFERENTIAL perturb the shadow's VALUES, re-run the real serving
* code, and require the served output to be byte-
* identical. Perturbation, not absence: switching a
* shadow OFF only proves its presence is harmless;
* changing what it says proves its values cannot travel.
* 2. STRUCTURAL the served letter and confidence must be exactly the
* functions of the served p_win that they are declared
* to be. A leak that moved the letter without moving the
* number is caught here and nowhere else.
* 3. CO-DETERMINATION rows with an IDENTICAL champion feature vector but a
* DIFFERENT shadow value must carry an identical p_win.
* This is the arm that cannot be satisfied by agreement-
* with-itself: it needs the shadow to actually vary.
* 4. CONTAINMENT every stored shadow block declares itself unservable,
* and no served column carries a shadow-only key.
*
* ── WHAT IS DELIBERATELY *NOT* AN ARM ────────────────────────────────────
* `chain_shadow.counter_p` equals the served `p_win` on every row. That looks
* like a perfect freeze check and it is worthless: `mergeChainShadow` COPIES
* `p_win` into `counter_p`, so the equality is an identity, not evidence. It is
* named here so nobody adds it later mistaking a tautology for a test.
*
* PURE. Takes rows, returns verdicts. No I/O, no database, no clock.
*/
const { VERDICT, arm } = require('./verdict');
/**
* ── FROZEN ANCHOR: the served-grade bands, restated independently ────────
* A deliberate second copy of `servedGrade.BANDS`. Re-deriving the letter from
* the module that produced it would be asking the generator whether it agrees
* with itself. The meta-guard asserts this copy still matches the live table —
* so a band change fails a test and has to be re-affirmed by a human, instead
* of silently redefining what "correct" means for every historical row.
*/
const BANDS = Object.freeze([
Object.freeze({ letter: 'B+', min: 0.780 }),
Object.freeze({ letter: 'B', min: 0.700 }),
Object.freeze({ letter: 'C+', min: 0.640 }),
Object.freeze({ letter: 'C', min: 0.560 }),
Object.freeze({ letter: 'C-', min: 0.480 }),
Object.freeze({ letter: 'D', min: 0.350 }),
Object.freeze({ letter: 'F', min: 0.000 }),
]);
/** Letters the served scale declares structurally unissuable. */
const UNISSUABLE = Object.freeze(['A+', 'A', 'A-']);
/**
* ── FROZEN ANCHOR: tolerances ────────────────────────────────────────────
* Every number a gate is allowed to forgive, in one place, so loosening one is
* a visible edit to a frozen constant rather than a `+ 0.001` buried in a
* comparison. The meta-guard asserts these EXACT values.
*
* MAX_* are counts of violations tolerated, and they are ZERO. A freeze that
* tolerates a leak is not a freeze. The only non-zero number here is the
* minimum evidence a PASS requires.
*/
const TOLERANCES = Object.freeze({
/** Served letter vs the band table: no mismatch is acceptable. */
MAX_LETTER_MISMATCH: 0,
/** Served confidence vs the exact-decimal re-derivation: no mismatch is
* acceptable. Pre-repair rows on the four known half-way values are counted
* separately (see CONFIDENCE_REPAIR), never forgiven as a tolerance here. */
MAX_CONFIDENCE_MISMATCH: 0,
/** Structurally unissuable letters actually emitted. */
MAX_UNISSUABLE_EMITTED: 0,
/** Co-determination groups where the shadow moved and p_win moved with it. */
MAX_CODETERMINATION_VIOLATIONS: 0,
/** Served fields that differ when the shadow's values are perturbed. */
MAX_DIFFERENTIAL_DIVERGENCES: 0,
/** Stored shadow blocks failing to declare themselves unservable. */
MAX_UNDECLARED_SHADOW_BLOCKS: 0,
/**
* A structural arm needs rows to be evidence of anything. Below this the arm
* reports INSUFFICIENT — not PASS. "We checked 3 rows and found nothing" is
* not the same claim as "this holds".
*/
MIN_STRUCTURAL_ROWS: 100,
/**
* The co-determination arm needs the shadow to have ACTUALLY VARIED while the
* champion inputs held still. One such group is a real observation; the arm
* reports how thin it is rather than rounding thinness up to confidence.
*/
MIN_CODETERMINATION_GROUPS: 1,
/**
* The differential arm needs the perturbation to have actually changed the
* shadow. If perturbing produced an identical shadow, the arm proved nothing
* about propagation and must say so.
*/
MIN_PERTURBED_ROWS: 1,
});
/**
* The fields a user, the ledger, or the public record actually consumes.
* A leak matters exactly insofar as it reaches one of these.
*/
const SERVED_FIELDS = Object.freeze([
'p_win', 'grade', 'engine_grade', 'confidence', 'confidence_basis',
'ev_pct', 'projection', 'edge_pct', 'takeable', 'value',
'kelly', 'model_odds', 'fair_odds', 'fair_prob', 'alt_lines',
'refused', 'insufficient_data', 'refusal_reason',
]);
/**
* Keys that exist only to carry shadow evidence. None may appear as, or feed,
* a served field. Listed explicitly so adding a new shadow column is a
* conscious act with a test attached.
*/
const SHADOW_KEYS = Object.freeze([
'factor_inputs', 'chain_shadow', 'would_fire', 'shadow_inputs',
'p_win_isotonic_shadow', 'p_win_calibrated', 'p_win_lowparam',
'chain_p', 'chain_p_over', 'chain_projected_value',
]);
const num = (v) => {
if (v === null || v === undefined || v === '') return null;
const n = Number(v);
return Number.isFinite(n) ? n : null;
};
/** The letter the served number implies, re-derived from the frozen band table. */
function expectedLetter(pWin) {
const p = num(pWin);
if (p === null) return null;
const band = BANDS.find((b) => p >= b.min);
return band ? band.letter : BANDS[BANDS.length - 1].letter;
}
/**
* The confidence the served number implies.
*
* ── THE VERIFIER MIRRORS THE PRODUCER, AND MOVED WITH IT ─────────────────
* This originally reproduced the producer's IEEE-754 `Math.round(p_win * 100)`
* on purpose: re-deriving in exact decimal reported 465 real production rows as
* violations of an invariant they satisfied, because the checker was using
* BETTER arithmetic than the thing it checked. Mirroring is the rule.
*
* On 2026-08-17 the PRODUCER was corrected to exact-decimal round-half-up
* (analyzeViaEngine1.js — 0.565 now yields 57, as a reader intends), so this
* moved in the same commit. If only one of the two had changed, the gate would
* have reported the corrected rows as a leak — which is precisely the failure
* mode this pairing exists to avoid.
*
* The implementation is deliberately a SECOND, INDEPENDENT copy rather than an
* import from the engine: asking the producer to confirm the producer proves
* nothing. A unit test cross-checks the two agree on every boundary case.
*/
function expectedConfidence(pWin) {
const p = num(pWin);
if (p === null) return null;
// p_win is quantised to three decimals by the producer, so its exact value is
// an integer number of thousandths and the rounding becomes integer maths.
const thousandths = Math.round(p * 1000);
return Math.floor((thousandths + 5) / 10);
}
/**
* ── FROZEN ANCHOR: the rounding-repair allowance ─────────────────────────
* Rows graded BEFORE the exact-decimal repair carry the float value and are not
* wrong — they were correct under the arithmetic that produced them. Rows
* graded after carry the repaired value. Both are legitimate, so the check
* accepts either, and ONLY on the four half-way values where the two can
* differ at all. Every other deviation still fails.
*
* This is not a blanket +/-1 tolerance: outside these four p_win values a
* one-off confidence is a real violation and blocks. The allowance is retired
* once no pre-repair rows remain in the scored window, which is a deliberate
* later edit, not a decay.
*/
const CONFIDENCE_REPAIR = Object.freeze({
repaired_at: '2026-08-17',
producer: 'src/services/intelligence/analyzeViaEngine1.js',
reason: 'Math.round(p_win*100) rounded the IEEE-754 artifact (0.565*100 = 56.49999999999999) '
+ 'rather than the value; replaced with exact-decimal round-half-up over integer thousandths',
measured: Object.freeze({
rows_affected: 488,
letters_changed: 0,
p_win_values: Object.freeze([0.145, 0.285, 0.565, 0.575]),
}),
});
/** Is this the known pre-repair value for one of the four half-way cases? */
function isPreRepairConfidence(pWin, stored) {
const p = num(pWin);
const c = num(stored);
if (p === null || c === null) return false;
if (!CONFIDENCE_REPAIR.measured.p_win_values.includes(p)) return false;
// The float result is exactly one below the exact-decimal result on these.
return c === expectedConfidence(p) - 1;
}
/**
* ARM 2 — STRUCTURAL IDENTITY.
*
* The served letter and confidence are declared to be pure functions of the
* served p_win. If a shadow ever reaches the grade surface without passing
* through p_win, this is where it shows up.
*
* @param {Array} rows { p_win, grade, confidence, confidence_basis, refused }
*/
function structuralIdentity(rows = [], opts = {}) {
const scoped = rows.filter((r) => num(r.p_win) !== null && !r.refused && r.grade);
const letterMismatch = [];
const confMismatch = [];
const unissuable = [];
// Rows that fall inside a declared, dated deploy boundary. Counted and
// reported, never silently dropped — see `opts.boundary` below.
let boundaryRows = 0;
const boundaryMismatch = [];
// Rows still carrying the pre-repair float confidence — legitimate, counted.
let preRepairConfidence = 0;
const onBoundary = (r) => {
if (!opts.boundary) return false;
const stamp = String(r[opts.boundary.field] || '');
return stamp.slice(0, 10) === opts.boundary.date;
};
for (const r of scoped) {
const inBoundary = onBoundary(r);
if (inBoundary) boundaryRows += 1;
const exp = expectedLetter(r.p_win);
if (exp !== r.grade) {
const rec = { id: r.id, p_win: num(r.p_win), grade: r.grade, expected: exp };
if (inBoundary) boundaryMismatch.push(rec); else letterMismatch.push(rec);
}
if (UNISSUABLE.includes(r.grade)) {
const rec = { id: r.id, grade: r.grade, p_win: num(r.p_win) };
if (inBoundary) boundaryMismatch.push(rec); else unissuable.push(rec);
}
if (r.confidence_basis === 'p_win') {
const expC = expectedConfidence(r.p_win);
if (num(r.confidence) !== expC) {
// A pre-repair row carrying the float value on one of the four known
// half-way p_win values is counted separately, not as a violation.
if (isPreRepairConfidence(r.p_win, r.confidence)) preRepairConfidence += 1;
else confMismatch.push({ id: r.id, p_win: num(r.p_win), confidence: num(r.confidence), expected: expC });
}
}
}
const numbers = {
rows_scoped: scoped.length,
letter_mismatches: letterMismatch.length,
confidence_mismatches: confMismatch.length,
unissuable_emitted: unissuable.length,
boundary_date: opts.boundary ? opts.boundary.date : null,
boundary_rows: boundaryRows,
boundary_anomalies: boundaryMismatch.length,
pre_repair_confidence_rows: preRepairConfidence,
};
// A BOUNDARY THAT GREW IS NOT A BOUNDARY. The declared exception carries the
// exact count measured when it was declared; anything above it is a new
// violation wearing an old excuse, and it blocks.
if (opts.boundary && boundaryMismatch.length > opts.boundary.max_anomalies) {
return arm('grade_freeze.structural', VERDICT.FAIL,
`the declared ${opts.boundary.date} deploy boundary holds ${boundaryMismatch.length} anomalies, `
+ `above the ${opts.boundary.max_anomalies} recorded when it was declared — this is new, not historical`,
numbers, { samples: { boundary: boundaryMismatch.slice(0, 5) } });
}
if (scoped.length < TOLERANCES.MIN_STRUCTURAL_ROWS) {
return arm('grade_freeze.structural', VERDICT.INSUFFICIENT,
`only ${scoped.length} graded rows in scope (need ${TOLERANCES.MIN_STRUCTURAL_ROWS}) — too few to call this held`,
numbers, { samples: { letterMismatch: letterMismatch.slice(0, 5) } });
}
const violations = letterMismatch.length + confMismatch.length + unissuable.length;
if (violations > 0) {
return arm('grade_freeze.structural', VERDICT.FAIL,
`${letterMismatch.length} letters, ${confMismatch.length} confidences and ${unissuable.length} unissuable grades do not follow from the served p_win`,
numbers, {
samples: {
letterMismatch: letterMismatch.slice(0, 5),
confMismatch: confMismatch.slice(0, 5),
unissuable: unissuable.slice(0, 5),
},
});
}
return arm('grade_freeze.structural', VERDICT.PASS,
`${scoped.length - boundaryRows} graded rows outside the declared boundary: every letter and confidence `
+ `follows exactly from the served p_win`
+ (boundaryRows ? ` (${boundaryRows} rows on the declared ${opts.boundary.date} boundary carry ${boundaryMismatch.length} known anomalies)` : ''),
numbers, { ...(opts.extra || {}) });
}
/**
* ARM 3 — CO-DETERMINATION.
*
* Group rows by the CHAMPION's inputs (the identity plus a hash of the feature
* vector). Within a group the champion cannot produce two answers. So if the
* shadow value varies inside a group and `p_win` also varies, the served number
* moved with something other than the champion — which is a leak.
*
* This arm is the one that cannot be passed by construction. Its power comes
* entirely from groups where the shadow genuinely varied; with none, it reports
* INSUFFICIENT and says how many groups it had.
*
* @param {Array} rows { fingerprint, shadow_value, p_win }
*/
function coDetermination(rows = []) {
const groups = new Map();
for (const r of rows) {
const fp = r.fingerprint;
const sv = num(r.shadow_value);
const pw = num(r.p_win);
if (!fp || sv === null || pw === null) continue;
if (!groups.has(fp)) groups.set(fp, { shadow: new Set(), pwin: new Set(), n: 0 });
const g = groups.get(fp);
g.shadow.add(sv);
g.pwin.add(pw);
g.n += 1;
}
const discriminating = [];
const violations = [];
for (const [fp, g] of groups) {
if (g.shadow.size < 2) continue;
const rec = {
fingerprint: fp,
rows: g.n,
distinct_shadow: g.shadow.size,
distinct_p_win: g.pwin.size,
shadow_spread: Math.max(...g.shadow) - Math.min(...g.shadow),
};
discriminating.push(rec);
if (g.pwin.size > 1) violations.push(rec);
}
const numbers = {
rows: rows.length,
groups: groups.size,
discriminating_groups: discriminating.length,
violations: violations.length,
max_shadow_spread: discriminating.length
? Math.max(...discriminating.map((d) => d.shadow_spread)) : null,
};
if (violations.length > TOLERANCES.MAX_CODETERMINATION_VIOLATIONS) {
return arm('grade_freeze.codetermination', VERDICT.FAIL,
`${violations.length} group(s) where the champion inputs were identical, the shadow moved, and the served p_win moved with it`,
numbers, { samples: violations.slice(0, 5) });
}
if (discriminating.length < TOLERANCES.MIN_CODETERMINATION_GROUPS) {
return arm('grade_freeze.codetermination', VERDICT.INSUFFICIENT,
`${groups.size} groups examined but the shadow never varied inside one — nothing here can distinguish a frozen grade from a leaking one`,
numbers);
}
return arm('grade_freeze.codetermination', VERDICT.PASS,
`${discriminating.length} group(s) where the shadow moved (spread up to ${numbers.max_shadow_spread}) on identical champion inputs: served p_win held constant in every one`,
numbers, { samples: discriminating.slice(0, 5) });
}
/**
* ARM 4 — CONTAINMENT.
*
* A stored shadow block must declare itself unservable IN THE PAYLOAD, and no
* served field may carry a shadow-only key. The declaration matters because the
* data outlives the comment: a later query joining these columns has only what
* the row says about itself.
*
* @param {Array} blocks { id, block, served } block = the stored shadow jsonb
*/
/**
* ── WHAT BLOCKS, AND WHAT IS ONLY REPORTED ───────────────────────────────
* The first live run of this arm found 2,043 A6 `would_fire` blocks carrying a
* multiplier with no `servable:false` — because the two shadow layers were
* built in different sessions and only `chainShadow` adopted the self-declaring
* convention. That is a real gap and it is reported with its exact count.
*
* It does not block, and the distinction is substantive rather than
* convenient: a stored PROBABILITY is the thing that can be mistaken for a
* forecast when someone queries the column a year from now, which is the harm
* the declaration exists to prevent. A multiplier is not a forecast — it cannot
* be read as one — and its own inputs are stored beside it, so it stays
* re-derivable either way. Blocking on it would fail every build from now until
* a writer changes, which is how a real gate gets switched off.
*/
const PROBABILITY_KEYS = Object.freeze(['chain_p', 'chain_p_over', 'p_calibrated']);
const MULTIPLIER_KEYS = Object.freeze(['multiplier']);
function containment(blocks = []) {
const undeclared = [];
const undeclaredMultiplier = [];
const contaminated = [];
for (const b of blocks) {
const block = b.block || {};
const carriesProbability = PROBABILITY_KEYS.some((k) => num(block[k]) !== null);
const carriesMultiplier = MULTIPLIER_KEYS.some((k) => num(block[k]) !== null);
if (carriesProbability && block.servable !== false) {
undeclared.push({ id: b.id, servable: block.servable ?? null, keys: Object.keys(block).slice(0, 8) });
} else if (carriesMultiplier && block.servable !== false) {
undeclaredMultiplier.push({ id: b.id, keys: Object.keys(block).slice(0, 8) });
}
const served = b.served || {};
for (const k of SHADOW_KEYS) {
if (Object.prototype.hasOwnProperty.call(served, k)) {
contaminated.push({ id: b.id, key: k });
}
}
}
const numbers = {
blocks: blocks.length,
undeclared_shadow_blocks: undeclared.length,
served_fields_carrying_shadow_keys: contaminated.length,
// REPORTED, NOT BLOCKING — the standing gap named above.
undeclared_multiplier_blocks: undeclaredMultiplier.length,
};
if (!blocks.length) {
return arm('grade_freeze.containment', VERDICT.INSUFFICIENT,
'no shadow blocks present to inspect — the invariant is not violated, but it is also not demonstrated', numbers);
}
if (undeclared.length > TOLERANCES.MAX_UNDECLARED_SHADOW_BLOCKS || contaminated.length > 0) {
return arm('grade_freeze.containment', VERDICT.FAIL,
`${undeclared.length} shadow block(s) carry a probability without declaring servable:false; `
+ `${contaminated.length} served field(s) carry a shadow-only key`,
numbers, { samples: { undeclared: undeclared.slice(0, 5), contaminated: contaminated.slice(0, 5) } });
}
return arm('grade_freeze.containment', VERDICT.PASS,
`${blocks.length} shadow block(s): none carries an undeclared probability and no served field carries a shadow key`
+ (undeclaredMultiplier.length
? `. STANDING FINDING: ${undeclaredMultiplier.length} A6 would_fire block(s) carry a multiplier without the servable:false declaration chainShadow adopted — reported, not blocking`
: ''),
numbers, undeclaredMultiplier.length ? { samples: { undeclared_multiplier: undeclaredMultiplier.slice(0, 3) } } : {});
}
/**
* ARM 1 — DIFFERENTIAL.
*
* Compare the served output of two runs of the REAL serving code that differ
* ONLY in what the shadow said. Any difference is a leak, by definition.
*
* @param {Array} pairs { id, base, perturbed, shadow_changed }
* base/perturbed are the two served result objects.
*/
function differential(pairs = [], fields = SERVED_FIELDS) {
const divergences = [];
let perturbedCount = 0;
for (const p of pairs) {
if (p.shadow_changed) perturbedCount += 1;
for (const f of fields) {
const a = p.base ? p.base[f] : undefined;
const b = p.perturbed ? p.perturbed[f] : undefined;
const sa = JSON.stringify(a === undefined ? null : a);
const sb = JSON.stringify(b === undefined ? null : b);
if (sa !== sb) divergences.push({ id: p.id, field: f, base: a, perturbed: b });
}
}
const numbers = {
pairs: pairs.length,
fields_compared: fields.length,
shadow_actually_perturbed: perturbedCount,
divergences: divergences.length,
};
if (divergences.length > TOLERANCES.MAX_DIFFERENTIAL_DIVERGENCES) {
return arm('grade_freeze.differential', VERDICT.FAIL,
`${divergences.length} served field(s) changed when only the shadow's values changed — the shadow reaches the served path`,
numbers, { samples: divergences.slice(0, 8) });
}
if (perturbedCount < TOLERANCES.MIN_PERTURBED_ROWS) {
// The perturbation did not take. Reporting PASS here would be the exact
// failure this layer exists to catch, one level up: a check that ran, found
// nothing, and had nothing to find.
return arm('grade_freeze.differential', VERDICT.INSUFFICIENT,
`the perturbation changed the shadow on ${perturbedCount} of ${pairs.length} rows — with the shadow unchanged, identical output proves nothing`,
numbers);
}
return arm('grade_freeze.differential', VERDICT.PASS,
`${pairs.length} rows re-run with the shadow's values perturbed on ${perturbedCount} of them: all ${fields.length} served fields byte-identical`,
numbers);
}
module.exports = {
BANDS, UNISSUABLE, TOLERANCES, SERVED_FIELDS, SHADOW_KEYS,
expectedLetter, expectedConfidence, isPreRepairConfidence, CONFIDENCE_REPAIR,
structuralIdentity, coDetermination, containment, differential,
};
+101
View File
@@ -0,0 +1,101 @@
'use strict';
/**
* verdict — the evaluator's shared vocabulary.
*
* ── WHY THIS IS ITS OWN MODULE ────────────────────────────────────────────
* Every gate in this layer has to answer the same question the same way, and
* the single most dangerous thing a gate can do is turn "I could not check
* this" into "this is fine". That is the shape of the 2026-08-01 settlement
* outage (`{settled:0,pending:0}` was byte-identical to a healthy nothing-to-do)
* and of the `fielding_oaa: 0 rows` wiring bug (a failed feed degrading to an
* honest-looking absence). So INSUFFICIENT is a first-class verdict here, it is
* NEVER folded into PASS, and `combine()` refuses to report an overall PASS
* when any arm could not be measured.
*
* PURE. No I/O, no imports. The gate scripts and the unit tests share it.
*/
const VERDICT = Object.freeze({
/** Measured, and the invariant held. */
PASS: 'PASS',
/** Measured, and the invariant was violated. Blocking. */
FAIL: 'FAIL',
/** Not measurable right now (no DB, no rows, no discriminating sample).
* Not a pass. Blocking only where the caller says the arm is required. */
INSUFFICIENT: 'INSUFFICIENT',
/** The check itself broke. An error is an error — never an empty result. */
ERROR: 'ERROR',
});
/** A verdict that may not be treated as success under any circumstance. */
const NOT_PASSING = Object.freeze([VERDICT.FAIL, VERDICT.INSUFFICIENT, VERDICT.ERROR]);
/**
* Build one arm's result.
*
* `numbers` is mandatory in spirit: a verdict with no measured quantities
* behind it is an attestation, which is the thing this whole layer exists to
* replace. Callers pass what they counted.
*/
function arm(id, verdict, reason, numbers = {}, extra = {}) {
if (!Object.values(VERDICT).includes(verdict)) {
throw new Error(`unknown verdict '${verdict}' for arm '${id}' — add it explicitly rather than passing it through`);
}
return { id, verdict, reason, numbers, ...extra };
}
/**
* Combine arms into one blocking answer.
*
* @param {Array} arms
* @param {object} opts
* - requiredArms: ids that MUST reach PASS. An INSUFFICIENT on a required arm
* is a FAIL of the whole gate: a gate you cannot run is not a gate you
* passed. Arms not in this list may report INSUFFICIENT without blocking,
* and the summary says so out loud rather than rounding it up.
*/
function combine(arms, opts = {}) {
const required = new Set(opts.requiredArms || []);
const failed = arms.filter((a) => a.verdict === VERDICT.FAIL);
const errored = arms.filter((a) => a.verdict === VERDICT.ERROR);
const insufficient = arms.filter((a) => a.verdict === VERDICT.INSUFFICIENT);
const requiredNotPassing = arms.filter(
(a) => required.has(a.id) && a.verdict !== VERDICT.PASS,
);
let verdict = VERDICT.PASS;
let reason = `all ${arms.length} arms passed`;
if (errored.length) {
verdict = VERDICT.ERROR;
reason = `${errored.length} arm(s) errored: ${errored.map((a) => a.id).join(', ')}`;
} else if (failed.length) {
verdict = VERDICT.FAIL;
reason = `${failed.length} arm(s) FAILED: ${failed.map((a) => a.id).join(', ')}`;
} else if (requiredNotPassing.length) {
verdict = VERDICT.FAIL;
reason = `required arm(s) did not pass: ${requiredNotPassing.map((a) => `${a.id}=${a.verdict}`).join(', ')}`;
} else if (insufficient.length) {
// Not a pass, not a block — stated as itself so a green line never means
// "we checked" when what happened was "we could not check".
verdict = VERDICT.INSUFFICIENT;
reason = `${arms.length - insufficient.length} passed; ${insufficient.length} could not be measured: ${insufficient.map((a) => a.id).join(', ')}`;
}
return {
verdict,
reason,
blocking: verdict === VERDICT.FAIL || verdict === VERDICT.ERROR,
counts: {
arms: arms.length,
pass: arms.filter((a) => a.verdict === VERDICT.PASS).length,
fail: failed.length,
insufficient: insufficient.length,
error: errored.length,
},
arms,
};
}
module.exports = { VERDICT, NOT_PASSING, arm, combine };
+331
View File
@@ -0,0 +1,331 @@
'use strict';
/**
* CANONICAL EVENT IDENTITY — a sport-neutral contract, an MLB vertical slice.
*
* -- THE DEFECT THIS CLOSES ----------------------------------------------
* `ledgerService.gameIdFor` derives an event label as `sport:date:away@home`.
* That is a LABEL, not an identity: it carries no occurrence component, so the
* two halves of a doubleheader produce a BYTE-IDENTICAL string.
*
* VERIFIED against the authoritative source, MLB 2026 through 2026-08-27:
* 19 doubleheaders / 38 real games collapse into 19 derived labels.
* Real case 2026-08-17, St. Louis Cardinals @ Cincinnati Reds:
* gamePk 824514 (gameNumber 1, 17:40Z) and 824478 (gameNumber 2, 22:40Z).
* VYNDR recorded ONE id, `mlb:2026-08-17:St.LouisCardinals@CincinnatiReds`,
* holding 171 ledger rows and FOUR distinct starting pitchers (Pallante,
* Mathews, Lowder, Emanuel) — i.e. both games merged into one event.
*
* -- WHY gamePk AND NOT THE ESPN EVENT ID --------------------------------
* `gameBinder` binds through `scheduleService`, which is ESPN. ESPN ids are
* also per-game, but statsapi is the AUTHORITATIVE MLB source, is already
* fetched inside the snapshot pipeline (`snapshotService` calls
* `mlbStatsAdapter.getScheduleWithPitchers`), is free and unlimited, and
* carries `gameNumber` / `doubleHeader` alongside `gamePk` so the occurrence is
* self-describing. Using the source of record avoids a second identity space.
*
* -- SPORT-NEUTRAL CONTRACT, ONE IMPLEMENTED SPORT ------------------------
* The FIELD is generic. Only MLB has a verified resolver. Every other sport
* returns UNSUPPORTED rather than a guessed identifier — an invented id would
* be worse than none, because downstream code would trust it.
*/
const IDENTITY_VERSION = 'evid@1';
/** Where an identity came from. */
const IDENTITY_SOURCE = Object.freeze({
MLB_STATSAPI_GAMEPK: 'MLB_STATSAPI_GAMEPK',
LEGACY_DERIVED: 'LEGACY_DERIVED',
UNKNOWN: 'UNKNOWN',
});
/** How strongly it identifies the event. */
const IDENTITY_METHOD = Object.freeze({
CANONICAL: 'CANONICAL',
LEGACY_DERIVED: 'LEGACY_DERIVED',
UNRESOLVED: 'UNRESOLVED',
UNSUPPORTED_SPORT: 'UNSUPPORTED_SPORT',
});
const SUPPORTED = Object.freeze(['mlb']);
/** Namespaced so two sports' numeric ids can never collide. */
function canonicalEventId(sport, nativeId) {
if (!sport || nativeId === null || nativeId === undefined || nativeId === '') return null;
return `${String(sport).toLowerCase()}:gamepk:${nativeId}`;
}
const refusal = (method, reason) => Object.freeze({
canonical_event_id: null,
event_identity_source: IDENTITY_SOURCE.UNKNOWN,
event_identity_method: method,
event_identity_version: IDENTITY_VERSION,
event_occurrence: null,
reason: reason || null,
});
const norm = (v) => String(v || '').toLowerCase().replace(/[^a-z]/g, '');
/**
* TEAM RESOLUTION — canonical, not fuzzy.
*
* Real prop feeds spell one team several ways: `Cincinnati Reds`,
* `CINReds`, `CIN`. A prefix test fails on all the abbreviated forms
* (`newyorkmets`.startsWith(`nymets`) is false), which is why the first draft
* resolved nothing for them.
*
* So resolution is anchored on `environmentContext.NAME_TO_ABBR` — the
* repository's existing canonical 31-team map — and every spelling is reduced
* to an ABBREVIATION before comparison:
* 1. the flattened full name `cincinnatireds` -> CIN
* 2. the canonical nickname as SUFFIX `cinreds` -> CIN
* 3. the abbreviation itself `cin` -> CIN
*
* Nicknames are matched LONGEST-FIRST because `whitesox` and `redsox` both end
* in `sox`; shortest-first would map every Chicago White Sox prop to Boston.
*
* This narrows candidates only. Choosing BETWEEN two candidates is never done
* by name — that is decided by start time, below.
*/
const { NAME_TO_ABBR } = require('../environmentContext');
const { nameKey } = require('../../utils/playerName');
const TEAM_INDEX = (() => {
const byFlatName = new Map();
const byAbbr = new Map();
const nicknames = [];
for (const [full, abbr] of Object.entries(NAME_TO_ABBR || {})) {
byFlatName.set(norm(full), abbr);
byAbbr.set(norm(abbr), abbr);
const words = String(full).trim().split(/\s+/);
// The nickname is everything after the city. Two-word nicknames
// (`red sox`, `white sox`, `blue jays`) need both trailing words.
for (const take of [2, 1]) {
if (words.length > take) nicknames.push({ nick: norm(words.slice(-take).join('')), abbr });
}
}
nicknames.sort((a, b) => b.nick.length - a.nick.length);
return { byFlatName, byAbbr, nicknames };
})();
/** Any spelling -> canonical abbreviation, or null. */
function teamAbbr(name) {
const f = norm(name);
if (!f) return null;
const exact = TEAM_INDEX.byFlatName.get(f);
if (exact) return exact;
const asAbbr = TEAM_INDEX.byAbbr.get(f);
if (asAbbr) return asAbbr;
for (const { nick, abbr } of TEAM_INDEX.nicknames) {
if (f.endsWith(nick)) return abbr;
}
return null;
}
/**
* Does a schedule game match a prop's team pair?
* Both sides must resolve to a canonical abbreviation and both must agree. An
* unresolvable name matches NOTHING rather than falling back to a text compare.
*/
function teamsMatch(game, homeTeam, awayTeam) {
const gh = teamAbbr(game && game.home && game.home.team);
const ga = teamAbbr(game && game.away && game.away.team);
const ph = teamAbbr(homeTeam);
const pa = teamAbbr(awayTeam);
if (!gh || !ga || !ph || !pa) return false;
return gh === ph && ga === pa;
}
/**
* -- IMPOSSIBLE-BINDING INVARIANT -----------------------------------------
* A player proposition must never be bound to an event his team is not in.
*
* VERIFIED IN PRODUCTION, mlb 2026-08-26. Two props were bound to real games
* their players were not playing in:
*
* Joe Mack (statsapi id 691788, MIAMI MARLINS)
* -> mlb:2026-08-26:LosAngelesDodgers@AtlantaBraves (pinnacle, 03:01 UTC)
* Yandy Diaz (statsapi id 650490, TAMPA BAY RAYS)
* -> mlb:2026-08-26:TexasRangers@ChicagoWhiteSox (pinnacle, 01:00/03:01 UTC)
*
* ROOT CAUSE IS SOURCE MARKET DATA, NOT THE BINDER. `gameMatchesTeams` is
* correctly guarded (empty team names return false), and the stored `game_id`
* is built from the PROP'S OWN home/away fields — so the provider delivered
* these players nested inside the wrong event. Both players' real games had
* already started (Marlins 22:40Z, Rays 17:10Z) while the events they were
* mis-nested into were still open (Braves 23:15Z, White Sox 23:40Z), which is
* why it only appears in the 01:00 and 03:01 UTC cycles and only from one book.
*
* The prop's own team fields are therefore NOT independent evidence — they are
* the wrong event's teams. The check needs the PLAYER'S team from a source that
* cannot be wrong in the same way, and refuses when they disagree.
*
* UNKNOWN TEAM IS NOT A VIOLATION. A player we cannot place is left alone:
* refusing on absent evidence would drop real props for a data gap.
*/
function checkParticipantEvidence(prop, game, playerTeams) {
if (!playerTeams || typeof playerTeams.get !== 'function') return null;
const key = nameKey(prop && (prop.player || prop.player_name));
if (!key) return null;
const teams = playerTeams.get(key);
if (!teams || teams.size === 0) return null; // unknown team -> not a violation
const home = teamAbbr(game && game.home && game.home.team);
const away = teamAbbr(game && game.away && game.away.team);
if (!home || !away) return null; // unreadable event -> not a violation
for (const t of teams) if (t === home || t === away) return null;
return { violated: true, player_teams: [...teams], event_teams: [away, home] };
}
/**
* Build nameKey -> Set(team abbreviation) for the teams playing on a slate.
* Injectable; production supplies the statsapi roster reader. A team whose
* roster cannot be read contributes nothing rather than an empty answer that
* would look like 'this player has no team'.
*/
async function buildPlayerTeamIndex(games, deps) {
const d = deps || {};
const getRoster = d.getTeamRoster;
const index = new Map();
const out = { teams: 0, players: 0, failed: 0 };
if (typeof getRoster !== 'function') return { index, stats: out };
const ids = new Map();
for (const g of Array.isArray(games) ? games : []) {
for (const side of ['home', 'away']) {
const t = g && g[side];
if (t && t.teamId != null) ids.set(t.teamId, t.team);
}
}
for (const [id, name] of ids) {
let roster = null;
try { roster = await getRoster(id); } catch { roster = null; }
if (!Array.isArray(roster)) { out.failed += 1; continue; }
out.teams += 1;
const abbr = teamAbbr(name);
if (!abbr) continue;
for (const p of roster) {
const k = nameKey(p && (p.name || p.fullName || p.player_name));
if (!k) continue;
if (!index.has(k)) index.set(k, new Set());
index.get(k).add(abbr);
out.players += 1;
}
}
return { index, stats: out };
}
/**
* Resolve ONE prop to a canonical MLB event.
*
* -- THE DISAMBIGUATOR IS TIME, AND IT IS REQUIRED -----------------------
* When teams match TWO games the pair alone cannot choose, and choosing wrong
* attributes a claim to the game that did not happen. `game_time` separates
* them (17:40Z vs 22:40Z in the verified case). Without a usable time the
* result is UNRESOLVED — a fracture, never a guess. A false fracture is
* recoverable; a false merge destroys the receipt.
*
* @param {object} prop needs home_team, away_team, optionally game_time
* @param {Array} games mlbStatsAdapter.getScheduleWithPitchers() output
*/
function resolveMlbEvent(prop, games, playerTeams) {
const p = prop || {};
const list = Array.isArray(games) ? games.filter((g) => g && g.gamePk != null) : [];
if (list.length === 0) return refusal(IDENTITY_METHOD.UNRESOLVED, 'no_schedule');
const matches = list.filter((g) => teamsMatch(g, p.home_team, p.away_team));
if (matches.length === 0) return refusal(IDENTITY_METHOD.UNRESOLVED, 'no_team_match');
// Refuse an event the player's own team is not in, BEFORE any other
// disambiguation. Time proximity must never override an impossible
// team/event relationship.
const possible = matches.filter((g) => !checkParticipantEvidence(prop, g, playerTeams));
if (possible.length === 0) {
const ev = checkParticipantEvidence(prop, matches[0], playerTeams);
return Object.freeze({
canonical_event_id: null,
event_identity_source: IDENTITY_SOURCE.UNKNOWN,
event_identity_method: IDENTITY_METHOD.UNRESOLVED,
event_identity_version: IDENTITY_VERSION,
event_occurrence: null,
reason: 'impossible_binding',
evidence: ev || null,
});
}
if (possible.length !== matches.length) matches.length = 0, matches.push(...possible);
let chosen = matches[0];
if (matches.length > 1) {
const t = Date.parse(p.game_time || '');
if (!Number.isFinite(t)) {
return refusal(IDENTITY_METHOD.UNRESOLVED, 'ambiguous_no_game_time');
}
let best = null; let bestDelta = Infinity;
for (const g of matches) {
const gt = Date.parse(g.gameDate || '');
if (!Number.isFinite(gt)) continue;
const d = Math.abs(gt - t);
if (d < bestDelta) { bestDelta = d; best = g; }
}
if (!best) return refusal(IDENTITY_METHOD.UNRESOLVED, 'ambiguous_no_schedule_time');
// A start time that matches neither game within an hour is not a match at
// all. Refusing beats attributing the claim to the nearer wrong game.
if (bestDelta > 60 * 60 * 1000) {
return refusal(IDENTITY_METHOD.UNRESOLVED, 'ambiguous_time_too_far');
}
chosen = best;
}
return Object.freeze({
canonical_event_id: canonicalEventId('mlb', chosen.gamePk),
event_identity_source: IDENTITY_SOURCE.MLB_STATSAPI_GAMEPK,
event_identity_method: IDENTITY_METHOD.CANONICAL,
event_identity_version: IDENTITY_VERSION,
// Self-describing occurrence, straight from the source. Never inferred.
event_occurrence: chosen.gameNumber ?? null,
reason: null,
});
}
/**
* Sport-neutral entry point. Only MLB resolves; everything else is
* UNSUPPORTED_SPORT, which is a truthful state and not a failure.
*/
function resolveEvent(sport, prop, games, playerTeams) {
const sp = String(sport || '').toLowerCase();
if (!SUPPORTED.includes(sp)) return refusal(IDENTITY_METHOD.UNSUPPORTED_SPORT, `sport_${sp || 'unknown'}`);
return resolveMlbEvent(prop, games, playerTeams);
}
/**
* Attach identity to a slate's props IN PLACE, returning counts so a silent
* resolution collapse is visible rather than quietly reintroducing collisions.
*/
function attachEventIdentity(sport, props, games, playerTeams) {
const list = Array.isArray(props) ? props : [];
const out = { total: list.length, canonical: 0, unresolved: 0, unsupported: 0, impossible: 0, reasons: {} };
for (const p of list) {
if (!p) continue;
const r = resolveEvent(sport, p, games, playerTeams);
// Declared on EVERY prop, resolved or not.
p.canonical_event_id = r.canonical_event_id;
p.event_identity_source = r.event_identity_source;
p.event_identity_method = r.event_identity_method;
p.event_identity_version = r.event_identity_version;
p.event_occurrence = r.event_occurrence;
if (r.event_identity_method === IDENTITY_METHOD.CANONICAL) out.canonical += 1;
else if (r.event_identity_method === IDENTITY_METHOD.UNSUPPORTED_SPORT) out.unsupported += 1;
else {
out.unresolved += 1;
if (r.reason === 'impossible_binding') out.impossible += 1;
if (r.reason) out.reasons[r.reason] = (out.reasons[r.reason] || 0) + 1;
}
}
return out;
}
module.exports = {
resolveEvent, resolveMlbEvent, attachEventIdentity, canonicalEventId, teamsMatch, teamAbbr,
checkParticipantEvidence, buildPlayerTeamIndex,
IDENTITY_SOURCE, IDENTITY_METHOD, IDENTITY_VERSION, SUPPORTED,
};
+89 -1
View File
@@ -77,13 +77,86 @@ const DEFAULT_TTL = 7200; // 2 hours — matches the spec's grades-cache TTL.
// and before the limit — which makes the graded set byte-identical to what it
// was before the widening. This gate lifts only when the MLB calibration is
// re-run on the consensus ruler and v2 is promoted.
/**
* WHAT DUPLICATES IS THIS FUNCTION SUPPOSED TO REMOVE?
*
* ONE: the same proposition offered by SEVERAL BOOKS. Provider is collapsed on
* purpose — grading the same player/stat/line once per book would multiply the
* work with no new information, and the price anchor is chosen later.
*
* WHAT IT WAS ALSO REMOVING, WRONGLY
*
* The key was `player::stat_type::line` with NO event component, so the SAME
* proposition in TWO DIFFERENT GAMES also collapsed — and first-row-wins
* silently discarded the second real game.
*
* VERIFIED on production: 2026-08-17 St. Louis @ Cincinnati was a doubleheader
* (gamePk 824514 / 824478). VYNDR recorded ONE derived id holding FOUR distinct
* starting pitchers, i.e. both games merged. 19 doubleheaders / 38 games are
* affected across the 2026 season to date.
*
* The key therefore adds the EVENT, and nothing else. Provider stays collapsed
* because collapsing books is the function's actual job.
*
* `canonical_event_id` is used when identity resolved. When it did not, the
* fallback is the legacy derived label, which reproduces the previous behaviour
* exactly — it does not pretend to distinguish what it cannot.
*/
let fractureSeq = 0;
/**
* The EVENT component of the dedupe key.
*
* -- FAIL CLOSED WHERE CANONICAL IDENTITY IS EXPECTED --------------------
* A first version fell back to `legacy:{date}:{away}@{home}` whenever identity
* did not resolve. That is precisely the collision-prone key this whole change
* exists to remove: on a doubleheader it is byte-identical for two real games,
* so an unresolved MLB prop could still merge two events.
*
* So for a sport where canonical identity is EXPECTED (MLB), an unresolved prop
* FRACTURES: it gets a key unique to that row, which cannot merge with anything.
* The cost is that books stop collapsing for that prop, so it may be graded more
* than once. That is a coverage/efficiency cost, and it is the right direction:
*
* UNKNOWN EVENT IDENTITY MAY LOSE COVERAGE.
* IT MAY NOT MERGE TWO REAL EVENTS.
*
* MEASURED: over MLB 2026-08-20..27, 101 of 101 team-pairs resolve from teams
* alone; only a doubleheader pair needs the start time to disambiguate. So the
* fracture path is rare in normal operation and bites only where identity is
* genuinely unknown.
*
* -- SPORTS WITHOUT A RESOLVER ARE UNCHANGED ----------------------------
* WNBA, NBA and the rest have no canonical adapter and never did. Fracturing
* them would multiply their grading load for no safety gain, because their
* identity was never canonical to begin with. They keep the legacy label, which
* is the exact behaviour that existed before this change, and their lineage is
* excluded from the canary for the same reason.
*/
function propositionEventKey(p) {
if (p && p.canonical_event_id) return p.canonical_event_id;
const method = p && p.event_identity_method;
const expectsCanonical = method && method !== 'UNSUPPORTED_SPORT';
if (expectsCanonical) {
// Unique per row: this prop can never share an event key with another.
fractureSeq += 1;
return `unresolved-event:${method}:${fractureSeq}`;
}
const away = String((p && p.away_team) || '').replace(/\s+/g, '');
const home = String((p && p.home_team) || '').replace(/\s+/g, '');
const date = (p && p.game_date) || '';
return `legacy:${date}:${away}@${home}`;
}
function dedupeProps(props, limit) {
const seen = new Set();
const out = [];
for (const p of props || []) {
if (!p || !p.player || !p.stat_type || p.line == null) continue;
if (!isModelBook(p.book)) continue;
const key = `${p.player}::${p.stat_type}::${p.line}`;
const key = `${propositionEventKey(p)}::${p.player}::${p.stat_type}::${p.line}`;
if (seen.has(key)) continue;
seen.add(key);
out.push(p);
@@ -128,6 +201,13 @@ async function gradeBestSide(grade, prop, sport, opts = {}) {
game_time: prop.game_time ?? null,
home_team: prop.home_team ?? null,
away_team: prop.away_team ?? null,
// CANONICAL EVENT IDENTITY rides with the grade so retention and lineage
// key on the real event rather than re-deriving a label from team names.
canonical_event_id: prop.canonical_event_id ?? null,
event_identity_source: prop.event_identity_source ?? null,
event_identity_method: prop.event_identity_method ?? null,
event_identity_version: prop.event_identity_version ?? null,
event_occurrence: prop.event_occurrence ?? null,
};
const sides = await Promise.all([
Promise.resolve()
@@ -155,6 +235,14 @@ async function gradeBestSide(grade, prop, sport, opts = {}) {
// Strip the internal retention fields so they never reach a cache or payload.
delete winner._features;
delete winner._grade_11;
// PUBLICATION SIGNAL. `onGraded` above fired with both sides before any
// filtering; this fires ONLY for the side that becomes the served Read, so
// retention can tell a published claim from a captured model state. Measured:
// 64.9% of captured rows describe a state no user ever saw.
if (typeof opts.onPublished === 'function') {
try { opts.onPublished(base, winner); } catch { /* never breaks the slate */ }
}
// CARRY THE GAME (2026-08-01). The legacy grade shape drops home/away, so by
// the time the challenger runs, nothing on the grade says WHICH GAME it is —
// measured: `team` was null on 416/416 stored grades, so the park/weather
+11 -2
View File
@@ -244,7 +244,7 @@ function archetypeVectorOf(g) {
return null;
}
function rowsFromSnapshot(sport, grades, oddsProps, nowIso) {
function rowsFromSnapshot(sport, grades, oddsProps, nowIso, lineageIndex) {
const sp = String(sport || '').toLowerCase();
let skippedUnbound = 0;
// TWO indexes, two roles — see indexProps. `byKey` supplies GAME FACTS from
@@ -375,6 +375,15 @@ function rowsFromSnapshot(sport, grades, oddsProps, nowIso) {
ruler_version: CURRENT_RULER_VERSION,
game_id: gameIdFor(sp, prop, gameDate),
game_date: gameDate,
// SHADOW LINEAGE LINKAGE. Which logical Read this receipt belongs to.
// Declared on EVERY row (PostgREST builds the bulk insert from the first
// row's shape), null when lineage did not resolve — unknown, never
// guessed. Read by nothing; settlement is unaffected.
read_id: lineageIndex
? (lineageIndex.get([sp, gameDate, nameKey(player), stat, side,
Number(line).toFixed(6)].join('|')) || null)
: null,
read_revision_id: null,
});
}
if (skippedUnbound > 0) {
@@ -392,7 +401,7 @@ async function recordPipelineGrades(sport, grades, oddsProps, opts = {}) {
if (!opts.sb && !isConfigured()) return { skipped: 'supabase not configured', written: 0 };
const sb = opts.sb || defaultClient();
const nowIso = (opts.now || (() => new Date().toISOString()))();
const rows = rowsFromSnapshot(sport, grades, oddsProps, nowIso);
const rows = rowsFromSnapshot(sport, grades, oddsProps, nowIso, opts.lineageIndex || null);
if (rows.length === 0) return { written: 0 };
let written = 0;
for (let i = 0; i < rows.length; i += UPSERT_CHUNK) {
+602
View File
@@ -0,0 +1,602 @@
'use strict';
/**
* READ LINEAGE — the append-only chronology of what VYNDR published.
*
* ── THE FINDING THIS IS BUILT ON ─────────────────────────────────────────
* VYNDR already stores immutable published claims. `model_snapshots` is
* upserted with `ignoreDuplicates: true` on
* (snapshot_id, player_key, stat, line, side), so a row is written once per
* snapshot cycle and NEVER overwritten; `captured_at` is NOT NULL, and
* `model_version` / `code_sha` / `ruler_version` record provenance. The only
* post-insert write is settlement (outcome, actual_value, settled_at), which is
* not part of the claim.
*
* MEASURED (mlb, 2026-08-26): 5,376 props, 4,286 with more than one row, avg
* 2.42 rows per prop, and 183 props whose GRADE CHANGED across them.
*
* So the immutable revisions exist. What was missing is LINEAGE:
* - a stable logical Read identity linking a prop's rows across cycles,
* - an explicit supersession pointer,
* - a way for a receipt to name the exact revision it settles.
*
* This module is that linkage. It does NOT create a second claim store, and it
* does not create a grade-history table — the claim already lives on the row.
*
* ── CAPTURE IS NOT REVISION, AND CONFLATING THEM WOULD FABRICATE HISTORY ─
* Most cycles re-capture an unchanged claim. On the measured date only 183 of
* 5,376 props actually changed. If every capture became a "revision", the
* chronology would report ~13,000 published claims where roughly 183 occurred.
* A row is a REVISION only when its claim DIGEST differs from the standing one;
* otherwise it is a RECAPTURE of the same published state and is labelled so.
*
* ── IDENTITY MUST SURVIVE UNSTABLE game_id SPELLING ──────────────────────
* The same real game appears as `mlb:2026-08-22:DetroitTigers@KansasCityRoyals`
* from one feed and as an abbreviated spelling from another. Logical Read
* identity therefore does NOT include the raw game_id.
*
* MEASURED over 30,746 identity groups: 1,328 carried more than one game_id
* spelling, and only 18 carried genuinely different team pairs (real distinct
* events — e.g. `jac caglianone` on 2026-08-22 appearing under both
* DetroitTigers@KansasCityRoyals and MinnesotaTwins@SanDiegoPadres). The prefix
* rule below reconciles 1,310 of the 1,328 and isolates exactly those 18, which
* MUST stay separate Reads.
*/
const crypto = require('crypto');
const { SERVED_FIELDS } = require('../evaluator/gradeFreeze');
const LINEAGE_VERSION = 'lin@1';
/**
* ── CLAIM SCHEMA VERSIONING ──────────────────────────────────────────────
* A digest is only interpretable against the claim definition that produced
* it. `SERVED_FIELDS` will change; when it does, digests computed afterwards
* mean something different from digests computed before, and nothing in the
* bytes says so.
*
* This bit us already: the first digest listed only `locked_odds` while
* `model_snapshots` stores `over_odds`/`under_odds`, so it was blind to price
* on exactly the rows it digested. That was caught by a replay crashing on a
* missing column — i.e. by luck. A version stamp is what makes the next such
* change visible instead of silent.
*
* Both are stored on every row. A historical digest is NEVER recomputed under a
* newer definition and claimed to be original.
*/
const CLAIM_SCHEMA_VERSION = 'claim@1';
const DIGEST_ALGORITHM_VERSION = 'sha256-json-sorted@1';
/** What a row is, relative to the chronology. */
const LINEAGE_ACTION = Object.freeze({
ORIGIN: 'ORIGIN', // the first published state of a logical Read
REVISION: 'REVISION', // a materially different claim; supersedes the prior
RECAPTURE: 'RECAPTURE', // same claim, observed again — NOT a new revision
FORK_DETECTED: 'FORK_DETECTED', // two writers superseded the same parent
});
/** Honest states for rows whose provenance predates lineage. */
const LINEAGE_STATE = Object.freeze({
LIVE: 'LIVE',
LEGACY_UNVERIFIED: 'LEGACY_UNVERIFIED',
});
/**
* Normalise a derived game_id into a comparable team pair.
* Returns null when the string is not in the derived shape — a null is honest;
* a guessed pair would silently merge two events.
*/
function eventFingerprint(gameId) {
const s = String(gameId || '');
const tail = s.split(':').pop();
if (!tail || !tail.includes('@')) return null;
const [awayRaw, homeRaw] = tail.split('@');
const norm = (v) => String(v || '').toLowerCase().replace(/[^a-z]/g, '');
const away = norm(awayRaw);
const home = norm(homeRaw);
if (!away || !home) return null;
return { away, home };
}
/** One side matches if either spelling is a prefix of the other. */
function sideMatches(a, b) {
if (!a || !b) return false;
return a === b || a.startsWith(b) || b.startsWith(a);
}
/**
* Are two game_id spellings the same real event?
* Unknown fingerprints are NOT assumed equal — an unparseable id cannot be
* proven to match anything, and merging on a guess is the failure this exists
* to prevent.
*/
function sameEvent(gameIdA, gameIdB) {
if (gameIdA === gameIdB) return true;
const a = eventFingerprint(gameIdA);
const b = eventFingerprint(gameIdB);
if (!a || !b) return false;
return sideMatches(a.away, b.away) && sideMatches(a.home, b.home);
}
const norm = (v) => (v === undefined ? null : v);
/**
* The digest of the MATERIAL CLAIM.
*
* Built from `gradeFreeze.SERVED_FIELDS` — the repository's existing canonical
* definition of what a user, the ledger and the public record actually consume.
* Reusing it means the claim the lineage tracks cannot drift from the claim the
* evaluator polices; a new served field joins both at once.
*
* The market terms the claim is expressed against (line, side, book, odds) are
* included because changing any of them changes what the user was told.
*/
/**
* Both odds spellings are listed because the two stores name price differently:
* `ledger_entries` has `locked_odds`, `model_snapshots` has `over_odds` /
* `under_odds`. Listing only one made the digest BLIND TO PRICE on exactly the
* rows it digests — a price move would have read as an unchanged claim. A row
* supplies whichever fields it has; the rest are null on both sides of any
* comparison, so the digest stays stable within a store.
*/
const CLAIM_MARKET_FIELDS = Object.freeze([
'line', 'side', 'book', 'locked_odds', 'over_odds', 'under_odds',
]);
/**
* ── THE CLAIM, CLASSIFIED ────────────────────────────────────────────────
* "Did price count as a revision?" is the wrong question because it has no
* single answer. VYNDR's BELIEF can hold still while the MARKET moves, and the
* composed PUBLISHED state changes either way. A chronology that cannot tell
* those apart reports 4,206 price moves and 301 belief moves as one undifferen-
* tiated number.
*
* So every claim field carries a class, and the change type is derived from
* WHICH classes moved.
*/
const FIELD_CLASS = Object.freeze({
BELIEF: 'BELIEF',
MARKET: 'MARKET',
COMPARISON: 'COMPARISON',
PROVENANCE: 'PROVENANCE',
PRESENTATION: 'PRESENTATION',
});
const CLAIM_FIELD_CLASSES = Object.freeze({
// What VYNDR independently believes will happen.
p_win: FIELD_CLASS.BELIEF,
projection: FIELD_CLASS.BELIEF,
grade: FIELD_CLASS.BELIEF,
engine_grade: FIELD_CLASS.BELIEF,
confidence: FIELD_CLASS.BELIEF,
fair_prob: FIELD_CLASS.BELIEF,
fair_odds: FIELD_CLASS.BELIEF,
model_odds: FIELD_CLASS.BELIEF,
// NOTE: canonical_event_id is deliberately NOT a claim field. Correcting a
// previously-derived event label is an IDENTITY correction, not a change of
// opinion or of price — classifying it as BELIEF would report every identity
// repair as VYNDR changing its mind.
//
// What the market offered.
line: FIELD_CLASS.MARKET,
side: FIELD_CLASS.MARKET,
book: FIELD_CLASS.MARKET,
locked_odds: FIELD_CLASS.MARKET,
over_odds: FIELD_CLASS.MARKET,
under_odds: FIELD_CLASS.MARKET,
// Belief measured against the market — moves when either side moves, which
// is why it is its own class and not evidence of a belief change.
ev_pct: FIELD_CLASS.COMPARISON,
edge_pct: FIELD_CLASS.COMPARISON,
value: FIELD_CLASS.COMPARISON,
kelly: FIELD_CLASS.COMPARISON,
takeable: FIELD_CLASS.COMPARISON,
// How the claim was produced.
confidence_basis: FIELD_CLASS.PROVENANCE,
// What the user was shown about refusal / alternatives.
alt_lines: FIELD_CLASS.PRESENTATION,
refused: FIELD_CLASS.PRESENTATION,
insufficient_data: FIELD_CLASS.PRESENTATION,
refusal_reason: FIELD_CLASS.PRESENTATION,
});
/** Why the published state advanced. Ordinal says WHEN; this says WHY. */
const CHANGE_TYPE = Object.freeze({
INITIAL_PUBLICATION: 'INITIAL_PUBLICATION',
BELIEF_CHANGE: 'BELIEF_CHANGE',
MARKET_REPRICE: 'MARKET_REPRICE',
MARKET_LINE_CHANGE: 'MARKET_LINE_CHANGE',
COMPARISON_CHANGE: 'COMPARISON_CHANGE',
PROVENANCE_CHANGE: 'PROVENANCE_CHANGE',
PRESENTATION_CHANGE: 'PRESENTATION_CHANGE',
MIXED_CHANGE: 'MIXED_CHANGE',
NO_MATERIAL_PUBLISHED_CHANGE: 'NO_MATERIAL_PUBLISHED_CHANGE',
});
/** Which claim fields differ between two published states, and their classes. */
function classifyChange(prev, next) {
const changed = [];
const classes = new Set();
for (const f of [...CLAIM_MARKET_FIELDS, ...SERVED_FIELDS]) {
const a = norm(prev ? prev[f] : null);
const b = norm(next ? next[f] : null);
const same = (typeof a === 'number' && typeof b === 'number')
? a === b : String(a) === String(b);
if (!same) {
changed.push(f);
classes.add(CLAIM_FIELD_CLASSES[f] || FIELD_CLASS.PRESENTATION);
}
}
if (changed.length === 0) {
return Object.freeze({
change_type: CHANGE_TYPE.NO_MATERIAL_PUBLISHED_CHANGE,
changed_fields: Object.freeze([]), changed_classes: Object.freeze([]),
belief_changed: false, market_changed: false,
});
}
const has = (c) => classes.has(c);
const beliefChanged = has(FIELD_CLASS.BELIEF);
const marketChanged = has(FIELD_CLASS.MARKET);
// COMPARISON alone is derived, so it is not promoted to a belief change.
const substantive = [...classes].filter((c) => c !== FIELD_CLASS.COMPARISON);
let type;
if (substantive.length === 0) type = CHANGE_TYPE.COMPARISON_CHANGE;
else if (substantive.length > 1) type = CHANGE_TYPE.MIXED_CHANGE;
else if (beliefChanged) type = CHANGE_TYPE.BELIEF_CHANGE;
else if (marketChanged) {
// A line move and a reprice are different market events and are not merged.
type = changed.includes('line') ? CHANGE_TYPE.MARKET_LINE_CHANGE : CHANGE_TYPE.MARKET_REPRICE;
} else if (has(FIELD_CLASS.PROVENANCE)) type = CHANGE_TYPE.PROVENANCE_CHANGE;
else type = CHANGE_TYPE.PRESENTATION_CHANGE;
return Object.freeze({
change_type: type,
changed_fields: Object.freeze(changed),
changed_classes: Object.freeze([...classes]),
belief_changed: beliefChanged,
market_changed: marketChanged,
});
}
function claimDigest(row) {
const r = row || {};
const material = {};
for (const f of [...CLAIM_MARKET_FIELDS, ...SERVED_FIELDS].sort()) {
const v = norm(r[f]);
// Numbers are stringified at fixed precision so 1.5 and "1.5" cannot
// produce two digests for one claim.
material[f] = typeof v === 'number' ? v.toFixed(6) : v;
}
return crypto.createHash('sha256').update(JSON.stringify(material)).digest('hex').slice(0, 32);
}
/**
* The natural key of a logical Read, EXCLUDING event spelling.
* Event compatibility is decided separately by `sameEvent`, because it is a
* fuzzy match and a hash cannot express one.
*/
/**
* -- EVENT IDENTITY METHOD, RECORDED NOT ASSUMED --------------------------
* A heuristic identity must carry its method, version and a refusal state
* rather than passing as truth.
*
* CANONICAL_EVENT_ID a provider-native event id was available and used.
* DERIVED_TEAM_PAIR identity came from `sport:date:away@home` plus the
* prefix rule. This is the CURRENT state for every row.
*/
const IDENTITY_METHOD = Object.freeze({
CANONICAL_EVENT_ID: 'CANONICAL_EVENT_ID',
DERIVED_TEAM_PAIR: 'DERIVED_TEAM_PAIR',
});
const IDENTITY_VERSION = 'ident@1';
/**
* A provider-native event occurrence, if the row carries one.
*
* -- WHY THIS HOOK EXISTS AND IS CURRENTLY ALWAYS NULL --------------------
* `ledgerService.gameIdFor` emits `sport:date:away@home` with NO game number,
* so the two halves of a doubleheader produce a BYTE-IDENTICAL id. No rule
* over that string can separate them, and `model_snapshots` carries no other
* event column.
*
* `mlbStatsAdapter` DOES expose `gamePk` (the canonical MLB game id), but it
* is not threaded onto the prop or the retention row. When it is, this reads
* it and the natural key separates the two games automatically.
*
* Until then identity is DERIVED_TEAM_PAIR and the doubleheader case is a
* reported limit, never a silent merge.
*/
function eventOccurrence(row) {
const r = row || {};
const v = r.canonical_event_id ?? r.game_pk ?? r.gamePk ?? r.event_id ?? null;
return v === null || v === undefined || v === '' ? null : String(v);
}
function identityMethod(row) {
return eventOccurrence(row)
? IDENTITY_METHOD.CANONICAL_EVENT_ID
: IDENTITY_METHOD.DERIVED_TEAM_PAIR;
}
function readNaturalKey(row) {
const r = row || {};
// A canonical occurrence, when present, is PART OF IDENTITY -- that is what
// makes two doubleheader games two Reads rather than one.
const occ = eventOccurrence(r);
const parts = [
String(r.sport || '').toLowerCase(),
String(r.game_date || ''),
String(r.player_key || ''),
String(r.stat || '').toLowerCase(),
String(r.side || '').toLowerCase(),
r.line === null || r.line === undefined ? '' : Number(r.line).toFixed(6),
];
if (parts.some((p) => p === '')) return null; // incomplete identity ⇒ no key
return occ ? `${parts.join('|')}|#${occ}` : parts.join('|');
}
/**
* Resolve where a candidate row belongs in the chronology.
*
* @param {object} spec
* candidate the row being written (needs sport/game_date/player_key/stat/
* side/line/game_id/captured_at + the claim fields)
* existing rows already in this logical Read's lineage, any order. Each
* needs: id, read_id, game_id, claim_digest, revision_ordinal,
* lineage_action, captured_at.
* @returns {object} never throws
*/
function resolveLineage(spec) {
const s = spec || {};
const cand = s.candidate || {};
const key = readNaturalKey(cand);
if (!key) {
return Object.freeze({
ok: false, refused: 'incomplete_read_identity',
lineage_version: LINEAGE_VERSION,
});
}
if (!cand.captured_at) {
return Object.freeze({ ok: false, refused: 'no_capture_clock', lineage_version: LINEAGE_VERSION });
}
// ── PUBLICATION GATE ────────────────────────────────────────────────────
// A capture is not a publication. `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: 8,443 of 13,012
// captured rows (64.9%) describe a state no user was ever shown.
//
// Publication chronology therefore advances ONLY for rows the collector's
// publication signal marked. An unmarked row is retained as model history
// with NULL lineage — truthful, and not a published claim.
if (cand.published !== true) {
return Object.freeze({
ok: false, refused: 'not_published', lineage_version: LINEAGE_VERSION,
});
}
const digest = claimDigest(cand);
const stamp = {
claim_schema_version: CLAIM_SCHEMA_VERSION,
digest_algorithm_version: DIGEST_ALGORITHM_VERSION,
};
// Only rows whose EVENT is compatible belong to this Read. Two genuinely
// different team pairs on one date are two Reads, not one.
const candOcc = eventOccurrence(cand);
const family = (Array.isArray(s.existing) ? s.existing : []).filter((r) => {
if (!r) return false;
const rOcc = eventOccurrence(r.claim || r);
// When BOTH sides carry a canonical occurrence it is the authority, and
// the team-name heuristic is not consulted at all.
if (candOcc && rOcc) return candOcc === rOcc;
return sameEvent(r.game_id, cand.game_id);
});
if (family.length === 0) {
return Object.freeze({
ok: true,
lineage_version: LINEAGE_VERSION,
action: LINEAGE_ACTION.ORIGIN,
read_id: s.mintReadId ? s.mintReadId() : null,
read_natural_key: key,
claim_digest: digest,
revision_ordinal: 0,
supersedes_id: null,
lineage_state: LINEAGE_STATE.LIVE,
change_type: CHANGE_TYPE.INITIAL_PUBLICATION,
identity_method: identityMethod(cand),
identity_version: IDENTITY_VERSION,
...stamp,
});
}
const readId = family[0].read_id || null;
// The standing published state = the highest-ordinal REVISION/ORIGIN row.
// Recaptures do not advance the chain, so they are not candidates for it.
const chain = family
.filter((r) => r.lineage_action === LINEAGE_ACTION.ORIGIN || r.lineage_action === LINEAGE_ACTION.REVISION)
.sort((a, b) => (a.revision_ordinal - b.revision_ordinal)
|| String(a.captured_at).localeCompare(String(b.captured_at)));
if (chain.length === 0) {
// Rows exist but none carries a chain position — legacy data. Do not invent
// an ordering for it; attach the candidate as an origin and say the history
// before it is unverified.
return Object.freeze({
ok: true,
lineage_version: LINEAGE_VERSION,
action: LINEAGE_ACTION.ORIGIN,
read_id: readId || (s.mintReadId ? s.mintReadId() : null),
read_natural_key: key,
claim_digest: digest,
revision_ordinal: 0,
supersedes_id: null,
lineage_state: LINEAGE_STATE.LEGACY_UNVERIFIED,
change_type: CHANGE_TYPE.INITIAL_PUBLICATION,
...stamp,
});
}
const head = chain[chain.length - 1];
// IDEMPOTENCY + NO-OP, one rule. An identical claim is the SAME published
// state however many times it is observed or retried.
if (head.claim_digest === digest) {
return Object.freeze({
ok: true,
lineage_version: LINEAGE_VERSION,
action: LINEAGE_ACTION.RECAPTURE,
read_id: readId,
read_natural_key: key,
claim_digest: digest,
// The recapture belongs to the standing revision; it does not advance it.
revision_ordinal: head.revision_ordinal,
supersedes_id: null,
recaptures_id: head.id,
lineage_state: LINEAGE_STATE.LIVE,
change_type: CHANGE_TYPE.NO_MATERIAL_PUBLISHED_CHANGE,
...stamp,
});
}
// ── WHY THERE IS NO IN-MEMORY FORK CHECK HERE ──────────────────────────
// A first draft refused when `family` already contained a row superseding the
// head. That check is unreachable: the head is BY DEFINITION the highest
// ordinal, so anything superseding it would already be the head. A resolver
// that can see its rival is not racing it — it appends after it, correctly.
//
// The real race is two resolvers that each read the same `existing` before
// either wrote. Neither can see the other, so no amount of in-memory logic
// detects it. That is enforced in the DATABASE by a UNIQUE index on
// supersedes_id (migration 045): only one row may supersede a given parent,
// so the second writer's insert fails loudly instead of forking history.
// `chronology()` also reports a fork after the fact.
// WHY the published state advanced, not merely THAT it did.
const delta = classifyChange(head.claim || head, cand);
return Object.freeze({
ok: true,
lineage_version: LINEAGE_VERSION,
action: LINEAGE_ACTION.REVISION,
read_id: readId,
read_natural_key: key,
claim_digest: digest,
revision_ordinal: head.revision_ordinal + 1,
supersedes_id: head.id,
lineage_state: LINEAGE_STATE.LIVE,
change_type: delta.change_type,
identity_method: identityMethod(cand),
identity_version: IDENTITY_VERSION,
changed_fields: delta.changed_fields,
belief_changed: delta.belief_changed,
market_changed: delta.market_changed,
...stamp,
});
}
/**
* Walk a Read's rows into its published chronology.
* Recaptures are attached to the revision they re-observed rather than dropped
* — they are evidence of when a claim was still standing.
*/
function chronology(rows) {
const list = (Array.isArray(rows) ? rows : []).filter(Boolean);
if (list.length === 0) return Object.freeze({ ok: false, reason: 'no rows', revisions: Object.freeze([]) });
const ids = new Set(list.map((r) => String(r.read_id)));
if (ids.size > 1) {
return Object.freeze({ ok: false, reason: 'rows span multiple read_ids', revisions: Object.freeze([]) });
}
const chain = list
.filter((r) => r.lineage_action === LINEAGE_ACTION.ORIGIN || r.lineage_action === LINEAGE_ACTION.REVISION)
.sort((a, b) => a.revision_ordinal - b.revision_ordinal);
const origins = chain.filter((r) => r.lineage_action === LINEAGE_ACTION.ORIGIN);
if (origins.length !== 1) {
return Object.freeze({ ok: false, reason: `expected exactly 1 ORIGIN, found ${origins.length}`, revisions: Object.freeze([]) });
}
for (let i = 1; i < chain.length; i += 1) {
if (String(chain[i].supersedes_id) !== String(chain[i - 1].id)) {
return Object.freeze({ ok: false, reason: `broken supersession at ordinal ${chain[i].revision_ordinal}`, revisions: Object.freeze([]) });
}
}
const recaptures = list.filter((r) => r.lineage_action === LINEAGE_ACTION.RECAPTURE);
return Object.freeze({
ok: true, reason: null,
read_id: chain[0] ? chain[0].read_id : null,
revisions: Object.freeze(chain),
original: chain[0] || null,
current: chain[chain.length - 1] || null,
revision_count: Math.max(0, chain.length - 1),
recapture_count: recaptures.length,
});
}
/** The revision standing at an instant — "what were we saying at 7pm?" */
function revisionAt(rows, instant) {
const c = chronology(rows);
if (!c.ok) return null;
let found = null;
for (const r of c.revisions) {
if (String(r.captured_at) <= String(instant)) found = r;
else break;
}
return found;
}
/**
* ── SETTLEMENT / CLAIM OWNERSHIP BOUNDARY ────────────────────────────────
* `model_snapshots` legitimately receives post-hoc evaluation after the event.
* That is enrichment of a historical row, not a change to what was claimed.
*
* The boundary is stated here rather than left implicit, because both are
* columns on the same table and nothing else distinguishes them. A settle pass
* that touched a claim-owned column would silently rewrite a receipt.
*/
const POST_HOC_FIELDS = Object.freeze([
'outcome', 'actual_value', 'settled_at', 'settlement_source',
'settlement_version', 're_settled_at', 'settle_attempts',
]);
/** Every field the published claim owns. Immutable after publication. */
function claimOwnedFields() {
return Object.freeze([
...CLAIM_MARKET_FIELDS,
...SERVED_FIELDS,
'read_id', 'read_natural_key', 'claim_digest', 'revision_ordinal',
'supersedes_id', 'lineage_action', 'change_type',
'claim_schema_version', 'digest_algorithm_version',
'captured_at', 'model_version', 'code_sha', 'published',
]);
}
/**
* Is a proposed post-event update confined to evaluation fields?
* Returns the violating keys rather than a bare boolean — a guard that cannot
* name what it caught is hard to act on.
*/
function settlementUpdateIsSafe(patch) {
const keys = Object.keys(patch || {});
const owned = new Set(claimOwnedFields());
const violations = keys.filter((k) => owned.has(k));
return Object.freeze({
safe: violations.length === 0,
violations: Object.freeze(violations),
// A key that is neither claim-owned nor a known evaluation field is
// reported: unrecognised is not the same as safe.
unrecognised: Object.freeze(keys.filter((k) => !owned.has(k) && !POST_HOC_FIELDS.includes(k))),
});
}
module.exports = {
resolveLineage, chronology, revisionAt, classifyChange,
claimDigest, readNaturalKey, eventFingerprint, sameEvent,
LINEAGE_ACTION, LINEAGE_STATE, LINEAGE_VERSION, CLAIM_MARKET_FIELDS,
CLAIM_SCHEMA_VERSION, DIGEST_ALGORITHM_VERSION,
CLAIM_FIELD_CLASSES, FIELD_CLASS, CHANGE_TYPE,
IDENTITY_METHOD, IDENTITY_VERSION, eventOccurrence, identityMethod,
POST_HOC_FIELDS, claimOwnedFields, settlementUpdateIsSafe,
};
+436 -1
View File
@@ -88,6 +88,10 @@ function boolOrNull(v) {
*/
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;
@@ -143,6 +147,17 @@ function rowsFromSides(base, sides, ctx = {}) {
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,
@@ -171,16 +186,397 @@ function rowsFromSides(base, sides, ctx = {}) {
/** 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 {
rows.push(...rowsFromSides(base, sides, ctx));
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));
if (hit) { hit.published = true; hit.published_side = 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
@@ -193,6 +589,15 @@ async function persist(rows, deps = {}) {
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);
@@ -253,6 +658,29 @@ function mergeEnrichment(rows, enrichedGrades) {
});
}
/**
* 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.
@@ -303,6 +731,13 @@ module.exports = {
mergeEnrichment,
mergeChainShadow,
persist,
attachLineage,
commitPublication,
recoverFromFork,
isSupersedesConflict,
FORK_RETRY_LIMIT,
LINEAGE_KEYS,
LINEAGE_FETCH_CLAIM,
newSnapshotId,
__internals: { numOrNull, intOrNull, boolOrNull },
};
+135 -1
View File
@@ -299,6 +299,25 @@ async function loadPitcherArsenals(sport) {
* flatten it. Both now serve RAW.
*/
const CALIBRATION_DEPLOYED = Object.freeze([]);
/**
* Which sports record publication lineage. MLB only by default — it is the only
* sport with a verified canonical event identity.
*/
// DARK BY DEFAULT. This release deploys the event-integrity and
// publication-commit corrections with the lineage OBSERVER OFF, so the repaired
// pipeline is exposed to production before anything starts writing history.
// Enabling is a deliberate configuration act: LINEAGE_CANARY_SPORTS=mlb.
//
// The flag governs the observer only. It does NOT control canonical event
// identity, the impossible-binding refusal, event-aware dedupe, or
// publication-commit ordering — those are integrity repairs and ship active.
const LINEAGE_CANARY_SPORTS = String(process.env.LINEAGE_CANARY_SPORTS || '')
.split(',').map((x) => x.trim().toLowerCase()).filter(Boolean);
function lineageCanaryEnabled(sport) {
return LINEAGE_CANARY_SPORTS.includes(String(sport || '').toLowerCase());
}
/**
* NOTHING IS SERVED CALIBRATED, AND THIS IS DELIBERATE.
*
@@ -412,6 +431,52 @@ async function runSnapshot(sport, opts = {}) {
const binder = deps.gameBinder || require('./gameBinder');
const b = await binder.attachGameTimes(sp, props, { gradedAt: ts });
console.log(`[snapshot] game binding ${sp}: ${b.bound} bound, ${b.alreadyHad} already had times, ${b.unresolved} UNRESOLVED, ${b.ambiguous} ambiguous(doubleheader)`);
// CANONICAL EVENT IDENTITY (MLB). The bound game_time above is what makes
// this resolvable: two halves of a doubleheader share teams and date and
// differ only by start time, so identity is decided by the clock, never by
// the team pair. Attached BEFORE dedupe, because dedupe is where two real
// games were previously collapsed into one.
//
// Sport-neutral field, MLB-only resolver: every other sport records
// UNSUPPORTED_SPORT rather than a guessed id.
if (sp === 'mlb') {
try {
const evid = deps.eventIdentity || require('./event/eventIdentity');
const mlbAdapter = deps.mlbAdapter || require('./adapters/mlbStatsAdapter');
const dates = [...new Set(props.map((p) => p && p.game_date).filter(Boolean))];
const games = [];
for (const d of dates) {
// eslint-disable-next-line no-await-in-loop
const g = await mlbAdapter.getScheduleWithPitchers(d);
if (Array.isArray(g)) games.push(...g);
}
// PARTICIPANT EVIDENCE. The prop's own team fields cannot police this:
// when the provider nests a player under the wrong event, those fields
// ARE the wrong event's teams. The player's real team has to come from
// a source that cannot be wrong the same way — the slate's rosters.
// Best-effort: no index means no refusals, i.e. the previous behaviour.
let playerTeams = null;
try {
const built = await evid.buildPlayerTeamIndex(games, {
getTeamRoster: (id) => mlbAdapter.getTeamRoster(id),
});
playerTeams = built.index;
console.log(`[snapshot] participant evidence ${sp}: ${built.stats.players} players`
+ ` across ${built.stats.teams} rosters${built.stats.failed ? `, ${built.stats.failed} unreadable` : ''}`);
} catch (e3) {
console.warn(`[snapshot] participant evidence unavailable for ${sp}:`, e3.message);
}
const e = evid.attachEventIdentity(sp, props, games, playerTeams);
console.log(`[snapshot] event identity ${sp}: ${e.canonical}/${e.total} canonical, `
+ `${e.unresolved} unresolved, ${e.impossible} impossible-binding refusals`
+ `${Object.keys(e.reasons).length ? ' ' + JSON.stringify(e.reasons) : ''}`);
} catch (e2) {
// Identity is additive: a failure leaves props on legacy identity, the
// exact behaviour that existed before. It must never cost a slate.
console.warn(`[snapshot] event identity failed for ${sp} (slate continues):`, e2.message);
}
}
if (b.unresolved > 0 && b.bound === 0 && b.alreadyHad === 0) {
await deps.notify(`Game binding produced NOTHING for ${sp.toUpperCase()} — ${b.unresolved} props could not be tied to a scheduled game. Their rows will be skipped rather than mis-dated.`, {
title: 'VYNDR pipeline', priority: 'high', tags: ['rotating_light'],
@@ -539,6 +604,7 @@ async function runSnapshot(sport, opts = {}) {
now: deps.now,
cacheSet: async (_k, v) => { envelope = v; },
onGraded: collector ? collector.onGraded : undefined,
onPublished: collector ? collector.onPublished : undefined,
});
// Retention is COLLECTED here (grade time — features must be exactly what the
@@ -547,6 +613,13 @@ async function runSnapshot(sport, opts = {}) {
// empty-slate early return: a slate that graded nothing but refused
// everything is exactly the case worth recording.
let retentionRows = 0;
// SHADOW LEDGER LINKAGE. retention persists BEFORE the ledger write, so the
// read_ids it assigned are already in memory — the join costs no query.
// read_revision_id is deliberately NOT resolved here: the upsert does not
// return row ids, and issuing a second query in the hot path to obtain one
// would be a real cost for a column nothing reads yet.
let lineageIndex = null;
let persistedRows = null;
const persistRetention = async (enrichedGrades, chainShadow = null) => {
if (!retention || !collector || !collector.rows.length) return;
try {
@@ -563,6 +636,13 @@ async function runSnapshot(sport, opts = {}) {
}
const r = await retention.persist(rows);
retentionRows = r.written || 0;
// Held for the publication commit, which happens only after the
// authoritative slate write succeeds.
persistedRows = rows;
// NOTE: no lineage index is built here. Lineage is assigned by
// commitPublication AFTER the slate write, so at this point no row has a
// read_id yet — building an index here would produce an empty one and
// quietly leave every ledger row unlinked.
console.log(`[snapshot] retention ${sp}: ${r.written}/${r.attempted} rows${r.skipped ? ' (skipped — no supabase env)' : ''}${r.error ? ` ERROR: ${r.error}` : ''}`);
} catch (e) {
console.warn(`[snapshot] retention write failed for ${sp} (snapshot continues):`, e.message);
@@ -1013,6 +1093,7 @@ async function runSnapshot(sport, opts = {}) {
}
}
// LINEUP + BASERUNNER CONTEXT — the input RBI and runs have always needed.
// Best-effort and dated: a context failure must never break a snapshot, and a
// lineup is a PRE-GAME fact that changes by the hour, so what we knew at grade
@@ -1046,6 +1127,57 @@ async function runSnapshot(sport, opts = {}) {
// what produced the "SIGNAL LIVE vs STALE 8h" contradiction.
const snapshot = { sport: sp, updated_at: ts, refreshed_at: ts, grades: enriched, deltas, gradeCount: enriched.length };
await deps.cacheSet(`snapshot:${sp}:latest`, snapshot, SNAP_TTL);
// ── PUBLICATION COMMITTED ───────────────────────────────────────────────
// The line above is the authoritative act: it is the single atomic write that
// makes the slate available to `GET /api/snapshot/:sport`. Only now is a
// claim "published", so only now may lineage record one.
//
// If the cacheSet above throws, execution never reaches here and no row is
// marked published — which is the entire point of the ordering. If THIS block
// fails, rows stay unpublished: the record understates what was served, which
// is observable and recoverable.
// ── CANARY SCOPE ────────────────────────────────────────────────────────
// Lineage dual-write runs for MLB only, because MLB is the one sport with a
// VERIFIED canonical event identity (statsapi gamePk). A sport without one
// cannot separate a doubleheader, and recording publication history on an
// identity that can merge two real games is the failure this whole layer
// exists to prevent.
//
// Env-gated both ways, repo-native (same shape as SNAPSHOT_CRON /
// INTRADAY_REFRESH): `LINEAGE_CANARY_SPORTS=` disables it entirely,
// `LINEAGE_CANARY_SPORTS=mlb,wnba` would widen it. Disabling needs no
// migration revert and touches no historical row.
if (retention && retention.commitPublication && persistedRows && lineageCanaryEnabled(sp)) {
try {
const pub = await retention.commitPublication({
rows: persistedRows,
publicationId: retentionCtx.snapshotId,
publishedAt: deps.now(),
});
const lin = pub.lineage || {};
console.log(`[publication] ${sp}: ${pub.published}/${pub.candidates} served rows committed`
+ `${pub.skipped ? ' (skipped — no supabase env)' : ''}${pub.error ? ` ERROR: ${pub.error}` : ''}`
+ ` · lineage origins=${lin.origins || 0} revisions=${lin.revisions || 0}`
+ ` recaptures=${lin.recaptures || 0} not_published=${lin.not_published || 0}`
+ ` forks=${lin.forks || 0}${lin.error ? ` lineage_error=${lin.error}` : ''}`);
if (pub.error || (lin && lin.error)) {
// A parity gap must be loud: the product published, the record did not.
await deps.notify(
`Lineage parity gap on ${sp.toUpperCase()}: slate published but lineage did not record it `
+ `(${pub.error || lin.error}). The product is unaffected; the historical record is incomplete.`,
{ title: 'VYNDR pipeline', priority: 'high', tags: ['rotating_light'] },
);
}
const idx = new Map();
for (const row of persistedRows) {
if (row && row.read_natural_key && row.read_id) idx.set(row.read_natural_key, row.read_id);
}
if (idx.size > 0) lineageIndex = idx;
} catch (e) {
console.warn(`[publication] commit failed for ${sp} (product unaffected):`, e.message);
}
}
// Session 59 — grades:{sport} must outlive the gap between cron runs (up to
// 5h) or team rosters / Explore / leaders go dark mid-day. SNAP_TTL (6h),
// NOT the legacy 2h gradeSlateService TTL — that gap was why /team showed
@@ -1063,7 +1195,7 @@ async function runSnapshot(sport, opts = {}) {
// Best-effort: the ledger must never break the snapshot.
let ledgerWritten = 0;
try {
const rec = await deps.ledger.recordPipelineGrades(sp, withChallenger, props, { now: deps.now });
const rec = await deps.ledger.recordPipelineGrades(sp, withChallenger, props, { now: deps.now, lineageIndex });
ledgerWritten = rec.written || 0;
await deps.ledger.captureClosing(sp, props);
} catch (e) {
@@ -1134,6 +1266,8 @@ async function runAllSnapshots(opts = {}) {
}
module.exports = {
lineageCanaryEnabled,
LINEAGE_CANARY_SPORTS,
runSnapshot,
runAllSnapshots,
computeLineDeltas,