diff --git a/scripts/generate-schema-contract.js b/scripts/generate-schema-contract.js new file mode 100644 index 0000000..db1768a --- /dev/null +++ b/scripts/generate-schema-contract.js @@ -0,0 +1,56 @@ +'use strict'; + +/** + * Regenerate supabase/schema/model_snapshots.columns.json from the MIGRATION + * CHAIN, by applying every committed migration to a disposable postgres and + * reading information_schema. + * + * The contract must be DERIVED, never hand-maintained: a hand-written list is + * a second opinion about the schema, and a second opinion is exactly what let + * `published_side` reach production. Deriving it from the same migrations the + * release verification applies means the list cannot drift from the release. + * + * docker run -d --name schema-gen -e POSTGRES_PASSWORD=test -e POSTGRES_DB=v postgres:15-alpine + * node scripts/generate-schema-contract.js schema-gen + * + * Drift against a LIVE database is a separate question — see + * scripts/verify-schema-contract.js. + */ + +const { execFileSync } = require('child_process'); +const fs = require('fs'); +const path = require('path'); + +const CONTAINER = process.argv[2]; +const TABLE = process.env.SCHEMA_TABLE || 'model_snapshots'; +const OUT = path.join(__dirname, '..', 'supabase', 'schema', `${TABLE}.columns.json`); + +function psql(sql) { + return execFileSync('docker', ['exec', CONTAINER, 'psql', '-U', 'postgres', '-d', 'v', '-tAc', sql], + { encoding: 'utf8' }).trim(); +} + +function main() { + if (!CONTAINER) throw new Error('usage: node scripts/generate-schema-contract.js '); + const cols = psql( + `select string_agg(column_name, ',' order by column_name)` + + ` from information_schema.columns` + + ` where table_schema='public' and table_name='${TABLE}'` + + ` and is_generated='NEVER' and identity_generation is null;`, + ).split(',').filter(Boolean); + if (cols.length === 0) throw new Error(`no columns found for ${TABLE} — was the migration chain applied?`); + + const prev = fs.existsSync(OUT) ? JSON.parse(fs.readFileSync(OUT, 'utf8')) : {}; + const doc = { + ...prev, + table: TABLE, + generated_by: 'scripts/generate-schema-contract.js', + derived_from: 'supabase/migrations/*.sql applied in order to a disposable postgres:15-alpine', + column_count: cols.length, + columns: cols, + }; + fs.writeFileSync(OUT, `${JSON.stringify(doc, null, 2)}\n`); + console.log(`wrote ${OUT} — ${cols.length} columns`); +} + +if (require.main === module) main(); diff --git a/scripts/verify-schema-contract.js b/scripts/verify-schema-contract.js new file mode 100644 index 0000000..c85d60b --- /dev/null +++ b/scripts/verify-schema-contract.js @@ -0,0 +1,49 @@ +'use strict'; + +/** + * DRIFT CHECK — compare the committed schema contract against a LIVE database. + * + * The unit guard proves the retention writer stays inside the contract. This + * proves the CONTRACT still describes reality. Both are needed: a correct + * writer against a stale contract is the same failure wearing a different hat. + * + * SUPABASE_URL=... SUPABASE_SERVICE_KEY=... node scripts/verify-schema-contract.js + * + * Exit 0 when the live table is a superset of the contract (safe: every column + * the writer may emit exists). Exit 1 when a contract column is MISSING live — + * that is the condition that breaks retention. + */ + +const fs = require('fs'); +const path = require('path'); +const { createClient } = require('@supabase/supabase-js'); + +const TABLE = process.env.SCHEMA_TABLE || 'model_snapshots'; + +async function main() { + const contract = JSON.parse(fs.readFileSync( + path.join(__dirname, '..', 'supabase', 'schema', `${TABLE}.columns.json`), 'utf8')); + const url = process.env.SUPABASE_URL; + const key = process.env.SUPABASE_SERVICE_KEY || process.env.SUPABASE_SERVICE_ROLE_KEY; + if (!url || !key) throw new Error('SUPABASE_URL + service key required'); + const sb = createClient(url, key, { auth: { persistSession: false } }); + + // One row is enough to learn the live column set from the response shape. + const { data, error } = await sb.from(TABLE).select('*').limit(1); + if (error) throw new Error(`live read failed: ${error.message}`); + if (!data || data.length === 0) throw new Error(`${TABLE} is empty — cannot infer live columns`); + + const live = new Set(Object.keys(data[0])); + const missing = contract.columns.filter((c) => !live.has(c)); + const extra = [...live].filter((c) => !contract.columns.includes(c)).sort(); + + console.log(`contract columns : ${contract.columns.length}`); + console.log(`live columns : ${live.size}`); + console.log(`live-only (informational, safe): ${extra.length ? extra.join(', ') : 'none'}`); + console.log(`MISSING LIVE (breaks retention): ${missing.length ? missing.join(', ') : 'none'}`); + process.exit(missing.length ? 1 : 0); +} + +if (require.main === module) { + main().catch((e) => { console.error('FAILED:', e.message); process.exit(1); }); +} diff --git a/src/services/retentionService.js b/src/services/retentionService.js index aa88573..f2b7bca 100644 --- a/src/services/retentionService.js +++ b/src/services/retentionService.js @@ -221,7 +221,12 @@ function createCollector(ctx) { const made = rowsFromSides(base, [winner], ctx); for (const r of made) { const hit = index.get(keyOf(r)); - if (hit) { hit.published = true; hit.published_side = true; } + // `published` alone. A companion `published_side` was written here + // and is NOT a model_snapshots column — it made PostgREST reject the + // ENTIRE retention batch with a 400, silently, for every sport. The + // side is already on the row (`side`), so the flag was redundant as + // well as invalid. + if (hit) { hit.published = true; } } } catch { /* signalling never affects grading */ } }, diff --git a/src/services/snapshotService.js b/src/services/snapshotService.js index 8d6cc27..46215cd 100644 --- a/src/services/snapshotService.js +++ b/src/services/snapshotService.js @@ -659,6 +659,25 @@ async function runSnapshot(sport, opts = {}) { // read_id yet — building an index here would produce an empty one and // quietly leave every ledger row unlinked. console.log(`[snapshot] retention ${sp}: ${r.written}/${r.attempted} rows${r.skipped ? ' (skipped — no supabase env)' : ''}${r.error ? ` ERROR: ${r.error}` : ''}`); + // RETENTION FAILURE MUST NOT BE SILENT. + // + // Retention stays best-effort — the product keeps publishing to Redis and + // that is deliberate. But a failed batch previously reached only this log + // line, and a PostgREST 400 (`published_side` is not a column) killed + // every batch for every sport with nothing surfacing anywhere. The + // measurement record died quietly while the product looked healthy. + // + // `skipped` is NOT a failure: it means no database is configured, which is + // the normal state in dev and test. + if (r.error || (!r.skipped && r.attempted > 0 && r.written === 0)) { + await deps.notify( + `Retention write FAILED for ${sp.toUpperCase()} — ${r.written}/${r.attempted} rows persisted. ` + + `stage=model_snapshots snapshot_id=${retentionCtx.snapshotId} code_sha=${retention.codeSha() || 'unknown'} ` + + `at=${deps.now()} error=${r.error || 'zero rows written with candidates present'}. ` + + 'The product is unaffected; the historical record for this cycle is missing.', + { title: 'VYNDR pipeline', priority: 'high', tags: ['rotating_light'] }, + ); + } } catch (e) { console.warn(`[snapshot] retention write failed for ${sp} (snapshot continues):`, e.message); } diff --git a/supabase/schema/model_snapshots.columns.json b/supabase/schema/model_snapshots.columns.json new file mode 100644 index 0000000..5119314 --- /dev/null +++ b/supabase/schema/model_snapshots.columns.json @@ -0,0 +1,81 @@ +{ + "table": "model_snapshots", + "generated_by": "scripts/generate-schema-contract.js", + "derived_from": "supabase/migrations/*.sql applied in order to a disposable postgres:15-alpine", + "note": "Insertable columns only (no generated/identity columns). This is the contract the retention writer must satisfy. Regenerate after any migration that touches model_snapshots; scripts/verify-schema-contract.js detects drift against a live database.", + "known_production_drift": { + "columns_in_production_not_in_migrations": [ + "quarantine_reason", + "re_settled_at", + "settlement_source" + ], + "explanation": "Present in production but not created by any committed migration - added out of band, same class as the migration-014 debt. The contract deliberately uses the MIGRATION-derived set, which is the stricter of the two: a writer that stays within it is valid against both." + }, + "column_count": 64, + "columns": [ + "actual_value", + "archetype", + "book", + "book_odds", + "canonical_event_id", + "captured_at", + "chain_shadow", + "change_type", + "claim_digest", + "claim_schema_version", + "code_sha", + "confidence", + "confidence_basis", + "created_at", + "cycle_hour_utc", + "devig_method", + "digest_algorithm_version", + "edge_pct", + "ev_pct", + "event_identity_method", + "event_identity_source", + "event_identity_version", + "event_occurrence", + "factor_inputs", + "fair_odds", + "fair_prob", + "features", + "game_date", + "game_id", + "grade", + "grade_11", + "id", + "line", + "lineage_action", + "lineage_state", + "lineage_version", + "model_version", + "opponent", + "outcome", + "over_odds", + "overround", + "p_win", + "player_key", + "player_name", + "projection", + "publication_id", + "published", + "published_at", + "read_id", + "read_natural_key", + "recaptures_id", + "refusal_reason", + "refused", + "revision_ordinal", + "settled_at", + "side", + "snapshot_id", + "sport", + "stat", + "supersedes_id", + "takeable", + "team", + "under_odds", + "value" + ] +} diff --git a/tests/unit/retentionSchemaContract.test.js b/tests/unit/retentionSchemaContract.test.js new file mode 100644 index 0000000..0100b1e --- /dev/null +++ b/tests/unit/retentionSchemaContract.test.js @@ -0,0 +1,217 @@ +'use strict'; + +/** + * RETENTION SCHEMA CONTRACT — the guard that would have caught the outage. + * + * WHAT HAPPENED. `createCollector.onPublished` set `published_side` alongside + * `published`. `published_side` is not a `model_snapshots` column, supabase-js + * declares the UNION of row keys in the `columns=` parameter, and PostgREST + * rejected the entire batch with a 400 — for every sport, silently, because + * retention is best-effort. Verified in production edge logs. + * + * WHY 381 GREEN SUITES MISSED IT. Every retention test injects a permissive + * fake client — `{ from: () => ({ upsert: async () => ({ error: null }) }) }` — + * which accepts any column set. The suite proved the logic and never once + * compared a row against the database contract. + * + * WHY THE FIRST MANUAL CHECK ALSO MISSED IT. It sampled the collector after + * `onGraded` only and never called `onPublished`, so the offending key was not + * yet on the row. It inspected a PRE-PUBLICATION shape and reported the FINAL + * outbound shape as clean. + * + * So this test does two things the old ones could not: + * 1. it validates the payload ACTUALLY HANDED TO `.upsert()`, captured by a + * spy, after the full production call order including `onPublished`; + * 2. it validates against a contract DERIVED from the migration chain, not a + * hand-maintained list. + */ + +const path = require('path'); +const fs = require('fs'); +const retention = require('../../src/services/retentionService'); + +const ROOT = path.resolve(__dirname, '..', '..'); +const CONTRACT = JSON.parse(fs.readFileSync( + path.join(ROOT, 'supabase/schema/model_snapshots.columns.json'), 'utf8', +)); + +/** Captures the exact array supabase-js would send. */ +function spyClient() { + const seen = []; + return { + seen, + client: { + from: (table) => ({ + upsert: async (rows, opts) => { seen.push({ table, rows, opts }); return { error: null }; }, + }), + }, + }; +} + +const CTX = { + snapshotId: '00000000-0000-0000-0000-000000000000', + capturedAt: '2026-08-27T22:00:00Z', + gameDate: '2026-08-27', + gameIdFor: () => 'mlb:2026-08-27:BostonRedSox@MiamiMarlins', +}; +const BASE = { + player: 'Aaron Judge', stat_type: 'hits', line: 0.5, sport: 'mlb', book: 'draftkings', + canonical_event_id: 'mlb:gamepk:823825', event_identity_source: 'MLB_STATSAPI_GAMEPK', + event_identity_method: 'CANONICAL', event_identity_version: 'evid@1', event_occurrence: 1, +}; +const OVER = { ...BASE, direction: 'over', grade: 'B', confidence: 61, p_win: 0.61 }; +const UNDER = { ...BASE, direction: 'under', grade: 'C', confidence: 39, p_win: 0.39 }; + +/** The FULL production call order — onGraded for both sides, then onPublished. */ +function finalOutboundRows() { + const c = retention.createCollector(CTX); + c.onGraded(BASE, [OVER, UNDER]); + c.onPublished(BASE, OVER); + return c.rows; +} + +describe('the contract itself is derived, not hand-written', () => { + test('it declares its generator and its derivation', () => { + expect(CONTRACT.table).toBe('model_snapshots'); + expect(CONTRACT.generated_by).toBe('scripts/generate-schema-contract.js'); + expect(CONTRACT.derived_from).toMatch(/supabase\/migrations/); + expect(fs.existsSync(path.join(ROOT, 'scripts/generate-schema-contract.js'))).toBe(true); + expect(fs.existsSync(path.join(ROOT, 'scripts/verify-schema-contract.js'))).toBe(true); + }); + + test('it is non-trivial and self-consistent', () => { + expect(CONTRACT.columns.length).toBe(CONTRACT.column_count); + expect(CONTRACT.columns.length).toBeGreaterThan(50); + expect(CONTRACT.columns).toContain('published'); + expect(CONTRACT.columns).toContain('canonical_event_id'); + // The offending field must NOT be in the contract — that is the fact. + expect(CONTRACT.columns).not.toContain('published_side'); + }); + + test('known production drift is recorded rather than hidden', () => { + const d = CONTRACT.known_production_drift; + expect(d.columns_in_production_not_in_migrations).toEqual( + expect.arrayContaining(['quarantine_reason', 're_settled_at', 'settlement_source']), + ); + // The migration-derived set is the stricter of the two, so a writer inside + // it is valid against production as well. + expect(d.explanation).toMatch(/stricter/); + }); +}); + +describe('FINAL OUTBOUND payload is within the schema contract', () => { + test('the captured upsert payload uses only real columns', async () => { + const spy = spyClient(); + await retention.persist(finalOutboundRows(), { getClient: () => spy.client }); + + expect(spy.seen.length).toBeGreaterThan(0); + const call = spy.seen[0]; + expect(call.table).toBe('model_snapshots'); + + // supabase-js sends the UNION of keys across the batch — validate the union. + const union = new Set(); + for (const r of call.rows) for (const k of Object.keys(r)) union.add(k); + const invalid = [...union].filter((k) => !CONTRACT.columns.includes(k)).sort(); + expect(invalid).toEqual([]); + }); + + test('the payload is captured AFTER onPublished, not before', async () => { + // The pre-publication shape is what the earlier manual check inspected. + const pre = retention.createCollector(CTX); + pre.onGraded(BASE, [OVER, UNDER]); + const preKeys = new Set(pre.rows.flatMap((r) => Object.keys(r))); + + const post = finalOutboundRows(); + const postKeys = new Set(post.flatMap((r) => Object.keys(r))); + + // Both must be valid; the point is that this test exercises the later one. + for (const set of [preKeys, postKeys]) { + expect([...set].filter((k) => !CONTRACT.columns.includes(k))).toEqual([]); + } + // And onPublished must actually have marked a row, or the test is vacuous. + expect(post.filter((r) => r.published === true)).toHaveLength(1); + expect(post.filter((r) => r.published === false)).toHaveLength(1); + }); + + test('every declared key is present on EVERY row', async () => { + // A key on only some rows is dropped for the whole batch by PostgREST, so + // a ragged batch is its own defect class. + const rows = finalOutboundRows(); + const union = new Set(rows.flatMap((r) => Object.keys(r))); + for (const r of rows) { + expect(new Set(Object.keys(r))).toEqual(union); + } + }); + + test('published_side is gone from the source entirely', () => { + const src = fs.readFileSync(path.join(ROOT, 'src/services/retentionService.js'), 'utf8') + .replace(/\/\*[\s\S]*?\*\//g, '') + .replace(/(^|[^:])\/\/.*$/gm, '$1'); + expect(src).not.toMatch(/published_side/); + }); +}); + +describe('RETENTION FAILURE IS NOT SILENT', () => { + test('a PostgREST 400 is reported, never swallowed as success', async () => { + const failing = { + from: () => ({ + upsert: async () => ({ + error: { message: "Could not find the 'published_side' column of 'model_snapshots'", code: 'PGRST204' }, + }), + }), + }; + const out = await retention.persist(finalOutboundRows(), { getClient: () => failing }); + expect(out.written).toBe(0); + expect(out.error).toMatch(/published_side/); + // Never reported as a success. + expect(out.skipped).toBe(false); + }); + + test('the pipeline emits a high-severity structured event on that failure', () => { + const src = fs.readFileSync(path.join(ROOT, 'src/services/snapshotService.js'), 'utf8'); + const block = src.slice(src.indexOf('Retention write FAILED') - 400, + src.indexOf('Retention write FAILED') + 700); + for (const field of ['stage=model_snapshots', 'snapshot_id=', 'code_sha=', 'at=', 'error=']) { + expect(block).toContain(field); + } + expect(block).toMatch(/priority: 'high'/); + // Best-effort semantics preserved: the product is explicitly unaffected. + expect(block).toMatch(/product is unaffected/); + }); + + test('a skipped write (no database configured) is NOT alerted', async () => { + const out = await retention.persist(finalOutboundRows(), { getClient: () => null }); + expect(out.skipped).toBe(true); + expect(out.error).toBeNull(); + const src = fs.readFileSync(path.join(ROOT, 'src/services/snapshotService.js'), 'utf8'); + expect(src).toMatch(/!r\.skipped/); + }); + + test('zero rows written with candidates present also alerts', () => { + const src = fs.readFileSync(path.join(ROOT, 'src/services/snapshotService.js'), 'utf8'); + expect(src).toMatch(/r\.attempted > 0 && r\.written === 0/); + }); +}); + +describe('CYCLE PERSISTENCE MECHANISM', () => { + test('a cycle is written in chunks, and a failed chunk aborts the rest', async () => { + // Documented for the dark-cycle gate: persistence is CHUNKED (250), so a + // partial cycle is possible in principle — the loop breaks on the first + // error rather than continuing. There is no terminal completion marker. + const src = fs.readFileSync(path.join(ROOT, 'src/services/retentionService.js'), 'utf8'); + const body = src.slice(src.indexOf('async function persist(')); + expect(body).toMatch(/const CHUNK = 250/); + expect(body).toMatch(/if \(error\) \{ out\.error = error\.message; break; \}/); + // `written` is the exact count of rows in successfully committed chunks, so + // written === attempted is the completion signal available today. + expect(body).toMatch(/out\.written \+= chunk\.length/); + }); + + test('written vs attempted distinguishes a complete cycle from a partial one', async () => { + const spy = spyClient(); + const rows = finalOutboundRows(); + const out = await retention.persist(rows, { getClient: () => spy.client }); + expect(out.attempted).toBe(rows.length); + expect(out.written).toBe(rows.length); + }); +});