Retention hotfix: drop published_side, derive the schema contract, break the silence

`createCollector.onPublished` set `published_side` beside `published`.
`published_side` is not a model_snapshots column. supabase-js declares the
UNION of row keys in the `columns=` parameter, so one invalid key made
PostgREST reject the ENTIRE batch with a 400 — every sport, every cycle.
Retention is best-effort, so nothing surfaced. Confirmed in edge logs.

The field was redundant as well as invalid: `side` is already on the row.
Deleted rather than added to the schema — a column would preserve an
accidental artifact.

Three things missed it, and each is now closed:

1. WRONG SHAPE INSPECTED. The manual check sampled the collector after
   onGraded only and never called onPublished, so the offending key was
   not yet on the row. It read a pre-publication shape and reported the
   final outbound shape as clean. The new test captures the array actually
   handed to .upsert(), after the full production call order.

2. NO CONTRACT. Every retention test injects a permissive fake client that
   accepts any column set, so 381 suites proved the logic and never once
   compared a row against the database. The contract is now DERIVED — the
   migration chain applied to a disposable postgres, read out of
   information_schema (scripts/generate-schema-contract.js). A
   hand-maintained list would be a second opinion about the schema, and a
   second opinion is what let this through. scripts/verify-schema-contract.js
   checks the contract still describes a live database.

3. SILENT FAILURE. A failed batch reached one console.log. It now emits a
   high-severity structured event carrying sport, snapshot id, stage,
   error, code_sha and timestamp. Best-effort semantics are unchanged —
   the product continues and says so — but the failure is observable.
   `skipped` (no database configured) is not a failure and does not alert.

Teeth, each with the injection verified present before the run:
  - published_side back into the final payload -> 4 tests fail; restored
    byte-identically (sha 6a0ced7c52134135 both sides)
  - settled_at (a REAL contract column) -> accepted, so the guard
    discriminates by contract membership, not by novelty
  - alert block deleted -> 3 tests fail; restored byte-identically

Model and product behaviour untouched: analyzeViaEngine1,
probabilityEstimator, gradeSlateService, lineageCanaryConfig all unchanged.
Lineage stays OFF. Net source change is one behavioural line plus the alert.

382 suites / 5,094 tests pass. web tsc exit 0 (zero web paths touched).

Measurement blackout recorded, NOT backfilled: last good retention write
2026-08-27T19:08:32Z; ceaa896 started 21:16:41Z; the 22:00 UTC cycle ran
(ledger wrote 22:05:02) and persisted zero snapshot rows.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01CQJeAG8vcDoL5zkiaJyVb8
This commit is contained in:
Kev
2026-08-27 18:53:37 -04:00
parent ceaa896f77
commit 35da190f2c
6 changed files with 428 additions and 1 deletions
+56
View File
@@ -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 <container>');
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();
+49
View File
@@ -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); });
}
+6 -1
View File
@@ -221,7 +221,12 @@ function createCollector(ctx) {
const made = rowsFromSides(base, [winner], ctx); const made = rowsFromSides(base, [winner], ctx);
for (const r of made) { for (const r of made) {
const hit = index.get(keyOf(r)); 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 */ } } catch { /* signalling never affects grading */ }
}, },
+19
View File
@@ -659,6 +659,25 @@ async function runSnapshot(sport, opts = {}) {
// read_id yet — building an index here would produce an empty one and // read_id yet — building an index here would produce an empty one and
// quietly leave every ledger row unlinked. // 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}` : ''}`); 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) { } catch (e) {
console.warn(`[snapshot] retention write failed for ${sp} (snapshot continues):`, e.message); console.warn(`[snapshot] retention write failed for ${sp} (snapshot continues):`, e.message);
} }
@@ -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"
]
}
+217
View File
@@ -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);
});
});