From 8d741e39321acbd2d6a7872cd2a7b2fa760f28af Mon Sep 17 00:00:00 2001 From: Kev Date: Fri, 28 Aug 2026 00:41:34 -0400 Subject: [PATCH] Pre-grading stage trace: record which branch loses the cohort MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The acquisition recorder proved the previous diagnosis wrong: the 03:00 MLB attempt acquired 6,365 props from a propline cache hit and CONTINUED. Zero reached the first grading callback, and nothing durable says why. WHY REPLAY WAS IMPOSSIBLE. The 03:00 input was the cache object written 02:20:24.889Z. ODDS_CACHE_TTL defaults to 3600s, so it expired ~03:20; it is now gone. No historical copy exists: closing_captures holds no MLB rows past 01:40:52, model_snapshots none, and `bookprices:mlb` carries no timestamp the book-comparison route exposes. Calling /api/odds/mlb was REFUSED — a cold cache would fetch and fire recordDownstream -> gradeAndCacheSlate, writing product state and manufacturing a false recovery. So the input is NOT RECOVERABLE and no replay was attempted. A SECOND CLAIM IS WITHDRAWN. The previous tranche concluded "all 6,365 props were rejected at event admission", reasoning that dedupe cannot empty a non-empty list. That is FALSE: `dedupeProps` also FILTERS — it drops props with no player/stat_type/line and any prop whose book is not a MODEL book. A fully admitted cohort can still dedupe to zero. Demonstrated in test. So the failing branch was never established, only assumed — which is exactly what the order forbade, and why DEDUPE_EMPTY is a first-class outcome here. THE REGION IS SMALL AND EVERY BRANCH LOOKS IDENTICAL FROM OUTSIDE: identity annotation (runSnapshot, mlb only, own catch) -> admitForGrading (pure, can throw) -> dedupeProps (pure, can throw, ALSO filters) -> mapLimit(gradeBestSide) <- first onGraded-capable call gradeAndCacheSlate swallows every throw in that region and returns the same {written:false,count:0} it returns for an honest zero. The recorder distinguishes them: READY_FOR_GRADING · ALL_REJECTED · IDENTITY_STAGE_ERROR · ADMISSION_STAGE_ERROR · DEDUPE_STAGE_ERROR · DEDUPE_EMPTY · NO_INPUT_PROPS · OTHER_PREGRADING_ERROR. An exception is never folded into ALL_REJECTED, and with no admission evidence the classifier refuses to classify at all — a test pins that. EXCEPTION SEMANTICS UNCHANGED. Each stage is wrapped to record and then RETHROW the identical error, so the enclosing best-effort catch still handles it exactly as before: nothing caught that was not caught, nothing swallowed that was not swallowed. Admission and dedupe rules, the impossible-binding invariant, event identity, model books and the first-row-wins cap are untouched — the only behavioural lines in the diff are `const gate/unique` becoming `let`. Correlated to the SAME snapshot_attempt_id the acquisition recorder minted — not a new run id — and stored under its own key `ops:pregrade:{sport}` so it can never displace the acquisition record. Same atomic LPUSH/LTRIM pattern, bounded per sport, scheduled-only, best-effort at the call site, auto-disabled under test. CAUGHT DURING BUILD: `savePgTrace` was defined and NEVER CALLED — the recorder would have persisted nothing, the same "built, correct, never invoked" failure this programme has hit before. A test now drives the real runSnapshot and asserts a trace is persisted carrying the acquisition's attempt id. Ten teeth, injections verified present, against a green baseline: empty collector read as ALL_REJECTED (2) · admission throw reported as admitted=0 (1) · dedupe throw reported as output=0 (1) · dedupe-to-zero mislabeled as rejection (1) · identity exception hidden (1) · observer reorders the candidate array (1) · trace failure changes product outcome (1) · intraday overwrites scheduled (1) · different attempt id (2) · secret leak (1). Restored byte-identically. THREE OF THOSE LANDED AND PASSED FIRST TIME — coverage holes, not safe defects: the dedupe-throw and identity-throw tests only exercised the recorder directly, never the real path, and nothing asserted the CANDIDATE array is not reordered. All three closed with real-path tests, then re-run failing. Two stale source assertions updated with the reason recorded: both pinned the exact `const gate = …` / opts-key order that the recording wrapper changed. 387 suites / 5,237 tests pass. web tsc exit 0. Lineage stays OFF. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_01CQJeAG8vcDoL5zkiaJyVb8 --- src/routes/internal.js | 27 +++ src/services/gradeSlateService.js | 42 +++- src/services/ops/acquisitionTrace.js | 145 +++++++++++- src/services/snapshotService.js | 32 +++ tests/unit/pregradeTrace.test.js | 335 +++++++++++++++++++++++++++ tests/unit/retentionIdentity.test.js | 5 +- tests/unit/shadowAccrual.test.js | 7 +- 7 files changed, 579 insertions(+), 14 deletions(-) create mode 100644 tests/unit/pregradeTrace.test.js diff --git a/src/routes/internal.js b/src/routes/internal.js index 0d2e3ac..6d44f12 100644 --- a/src/routes/internal.js +++ b/src/routes/internal.js @@ -46,6 +46,33 @@ router.use(requireInternalAuth({ loopbackOnly: false })); * runs no part of the pipeline. Counts, statuses and sanitized error classes * only: no keys, no URLs, no payloads, no prop arrays. */ +/** + * GET /api/internal/pregrading/:sport (pre-grading stage trace) + * + * READ-ONLY. The stage-by-stage record of cohort construction between + * acquisition and the first grading call — identity, admission, dedupe — for + * SCHEDULED attempts, correlated to the same snapshot_attempt_id the + * acquisition recorder minted. Counts and sanitized error classes only. Runs no + * pipeline and makes no fetch. + */ +router.get('/pregrading/:sport', async (req, res) => { + try { + const acq = require('../services/ops/acquisitionTrace'); + const attempts = await acq.pregradeHistory(req.params.sport); + res.set('Cache-Control', 'no-store'); + return res.json({ + ok: true, + sport: String(req.params.sport || '').toLowerCase(), + trigger_scope: acq.TRIGGER.SCHEDULED, + history_cap: acq.HISTORY_CAP, + count: attempts.length, + attempts, + }); + } catch (err) { + return res.status(500).json({ ok: false, error: err && err.message }); + } +}); + router.get('/acquisition/:sport', async (req, res) => { try { const acq = require('../services/ops/acquisitionTrace'); diff --git a/src/services/gradeSlateService.js b/src/services/gradeSlateService.js index fb00f90..923ab8e 100644 --- a/src/services/gradeSlateService.js +++ b/src/services/gradeSlateService.js @@ -58,6 +58,8 @@ const { isModelBook } = require('../config/bookRoles'); // statsapi is free and unlimited. Concurrency stays at 5 — one variable at a time. // // Env-tunable so the ceiling can move without a deploy: GRADE_SLATE_LIMIT. +const NO_PREGRADE = Object.freeze({ identity() {}, admission() {}, dedupe() {}, gradeLoop() {} }); + const DEFAULT_LIMIT = Number(process.env.GRADE_SLATE_LIMIT) > 0 ? Number(process.env.GRADE_SLATE_LIMIT) : 1500; @@ -343,6 +345,9 @@ async function gradeAndCacheSlate(sport, props, opts = {}) { const ttl = opts.ttl || DEFAULT_TTL; const concurrency = opts.concurrency || DEFAULT_CONCURRENCY; const now = opts.now || (() => new Date().toISOString()); + // Optional pre-grading recorder; a frozen no-op by default, so every existing + // caller keeps byte-identical behaviour. + const pregrade = opts.pregrade || NO_PREGRADE; if (!Array.isArray(props) || props.length === 0) { return { written: false, count: 0 }; @@ -351,12 +356,45 @@ async function gradeAndCacheSlate(sport, props, opts = {}) { try { // ADMISSION BEFORE DEDUPE: a prop that cannot name its game never reaches // the grader, and a contradicted one never reaches it either. - const gate = admitForGrading(props, sport); + // + // RECORDING ONLY. Each stage is wrapped to record its outcome and then + // RETHROW the identical error, so the enclosing catch below handles it + // exactly as before: nothing is caught that was not caught, nothing is + // swallowed that was not swallowed. Without this an exception here and a + // deliberate rejection both surface as {written:false,count:0}, which is + // why ALL_REJECTED must never be inferred from an empty collector. + let gate; + try { + gate = admitForGrading(props, sport); + } catch (eAdm) { + pregrade.admission({ started: true, completed: false, threw: true, input_count: props.length, error: eAdm && eAdm.message }); + throw eAdm; + } + pregrade.admission({ + started: true, completed: true, threw: false, + input_count: props.length, + admitted_count: gate.admitted.length, + rejected_count: gate.rejected.length, + rejection_reason_counts: { ...gate.reasons }, + }); if (gate.rejected.length > 0) { console.log(`[grade-slate] ${sport} event admission: ${gate.admitted.length} admitted, ` + `${gate.rejected.length} rejected ${JSON.stringify(gate.reasons)}`); } - const unique = dedupeProps(gate.admitted, limit); + let unique; + try { + unique = dedupeProps(gate.admitted, limit); + } catch (eDed) { + pregrade.dedupe({ started: true, completed: false, threw: true, input_count: gate.admitted.length, error: eDed && eDed.message }); + throw eDed; + } + pregrade.dedupe({ + started: true, completed: true, threw: false, + input_count: gate.admitted.length, + output_count: unique.length, + removed_count: gate.admitted.length - unique.length, + }); + pregrade.gradeLoop({ candidate_count: unique.length, reached: unique.length > 0 }); if (unique.length === 0) return { written: false, count: 0 }; const graded = (await mapLimit(unique, concurrency, (p) => gradeBestSide(grade, p, sport, opts))) diff --git a/src/services/ops/acquisitionTrace.js b/src/services/ops/acquisitionTrace.js index 19575a2..c70c0d5 100644 --- a/src/services/ops/acquisitionTrace.js +++ b/src/services/ops/acquisitionTrace.js @@ -165,11 +165,8 @@ function finish(trace, { final, propsCount, error, outcome, now } = {}) { * snapshot — that would make the observer more dangerous than the blindness it * exists to cure. */ -async function persist(trace, deps = {}) { - if (!trace || !trace.sport) return { stored: false, reason: 'no_sport' }; - // Only SCHEDULED attempts are retained; an intraday success must not be able - // to displace a scheduled failure. - if (trace.trigger !== TRIGGER.SCHEDULED) return { stored: false, reason: 'not_scheduled' }; +/** Shared bounded, ATOMIC append. Best-effort: never throws to the caller. */ +async function pushBounded(k, obj, deps = {}) { // Auto-disabled under test unless a client is injected — the `opsNotify` // precedent. Without this the default path constructs a real ioredis client // inside every suite that drives runSnapshot, which blocks on connect. @@ -180,8 +177,7 @@ async function persist(trace, deps = {}) { const getClient = deps.getRedisClient || require('../../utils/redis').getRedisClient; const client = getClient(); if (!client) return { stored: false, reason: 'no_redis' }; - const k = key(trace.sport); - await client.lpush(k, JSON.stringify(trace)); + await client.lpush(k, JSON.stringify(obj)); await client.ltrim(k, 0, HISTORY_CAP - 1); await client.expire(k, TTL_SECONDS); return { stored: true, key: k }; @@ -190,20 +186,149 @@ async function persist(trace, deps = {}) { } } -async function history(sport, deps = {}) { +async function readBounded(k, deps = {}) { try { const getClient = deps.getRedisClient || require('../../utils/redis').getRedisClient; const client = getClient(); if (!client) return []; - const raw = await client.lrange(key(sport), 0, HISTORY_CAP - 1); + const raw = await client.lrange(k, 0, HISTORY_CAP - 1); return (raw || []).map((r) => { try { return JSON.parse(r); } catch { return null; } }).filter(Boolean); } catch { return []; } } +async function persist(trace, deps = {}) { + if (!trace || !trace.sport) return { stored: false, reason: 'no_sport' }; + // Only SCHEDULED attempts are retained; an intraday success must not be able + // to displace a scheduled failure. + if (trace.trigger !== TRIGGER.SCHEDULED) return { stored: false, reason: 'not_scheduled' }; + return pushBounded(key(trace.sport), trace, deps); +} + +async function history(sport, deps = {}) { + return readBounded(key(sport), deps); +} + +/* ------------------------------------------------------------------ * + * PRE-GRADING STAGE TRACE + * + * The acquisition recorder proved the previous diagnosis wrong: the 03:00 MLB + * attempt acquired 6,365 props and CONTINUED. Zero of them reached the first + * grading callback, and nothing durable says why. + * + * The failing region is PRE-GRADING COHORT CONSTRUCTION, and it is small: + * identity annotation (runSnapshot, mlb only, own catch) + * -> admitForGrading (pure, can throw, caught by gradeAndCacheSlate) + * -> dedupeProps (pure, can throw, caught by gradeAndCacheSlate) + * -> mapLimit(gradeBestSide) <- the first onGraded-capable call + * + * `gradeAndCacheSlate` swallows every throw in that region and returns the SAME + * `{written:false,count:0}` it returns for an honest zero — so an exception and + * a deliberate rejection are indistinguishable from outside. That is exactly + * what this records, and it is why `ALL_REJECTED` must never be inferred from an + * empty collector. + * + * Correlated to the SAME `snapshot_attempt_id` the acquisition recorder minted. + * Stored under its own key so a pre-grading record can never displace the + * acquisition record for the same attempt. + * ------------------------------------------------------------------ */ + +const PREGRADE_PREFIX = 'ops:pregrade:'; + +const PREGRADE_OUTCOME = Object.freeze({ + READY_FOR_GRADING: 'READY_FOR_GRADING', + ALL_REJECTED: 'ALL_REJECTED', + IDENTITY_STAGE_ERROR: 'IDENTITY_STAGE_ERROR', + ADMISSION_STAGE_ERROR: 'ADMISSION_STAGE_ERROR', + DEDUPE_STAGE_ERROR: 'DEDUPE_STAGE_ERROR', + DEDUPE_EMPTY: 'DEDUPE_EMPTY', + NO_INPUT_PROPS: 'NO_INPUT_PROPS', + OTHER_PREGRADING_ERROR: 'OTHER_PREGRADING_ERROR', + INCOMPLETE: 'INCOMPLETE', +}); + +function pregradeKey(sport) { + return `${PREGRADE_PREFIX}${String(sport || 'unknown').toLowerCase()}`; +} + +function beginPregrade({ attemptId, sport, trigger, codeSha, processStartedAt, propsCount, now } = {}) { + return { + // The SAME identity the acquisition recorder minted — not a new run id. + snapshot_attempt_id: attemptId || null, + sport: sport ? String(sport).toLowerCase() : null, + trigger: trigger || TRIGGER.SCHEDULED, + code_sha: codeSha || null, + process_generation: processStartedAt || null, + started_at: now || new Date().toISOString(), + input_props_count: Number.isFinite(propsCount) ? propsCount : null, + identity: null, + admission: null, + dedupe: null, + grade_loop: null, + outcome: PREGRADE_OUTCOME.INCOMPLETE, + completed_at: null, + }; +} + +/** + * Pure writes into the trace. Nothing here may touch a prop object, reorder or + * filter an array, or influence a branch. + */ +function pregradeRecorder(trace) { + const t = trace || {}; + return { + identity(info) { t.identity = { ...info, error: sanitize(info && info.error) }; }, + admission(info) { t.admission = { ...info, error: sanitize(info && info.error) }; }, + dedupe(info) { t.dedupe = { ...info, error: sanitize(info && info.error) }; }, + gradeLoop(info) { t.grade_loop = { ...info }; }, + }; +} + +/** + * Classify from the recorded stages ONLY. An exception is never folded into + * ALL_REJECTED, and a dedupe that empties a non-empty admitted set is its own + * state rather than being blamed on admission. + */ +function classifyPregrade(trace) { + const t = trace || {}; + if (t.identity && t.identity.threw) return PREGRADE_OUTCOME.IDENTITY_STAGE_ERROR; + if (t.admission && t.admission.threw) return PREGRADE_OUTCOME.ADMISSION_STAGE_ERROR; + if (t.dedupe && t.dedupe.threw) return PREGRADE_OUTCOME.DEDUPE_STAGE_ERROR; + if (t.input_props_count === 0) return PREGRADE_OUTCOME.NO_INPUT_PROPS; + if (t.grade_loop && t.grade_loop.reached) return PREGRADE_OUTCOME.READY_FOR_GRADING; + if (t.admission && t.admission.completed && t.admission.admitted_count === 0) { + return PREGRADE_OUTCOME.ALL_REJECTED; + } + if (t.dedupe && t.dedupe.completed && t.dedupe.output_count === 0 + && t.admission && t.admission.admitted_count > 0) { + return PREGRADE_OUTCOME.DEDUPE_EMPTY; + } + return PREGRADE_OUTCOME.OTHER_PREGRADING_ERROR; +} + +function finishPregrade(trace, { now } = {}) { + const t = trace; + t.outcome = classifyPregrade(t); + t.completed_at = now || new Date().toISOString(); + return t; +} + +async function persistPregrade(trace, deps = {}) { + if (!trace || !trace.sport) return { stored: false, reason: 'no_sport' }; + if (trace.trigger !== TRIGGER.SCHEDULED) return { stored: false, reason: 'not_scheduled' }; + return pushBounded(pregradeKey(trace.sport), trace, deps); +} + +async function pregradeHistory(sport, deps = {}) { + return readBounded(pregradeKey(sport), deps); +} + + module.exports = { - TRIGGER, SOURCE_OUTCOME, FINAL, OUTCOME, + TRIGGER, SOURCE_OUTCOME, FINAL, OUTCOME, PREGRADE_OUTCOME, + pregradeKey, beginPregrade, pregradeRecorder, classifyPregrade, finishPregrade, + persistPregrade, pregradeHistory, HISTORY_CAP, TTL_SECONDS, KEY_PREFIX, key, sanitize, newAttemptId, begin, beginAttempt, currentAttempt, recorder, finishAttempt, finish, diff --git a/src/services/snapshotService.js b/src/services/snapshotService.js index 5d4e963..0c171b7 100644 --- a/src/services/snapshotService.js +++ b/src/services/snapshotService.js @@ -379,6 +379,9 @@ async function runSnapshot(sport, opts = {}) { persistAcquisitionTrace: opts.persistAcquisitionTrace || (async (t) => { try { return await acq.persist(t); } catch { return { stored: false }; } }), + persistPregradeTrace: opts.persistPregradeTrace || (async (t) => { + try { return await acq.persistPregrade(t); } catch { return { stored: false }; } + }), processStartedAt: opts.processStartedAt || null, scheduledHourUtc: opts.scheduledHourUtc, trigger: opts.trigger, @@ -431,6 +434,21 @@ async function runSnapshot(sport, opts = {}) { now: deps.now(), }); const acqRec = acq.recorder(acqTrace); + // Pre-grading stage trace, correlated to the SAME attempt identity. + const pgTrace = acq.beginPregrade({ + attemptId: acqTrace.snapshot_attempt_id, + sport: sp, + trigger: deps.trigger || acq.TRIGGER.SCHEDULED, + codeSha: (deps.retention && deps.retention.codeSha) ? deps.retention.codeSha() : null, + processStartedAt: deps.processStartedAt, + now: deps.now(), + }); + const pgRec = acq.pregradeRecorder(pgTrace); + const savePgTrace = async () => { + try { await deps.persistPregradeTrace(acq.finishPregrade(pgTrace, { now: deps.now() })); } catch (e) { + console.warn(`[pregrade-trace] store failed for ${sp} (snapshot continues):`, e && e.message); + } + }; // Best-effort AT THE CALL SITE, so ANY implementation — default or injected — // is safe. Guarding only the default dep left the product one bad injection // away from a telemetry write failing a healthy snapshot. @@ -486,6 +504,7 @@ async function runSnapshot(sport, opts = {}) { outcome: acq.OUTCOME.CONTINUED, now: deps.now(), }); await saveAcqTrace(acqTrace); + pgTrace.input_props_count = props.length; // Session 64 (Order 1.5) — BIND EVERY PROP TO ITS REAL GAME before anything // downstream dates it. PropLine emits no commence_time, so ledgerService's @@ -546,10 +565,19 @@ async function runSnapshot(sport, opts = {}) { // Without date-valid evidence a conflict can only be UNRESOLVED, never // CONTRADICTED — uncertainty must not become an accusation. const e = evid.attachEventIdentity(sp, props, games, playerTeams, evidenceDateValid); + pgRec.identity({ + started: true, completed: true, threw: false, + schedule_games: games.length, dates_requested: dates.length, + evidence_date_valid: evidenceDateValid, + total: e.total, canonical: e.canonical, unresolved: e.unresolved, + unsupported: e.unsupported, contradicted: e.impossible, + reasons: { ...e.reasons }, + }); console.log(`[snapshot] event identity ${sp}: ${e.canonical}/${e.total} canonical, ` + `${e.unresolved} unresolved, ${e.impossible} contradicted` + `${Object.keys(e.reasons).length ? ' ' + JSON.stringify(e.reasons) : ''}`); } catch (e2) { + pgRec.identity({ started: true, completed: false, threw: true, error: e2 && e2.message }); // 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); @@ -671,6 +699,7 @@ async function runSnapshot(sport, opts = {}) { } await deps.gradeAndCacheSlate(sp, props, { + pregrade: pgRec, factorContext, matchupKeys, // Bisect hook (2026-08-01): lets the internal trigger run a bounded slate @@ -684,6 +713,9 @@ async function runSnapshot(sport, opts = {}) { onGraded: collector ? collector.onGraded : undefined, onPublished: collector ? collector.onPublished : undefined, }); + // Persist the pre-grading stage trace. Best-effort at the call site: a + // telemetry failure must never fail a healthy slate. + await savePgTrace(); // Retention is COLLECTED here (grade time — features must be exactly what the // model saw) but PERSISTED after enrichment below, so archetype/team/opponent diff --git a/tests/unit/pregradeTrace.test.js b/tests/unit/pregradeTrace.test.js new file mode 100644 index 0000000..81e7f84 --- /dev/null +++ b/tests/unit/pregradeTrace.test.js @@ -0,0 +1,335 @@ +'use strict'; + +/** + * PRE-GRADING COHORT CONSTRUCTION — the stage between acquisition and grading. + * + * The acquisition recorder proved the earlier diagnosis wrong: the 03:00 MLB + * attempt acquired 6,365 props and CONTINUED. Zero reached the first grading + * callback, and nothing durable said why. + * + * The region is small and every branch in it returns the SAME + * `{written:false,count:0}`: + * admitForGrading (pure, can throw -> swallowed by gradeAndCacheSlate) + * dedupeProps (pure, can throw -> swallowed; ALSO FILTERS by model book) + * mapLimit(gradeBestSide) <- the first onGraded-capable call + * + * So an exception, a mass rejection and a dedupe-to-zero are indistinguishable + * from outside. `ALL_REJECTED` must never be inferred from an empty collector. + */ + +const fs = require('fs'); +const path = require('path'); +const acq = require('../../src/services/ops/acquisitionTrace'); +const gs = require('../../src/services/gradeSlateService'); +const { MODEL_BOOKS } = require('../../src/config/bookRoles'); +const ROOT = path.resolve(__dirname, '..', '..'); + +const BOOK = [...MODEL_BOOKS][0]; +const props = (n, over = {}) => Array.from({ length: n }, (_, i) => ({ + player: `Player ${i}`, stat_type: 'hits', line: 0.5 + i, sport: 'mlb', book: BOOK, + home_team: 'Reds', away_team: 'Cardinals', ...over, +})); + +async function runStage(input, opts = {}) { + const t = acq.beginPregrade({ attemptId: 'acq_test', sport: 'mlb', propsCount: input.length }); + const res = await gs.gradeAndCacheSlate('mlb', input, { + pregrade: acq.pregradeRecorder(t), + cacheSet: async () => {}, grade: async () => null, + ...opts, + }); + return { trace: acq.finishPregrade(t, {}), res }; +} + +describe('THE STAGE MATRIX', () => { + test('READY_FOR_GRADING — identity + admission + dedupe all succeed', async () => { + const { trace } = await runStage(props(5)); + expect(trace.admission.completed).toBe(true); + expect(trace.admission.admitted_count).toBe(5); + expect(trace.admission.rejected_count).toBe(0); + expect(trace.dedupe.completed).toBe(true); + expect(trace.dedupe.output_count).toBe(5); + expect(trace.grade_loop).toEqual({ candidate_count: 5, reached: true }); + expect(trace.outcome).toBe(acq.PREGRADE_OUTCOME.READY_FOR_GRADING); + }); + + test('ALL_REJECTED — with EXACT reason accounting', async () => { + const { trace } = await runStage(props(6, { event_binding_status: 'UNRESOLVED' })); + expect(trace.admission.admitted_count).toBe(0); + expect(trace.admission.rejected_count).toBe(6); + // admitted + rejected must reconcile to the admission input exactly. + expect(trace.admission.admitted_count + trace.admission.rejected_count) + .toBe(trace.admission.input_count); + const reasons = trace.admission.rejection_reason_counts; + expect(Object.values(reasons).reduce((a, b) => a + b, 0)).toBe(6); + expect(reasons.EVENT_UNRESOLVED).toBe(6); + expect(trace.outcome).toBe(acq.PREGRADE_OUTCOME.ALL_REJECTED); + }); + + test('MIXED admission reconciles admitted + rejected', async () => { + const { trace } = await runStage([ + ...props(4), + ...props(3, { event_binding_status: 'CONTRADICTED' }), + ]); + expect(trace.admission.input_count).toBe(7); + expect(trace.admission.admitted_count).toBe(4); + expect(trace.admission.rejected_count).toBe(3); + expect(trace.admission.rejection_reason_counts.EVENT_PLAYER_TEAM_CONTRADICTION).toBe(3); + expect(trace.outcome).toBe(acq.PREGRADE_OUTCOME.READY_FOR_GRADING); + }); + + test('DEDUPE_EMPTY — admitted > 0 but dedupe filters to zero, NOT ALL_REJECTED', async () => { + // dedupeProps FILTERS by model book. A fully-admitted set can still reach + // the grader with nothing, and that is a different fact from a rejection. + const { trace } = await runStage(props(5, { book: 'not_a_model_book' })); + expect(trace.admission.admitted_count).toBe(5); + expect(trace.dedupe.input_count).toBe(5); + expect(trace.dedupe.output_count).toBe(0); + expect(trace.outcome).toBe(acq.PREGRADE_OUTCOME.DEDUPE_EMPTY); + expect(trace.outcome).not.toBe(acq.PREGRADE_OUTCOME.ALL_REJECTED); + }); + + test('ADMISSION_STAGE_ERROR — a throw is never folded into ALL_REJECTED', async () => { + const t = acq.beginPregrade({ attemptId: 'a', sport: 'mlb', propsCount: 3 }); + const rec = acq.pregradeRecorder(t); + rec.admission({ started: true, completed: false, threw: true, input_count: 3, error: 'boom' }); + expect(acq.finishPregrade(t, {}).outcome).toBe(acq.PREGRADE_OUTCOME.ADMISSION_STAGE_ERROR); + expect(t.outcome).not.toBe(acq.PREGRADE_OUTCOME.ALL_REJECTED); + }); + + test('DEDUPE_STAGE_ERROR — a dedupe throw is not reported as output zero', async () => { + const t = acq.beginPregrade({ attemptId: 'a', sport: 'mlb', propsCount: 3 }); + const rec = acq.pregradeRecorder(t); + rec.admission({ started: true, completed: true, admitted_count: 3, rejected_count: 0, input_count: 3 }); + rec.dedupe({ started: true, completed: false, threw: true, input_count: 3, error: 'boom' }); + expect(acq.finishPregrade(t, {}).outcome).toBe(acq.PREGRADE_OUTCOME.DEDUPE_STAGE_ERROR); + }); + + test('a REAL dedupe throw is recorded as DEDUPE_STAGE_ERROR, not output zero', async () => { + // Drive the real gradeAndCacheSlate: admitted passes, then dedupeProps + // throws while reading the prop's book. + const bad = props(2); + Object.defineProperty(bad[1], 'book', { enumerable: true, get() { throw new Error('book poison'); } }); + const { res, trace } = await runStage(bad); + expect(res.written).toBe(false); + expect(trace.admission.completed).toBe(true); + expect(trace.admission.admitted_count).toBe(2); + expect(trace.dedupe.threw).toBe(true); + expect(trace.dedupe).not.toHaveProperty('output_count'); + expect(trace.outcome).toBe(acq.PREGRADE_OUTCOME.DEDUPE_STAGE_ERROR); + expect(trace.outcome).not.toBe(acq.PREGRADE_OUTCOME.DEDUPE_EMPTY); + }); + + test('a REAL identity throw reaches the trace through runSnapshot', async () => { + // The mlb identity block has its own catch. Recording must survive it. + const snap2 = require('../../src/services/snapshotService'); + const pg = []; + await snap2.runSnapshot('mlb', { + getOdds: async () => ({ sport: 'mlb', props: props(2), provider: 'propline', source: 'cache' }), + gradeAndCacheSlate: async () => ({ written: false, count: 0 }), + eventIdentity: { buildPlayerTeamIndex: async () => { throw new Error('roster down'); }, + evidenceIsDateValid: () => false, + attachEventIdentity: () => { throw new Error('statsapi shape changed'); } }, + mlbAdapter: { getScheduleWithPitchers: async () => [] }, + notify: async () => {}, sleep: async () => {}, retryDelayMs: 0, + now: () => '2026-08-28T14:00:00.000Z', scheduledHourUtc: 14, processStartedAt: 'p', + retention: null, ledger: { recordPipelineGrades: async () => {}, captureClosing: async () => {}, gameDateFor: () => '2026-08-28' }, + captureBookPrices: async () => ({}), buildEspnIndex: async () => ({}), + cacheGet: async () => null, cacheSet: async () => {}, + persistAcquisitionTrace: async () => {}, + persistPregradeTrace: async (t) => { pg.push(t); }, + }); + expect(pg).toHaveLength(1); + expect(pg[0].identity).toBeTruthy(); + expect(pg[0].identity.threw).toBe(true); + expect(pg[0].identity.error).toMatch(/statsapi shape changed|roster down/); + expect(pg[0].outcome).toBe(acq.PREGRADE_OUTCOME.IDENTITY_STAGE_ERROR); + }); + + test('IDENTITY_STAGE_ERROR outranks everything downstream', async () => { + const t = acq.beginPregrade({ attemptId: 'a', sport: 'mlb', propsCount: 3 }); + acq.pregradeRecorder(t).identity({ started: true, completed: false, threw: true, error: 'statsapi down' }); + expect(acq.finishPregrade(t, {}).outcome).toBe(acq.PREGRADE_OUTCOME.IDENTITY_STAGE_ERROR); + }); + + test('an empty collector alone can NEVER be classified', () => { + // The whole point: with no admission evidence the classifier refuses to say + // ALL_REJECTED. This is the defect the previous tranche talked itself into. + const t = acq.beginPregrade({ attemptId: 'a', sport: 'mlb', propsCount: 6365 }); + expect(acq.finishPregrade(t, {}).outcome).toBe(acq.PREGRADE_OUTCOME.OTHER_PREGRADING_ERROR); + expect(t.outcome).not.toBe(acq.PREGRADE_OUTCOME.ALL_REJECTED); + }); +}); + +describe('EXCEPTION SEMANTICS ARE UNCHANGED', () => { + const src = fs.readFileSync(path.join(ROOT, 'src/services/gradeSlateService.js'), 'utf8'); + + test('each recorded stage RETHROWS the identical error', () => { + expect(src).toMatch(/catch \(eAdm\) \{[\s\S]{0,220}throw eAdm;/); + expect(src).toMatch(/catch \(eDed\) \{[\s\S]{0,220}throw eDed;/); + }); + + test('the enclosing best-effort catch is still the only handler', () => { + expect(src).toMatch(/\/\/ Best-effort — slate grading must never break odds delivery\./); + expect(src).toMatch(/return \{ written: false, count: 0, error: e\.message \};/); + }); + + test('a real admission throw still returns the legacy shape, not a crash', async () => { + const boom = { admitForGrading: null }; + void boom; + // Drive the real function with a prop set that makes admitForGrading throw + // by poisoning the array itself (a non-object entry is skipped, so use a + // getter that throws on read). + const bad = props(2); + Object.defineProperty(bad, '2', { enumerable: true, get() { throw new Error('poison'); } }); + bad.length = 3; + const { res, trace } = await runStage(bad); + expect(res.written).toBe(false); + expect(res.error).toMatch(/poison/); + expect(trace.outcome).toBe(acq.PREGRADE_OUTCOME.ADMISSION_STAGE_ERROR); + }); +}); + +describe('THE OBSERVER DOES NOT TOUCH THE PROPS', () => { + test('prop objects and array order are unchanged', async () => { + const input = props(4); + const before = JSON.parse(JSON.stringify(input)); + const order = input.map((p) => p.player); + await runStage(input); + expect(input.map((p) => p.player)).toEqual(order); + expect(input).toHaveLength(4); + expect(JSON.parse(JSON.stringify(input))).toEqual(before); + }); + + test('the recorder writes only into the trace object', () => { + const t = acq.beginPregrade({ attemptId: 'a', sport: 'mlb' }); + const rec = acq.pregradeRecorder(t); + const p = { player: 'x' }; + rec.admission({ started: true, completed: true, admitted_count: 1, rejected_count: 0, input_count: 1 }); + expect(Object.keys(p)).toEqual(['player']); + expect(t.admission.admitted_count).toBe(1); + }); + + test('the CANDIDATE array is never reordered or filtered by the observer', () => { + // Reordering `unique` would change grading order and the first-row-wins cap + // — a behaviour change wearing an observer's clothes. + const src = fs.readFileSync(path.join(ROOT, 'src/services/gradeSlateService.js'), 'utf8'); + const from = src.indexOf('let unique;'); + const to = src.indexOf('const graded = (await mapLimit', from); + expect(from).toBeGreaterThan(-1); + expect(to).toBeGreaterThan(from); + const region = src.slice(from, to); + for (const mutator of ['.sort(', '.reverse(', '.splice(', '.filter(', '.push(', '.pop(', '.shift(']) { + expect(region).not.toContain(mutator); + } + }); + + test('the recorder default is a frozen no-op, so existing callers are identical', () => { + const src = fs.readFileSync(path.join(ROOT, 'src/services/gradeSlateService.js'), 'utf8'); + expect(src).toMatch(/const NO_PREGRADE = Object\.freeze\(\{ identity\(\) \{\}, admission\(\) \{\}, dedupe\(\) \{\}, gradeLoop\(\) \{\} \}\)/); + expect(src).toMatch(/const pregrade = opts\.pregrade \|\| NO_PREGRADE;/); + }); +}); + +describe('IDENTITY CORRELATION AND STORAGE', () => { + test('the pre-grading trace reuses the ACQUISITION attempt id', () => { + const t = acq.beginPregrade({ attemptId: 'acq_b01f75cc', sport: 'mlb' }); + expect(t.snapshot_attempt_id).toBe('acq_b01f75cc'); + const src = fs.readFileSync(path.join(ROOT, 'src/services/snapshotService.js'), 'utf8'); + expect(src).toMatch(/attemptId: acqTrace\.snapshot_attempt_id/); + }); + + test('it carries no product identity', () => { + const t = acq.beginPregrade({ attemptId: 'a', sport: 'mlb' }); + for (const f of ['snapshot_id', 'read_id', 'claim_digest', 'canonical_event_id', 'publication_id']) { + expect(t).not.toHaveProperty(f); + } + }); + + test('it stores under its OWN key, so it cannot displace the acquisition record', () => { + expect(acq.pregradeKey('mlb')).toBe('ops:pregrade:mlb'); + expect(acq.pregradeKey('mlb')).not.toBe(acq.key('mlb')); + }); + + test('per sport, atomic append, scheduled-only', async () => { + const keys = []; + const client = { lpush: async (k, v) => keys.push([k, JSON.parse(v).snapshot_attempt_id]), ltrim: async () => {}, expire: async () => {} }; + for (const sp of ['mlb', 'wnba']) { + await acq.persistPregrade(acq.finishPregrade(acq.beginPregrade({ attemptId: `a_${sp}`, sport: sp }), {}), { getRedisClient: () => client }); + } + expect(keys.map((k) => k[0])).toEqual(['ops:pregrade:mlb', 'ops:pregrade:wnba']); + const intraday = acq.beginPregrade({ attemptId: 'x', sport: 'mlb', trigger: acq.TRIGGER.INTRADAY }); + const out = await acq.persistPregrade(acq.finishPregrade(intraday, {}), { + getRedisClient: () => ({ lpush: async () => { throw new Error('must not be called'); }, ltrim: async () => {}, expire: async () => {} }), + }); + expect(out.reason).toBe('not_scheduled'); + }); + + test('two schedulers on the same slot are retained separately', async () => { + const stored = []; + const client = { lpush: async (k, v) => stored.push(JSON.parse(v)), ltrim: async () => {}, expire: async () => {} }; + for (const gen of ['proc-A', 'proc-B']) { + const t = acq.beginPregrade({ attemptId: `acq_${gen}`, sport: 'mlb', processStartedAt: gen }); + await acq.persistPregrade(acq.finishPregrade(t, {}), { getRedisClient: () => client }); + } + expect(new Set(stored.map((t) => t.process_generation)).size).toBe(2); + }); + + test('no credential or raw payload reaches the trace', async () => { + const t = acq.beginPregrade({ attemptId: 'a', sport: 'mlb' }); + acq.pregradeRecorder(t).admission({ started: true, completed: false, threw: true, error: 'boom https://x.io/y?apiKey=SUPERSECRETVALUE1234567' }); + const blob = JSON.stringify(t); + expect(blob).not.toMatch(/SUPERSECRETVALUE/); + expect(blob).not.toMatch(/x\.io/); + }); +}); + +describe('THE TRACE IS ACTUALLY PERSISTED BY THE REAL PATH', () => { + // A recorder that is built, correct, and never invoked is the exact failure + // this whole programme has hit before. Drive the REAL runSnapshot. + const snap = require('../../src/services/snapshotService'); + + test('runSnapshot persists a pre-grading trace correlated to the acquisition id', async () => { + const acqStored = []; const pgStored = []; + await snap.runSnapshot('wnba', { + getOdds: async () => ({ sport: 'wnba', props: props(3), provider: 'propline', source: 'cache' }), + gradeAndCacheSlate: require('../../src/services/gradeSlateService').gradeAndCacheSlate, + grade: async () => null, + notify: async () => {}, sleep: async () => {}, retryDelayMs: 0, + now: () => '2026-08-28T14:00:00.000Z', scheduledHourUtc: 14, + processStartedAt: 'proc-1', + retention: null, ledger: { recordPipelineGrades: async () => {}, captureClosing: async () => {}, gameDateFor: () => '2026-08-28' }, + captureBookPrices: async () => ({}), buildEspnIndex: async () => ({}), + cacheGet: async () => null, cacheSet: async () => {}, + persistAcquisitionTrace: async (t) => { acqStored.push(t); }, + persistPregradeTrace: async (t) => { pgStored.push(t); }, + }); + expect(acqStored).toHaveLength(1); + expect(pgStored).toHaveLength(1); + // SAME identity, not a second run id. + expect(pgStored[0].snapshot_attempt_id).toBe(acqStored[0].snapshot_attempt_id); + expect(pgStored[0].input_props_count).toBe(3); + expect(pgStored[0].admission.input_count).toBe(3); + expect(pgStored[0].outcome).toBe(acq.PREGRADE_OUTCOME.READY_FOR_GRADING); + }); +}); + +describe('TELEMETRY IS BEST-EFFORT', () => { + test('a trace-store failure does not change the product outcome', async () => { + const client = { lpush: async () => { throw new Error('redis down'); }, ltrim: async () => {}, expire: async () => {} }; + const out = await acq.persistPregrade(acq.finishPregrade(acq.beginPregrade({ attemptId: 'a', sport: 'mlb' }), {}), { getRedisClient: () => client }); + expect(out.stored).toBe(false); + expect(out.reason).toMatch(/redis down/); + }); + + test('the call site guards the write so it can never throw into runSnapshot', () => { + const src = fs.readFileSync(path.join(ROOT, 'src/services/snapshotService.js'), 'utf8'); + const i = src.indexOf('const savePgTrace = async ()'); + expect(i).toBeGreaterThan(-1); + expect(src.slice(i, i + 320)).toMatch(/try \{[\s\S]*catch/); + }); + + test('the default path opens no redis client under test', async () => { + const out = await acq.persistPregrade(acq.finishPregrade(acq.beginPregrade({ attemptId: 'a', sport: 'mlb' }), {})); + expect(out.reason).toBe('test_env'); + }); +}); diff --git a/tests/unit/retentionIdentity.test.js b/tests/unit/retentionIdentity.test.js index c97b543..6519d1a 100644 --- a/tests/unit/retentionIdentity.test.js +++ b/tests/unit/retentionIdentity.test.js @@ -102,7 +102,10 @@ describe('CANONICAL EVENT AVAILABILITY BY ROW CLASS', () => { test('UNRESOLVED / AMBIGUOUS / CONTRADICTED never reach retention at all', () => { // They are rejected BEFORE grading, and onGraded is what feeds retention. - expect(gs).toMatch(/const gate = admitForGrading\(props, sport\)/); + // Matches the CALL, not the assignment form: the pre-grading recorder + // wraps it in try/catch to record and rethrow, which changed `const gate =` + // into `gate = `. The invariant is that admission runs on (props, sport). + expect(gs).toMatch(/gate = admitForGrading\(props, sport\)/); expect(gs).toMatch(/dedupeProps\(gate\.admitted, limit\)/); expect(gs).not.toMatch(/dedupeProps\(gate\.rejected/); }); diff --git a/tests/unit/shadowAccrual.test.js b/tests/unit/shadowAccrual.test.js index 848b733..ba99833 100644 --- a/tests/unit/shadowAccrual.test.js +++ b/tests/unit/shadowAccrual.test.js @@ -151,7 +151,12 @@ describe('the scheduled path carries the shadow (not just a hand-run script)', ( it('runSnapshot builds the key index and passes it to the grader', () => { expect(src).toMatch(/matchupKeys/); - expect(src).toMatch(/gradeAndCacheSlate\(sp, props, \{\s*\n\s*factorContext,\s*\n\s*matchupKeys,/); + // The opts object now leads with the diagnostic `pregrade` recorder, so + // assert both keys are passed rather than pinning their exact order. + const call = src.slice(src.indexOf('await deps.gradeAndCacheSlate(sp, props, {')); + const opts = call.slice(0, call.indexOf('});')); + expect(opts).toMatch(/factorContext,/); + expect(opts).toMatch(/matchupKeys,/); }); it('the index is built BEFORE the grade call, not after', () => {