diff --git a/src/routes/internal.js b/src/routes/internal.js index 14f4582..0d2e3ac 100644 --- a/src/routes/internal.js +++ b/src/routes/internal.js @@ -37,6 +37,33 @@ router.use(requireInternalAuth({ loopbackOnly: false })); * Consumed by the admin dashboard's "Provider Quotas" tile. Cached * for 5s so a refresh-button mash doesn't flood Redis. */ +/** + * GET /api/internal/acquisition/:sport (scheduled acquisition trace) + * + * READ-ONLY. Returns the bounded history of SCHEDULED snapshot acquisition + * attempts for one sport — the primary/retry/fallback decision chain that + * production previously retained nothing about. Makes no provider request and + * runs no part of the pipeline. Counts, statuses and sanitized error classes + * only: no keys, no URLs, no payloads, no prop arrays. + */ +router.get('/acquisition/:sport', async (req, res) => { + try { + const acq = require('../services/ops/acquisitionTrace'); + const attempts = await acq.history(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('/quota', async (req, res) => { try { const providers = await quotaTracker.getAllQuotaStatuses(); diff --git a/src/services/oddsService.js b/src/services/oddsService.js index 877d795..d9a51f1 100644 --- a/src/services/oddsService.js +++ b/src/services/oddsService.js @@ -375,7 +375,15 @@ function shouldGradeSlate() { return process.env.NODE_ENV !== 'test'; } -async function getOdds(sport) { +// RECORDING ONLY. `trace` is an optional recorder supplied by the scheduled +// snapshot path so the acquisition decision chain becomes observable. Every +// method is a pure write; none can influence a branch, and nothing here makes +// an extra provider request. Callers that pass nothing keep the exact prior +// behaviour through this no-op. +const NO_TRACE = Object.freeze({ cache() {}, primary() {}, fallback() {} }); + +async function getOdds(sport, opts = {}) { + const trace = (opts && opts.trace) || NO_TRACE; const redis = getRedisClient(); const apiKey = process.env.ODDS_API_KEY; const cacheKey = getCacheKey(sport); @@ -385,6 +393,12 @@ async function getOdds(sport) { if (cached) { const data = JSON.parse(cached); const quota = await getQuotaRemaining(redis); + trace.cache({ + considered: true, hit: true, updated_at: data.updated_at || null, + provider: data.provider || 'odds-api', + count: Array.isArray(data.props) ? data.props.length : null, + decision: 'served_from_cache', + }); return { sport, updated_at: data.updated_at, @@ -401,10 +415,17 @@ async function getOdds(sport) { // first; on empty/error fall through to the conserved odds-api path // below. Gated on hasKeys() so environments without PropLine keys (incl. // the test suite) keep the exact prior behavior. + trace.cache({ considered: true, hit: false, decision: 'miss' }); const propline = require('./adapters/proplineAdapter'); + if (!propline.hasKeys()) trace.primary({ provider: 'propline', attempted: false, outcome: 'NOT_ATTEMPTED', reason: 'no_keys' }); if (propline.hasKeys()) { try { const pl = await propline.getProps(sport); + const plCount = (pl && Array.isArray(pl.props)) ? pl.props.length : 0; + trace.primary({ + provider: 'propline', attempted: true, + outcome: plCount > 0 ? 'NONZERO' : 'ZERO', count: plCount, + }); if (pl && Array.isArray(pl.props) && pl.props.length > 0) { const now = new Date().toISOString(); const cacheData = { updated_at: now, props: pl.props, spreads: pl.spreads || [], provider: 'propline' }; @@ -422,6 +443,7 @@ async function getOdds(sport) { }; } } catch (e) { + trace.primary({ provider: 'propline', attempted: true, outcome: 'ERROR', count: null, error: e && e.message }); console.warn('[oddsService] PropLine failed, falling back to odds-api:', e.message); } } @@ -442,6 +464,20 @@ async function getOdds(sport) { // quotaTracker.js for the rationale. const quotaTracker = require('./quotaTracker'); const quotaStatus = await quotaTracker.getQuotaStatus('odds-api'); + // Captured AT INVOCATION. Reading provider quota hours later and calling it + // historical evidence is the exact mistake this trace exists to stop. + trace.fallback({ + provider: 'odds-api', considered: true, + allowed_at_invocation: !!quotaStatus.allowed, + blocked_reason: quotaStatus.allowed ? null : 'quota_ceiling', + quota_at_invocation: { + used: quotaStatus.used != null ? quotaStatus.used : null, + limit: quotaStatus.limit != null ? quotaStatus.limit : null, + pct: quotaStatus.pct != null ? quotaStatus.pct : null, + period: quotaStatus.period || null, + }, + attempted: false, + }); if (!quotaStatus.allowed) { const error = new Error('Odds data temporarily unavailable. Try again later.'); error.statusCode = 429; @@ -452,6 +488,11 @@ async function getOdds(sport) { // Fetch live data try { const { props, spreads, quotaRemaining, headers } = await fetchAllOdds(sport, apiKey); + trace.fallback({ + provider: 'odds-api', considered: true, allowed_at_invocation: true, + attempted: true, outcome: (Array.isArray(props) && props.length > 0) ? 'NONZERO' : 'ZERO', + count: Array.isArray(props) ? props.length : null, + }); // Update quota in Redis if (headers) { @@ -477,6 +518,10 @@ async function getOdds(sport) { scratchedPlayers, }; } catch (err) { + trace.fallback({ + provider: 'odds-api', considered: true, allowed_at_invocation: true, + attempted: true, outcome: 'ERROR', error: err && err.message, + }); // If API fails, try stale cache (no TTL check — any cached data) const stale = await redis.get(cacheKey); if (stale) { diff --git a/src/services/ops/acquisitionTrace.js b/src/services/ops/acquisitionTrace.js new file mode 100644 index 0000000..19575a2 --- /dev/null +++ b/src/services/ops/acquisitionTrace.js @@ -0,0 +1,211 @@ +'use strict'; + +/** + * SCHEDULED ACQUISITION TRACE — diagnostic only. + * + * WHY THIS EXISTS. On 2026-08-27 the MLB snapshot stopped producing anything at + * the 22:00 and 01:00 slots while NBA/WNBA ran normally. The differential + * narrowed it to exactly two `runSnapshot` exits — `getOdds` THREW, or it + * RETURNED ZERO PROPS — and production retained nothing that could tell them + * apart. A failed acquisition has no retention snapshot_id, writes no ledger + * row and updates no slate, so it left no durable evidence at all. + * + * WHAT IT IS NOT. It changes no acquisition behaviour: not provider order, not + * the retry, not the quota threshold, not the fallback, not cache policy. It + * records decisions the real scheduled execution has ALREADY made. It issues + * NO provider request of its own — a test asserts that. + * + * WHY REDIS AND NOT MEMORY. `server.js` arms the scheduler in every process and + * there is no lock or leader election, and a rolling deploy demonstrably serves + * two containers at once. Process-local evidence could therefore be written by + * a container nobody later probes. The store is `LPUSH` + `LTRIM`, which is + * atomic, so two schedulers racing the same slot both survive instead of one + * silently overwriting the other. + * + * WHY A BOUNDED HISTORY AND NOT "LAST ACQUISITION". The scheduled MLB attempt + * failed at the hour and an intraday attempt SUCCEEDED ~20 minutes later. A + * single last-value would have erased the failure with the success — which is + * precisely the evidence we need. Only SCHEDULED attempts are recorded, and + * each sport keeps its own list. + */ + +const crypto = require('crypto'); + +const KEY_PREFIX = 'ops:acquisition:'; +/** Latest N scheduled attempts per sport. Small on purpose. */ +const HISTORY_CAP = 12; +/** Long enough to survive an overnight gap between slots. */ +const TTL_SECONDS = 172800; // 48h + +const TRIGGER = Object.freeze({ + SCHEDULED: 'SCHEDULED_SNAPSHOT', + INTRADAY: 'INTRADAY_REFRESH', + MANUAL: 'MANUAL_API', +}); + +const SOURCE_OUTCOME = Object.freeze({ + NONZERO: 'NONZERO', + ZERO: 'ZERO', + ERROR: 'ERROR', + NOT_ATTEMPTED: 'NOT_ATTEMPTED', +}); + +const FINAL = Object.freeze({ + NONZERO: 'NONZERO', + ZERO: 'ZERO', + THREW: 'THREW', +}); + +const OUTCOME = Object.freeze({ + CONTINUED: 'CONTINUED', + EARLY_RETURN_ODDS_ERROR: 'EARLY_RETURN_ODDS_ERROR', + EARLY_RETURN_ZERO_PROPS: 'EARLY_RETURN_ZERO_PROPS', + INCOMPLETE: 'INCOMPLETE', +}); + +function key(sport) { return `${KEY_PREFIX}${String(sport || 'unknown').toLowerCase()}`; } + +/** + * Strip anything credential-shaped before a message is stored. Provider errors + * can carry the full request URL, and PropLine's key rides in it. + */ +function sanitize(message) { + if (message == null) return null; + let s = String(message); + s = s.replace(/\b(api[_-]?key|apikey|key|token|authorization|auth)\s*[=:]\s*[^\s&"']+/gi, '$1=[redacted]'); + s = s.replace(/https?:\/\/[^\s"']+/gi, '[url]'); + s = s.replace(/\b[A-Za-z0-9_-]{24,}\b/g, '[redacted]'); + return s.slice(0, 240); +} + +/** A diagnostic identity ONLY. Never a Read, retention, Ledger or lineage id. */ +function newAttemptId() { + return `acq_${crypto.randomUUID()}`; +} + +function begin({ sport, trigger, scheduledHourUtc, codeSha, processStartedAt, now } = {}) { + return { + snapshot_attempt_id: newAttemptId(), + sport: sport ? String(sport).toLowerCase() : null, + trigger: trigger || TRIGGER.SCHEDULED, + scheduled_hour_utc: Number.isFinite(scheduledHourUtc) ? scheduledHourUtc : null, + started_at: now || new Date().toISOString(), + code_sha: codeSha || null, + // Distinguishes two containers running the same slot. + process_generation: processStartedAt || null, + attempts: [], + final: null, + final_props_count: null, + final_error: null, + outcome: OUTCOME.INCOMPLETE, + completed_at: null, + }; +} + +/** + * One `getOdds` call. `index` 0 is the primary call, 1 is the EXISTING + * runSnapshot retry — no retry is added here, the existing one is observed. + */ +function beginAttempt(trace, { index, isRetry, delayMs, now } = {}) { + const a = { + index: Number.isFinite(index) ? index : (trace.attempts.length), + is_retry: !!isRetry, + configured_delay_ms: Number.isFinite(delayMs) ? delayMs : null, + started_at: now || new Date().toISOString(), + cache: null, + primary: null, + fallback: null, + result: null, + }; + trace.attempts.push(a); + return a; +} + +function currentAttempt(trace) { + return trace.attempts.length ? trace.attempts[trace.attempts.length - 1] : null; +} + +/** + * The recorder handed to `getOdds`. Every method is a pure write into the + * current attempt — it can never influence a branch. + */ +function recorder(trace) { + const at = () => currentAttempt(trace); + return { + cache(info) { const a = at(); if (a) a.cache = { ...info }; }, + primary(info) { const a = at(); if (a) a.primary = { ...info, error: sanitize(info && info.error) }; }, + fallback(info) { const a = at(); if (a) a.fallback = { ...info, error: sanitize(info && info.error) }; }, + }; +} + +function finishAttempt(trace, { result, propsCount, error, provider, source } = {}) { + const a = currentAttempt(trace); + if (!a) return; + a.result = { + result, + props_count: Number.isFinite(propsCount) ? propsCount : null, + provider: provider || null, + source: source || null, + error: sanitize(error), + }; + a.completed_at = new Date().toISOString(); +} + +function finish(trace, { final, propsCount, error, outcome, now } = {}) { + trace.final = final || null; + trace.final_props_count = Number.isFinite(propsCount) ? propsCount : null; + trace.final_error = sanitize(error); + trace.outcome = outcome || OUTCOME.INCOMPLETE; + trace.completed_at = now || new Date().toISOString(); + return trace; +} + +/** + * Best-effort persist. A telemetry failure must NEVER fail a healthy product + * 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' }; + // 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. + if (!deps.getRedisClient && process.env.NODE_ENV === 'test') { + return { stored: false, reason: 'test_env' }; + } + try { + 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.ltrim(k, 0, HISTORY_CAP - 1); + await client.expire(k, TTL_SECONDS); + return { stored: true, key: k }; + } catch (e) { + return { stored: false, reason: sanitize(e && e.message) }; + } +} + +async function history(sport, 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); + return (raw || []).map((r) => { try { return JSON.parse(r); } catch { return null; } }).filter(Boolean); + } catch { + return []; + } +} + +module.exports = { + TRIGGER, SOURCE_OUTCOME, FINAL, OUTCOME, + HISTORY_CAP, TTL_SECONDS, KEY_PREFIX, + key, sanitize, newAttemptId, + begin, beginAttempt, currentAttempt, recorder, finishAttempt, finish, + persist, history, +}; diff --git a/src/services/snapshotService.js b/src/services/snapshotService.js index 1133251..5d4e963 100644 --- a/src/services/snapshotService.js +++ b/src/services/snapshotService.js @@ -29,6 +29,7 @@ const DELTA_MOVE = 1.0; // ticker MOVE threshold const STATS_CONCURRENCY = 5; const { nameKey, normalizeName } = require('../utils/playerName'); +const acq = require('./ops/acquisitionTrace'); // Session 46 — group/dedupe by the normalized name key so "A.J. Ewing" and // "AJ Ewing" (or "Jazz Chisholm" / "Jazz Chisholm Jr.") collapse to one player. const norm = (s) => nameKey(s); @@ -373,6 +374,14 @@ async function runSnapshot(sport, opts = {}) { notify: opts.notify || require('../utils/opsNotify').notify, sleep: opts.sleep || ((ms) => new Promise((r) => setTimeout(r, ms))), retryDelayMs: opts.retryDelayMs != null ? opts.retryDelayMs : 60_000, + // Diagnostic telemetry is best-effort BY CONSTRUCTION: a trace-store failure + // must never fail a healthy product snapshot. + persistAcquisitionTrace: opts.persistAcquisitionTrace || (async (t) => { + try { return await acq.persist(t); } catch { return { stored: false }; } + }), + processStartedAt: opts.processStartedAt || null, + scheduledHourUtc: opts.scheduledHourUtc, + trigger: opts.trigger, // Session 58 — Phase 1 truth infrastructure. ledger no-ops without // SUPABASE env, so tests / local dev never touch a database. ledger: opts.ledger || require('./ledgerService'), @@ -404,25 +413,79 @@ async function runSnapshot(sport, opts = {}) { // Session 56 — retry ONCE on a hard failure (thrown error / null response = a // transient PropLine/network blip). A successful-but-empty slate is NOT a // failure (off-hours), so it is not retried — that would waste quota + latency. + // SCHEDULED ACQUISITION TRACE — diagnostic only, recording only. + // + // A failed acquisition has no snapshot_id, writes no ledger row and updates no + // slate, so the 22:00 / 01:00 MLB failures left NOTHING able to tell `getOdds` + // threw from `getOdds` returned zero. This records decisions the existing + // control flow already makes: no extra provider request, no added retry, no + // branch changed. + const acqTrace = acq.begin({ + sport: sp, + trigger: deps.trigger || acq.TRIGGER.SCHEDULED, + scheduledHourUtc: deps.scheduledHourUtc, + // `retention` is bound from deps further down the function; reach through + // deps here rather than the not-yet-initialised local. + codeSha: (deps.retention && deps.retention.codeSha) ? deps.retention.codeSha() : null, + processStartedAt: deps.processStartedAt, + now: deps.now(), + }); + const acqRec = acq.recorder(acqTrace); + // 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. + const saveAcqTrace = async (t) => { + try { await deps.persistAcquisitionTrace(t); } catch (e) { + console.warn(`[acq-trace] store failed for ${sp} (snapshot continues):`, e && e.message); + } + }; let odds; try { - odds = await deps.getOdds(sp); + acq.beginAttempt(acqTrace, { index: 0, isRetry: false, now: deps.now() }); + odds = await deps.getOdds(sp, { trace: acqRec }); if (odds == null) throw new Error('null odds response'); + acq.finishAttempt(acqTrace, { + result: 'RETURNED', propsCount: Array.isArray(odds.props) ? odds.props.length : 0, + provider: odds.provider, source: odds.source, + }); } catch (e) { + acq.finishAttempt(acqTrace, { result: 'THREW', error: e && e.message }); try { + // The EXISTING retry. Observed, never added to. await deps.sleep(deps.retryDelayMs); - odds = await deps.getOdds(sp); + acq.beginAttempt(acqTrace, { index: 1, isRetry: true, delayMs: deps.retryDelayMs, now: deps.now() }); + odds = await deps.getOdds(sp, { trace: acqRec }); if (odds == null) throw new Error('null odds response (retry)'); + acq.finishAttempt(acqTrace, { + result: 'RETURNED', propsCount: Array.isArray(odds.props) ? odds.props.length : 0, + provider: odds.provider, source: odds.source, + }); } catch (e2) { + acq.finishAttempt(acqTrace, { result: 'THREW', error: e2 && e2.message }); + acq.finish(acqTrace, { + final: acq.FINAL.THREW, error: e2 && e2.message, + outcome: acq.OUTCOME.EARLY_RETURN_ODDS_ERROR, now: deps.now(), + }); + await saveAcqTrace(acqTrace); await deps.notify(`❌ ${sp.toUpperCase()} snapshot FAILED: ${e2.message}`, { title: 'VYNDR pipeline', priority: 'high', tags: ['x'] }); return { sport: sp, status: 'error', reason: e2.message, gradeCount: 0 }; } } const props = (odds && Array.isArray(odds.props)) ? odds.props : []; if (props.length === 0) { + acq.finish(acqTrace, { + final: acq.FINAL.ZERO, propsCount: 0, + outcome: acq.OUTCOME.EARLY_RETURN_ZERO_PROPS, now: deps.now(), + }); + await saveAcqTrace(acqTrace); await deps.notify(`⚠️ ${sp.toUpperCase()} snapshot: 0 props (odds unavailable)`, { title: 'VYNDR pipeline', priority: 'low', tags: ['warning'] }); return { sport: sp, status: 'skipped', reason: 'no props', gradeCount: 0 }; } + acq.finish(acqTrace, { + final: acq.FINAL.NONZERO, propsCount: props.length, + outcome: acq.OUTCOME.CONTINUED, now: deps.now(), + }); + await saveAcqTrace(acqTrace); // 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 diff --git a/src/snapshotScheduler.js b/src/snapshotScheduler.js index d3edfac..af0a96c 100644 --- a/src/snapshotScheduler.js +++ b/src/snapshotScheduler.js @@ -23,6 +23,10 @@ const HOURS_UTC = (process.env.SNAPSHOT_HOURS_UTC || '14,19,22,1,3') // PropLine-backed sports get the 20-min intraday refresh (soccer is odds-api // quota-gated). All of that lives in the config table, not here. const cadence = require('./config/sportCadence'); +const acq = require('./services/ops/acquisitionTrace'); +// One stable value per process lifetime — the generation marker that lets two +// containers running the same scheduled slot be told apart. +const PROCESS_STARTED_AT = new Date(Date.now() - Math.round(process.uptime() * 1000)).toISOString(); // Session 56 — missed-cron detection (pure, testable). // The latest scheduled hour:00 (UTC) at or before `date`. null if none in 48h. @@ -255,7 +259,16 @@ function startSnapshotScheduler(opts = {}) { // Job 1 — only the sports whose cadence includes this hour. MLB runs every // grid hour; WNBA skips 14/3 UTC (props not posted); soccer runs 14/19 only. const scheduled = cadence.sportsForHour(h); - const results = await runAll({ sports: scheduled }); + // Acquisition-trace context. Diagnostic only: it names WHICH slot and + // WHICH process made the attempt, so a scheduled failure can never be + // confused with an intraday call and two containers running the same + // slot stay distinguishable. + const results = await runAll({ + sports: scheduled, + trigger: acq.TRIGGER.SCHEDULED, + scheduledHourUtc: h, + processStartedAt: PROCESS_STARTED_AT, + }); const ok = results.filter((r) => r.status === 'ok'); console.log(`[snapshot] cron fired ${h}:00 UTC — ${ok.length}/${results.length} sports graded (scheduled: ${scheduled.join(',') || 'none'})`); // Session 63 (A1-S4) — after the day's FIRST slot, tell Kev the media diff --git a/tests/unit/acquisitionTrace.test.js b/tests/unit/acquisitionTrace.test.js new file mode 100644 index 0000000..746ced7 --- /dev/null +++ b/tests/unit/acquisitionTrace.test.js @@ -0,0 +1,391 @@ +'use strict'; + +/** + * SCHEDULED ACQUISITION OBSERVABILITY. + * + * The MLB snapshot stopped producing anything at the 22:00 and 01:00 slots on + * 2026-08-27 while NBA/WNBA ran normally. The differential narrowed it to two + * `runSnapshot` exits — `getOdds` THREW, or it RETURNED ZERO PROPS — and + * production retained nothing that could tell them apart: a failed acquisition + * has no snapshot_id, writes no ledger row and updates no slate. + * + * These tests prove the trace records the decision chain the existing control + * flow already makes, and that it changes nothing about acquisition. + */ + +const fs = require('fs'); +const path = require('path'); +const acq = require('../../src/services/ops/acquisitionTrace'); +const snap = require('../../src/services/snapshotService'); +const ROOT = path.resolve(__dirname, '..', '..'); + +jest.setTimeout(15000); + +// Deps that stop runSnapshot at the acquisition boundary: no grading, no +// retention, no network. The trace is finalized BEFORE any of this, so these +// only keep the test from running the whole pipeline. +const HALT = { + gradeAndCacheSlate: async () => ({ written: false, count: 0 }), + retention: null, + ledger: { recordPipelineGrades: async () => {}, captureClosing: async () => {}, gameDateFor: () => '2026-08-28' }, + captureBookPrices: async () => {}, + buildEspnIndex: async () => ({}), +}; + +const baseDeps = (over = {}) => ({ + notify: async () => {}, sleep: async () => {}, retryDelayMs: 0, + now: () => '2026-08-28T03:00:00.000Z', + scheduledHourUtc: 3, processStartedAt: '2026-08-28T00:31:17.657Z', + retention: require('../../src/services/retentionService'), + cacheGet: async () => null, cacheSet: async () => {}, + ...over, +}); + +/** Run runSnapshot and return the trace it persisted (if any). */ +async function runWithTrace(sport, getOdds, over = {}) { + const stored = []; + const result = await snap.runSnapshot(sport, baseDeps({ + getOdds, persistAcquisitionTrace: async (t) => { stored.push(t); return { stored: true }; }, ...over, + })); + return { result, trace: stored[0] || null, stored }; +} + +const threw = (msg, extra = {}) => async () => { const e = new Error(msg); Object.assign(e, extra); throw e; }; +const returns = (props, o = {}) => async () => ({ sport: 'mlb', props, provider: 'propline', source: 'live', ...o }); + +describe('ATTEMPT IDENTITY', () => { + test('every attempt gets a unique diagnostic id, distinct from every other identity', () => { + const a = acq.newAttemptId(); const b = acq.newAttemptId(); + expect(a).not.toBe(b); + expect(a.startsWith('acq_')).toBe(true); + }); + + test('the trace carries sport, trigger, slot, start, code_sha and process generation', async () => { + const { trace } = await runWithTrace('mlb', returns([])); + expect(trace.sport).toBe('mlb'); + expect(trace.trigger).toBe(acq.TRIGGER.SCHEDULED); + expect(trace.scheduled_hour_utc).toBe(3); + expect(trace.started_at).toBe('2026-08-28T03:00:00.000Z'); + expect(trace).toHaveProperty('code_sha'); + expect(trace.process_generation).toBe('2026-08-28T00:31:17.657Z'); + }); + + test('it does NOT reuse or alter any product identity', () => { + const t = acq.begin({ sport: 'mlb' }); + for (const forbidden of ['snapshot_id', 'read_id', 'claim_digest', 'canonical_event_id', 'ledger_id', 'publication_id']) { + expect(t).not.toHaveProperty(forbidden); + } + }); +}); + +describe('THE TEST MATRIX — one scheduled acquisition', () => { + test('PRIMARY NONZERO -> runSnapshot continues, no fallback invented', async () => { + // wnba: no MLB event-identity branch, so this stops at the acquisition + // boundary without reaching statsapi. + const { result, trace } = await runWithTrace('wnba', returns([{ player: 'x' }, { player: 'y' }]), HALT); + expect(trace.final).toBe(acq.FINAL.NONZERO); + expect(trace.final_props_count).toBe(2); + expect(trace.outcome).toBe(acq.OUTCOME.CONTINUED); + expect(trace.attempts).toHaveLength(1); + expect(result.status).not.toBe('error'); + expect(trace.sport).toBe('wnba'); + }); + + test('FINAL ZERO -> EARLY_RETURN_ZERO_PROPS, recorded distinctly from an error', async () => { + const { result, trace } = await runWithTrace('mlb', returns([])); + expect(trace.final).toBe(acq.FINAL.ZERO); + expect(trace.final_props_count).toBe(0); + expect(trace.outcome).toBe(acq.OUTCOME.EARLY_RETURN_ZERO_PROPS); + expect(trace.final).not.toBe(acq.FINAL.THREW); + expect(result.status).toBe('skipped'); + expect(result.reason).toBe('no props'); + }); + + test('FINAL THROW -> EARLY_RETURN_ODDS_ERROR, with BOTH attempts preserved', async () => { + const { result, trace } = await runWithTrace('mlb', threw('Odds data temporarily unavailable.', { statusCode: 429 })); + expect(trace.final).toBe(acq.FINAL.THREW); + expect(trace.outcome).toBe(acq.OUTCOME.EARLY_RETURN_ODDS_ERROR); + // The EXISTING retry — primary + retry, not an added one. + expect(trace.attempts).toHaveLength(2); + expect(trace.attempts[0].is_retry).toBe(false); + expect(trace.attempts[1].is_retry).toBe(true); + expect(trace.attempts.every((a) => a.result.result === 'THREW')).toBe(true); + expect(result.status).toBe('error'); + }); + + test('PRIMARY ERROR -> EXISTING RETRY SUCCESS: both attempts preserved, final NONZERO', async () => { + let n = 0; + const getOdds = async () => { + n += 1; + if (n === 1) throw new Error('transient'); + return { sport: 'wnba', props: [{ player: 'x' }], provider: 'propline', source: 'live' }; + }; + const { trace } = await runWithTrace('wnba', getOdds, HALT); + expect(trace.attempts).toHaveLength(2); + expect(trace.attempts[0].result.result).toBe('THREW'); + expect(trace.attempts[1].result.result).toBe('RETURNED'); + expect(trace.final).toBe(acq.FINAL.NONZERO); + expect(trace.outcome).toBe(acq.OUTCOME.CONTINUED); + }); + + test('the recorded retry delay is the CONFIGURED one, never an added retry', async () => { + const { trace } = await runWithTrace('mlb', threw('x'), { retryDelayMs: 60000 }); + expect(trace.attempts[1].configured_delay_ms).toBe(60000); + const src = fs.readFileSync(path.join(ROOT, 'src/services/snapshotService.js'), 'utf8'); + // Exactly two getOdds calls and exactly one retry sleep — unchanged. + expect((src.match(/deps\.getOdds\(/g) || [])).toHaveLength(2); + expect((src.match(/deps\.sleep\(deps\.retryDelayMs\)/g) || [])).toHaveLength(1); + }); +}); + +describe('THE DECISION CHAIN INSIDE getOdds', () => { + const chain = () => { + const t = acq.begin({ sport: 'mlb' }); + acq.beginAttempt(t, { index: 0 }); + return { t, rec: acq.recorder(t) }; + }; + + test('PRIMARY ZERO -> RETRY ZERO -> FALLBACK BLOCKED records the exact chain', () => { + const { t, rec } = chain(); + rec.cache({ considered: true, hit: false, decision: 'miss' }); + rec.primary({ provider: 'propline', attempted: true, outcome: 'ZERO', count: 0 }); + rec.fallback({ + provider: 'odds-api', considered: true, allowed_at_invocation: false, + blocked_reason: 'quota_ceiling', + quota_at_invocation: { used: 478, limit: 500, pct: 0.956, period: '2026-08' }, + attempted: false, + }); + acq.finishAttempt(t, { result: 'THREW', error: 'Odds data temporarily unavailable.' }); + const a = t.attempts[0]; + expect(a.cache.hit).toBe(false); + expect(a.primary.outcome).toBe('ZERO'); + expect(a.primary.count).toBe(0); + expect(a.fallback.allowed_at_invocation).toBe(false); + expect(a.fallback.blocked_reason).toBe('quota_ceiling'); + expect(a.fallback.quota_at_invocation.used).toBe(478); + expect(a.fallback.attempted).toBe(false); + }); + + test('PRIMARY ERROR is recorded distinctly from PRIMARY ZERO', () => { + const { t, rec } = chain(); + rec.primary({ provider: 'propline', attempted: true, outcome: 'ERROR', count: null, error: 'socket hang up' }); + expect(t.attempts[0].primary.outcome).toBe('ERROR'); + expect(t.attempts[0].primary.outcome).not.toBe('ZERO'); + expect(t.attempts[0].primary.count).toBeNull(); + }); + + test('FALLBACK SUCCESS is recorded when the existing behaviour reaches it', () => { + const { t, rec } = chain(); + rec.primary({ provider: 'propline', attempted: true, outcome: 'ZERO', count: 0 }); + rec.fallback({ provider: 'odds-api', considered: true, allowed_at_invocation: true, attempted: true, outcome: 'NONZERO', count: 120 }); + expect(t.attempts[0].fallback.outcome).toBe('NONZERO'); + expect(t.attempts[0].fallback.count).toBe(120); + }); + + test('quota is captured AT INVOCATION, inside getOdds, not read later', () => { + const src = fs.readFileSync(path.join(ROOT, 'src/services/oddsService.js'), 'utf8'); + const i = src.indexOf("const quotaStatus = await quotaTracker.getQuotaStatus('odds-api');"); + expect(i).toBeGreaterThan(-1); + const block = src.slice(i, i + 900); + expect(block).toMatch(/quota_at_invocation/); + expect(block).toMatch(/allowed_at_invocation: !!quotaStatus\.allowed/); + }); + + test('the recorder is a pure write — a no-op recorder leaves behaviour identical', async () => { + const src = fs.readFileSync(path.join(ROOT, 'src/services/oddsService.js'), 'utf8'); + expect(src).toMatch(/const NO_TRACE = Object\.freeze\(\{ cache\(\) \{\}, primary\(\) \{\}, fallback\(\) \{\} \}\)/); + expect(src).toMatch(/const trace = \(opts && opts\.trace\) \|\| NO_TRACE;/); + // No trace call may appear inside a conditional test. + expect(src).not.toMatch(/if\s*\([^)]*trace\.[a-z]/); + }); +}); + +describe('THE OBSERVER MAKES NO PROVIDER CALL', () => { + const oddsSrc = fs.readFileSync(path.join(ROOT, 'src/services/oddsService.js'), 'utf8'); + const traceSrc = fs.readFileSync(path.join(ROOT, 'src/services/ops/acquisitionTrace.js'), 'utf8'); + + test('the trace module contains no HTTP or provider call whatsoever', () => { + // Scan CODE, not prose or the sanitizer's own URL regex — a guard that + // flags its own redaction pattern is checking the wrong thing. + const code = traceSrc + .replace(/\/\*[\s\S]*?\*\//g, '') + .replace(/(^|[^:])\/\/.*$/gm, '$1') + .replace(/s\.replace\([\s\S]*?\);/g, ''); + for (const forbidden of ['fetch(', 'axios', 'getProps', 'fetchAllOdds', 'gateway', + "require('http", 'https://', 'getOdds', 'oddsService', 'proplineAdapter']) { + expect(code).not.toContain(forbidden); + } + }); + + test('instrumenting getOdds added no provider request', () => { + expect((oddsSrc.match(/await propline\.getProps\(/g) || [])).toHaveLength(1); + expect((oddsSrc.match(/await fetchAllOdds\(/g) || [])).toHaveLength(1); + }); + + test('the internal read endpoint runs no pipeline and no fetch', () => { + const src = fs.readFileSync(path.join(ROOT, 'src/routes/internal.js'), 'utf8'); + const i = src.indexOf("router.get('/acquisition/:sport'"); + expect(i).toBeGreaterThan(-1); + const block = src.slice(i, src.indexOf('});', src.indexOf('} catch', i))); + expect(block).not.toMatch(/getOdds|runSnapshot|diagnose|fetch\(/); + expect(block).toMatch(/acq\.history/); + }); +}); + +describe('SANITIZATION', () => { + test('credentials, tokens and URLs never reach the trace', () => { + const dirty = 'failed https://api.propline.io/v1/odds?apiKey=abcd1234abcd1234abcd1234 token=zzzzzzzzzzzzzzzzzzzzzzzzzz'; + const clean = acq.sanitize(dirty); + expect(clean).not.toMatch(/abcd1234abcd1234/); + expect(clean).not.toMatch(/api\.propline\.io/); + expect(clean).not.toMatch(/zzzzzzzzzzzzzzzzzzzzzz/); + expect(clean).toMatch(/\[url\]|\[redacted\]/); + }); + + test('a thrown provider error is sanitized on the way into the trace', async () => { + const { trace } = await runWithTrace('mlb', threw('boom https://x.io/y?key=SUPERSECRETKEYVALUE123456')); + const blob = JSON.stringify(trace); + expect(blob).not.toMatch(/SUPERSECRETKEYVALUE/); + expect(blob).not.toMatch(/x\.io/); + }); + + test('no prop payload is retained — counts only', async () => { + const { trace } = await runWithTrace('mlb', returns([{ player: 'Aaron Judge', stat: 'hits' }])); + expect(JSON.stringify(trace)).not.toMatch(/Aaron Judge/); + expect(trace.final_props_count).toBe(1); + }); +}); + +describe('HISTORY IS BOUNDED, PER SPORT, AND SCHEDULED-ONLY', () => { + test('storage is a bounded per-sport list with a TTL', () => { + expect(acq.key('mlb')).toBe('ops:acquisition:mlb'); + expect(acq.key('wnba')).not.toBe(acq.key('mlb')); + expect(acq.HISTORY_CAP).toBeGreaterThan(1); + expect(acq.HISTORY_CAP).toBeLessThanOrEqual(20); + expect(acq.TTL_SECONDS).toBeGreaterThanOrEqual(86400); + }); + + test('it uses ATOMIC append (lpush/ltrim), so two schedulers cannot overwrite each other', async () => { + const calls = []; + const client = { + lpush: async (k, v) => { calls.push(['lpush', k, JSON.parse(v).snapshot_attempt_id]); }, + ltrim: async (k, a, b) => calls.push(['ltrim', k, a, b]), + expire: async (k, t) => calls.push(['expire', k, t]), + }; + const a = acq.finish(acq.begin({ sport: 'mlb' }), { final: acq.FINAL.THREW }); + const b = acq.finish(acq.begin({ sport: 'mlb' }), { final: acq.FINAL.THREW }); + await acq.persist(a, { getRedisClient: () => client }); + await acq.persist(b, { getRedisClient: () => client }); + expect(calls.filter((c) => c[0] === 'lpush')).toHaveLength(2); + expect(calls.filter((c) => c[0] === 'ltrim')).toHaveLength(2); + // Read-modify-write would show a get; there is none. + expect(calls.some((c) => c[0] === 'get' || c[0] === 'set')).toBe(false); + const src = fs.readFileSync(path.join(ROOT, 'src/services/ops/acquisitionTrace.js'), 'utf8'); + expect(src).toMatch(/client\.lpush/); + expect(src).not.toMatch(/cacheSet\(/); + }); + + test('INTRADAY SUCCESS CANNOT OVERWRITE A SCHEDULED FAILURE', async () => { + // Structural: intraday never calls runSnapshot, and persist refuses any + // non-scheduled trigger outright. + const intraday = acq.finish(acq.begin({ sport: 'mlb', trigger: acq.TRIGGER.INTRADAY }), { final: acq.FINAL.NONZERO }); + const client = { lpush: async () => { throw new Error('must not be called'); }, ltrim: async () => {}, expire: async () => {} }; + const out = await acq.persist(intraday, { getRedisClient: () => client }); + expect(out.stored).toBe(false); + expect(out.reason).toBe('not_scheduled'); + const src = fs.readFileSync(path.join(ROOT, 'src/services/intradayRefreshService.js'), 'utf8'); + expect(src).not.toMatch(/runSnapshot/); + }); + + test('ANOTHER SPORT CANNOT OVERWRITE MLB EVIDENCE', async () => { + const keys = []; + const client = { lpush: async (k) => keys.push(k), ltrim: async () => {}, expire: async () => {} }; + for (const sp of ['mlb', 'nba', 'wnba']) { + await acq.persist(acq.finish(acq.begin({ sport: sp }), { final: acq.FINAL.ZERO }), { getRedisClient: () => client }); + } + expect(new Set(keys).size).toBe(3); + expect(keys).toContain('ops:acquisition:mlb'); + }); + + 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.begin({ sport: 'mlb', scheduledHourUtc: 3, processStartedAt: gen }); + await acq.persist(acq.finish(t, { final: acq.FINAL.ZERO }), { getRedisClient: () => client }); + } + expect(stored).toHaveLength(2); + expect(stored[0].snapshot_attempt_id).not.toBe(stored[1].snapshot_attempt_id); + expect(new Set(stored.map((t) => t.process_generation)).size).toBe(2); + }); +}); + +describe('TELEMETRY FAILURE MUST NOT BREAK THE PRODUCT', () => { + test('a trace-store failure does not fail a healthy acquisition', async () => { + const { result, stored } = await runWithTrace('wnba', returns([{ p: 1 }]), { + ...HALT, + persistAcquisitionTrace: async () => { throw new Error('redis down'); }, + }); + // The snapshot continued past acquisition; it did not return an odds error. + expect(result.status).not.toBe('error'); + expect(stored).toHaveLength(0); + }); + + test('the default path never constructs a redis client under test', async () => { + // The opsNotify precedent. Without it, every suite driving runSnapshot + // would open a real connection and block. + const out = await acq.persist(acq.finish(acq.begin({ sport: 'mlb' }), { final: acq.FINAL.ZERO })); + expect(out.stored).toBe(false); + expect(out.reason).toBe('test_env'); + }); + + test('persist swallows a store error and reports it rather than throwing', async () => { + const client = { lpush: async () => { throw new Error('redis down'); }, ltrim: async () => {}, expire: async () => {} }; + const out = await acq.persist(acq.finish(acq.begin({ sport: 'mlb' }), { final: acq.FINAL.ZERO }), { getRedisClient: () => client }); + expect(out.stored).toBe(false); + expect(out.reason).toMatch(/redis down/); + }); + + test('the default persist dep is wrapped so it can never throw into runSnapshot', () => { + const src = fs.readFileSync(path.join(ROOT, 'src/services/snapshotService.js'), 'utf8'); + const i = src.indexOf('persistAcquisitionTrace: opts.persistAcquisitionTrace'); + expect(i).toBeGreaterThan(-1); + expect(src.slice(i, i + 220)).toMatch(/try \{[\s\S]*catch/); + }); +}); + +describe('ACQUISITION BEHAVIOUR IS UNCHANGED', () => { + test('getOdds gained only an optional recorder argument', () => { + const src = fs.readFileSync(path.join(ROOT, 'src/services/oddsService.js'), 'utf8'); + expect(src).toMatch(/async function getOdds\(sport, opts = \{\}\)/); + }); + + test('provider order, quota threshold and cache policy are untouched', () => { + const src = fs.readFileSync(path.join(ROOT, 'src/services/oddsService.js'), 'utf8'); + // PropLine still first, still gated on hasKeys, still the >0 test. + expect(src).toMatch(/if \(propline\.hasKeys\(\)\) \{/); + expect(src).toMatch(/pl\.props\.length > 0/); + expect(src).toMatch(/if \(!quotaStatus\.allowed\) \{/); + }); + + test('runSnapshot HONOURS the caller trigger — it is not hardcoded scheduled', async () => { + // Without this, a non-scheduled acquisition would be stamped SCHEDULED and + // could displace real scheduled evidence. + const { trace } = await runWithTrace('mlb', returns([]), { trigger: acq.TRIGGER.INTRADAY }); + expect(trace.trigger).toBe(acq.TRIGGER.INTRADAY); + expect(trace.trigger).not.toBe(acq.TRIGGER.SCHEDULED); + const client = { lpush: async () => { throw new Error('must not be called'); }, ltrim: async () => {}, expire: async () => {} }; + expect((await acq.persist(trace, { getRedisClient: () => client })).reason).toBe('not_scheduled'); + }); + + test('the scheduler passes only diagnostic context, not behaviour', () => { + const src = fs.readFileSync(path.join(ROOT, 'src/snapshotScheduler.js'), 'utf8'); + const i = src.indexOf('const results = await runAll({'); + const block = src.slice(i, i + 320); + expect(block).toMatch(/sports: scheduled/); + expect(block).toMatch(/trigger: acq\.TRIGGER\.SCHEDULED/); + expect(block).toMatch(/scheduledHourUtc: h/); + expect(block).toMatch(/processStartedAt: PROCESS_STARTED_AT/); + expect(block).not.toMatch(/retryDelayMs|getOdds|quota/); + }); +});