Lineage productization: a live graph, an observer that doesn't trust it, and a surface that owns only what it knows
The authority review found there was nothing to switch: lineage is a
write-only graph with no living permanent writer and no product consumer.
Authority theater would have been a switch on a consumer that does not exist,
reading a store that is not being written. So: make the graph live, measure it
from outside, and expose the one category it alone owns.
WRITE-PATH FAILURE SEMANTICS, TRACED FIRST. persist(:857) -> the authoritative
cacheSet snapshot:latest(:1404) -> the lineage gate(:1426) -> ledger(:1473).
Lineage failure cannot fail base retention (earlier, separate upsert), cannot
fail publication (the product write precedes the gate), cannot fail the Ledger
(lineageIndex defaults null; the ledger has its own guard), and cannot create a
new partial row (blankLineage() first, atomicity sweep after the catch). The
isolation this tranche needed already existed; only a mode was missing.
THREE MODES, ONE EVALUATOR. `lineageWriteMode` distinguishes OFF /
CANARY_LEASED / PERSISTENT_SHADOW, and the write gate and the status surface
both read it, so they cannot disagree. The canary is CONSULTED, never
converted: its <=4h absolute expiry, its dynamic evaluation and its
fail-closed parse are untouched.
AN AMBIGUOUS CONFIGURATION FAILS CLOSED. If a sport is named by both persistent
mode and an active lease, the two instructions disagree about WHEN WRITING
STOPS — the lease says 22:45, persistent says never. The dangerous reading is
the quiet one: an operator sets a bounded lease believing writing will stop
while persistent keeps it going. We cannot know which they meant, so that sport
writes nothing until the configuration says one thing. The sport allowlist is
the canary's own, so persistent mode can never widen past it.
THE OBSERVER MAY NOT ASK THE WRITER HOW IT DID. settleLedger returned
{settled:0,pending:0} — byte-identical to a healthy "nothing to settle" — while
1,444 rows sat unprocessed, and the watchdog believed it. So `lineageCoverage`
reads durable retained state only, and THE DENOMINATOR MAY NOT CONSULT
lineage_action: eligibility is "the row was published AND a natural key is
derivable from its own identity columns", neither of which the lineage path
writes. If expectation were derived from whether lineage exists, coverage would
be 100% by construction and the metric would be decoration. Zero-expected and
zero-written are kept as different answers.
THE LEGACY BOUNDARY IS OBSERVED, NOT DECLARED. `publication_id` is stamped only
by commitPublication, so the row itself says whether lineage ran. Verified on
production: 5,353 rows carry it — 4,234 complete actions plus exactly the 1,119
historical partial rows — and zero actions exist without one. No epoch constant
is invented; a date would have been a guess about when the writer was on.
publication_id NULL -> LEGACY_UNVERIFIED. Stamped but incomplete ->
LINEAGE_UNAVAILABLE, which is the honest answer for the 1,119 and is never
quietly rewritten as legacy.
THE GRADE-SHIFT BADGE IS UNTOUCHED. revised_from_grade answers "did the letter
change"; lineage answers "which published claim superseded which". Different
questions, and a test now fails if either route learns the word lineage.
Suite 397/5,496/0 · tsc 0 · 15/15 teeth.
TWO OF MY OWN TESTS WERE VACUOUS AND A TOOTH FOUND IT. Tooth 2 came back green
because the isolation tests asserted the slate was published — true whether or
not the exception propagated — while never reaching the lineage gate at all:
the fake grader never fired `onGraded`, so the collector stayed empty and
persistedRows stayed null. Fixed by firing the hook and counting the commit.
A green teeth run means the test is missing.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01CQJeAG8vcDoL5zkiaJyVb8
This commit is contained in:
@@ -251,6 +251,16 @@ router.get('/snapshot/status', async (req, res) => {
|
||||
const lineageCanary = require('../services/lineageCanaryConfig');
|
||||
// Session 56 — surface the missed-cron signal in the health probe.
|
||||
const mlbTs = last_snapshot.mlb && last_snapshot.mlb.updated_at;
|
||||
// Best-effort: the health probe must never fail because the audit could not
|
||||
// run. An unavailable audit is reported as AUDIT_UNAVAILABLE, which is a
|
||||
// different answer from "coverage is zero".
|
||||
let lineageCoverage = null;
|
||||
try {
|
||||
lineageCoverage = await require('../services/lineageCoverage')
|
||||
.auditLatestCohort('mlb', { now: () => new Date().toISOString() });
|
||||
} catch (e) {
|
||||
lineageCoverage = { audit_available: false, reason: e.message, health: 'AUDIT_UNAVAILABLE' };
|
||||
}
|
||||
return res.json({
|
||||
runtime: {
|
||||
code_sha: codeSha(),
|
||||
@@ -262,6 +272,14 @@ router.get('/snapshot/status', async (req, res) => {
|
||||
// change without a new process, so `runtime.started_at` is a defensible
|
||||
// lower bound for how long this state has held.
|
||||
lineage_canary: lineageCanary.state(),
|
||||
// THE WRITE MODE, from the SAME evaluator the write gate calls. Reporting
|
||||
// it from a second parse is how a status surface and a gate come to
|
||||
// disagree, so there is only one.
|
||||
lineage_write_mode: require('../services/lineageWriteMode').state(),
|
||||
// INDEPENDENT COVERAGE. Derived from durable retained state, never from
|
||||
// the writer's own counters — an observer that reads the failing writer's
|
||||
// return value cannot see that writer fail.
|
||||
lineage_coverage: lineageCoverage,
|
||||
// Latest TERMINAL retention result per sport, exactly as the persistence
|
||||
// function reported it. Chunked writes stop at the first failed chunk, so
|
||||
// rows existing under a snapshot_id does not mean the cohort is complete —
|
||||
@@ -312,6 +330,54 @@ router.post('/snapshot/:sport', async (req, res) => {
|
||||
* is the self-learning loop's write path (the public /api/accuracy is read-only).
|
||||
* Registered BEFORE /outcomes/:sport so "all" isn't captured as a sport.
|
||||
*/
|
||||
/**
|
||||
* ADDITIVE READ ANCESTRY — non-authoritative, replaces nothing.
|
||||
*
|
||||
* The category lineage alone owns: which published Read came first, what
|
||||
* superseded what, and whether a later observation revised or recaptured the
|
||||
* same claim. It does NOT serve the grade-shift badge, which answers a
|
||||
* different question from a different source and is untouched.
|
||||
*
|
||||
* Internal-only on purpose. Nothing user-facing consumes lineage yet, and the
|
||||
* narrowest seam that exposes the truth without altering an existing response
|
||||
* is a protected route.
|
||||
*/
|
||||
router.get('/lineage/ancestry', async (req, res) => {
|
||||
const q = req.query || {};
|
||||
const required = ['sport', 'game_date', 'player_key', 'stat', 'side', 'line'];
|
||||
const missing = required.filter((k) => q[k] === undefined || q[k] === '');
|
||||
if (missing.length) {
|
||||
return res.status(400).json({ error: 'missing required identity fields', missing });
|
||||
}
|
||||
try {
|
||||
const ancestry = require('../services/read/readAncestry');
|
||||
const out = await ancestry.ancestryForRead({
|
||||
sport: String(q.sport).toLowerCase(),
|
||||
game_date: String(q.game_date),
|
||||
player_key: String(q.player_key),
|
||||
stat: String(q.stat),
|
||||
side: String(q.side),
|
||||
line: Number(q.line),
|
||||
canonical_event_id: q.canonical_event_id ? String(q.canonical_event_id) : null,
|
||||
});
|
||||
return res.json(out);
|
||||
} catch (e) {
|
||||
return res.status(500).json({ error: e.message });
|
||||
}
|
||||
});
|
||||
|
||||
/** Coverage audit for one sport, on demand — the same independent observer the
|
||||
* status probe reports, addressable when a specific cohort is in question. */
|
||||
router.get('/lineage/coverage/:sport', async (req, res) => {
|
||||
try {
|
||||
const out = await require('../services/lineageCoverage')
|
||||
.auditLatestCohort(req.params.sport, { now: () => new Date().toISOString() });
|
||||
return res.json(out);
|
||||
} catch (e) {
|
||||
return res.status(500).json({ error: e.message });
|
||||
}
|
||||
});
|
||||
|
||||
router.post('/outcomes/all', async (req, res) => {
|
||||
const outcomes = require('../services/outcomeService');
|
||||
try {
|
||||
|
||||
@@ -0,0 +1,232 @@
|
||||
/**
|
||||
* LINEAGE COVERAGE — an observer that does not trust the writer.
|
||||
*
|
||||
* The settlement outage and the retention outage taught the same lesson twice:
|
||||
* an observer that reads the failing writer's own return value cannot see that
|
||||
* writer fail. `settleLedger` returned `{settled:0, pending:0}` — byte-identical
|
||||
* to a healthy "nothing to settle" — while 1,444 rows sat unprocessed, and the
|
||||
* zero-settle watchdog believed it.
|
||||
*
|
||||
* So nothing here reads attachLineage's counters, its logs, or its terminal
|
||||
* record. Everything is derived from durable retained state.
|
||||
*
|
||||
* THE DENOMINATOR MAY NOT CONSULT `lineage_action`.
|
||||
*
|
||||
* If "was this claim expected to have lineage?" were answered by looking at
|
||||
* whether it HAS lineage, coverage would be 100% by construction and the metric
|
||||
* would be decoration. Eligibility is therefore derived from two durable facts
|
||||
* that exist independently of the lineage writer:
|
||||
*
|
||||
* 1. the row was PUBLISHED (`published === true`), which is what
|
||||
* commitPublication itself filters on, and
|
||||
* 2. a Read natural key is DERIVABLE from the row's own identity columns
|
||||
* (sport, game_date, player_key, stat, side, line, canonical event).
|
||||
*
|
||||
* Neither is written by the lineage path.
|
||||
*/
|
||||
|
||||
const readLineage = require('./read/readLineage');
|
||||
|
||||
const HEALTH = Object.freeze({
|
||||
HEALTHY: 'HEALTHY',
|
||||
NO_ELIGIBLE_CLAIMS: 'NO_ELIGIBLE_CLAIMS',
|
||||
MISSING_COVERAGE: 'MISSING_COVERAGE',
|
||||
PARTIAL_COVERAGE: 'PARTIAL_COVERAGE',
|
||||
INVALID_GRAPH: 'INVALID_GRAPH',
|
||||
STALE: 'STALE',
|
||||
AUDIT_UNAVAILABLE: 'AUDIT_UNAVAILABLE',
|
||||
});
|
||||
|
||||
/** The nine fields that together are a completed lineage action. Mirrors
|
||||
* retentionService.VALID_LINEAGE_ACTION_FIELDS; kept local so the observer
|
||||
* carries no dependency on the writer module, and cross-checked by a test. */
|
||||
const VALID_ACTION_FIELDS = Object.freeze([
|
||||
'read_id', 'read_natural_key', 'lineage_action', 'claim_digest',
|
||||
'revision_ordinal', 'lineage_state', 'lineage_version',
|
||||
'claim_schema_version', 'digest_algorithm_version',
|
||||
]);
|
||||
|
||||
function isValidAction(row) {
|
||||
if (!row) return false;
|
||||
return VALID_ACTION_FIELDS.every((f) => row[f] !== null && row[f] !== undefined);
|
||||
}
|
||||
|
||||
/**
|
||||
* ELIGIBILITY — pure, and deliberately blind to every lineage column.
|
||||
* Returns the set of natural keys that SHOULD carry lineage.
|
||||
*/
|
||||
function expectedKeys(rows) {
|
||||
const keys = new Set();
|
||||
let published = 0;
|
||||
let unkeyable = 0;
|
||||
for (const r of Array.isArray(rows) ? rows : []) {
|
||||
if (!r || r.published !== true) continue;
|
||||
published += 1;
|
||||
const k = readLineage.readNaturalKey(r);
|
||||
if (!k) { unkeyable += 1; continue; }
|
||||
keys.add(k);
|
||||
}
|
||||
return { keys, published_rows: published, unkeyable_rows: unkeyable };
|
||||
}
|
||||
|
||||
/** Keys that actually carry a complete, valid lineage action. */
|
||||
function coveredKeys(rows) {
|
||||
const keys = new Set();
|
||||
let actions = 0;
|
||||
let partial = 0;
|
||||
const byAction = { ORIGIN: 0, REVISION: 0, RECAPTURE: 0 };
|
||||
for (const r of Array.isArray(rows) ? rows : []) {
|
||||
if (!r) continue;
|
||||
const touched = VALID_ACTION_FIELDS.some((f) => r[f] !== null && r[f] !== undefined);
|
||||
if (!touched) continue;
|
||||
if (!isValidAction(r)) { partial += 1; continue; }
|
||||
actions += 1;
|
||||
if (byAction[r.lineage_action] !== undefined) byAction[r.lineage_action] += 1;
|
||||
keys.add(r.read_natural_key);
|
||||
}
|
||||
return { keys, actions, partial, byAction };
|
||||
}
|
||||
|
||||
/**
|
||||
* Graph defects detectable from the cohort alone: a read_id whose rows do not
|
||||
* form one coherent chain. `chronology` refuses rather than electing a head, so
|
||||
* its refusal reason IS the defect description.
|
||||
*/
|
||||
function graphDefects(rows) {
|
||||
const byRead = new Map();
|
||||
for (const r of Array.isArray(rows) ? rows : []) {
|
||||
if (!isValidAction(r)) continue;
|
||||
if (!byRead.has(r.read_id)) byRead.set(r.read_id, []);
|
||||
byRead.get(r.read_id).push(r);
|
||||
}
|
||||
const defects = [];
|
||||
for (const [readId, chain] of byRead) {
|
||||
// A cohort slice legitimately holds only part of a chain, so a chain whose
|
||||
// ORIGIN was written by an EARLIER cohort is not a defect. Only a chain
|
||||
// carrying more than one ORIGIN, or a duplicated ordinal, is.
|
||||
const origins = chain.filter((r) => r.lineage_action === 'ORIGIN').length;
|
||||
const ordinals = chain.map((r) => r.revision_ordinal);
|
||||
const dupOrdinal = ordinals.length !== new Set(ordinals).size;
|
||||
if (origins > 1) defects.push({ read_id: readId, defect: 'MULTIPLE_ORIGIN', origins });
|
||||
else if (dupOrdinal) defects.push({ read_id: readId, defect: 'DUPLICATE_ORDINAL', ordinals });
|
||||
}
|
||||
return defects;
|
||||
}
|
||||
|
||||
/** Pure classification, so the decision rule is testable without a database. */
|
||||
function classify(summary) {
|
||||
if (summary.audit_available === false) return HEALTH.AUDIT_UNAVAILABLE;
|
||||
if (summary.invalid_partial_actions > 0 || summary.graph_defects > 0 || summary.extra_keys > 0) {
|
||||
return HEALTH.INVALID_GRAPH;
|
||||
}
|
||||
// ZERO EXPECTED AND ZERO WRITTEN ARE DIFFERENT ANSWERS and must never be
|
||||
// collapsed: an off-hours slate with nothing to publish is healthy silence,
|
||||
// and a full slate with no lineage is the writer having never run.
|
||||
if (summary.expected_keys === 0) return HEALTH.NO_ELIGIBLE_CLAIMS;
|
||||
if (summary.covered_keys === 0) return HEALTH.MISSING_COVERAGE;
|
||||
if (summary.covered_keys < summary.expected_keys) return HEALTH.PARTIAL_COVERAGE;
|
||||
return HEALTH.HEALTHY;
|
||||
}
|
||||
|
||||
/** Everything a cohort audit reports. Pure over rows. */
|
||||
function auditRows(rows, meta = {}) {
|
||||
const exp = expectedKeys(rows);
|
||||
const cov = coveredKeys(rows);
|
||||
const missing = [...exp.keys].filter((k) => !cov.keys.has(k));
|
||||
const extra = [...cov.keys].filter((k) => !exp.keys.has(k));
|
||||
const defects = graphDefects(rows);
|
||||
const summary = {
|
||||
snapshot_id: meta.snapshot_id ?? null,
|
||||
sport: meta.sport ?? null,
|
||||
game_date: meta.game_date ?? null,
|
||||
code_sha: meta.code_sha ?? null,
|
||||
captured_at: meta.captured_at ?? null,
|
||||
retention_terminal: meta.retention_terminal ?? null,
|
||||
audit_available: true,
|
||||
rows: Array.isArray(rows) ? rows.length : 0,
|
||||
published_rows: exp.published_rows,
|
||||
unkeyable_published_rows: exp.unkeyable_rows,
|
||||
expected_keys: exp.keys.size,
|
||||
covered_keys: cov.keys.size,
|
||||
missing_keys: missing.length,
|
||||
extra_keys: extra.length,
|
||||
coverage_pct: exp.keys.size === 0 ? null
|
||||
: Math.round((cov.keys.size / exp.keys.size) * 10000) / 100,
|
||||
lineage_actions: cov.actions,
|
||||
origins: cov.byAction.ORIGIN,
|
||||
revisions: cov.byAction.REVISION,
|
||||
recaptures: cov.byAction.RECAPTURE,
|
||||
invalid_partial_actions: cov.partial,
|
||||
graph_defects: defects.length,
|
||||
};
|
||||
summary.health = classify(summary);
|
||||
return Object.freeze({
|
||||
...summary,
|
||||
missing_sample: Object.freeze(missing.slice(0, 10)),
|
||||
extra_sample: Object.freeze(extra.slice(0, 10)),
|
||||
defect_sample: Object.freeze(defects.slice(0, 10)),
|
||||
});
|
||||
}
|
||||
|
||||
const AUDIT_COLUMNS = [
|
||||
// identity + eligibility — none of these is written by the lineage path
|
||||
'id', 'snapshot_id', 'sport', 'game_date', 'player_key', 'stat', 'side', 'line',
|
||||
'canonical_event_id', 'event_occurrence', 'published', 'captured_at', 'code_sha',
|
||||
// lineage state — read to MEASURE coverage, never to decide eligibility
|
||||
...VALID_ACTION_FIELDS, 'supersedes_id', 'recaptures_id', 'change_type',
|
||||
].join(', ');
|
||||
|
||||
/**
|
||||
* Audit the most recent cohort that could have produced lineage, chosen from
|
||||
* durable retained state. This is the "writer never ran" detector: the cohort
|
||||
* is selected without reference to whether any lineage exists, so a slate that
|
||||
* published and recorded nothing surfaces as MISSING_COVERAGE rather than as
|
||||
* silence.
|
||||
*/
|
||||
async function auditLatestCohort(sport, deps = {}) {
|
||||
const sp = String(sport || '').toLowerCase();
|
||||
try {
|
||||
const getClient = deps.getClient || require('../utils/supabase').getSupabaseServiceClient;
|
||||
const supabase = getClient();
|
||||
if (!supabase) return Object.freeze({ audit_available: false, reason: 'no_supabase_env', sport: sp, health: HEALTH.AUDIT_UNAVAILABLE });
|
||||
|
||||
const { data: head, error: headErr } = await supabase
|
||||
.from('model_snapshots')
|
||||
.select('snapshot_id, game_date, captured_at, code_sha')
|
||||
.eq('sport', sp).eq('published', true)
|
||||
.order('captured_at', { ascending: false }).limit(1);
|
||||
if (headErr) return Object.freeze({ audit_available: false, reason: headErr.message, sport: sp, health: HEALTH.AUDIT_UNAVAILABLE });
|
||||
if (!head || head.length === 0) {
|
||||
return Object.freeze({ audit_available: true, sport: sp, expected_keys: 0, covered_keys: 0, health: HEALTH.NO_ELIGIBLE_CLAIMS, reason: 'no published rows' });
|
||||
}
|
||||
const h = head[0];
|
||||
|
||||
const { paginate } = require('../utils/safePaginate');
|
||||
const rows = await paginate(() => supabase
|
||||
.from('model_snapshots')
|
||||
.select(AUDIT_COLUMNS)
|
||||
.eq('snapshot_id', h.snapshot_id)
|
||||
.eq('sport', sp)
|
||||
.eq('game_date', h.game_date), { key: 'id', label: 'lineageCoverage.auditLatestCohort' });
|
||||
|
||||
const audit = auditRows(rows, {
|
||||
snapshot_id: h.snapshot_id, sport: sp, game_date: h.game_date,
|
||||
code_sha: h.code_sha, captured_at: h.captured_at,
|
||||
retention_terminal: deps.retentionTerminal ?? null,
|
||||
});
|
||||
const ageMs = deps.now ? (new Date(deps.now()).getTime() - new Date(h.captured_at).getTime()) : null;
|
||||
const staleAfter = deps.staleAfterMs ?? (6 * 60 * 60 * 1000);
|
||||
if (ageMs !== null && ageMs > staleAfter && audit.health === HEALTH.HEALTHY) {
|
||||
return Object.freeze({ ...audit, health: HEALTH.STALE, age_ms: ageMs });
|
||||
}
|
||||
return Object.freeze({ ...audit, age_ms: ageMs });
|
||||
} catch (e) {
|
||||
return Object.freeze({ audit_available: false, reason: e && e.message ? e.message : String(e), sport: sp, health: HEALTH.AUDIT_UNAVAILABLE });
|
||||
}
|
||||
}
|
||||
|
||||
module.exports = {
|
||||
HEALTH, VALID_ACTION_FIELDS, AUDIT_COLUMNS,
|
||||
isValidAction, expectedKeys, coveredKeys, graphDefects, classify,
|
||||
auditRows, auditLatestCohort,
|
||||
};
|
||||
@@ -0,0 +1,152 @@
|
||||
/**
|
||||
* LINEAGE WRITE MODE — one evaluator, three modes.
|
||||
*
|
||||
* CANARY PERMISSION AND PERSISTENT SHADOW WRITING ARE DIFFERENT THINGS.
|
||||
*
|
||||
* CANARY_LEASED temporary permission to shadow-write, bounded by an
|
||||
* absolute expiry of at most four hours, evaluated at the
|
||||
* write gate on every call. Proven to transition
|
||||
* ACTIVE -> EXPIRED on one unrestarted process generation.
|
||||
* Its purpose is a contained experiment.
|
||||
* PERSISTENT_SHADOW continuous NON-AUTHORITATIVE evidence collection. No
|
||||
* expiry, because the thing it is collecting is a
|
||||
* continuous history and a graph with holes in it is worth
|
||||
* less than no graph at all.
|
||||
* OFF no lineage attachment.
|
||||
*
|
||||
* The canary is NOT converted into a permanent mode and its contract is not
|
||||
* touched: `lineageCanaryConfig` is consulted, never modified.
|
||||
*
|
||||
* WHY AN AMBIGUOUS CONFIGURATION FAILS CLOSED.
|
||||
*
|
||||
* If a sport is named by BOTH persistent mode and an active lease, the two
|
||||
* instructions disagree about WHEN WRITING STOPS. The lease says "at 22:45";
|
||||
* persistent says "never". They agree on today and differ on tomorrow, and the
|
||||
* dangerous reading is the quiet one: an operator sets a bounded lease
|
||||
* believing writing will stop, while persistent mode keeps it going past the
|
||||
* expiry they were relying on. We cannot know which they meant, so that sport
|
||||
* writes nothing until the configuration says one thing.
|
||||
*/
|
||||
|
||||
const canary = require('./lineageCanaryConfig');
|
||||
|
||||
const WRITE_MODE = Object.freeze({
|
||||
OFF: 'OFF',
|
||||
CANARY_LEASED: 'CANARY_LEASED',
|
||||
PERSISTENT_SHADOW: 'PERSISTENT_SHADOW',
|
||||
});
|
||||
|
||||
const CONFLICT = Object.freeze({
|
||||
SPORT_IN_BOTH_MODES: 'SPORT_IN_BOTH_MODES',
|
||||
});
|
||||
|
||||
/**
|
||||
* Reuses the canary's allowlist rather than declaring a second one. A sport
|
||||
* becomes eligible for lineage in ONE place, so persistent mode can never
|
||||
* silently widen past what the canary was permitted to reach.
|
||||
*/
|
||||
const PERSISTABLE_SPORTS = canary.LEASABLE_SPORTS;
|
||||
|
||||
const CONFIGURATION_SOURCE = 'ENVIRONMENT';
|
||||
const ENV_VAR = 'LINEAGE_PERSISTENT_SPORTS';
|
||||
|
||||
/** Parse the persistent list. Repo-native shape: a comma-separated sport list,
|
||||
* the same as LINEAGE_CANARY_SPORTS minus the lease syntax. Never a decision. */
|
||||
function parsePersistent(raw) {
|
||||
const text = typeof raw === 'string' ? raw.trim() : '';
|
||||
if (!text) return Object.freeze({ configured: false, sports: Object.freeze([]), invalid_reason: null });
|
||||
const tokens = text.split(',').map((t) => t.trim().toLowerCase()).filter(Boolean);
|
||||
if (tokens.length === 0) {
|
||||
return Object.freeze({ configured: true, sports: Object.freeze([]), invalid_reason: 'INVALID_EMPTY_TOKENS' });
|
||||
}
|
||||
const bad = tokens.filter((t) => !PERSISTABLE_SPORTS.includes(t));
|
||||
if (bad.length > 0) {
|
||||
// Fail closed on the WHOLE list. Silently keeping the valid half would
|
||||
// enable a subset the operator never asked for on its own.
|
||||
return Object.freeze({ configured: true, sports: Object.freeze([]), invalid_reason: 'INVALID_SPORT_NOT_PERSISTABLE' });
|
||||
}
|
||||
return Object.freeze({ configured: true, sports: Object.freeze([...new Set(tokens)]), invalid_reason: null });
|
||||
}
|
||||
|
||||
/**
|
||||
* THE ONE EVALUATOR. Every write-gate decision and every status line comes
|
||||
* from here, so the gate and the status can never disagree about the mode.
|
||||
*
|
||||
* `now` is threaded because the canary lease is evaluated against the clock at
|
||||
* WRITE TIME, not at process start.
|
||||
*/
|
||||
function evaluate(now) {
|
||||
const persistent = parsePersistent(process.env[ENV_VAR]);
|
||||
const lease = canary.evaluate(now);
|
||||
// FIELD NAMES ARE THE CONTRACT. The lease exposes `effective_active_sports`
|
||||
// and `lease_state`; reading `active_sports`/`state` returns undefined and
|
||||
// would silently report every lease as inactive — disabling the canary path
|
||||
// through this evaluator while looking perfectly healthy.
|
||||
const leaseSports = Array.isArray(lease.effective_active_sports) ? lease.effective_active_sports : [];
|
||||
|
||||
const conflictSports = persistent.sports.filter((s) => leaseSports.includes(s));
|
||||
const persistentEffective = persistent.sports.filter((s) => !conflictSports.includes(s));
|
||||
const canaryEffective = leaseSports.filter((s) => !conflictSports.includes(s));
|
||||
|
||||
const modes = {};
|
||||
for (const s of canaryEffective) modes[s] = WRITE_MODE.CANARY_LEASED;
|
||||
// Persistent is assigned second only because the two sets are disjoint by
|
||||
// construction above; it is not a precedence rule.
|
||||
for (const s of persistentEffective) modes[s] = WRITE_MODE.PERSISTENT_SHADOW;
|
||||
|
||||
const sports = Object.keys(modes).sort();
|
||||
let mode = WRITE_MODE.OFF;
|
||||
if (persistentEffective.length > 0) mode = WRITE_MODE.PERSISTENT_SHADOW;
|
||||
else if (canaryEffective.length > 0) mode = WRITE_MODE.CANARY_LEASED;
|
||||
|
||||
return Object.freeze({
|
||||
mode,
|
||||
sports: Object.freeze(sports),
|
||||
modes: Object.freeze(modes),
|
||||
persistent_configured: persistent.configured,
|
||||
persistent_sports: persistent.sports,
|
||||
persistent_invalid_reason: persistent.invalid_reason,
|
||||
canary_lease_state: lease.lease_state,
|
||||
canary_sports: Object.freeze([...leaseSports]),
|
||||
conflict: conflictSports.length > 0 ? CONFLICT.SPORT_IN_BOTH_MODES : null,
|
||||
conflict_sports: Object.freeze(conflictSports),
|
||||
configuration_source: CONFIGURATION_SOURCE,
|
||||
});
|
||||
}
|
||||
|
||||
/** The write gate's only question. */
|
||||
function isEnabled(sport, now) {
|
||||
const sp = String(sport || '').toLowerCase();
|
||||
if (!sp) return false;
|
||||
return Object.prototype.hasOwnProperty.call(evaluate(now).modes, sp);
|
||||
}
|
||||
|
||||
/** Which mode a sport is writing under, for status and for the observer. */
|
||||
function modeFor(sport, now) {
|
||||
const sp = String(sport || '').toLowerCase();
|
||||
return evaluate(now).modes[sp] || WRITE_MODE.OFF;
|
||||
}
|
||||
|
||||
/** Status surface. Never returns the raw environment value. */
|
||||
function state(now) {
|
||||
const e = evaluate(now);
|
||||
return Object.freeze({
|
||||
mode: e.mode,
|
||||
effective_sports: e.sports,
|
||||
modes: e.modes,
|
||||
persistent_configured: e.persistent_configured,
|
||||
persistent_sports: e.persistent_sports,
|
||||
persistent_invalid_reason: e.persistent_invalid_reason,
|
||||
canary_lease_state: e.canary_lease_state,
|
||||
canary_sports: e.canary_sports,
|
||||
conflict: e.conflict,
|
||||
conflict_sports: e.conflict_sports,
|
||||
configuration_source: e.configuration_source,
|
||||
authority: 'NON_AUTHORITATIVE',
|
||||
});
|
||||
}
|
||||
|
||||
module.exports = {
|
||||
WRITE_MODE, CONFLICT, PERSISTABLE_SPORTS, ENV_VAR, CONFIGURATION_SOURCE,
|
||||
parsePersistent, evaluate, isEnabled, modeFor, state,
|
||||
};
|
||||
@@ -0,0 +1,205 @@
|
||||
/**
|
||||
* READ ANCESTRY — the additive, NON-AUTHORITATIVE history surface.
|
||||
*
|
||||
* This replaces nothing. `revised_from_grade` and the grade-shift badge answer
|
||||
* a DIFFERENT question — "did the model's letter change?" — and are untouched.
|
||||
* This answers the question only lineage can: which published Read came first,
|
||||
* what superseded what, and whether a later observation was a revision or a
|
||||
* recapture of the same claim.
|
||||
*
|
||||
* THE LEGACY BOUNDARY IS OBSERVED, NOT DECLARED.
|
||||
*
|
||||
* `publication_id` is stamped by `commitPublication`, which is the only code
|
||||
* that ever attempts lineage. So the row itself says whether lineage ran:
|
||||
*
|
||||
* publication_id NULL lineage never ran here LEGACY_UNVERIFIED
|
||||
* publication_id + valid action lineage ran and completed LINEAGE_AVAILABLE
|
||||
* publication_id + no valid action lineage ran and did not LINEAGE_UNAVAILABLE
|
||||
* valid actions that do not form one chain INVALID_LINEAGE
|
||||
*
|
||||
* Verified on production 2026-08-31: 5,353 rows carry publication_id — 4,234
|
||||
* complete actions plus exactly the 1,119 historical partial rows — and ZERO
|
||||
* lineage actions exist without one. No epoch constant is needed and none is
|
||||
* invented; a date-based boundary would have been a guess about when the
|
||||
* writer was on.
|
||||
*/
|
||||
|
||||
const readLineage = require('./readLineage');
|
||||
|
||||
const ANCESTRY_STATE = Object.freeze({
|
||||
LINEAGE_AVAILABLE: 'LINEAGE_AVAILABLE',
|
||||
LEGACY_UNVERIFIED: 'LEGACY_UNVERIFIED',
|
||||
LINEAGE_UNAVAILABLE: 'LINEAGE_UNAVAILABLE',
|
||||
INVALID_LINEAGE: 'INVALID_LINEAGE',
|
||||
NOT_FOUND: 'NOT_FOUND',
|
||||
});
|
||||
|
||||
const VALID_ACTION_FIELDS = Object.freeze([
|
||||
'read_id', 'read_natural_key', 'lineage_action', 'claim_digest',
|
||||
'revision_ordinal', 'lineage_state', 'lineage_version',
|
||||
'claim_schema_version', 'digest_algorithm_version',
|
||||
]);
|
||||
|
||||
const isValidAction = (r) => !!r && VALID_ACTION_FIELDS.every((f) => r[f] !== null && r[f] !== undefined);
|
||||
|
||||
/**
|
||||
* WHAT LINEAGE ACTUALLY KNOWS. Deliberately narrow.
|
||||
*
|
||||
* Lineage knows the shape of the published record: identity, order, parentage,
|
||||
* digest and the versions those were computed under. It does NOT know why the
|
||||
* model believed anything, whether a price was right, or how the bet settled,
|
||||
* and this projection refuses to imply otherwise by carrying those fields.
|
||||
*/
|
||||
function project(row) {
|
||||
return Object.freeze({
|
||||
snapshot_row_id: row.id ?? null,
|
||||
read_id: row.read_id,
|
||||
read_natural_key: row.read_natural_key,
|
||||
action: row.lineage_action,
|
||||
change_type: row.change_type ?? null,
|
||||
revision_ordinal: row.revision_ordinal,
|
||||
supersedes_row_id: row.supersedes_id ?? null,
|
||||
recaptures_row_id: row.recaptures_id ?? null,
|
||||
claim_digest: row.claim_digest,
|
||||
claim_schema_version: row.claim_schema_version,
|
||||
digest_algorithm_version: row.digest_algorithm_version,
|
||||
lineage_version: row.lineage_version,
|
||||
lineage_state: row.lineage_state,
|
||||
canonical_event_id: row.canonical_event_id ?? null,
|
||||
event_occurrence: row.event_occurrence ?? null,
|
||||
captured_at: row.captured_at ?? null,
|
||||
published_at: row.published_at ?? null,
|
||||
publication_id: row.publication_id ?? null,
|
||||
snapshot_id: row.snapshot_id ?? null,
|
||||
model_version: row.model_version ?? null,
|
||||
code_sha: row.code_sha ?? null,
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* Pure classification over the rows of ONE logical Read. No database, so the
|
||||
* decision rule is testable and cannot drift from what the route serves.
|
||||
*/
|
||||
function classifyAncestry(rows) {
|
||||
const list = (Array.isArray(rows) ? rows : []).filter(Boolean);
|
||||
if (list.length === 0) {
|
||||
return Object.freeze({ state: ANCESTRY_STATE.NOT_FOUND, reason: 'no retained rows for this Read', revisions: Object.freeze([]) });
|
||||
}
|
||||
const attempted = list.filter((r) => r.publication_id !== null && r.publication_id !== undefined);
|
||||
const actions = list.filter(isValidAction);
|
||||
|
||||
if (actions.length === 0) {
|
||||
if (attempted.length === 0) {
|
||||
// Lineage never ran for this Read. It is not missing history; there was
|
||||
// never a writer. Nothing is synthesized and no parent is inferred.
|
||||
return Object.freeze({
|
||||
state: ANCESTRY_STATE.LEGACY_UNVERIFIED,
|
||||
reason: 'published before lineage recording began',
|
||||
rows_retained: list.length,
|
||||
revisions: Object.freeze([]),
|
||||
});
|
||||
}
|
||||
// Lineage DID run and produced nothing usable. That is a real gap and it
|
||||
// must be visible: silently answering LEGACY here would rewrite a failure
|
||||
// as an absence.
|
||||
return Object.freeze({
|
||||
state: ANCESTRY_STATE.LINEAGE_UNAVAILABLE,
|
||||
reason: 'lineage was attempted for this Read and did not complete',
|
||||
rows_retained: list.length,
|
||||
rows_attempted: attempted.length,
|
||||
revisions: Object.freeze([]),
|
||||
});
|
||||
}
|
||||
|
||||
const byRead = new Map();
|
||||
for (const a of actions) {
|
||||
if (!byRead.has(a.read_id)) byRead.set(a.read_id, []);
|
||||
byRead.get(a.read_id).push(a);
|
||||
}
|
||||
if (byRead.size > 1) {
|
||||
return Object.freeze({
|
||||
state: ANCESTRY_STATE.INVALID_LINEAGE,
|
||||
reason: `rows span ${byRead.size} read_ids for one natural key`,
|
||||
read_ids: Object.freeze([...byRead.keys()]),
|
||||
revisions: Object.freeze([]),
|
||||
});
|
||||
}
|
||||
const chain = [...byRead.values()][0];
|
||||
const chron = readLineage.chronology(chain);
|
||||
if (!chron.ok) {
|
||||
// `chronology` REFUSES rather than electing a head from broken evidence,
|
||||
// and its refusal reason is the defect description.
|
||||
return Object.freeze({
|
||||
state: ANCESTRY_STATE.INVALID_LINEAGE,
|
||||
reason: chron.reason,
|
||||
revisions: Object.freeze([]),
|
||||
});
|
||||
}
|
||||
const recaptures = chain.filter((r) => r.lineage_action === readLineage.LINEAGE_ACTION.RECAPTURE);
|
||||
return Object.freeze({
|
||||
state: ANCESTRY_STATE.LINEAGE_AVAILABLE,
|
||||
reason: null,
|
||||
read_id: chron.read_id,
|
||||
read_natural_key: chain[0].read_natural_key,
|
||||
revision_count: chron.revision_count,
|
||||
recapture_count: chron.recapture_count,
|
||||
original: chron.original ? project(chron.original) : null,
|
||||
current: chron.current ? project(chron.current) : null,
|
||||
revisions: Object.freeze(chron.revisions.map(project)),
|
||||
recaptures: Object.freeze(recaptures.map(project)),
|
||||
// Said plainly on every response so no consumer can mistake this for a
|
||||
// source of product truth.
|
||||
authority: 'NON_AUTHORITATIVE',
|
||||
});
|
||||
}
|
||||
|
||||
const ANCESTRY_COLUMNS = [
|
||||
'id', 'snapshot_id', 'sport', 'game_date', 'player_key', 'stat', 'side', 'line',
|
||||
'canonical_event_id', 'event_occurrence', 'captured_at', 'published', 'published_at',
|
||||
'publication_id', 'model_version', 'code_sha',
|
||||
...VALID_ACTION_FIELDS, 'supersedes_id', 'recaptures_id', 'change_type',
|
||||
].join(', ');
|
||||
|
||||
/**
|
||||
* Look a Read up by its IDENTITY COLUMNS, never by `read_natural_key`.
|
||||
*
|
||||
* A row whose lineage failed has no natural key on it — that is the failure
|
||||
* atomicity contract — so querying by the lineage column would make exactly
|
||||
* the rows that need a LINEAGE_UNAVAILABLE answer invisible, and they would
|
||||
* come back NOT_FOUND instead.
|
||||
*/
|
||||
async function ancestryForRead(spec, deps = {}) {
|
||||
const s = spec || {};
|
||||
try {
|
||||
const getClient = deps.getClient || require('../../utils/supabase').getSupabaseServiceClient;
|
||||
const supabase = getClient();
|
||||
if (!supabase) return Object.freeze({ state: ANCESTRY_STATE.NOT_FOUND, reason: 'no supabase env', revisions: Object.freeze([]) });
|
||||
|
||||
const { paginate } = require('../../utils/safePaginate');
|
||||
const rows = await paginate(() => {
|
||||
let q = supabase.from('model_snapshots').select(ANCESTRY_COLUMNS)
|
||||
.eq('sport', String(s.sport || '').toLowerCase())
|
||||
.eq('game_date', s.game_date)
|
||||
.eq('player_key', s.player_key)
|
||||
.eq('stat', s.stat)
|
||||
.eq('side', s.side)
|
||||
.eq('line', s.line);
|
||||
if (s.canonical_event_id) q = q.eq('canonical_event_id', s.canonical_event_id);
|
||||
return q;
|
||||
}, { key: 'id', label: 'readAncestry.ancestryForRead' });
|
||||
|
||||
const out = classifyAncestry(rows);
|
||||
return Object.freeze({ ...out, query: Object.freeze({ ...s }) });
|
||||
} catch (e) {
|
||||
return Object.freeze({
|
||||
state: ANCESTRY_STATE.LINEAGE_UNAVAILABLE,
|
||||
reason: `ancestry lookup failed: ${e && e.message ? e.message : String(e)}`,
|
||||
revisions: Object.freeze([]),
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
module.exports = {
|
||||
ANCESTRY_STATE, ANCESTRY_COLUMNS, VALID_ACTION_FIELDS,
|
||||
isValidAction, project, classifyAncestry, ancestryForRead,
|
||||
};
|
||||
@@ -31,6 +31,7 @@ const STATS_CONCURRENCY = 5;
|
||||
const { nameKey, normalizeName } = require('../utils/playerName');
|
||||
const acq = require('./ops/acquisitionTrace');
|
||||
const participantIdentity = require('./model/participantIdentity');
|
||||
const lineageWriteMode = require('./lineageWriteMode');
|
||||
// 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);
|
||||
@@ -327,6 +328,19 @@ function lineageCanaryEnabled(sport, now) {
|
||||
return lineageCanaryConfig.isEnabled(sport, now);
|
||||
}
|
||||
|
||||
/**
|
||||
* THE WRITE GATE'S ONLY QUESTION.
|
||||
*
|
||||
* Delegates to the single write-mode evaluator, so the gate and the status
|
||||
* surface can never disagree about whether lineage is writing or why. The
|
||||
* canary lease keeps its own contract untouched and is one of the two inputs;
|
||||
* persistent shadow is the other. `now` is threaded because the lease is
|
||||
* evaluated against the clock at WRITE TIME.
|
||||
*/
|
||||
function lineageWriteEnabled(sport, now) {
|
||||
return lineageWriteMode.isEnabled(sport, now);
|
||||
}
|
||||
|
||||
/**
|
||||
* NOTHING IS SERVED CALIBRATED, AND THIS IS DELIBERATE.
|
||||
*
|
||||
@@ -1423,7 +1437,7 @@ async function runSnapshot(sport, opts = {}) {
|
||||
// INTRADAY_REFRESH): `LINEAGE_CANARY_SPORTS=` disables it entirely,
|
||||
// `LINEAGE_CANARY_SPORTS=mlb,wnba` would widen it. Disabling needs no
|
||||
// migration revert and touches no historical row.
|
||||
if (retention && retention.commitPublication && persistedRows && lineageCanaryEnabled(sp)) {
|
||||
if (retention && retention.commitPublication && persistedRows && lineageWriteEnabled(sp, deps.now && deps.now())) {
|
||||
try {
|
||||
const pub = await retention.commitPublication({
|
||||
rows: persistedRows,
|
||||
@@ -1542,6 +1556,7 @@ async function runAllSnapshots(opts = {}) {
|
||||
|
||||
module.exports = {
|
||||
lineageCanaryEnabled,
|
||||
lineageWriteEnabled,
|
||||
LINEAGE_CANARY_SPORTS,
|
||||
runSnapshot,
|
||||
runAllSnapshots,
|
||||
|
||||
Reference in New Issue
Block a user