Scheduled acquisition trace: record the fork instead of inferring it
MLB 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 exactly two runSnapshot exits — getOdds THREW, or getOdds RETURNED ZERO PROPS — and production retained nothing able to tell them apart. A failed acquisition has no snapshot_id, writes no ledger row and updates no slate, so it left no durable evidence at all. The alert channel could not fill the gap either: quota/test-alert reports sent:true while the ntfy topic replays 0 messages, so absence of alerts is not evidence. This records the decisions the existing control flow already makes. It changes no acquisition behaviour: provider order, the single existing retry, the quota threshold, the fallback and cache policy are all untouched. The only behavioural line in the diff is getOdds gaining an optional recorder argument, defaulted to a frozen no-op so every existing caller is byte-identical. IT MAKES NO PROVIDER CALL. Measured: zero added fetchAllOdds/getProps/gateway calls, and the trace module contains no HTTP of any kind. Tests assert both. WHY REDIS, 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 — so process-local evidence could be written by a container nobody later probes. Storage is LPUSH + LTRIM, which is atomic: two schedulers racing the same slot both survive instead of one silently overwriting the other, and each attempt carries a process_generation so they stay distinguishable. WHY A BOUNDED HISTORY, 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 — precisely the evidence needed. Only SCHEDULED attempts are retained (persist refuses any other trigger), each sport keeps its own key, and intraday structurally cannot write one because intradayRefreshService never calls runSnapshot. WHAT IS CAPTURED, per attempt: cache decision; PropLine outcome as NONZERO/ZERO/ERROR/NOT_ATTEMPTED with count; the odds-api fallback with allowed_at_invocation, blocked_reason and the quota AS OBSERVED AT THAT INVOCATION — reading provider quota hours later and calling it historical evidence is the exact mistake this exists to stop; then the final result and the runSnapshot terminal outcome. The EXISTING retry appears as a second attempt; no retry was added. Sanitized: keys, tokens, URLs and long opaque strings are redacted, and no prop payload is retained — counts only. Tests assert a dirty provider error and a real prop array both come out clean. BEST-EFFORT AT THE CALL SITE, not just in the default dep — a teeth proof showed an injected store could still throw into a healthy snapshot. Now any implementation is safe. persist also auto-disables under NODE_ENV=test unless a client is injected (the opsNotify precedent); without that the default path opened a real ioredis connection inside every suite driving runSnapshot. Read-only GET /api/internal/acquisition/:sport behind the existing internal auth. It runs no pipeline and makes no fetch. Nine teeth against a green baseline, injections verified present: intraday overwrites scheduled (1) · one global slot (2) · thrown getOdds with no terminal trace (2) · zero mislabeled as error (1) · quota not captured at invocation (1) · observer makes a provider call (1) · credential leak (2) · scheduled/intraday share an identity (1) · telemetry failure breaks the snapshot (1). Restored byte-identically. TWO OF THOSE LANDED AND PASSED FIRST TIME — coverage holes, not safe defects: the provider-call scan did not forbid getOdds, and nothing exercised a non-scheduled trigger through runSnapshot. Both closed, then re-run failing. Model, retention, event identity, admission, dedupe, ledger, lineage, cadence, quota tracker and the PropLine adapter are all UNCHANGED. Lineage stays OFF. 386 suites / 5,210 tests pass. web tsc exit 0. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01CQJeAG8vcDoL5zkiaJyVb8
This commit is contained in:
@@ -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();
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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,
|
||||
};
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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/);
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user