Pre-grading stage trace: record which branch loses the cohort
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) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01CQJeAG8vcDoL5zkiaJyVb8
This commit is contained in:
@@ -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');
|
||||
|
||||
@@ -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)))
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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');
|
||||
});
|
||||
});
|
||||
@@ -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/);
|
||||
});
|
||||
|
||||
@@ -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', () => {
|
||||
|
||||
Reference in New Issue
Block a user