From 9809626c99f8ffb8cb0bbc2b32dab8fdb4298cd0 Mon Sep 17 00:00:00 2001 From: Kev Date: Thu, 27 Aug 2026 19:12:59 -0400 Subject: [PATCH] Retention completion: a cohort is complete only when the writer says N of N MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The previous bug made the recorder write nothing. The dangerous successor is a recorder that writes half and looks healthy: persist() writes in chunks of 250 and STOPS AT THE FIRST FAILED CHUNK, so chunks committed before the failure are already durable. Rows exist under the snapshot_id, captured_at is uniform, Redis kept working — and the cohort is short. So row presence was never completion evidence, and neither was a matching timestamp. Completeness is now proven by the writer or not at all. TERMINAL RETENTION STATES (retentionService.classifyPersist): NOTHING_TO_PERSIST attempted 0 — a refusal-only slate is still a cycle SKIPPED_NO_DATABASE no database configured; not a failure COMPLETE attempted > 0, written === attempted, no error FAILED_ZERO_WRITE written === 0 — first chunk failed FAILED_PARTIAL 0 < written < attempted — a later chunk failed FAILED_UNRESOLVED_ERROR counts look complete but an error is unresolved; unreachable through today's loop, and kept because the alternative is reporting COMPLETE holding an error The invariant: any written < attempted with attempted > 0 is a FAILED cycle. A partial cohort is never degraded success. classifyPersist reads the EXACT persist() result and refuses anything else — it never recomputes attempted or written, because a second calculation could disagree with the writer and then the status would describe a cycle that did not happen. persist() itself is byte-identical to 35da190. `written` counts rows in COMMITTED CHUNKS, not database inserts: the upsert uses ignoreDuplicates, so a re-run legitimately inserts far fewer rows than it writes. Comparing written to count(*) will disagree by design. Documented, because that mismatch is exactly what would be misread as a partial write. VISIBILITY. The 35da190 alert condition was `r.error || (!r.skipped && r.attempted > 0 && r.written === 0)` — it could not see a partial cohort as a distinct state. It is now driven by terminal status, so FAILED_PARTIAL alerts as loudly as a total failure and is labelled INCOMPLETE and unusable as evidence. Best-effort is unchanged: the product continues and the alert says so. OBSERVABILITY. A successful cycle previously left only a console.log with no snapshot_id, no code_sha and no terminal status, so completion could not be established after the fact. `GET /api/internal/snapshot/status` now returns `last_retention` per sport — sport, snapshot_id, attempted, written, status, completed_at, code_sha, error_summary — taken verbatim from the persistence result. Existing internal auth, read-only, counts and status only, no payloads. No new table, no new route. RELEASE-AUTHORIZED INSERT CONTRACT. The migration-derived contract is the release authority; production is not. A prod-only column is DRIFT / RECORDED DEBT and never becomes permission by existing. Verifier classifies: release column missing in prod -> HARD FAILURE; prod-only -> drift warning; outbound key outside the contract -> contract failure (enforced against the real upsert payload). It is read-only and never rewrites the contract from live schema. Live: release 64, prod 67, prod-only 3, missing in prod 0. Six teeth, each with the injection verified present, against a green baseline: 1 written>0 as generic success -> 6 fail 2 later-chunk failure reports COMPLETE -> 5 fail 3 FAILED_PARTIAL does not alert -> 3 fail 4 status reports a recalculated count -> 1 fail 5 row presence treated as completion -> 1 fail 6 invalid outbound column reintroduced -> 4 fail Restored byte-identically (retention b341cf16c1baa992, snapshot 81ab1bd7730dee89). Two stale assertions updated rather than deleted, with the mechanism change recorded: the alert-shape tests described the superseded written===0 condition, and the runtime probe test pinned an exact import list. Model and product preserved: analyzeViaEngine1, probabilityEstimator, gradeSlateService, lineageCanaryConfig, eventIdentity, ledgerService, calibration and chain all UNCHANGED; zero lineage/publication files touched; zero cacheSet changes; zero web paths. Lineage stays OFF. 383 suites / 5,118 tests pass. web tsc exit 0. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_01CQJeAG8vcDoL5zkiaJyVb8 --- scripts/verify-schema-contract.js | 52 ++-- src/routes/internal.js | 8 +- src/services/retentionService.js | 99 +++++++ src/services/snapshotService.js | 38 ++- supabase/schema/model_snapshots.columns.json | 7 +- tests/unit/retentionCompletion.test.js | 264 +++++++++++++++++++ tests/unit/retentionSchemaContract.test.js | 22 +- tests/unit/runtimeObservability.test.js | 4 +- 8 files changed, 455 insertions(+), 39 deletions(-) create mode 100644 tests/unit/retentionCompletion.test.js diff --git a/scripts/verify-schema-contract.js b/scripts/verify-schema-contract.js index c85d60b..252d0f2 100644 --- a/scripts/verify-schema-contract.js +++ b/scripts/verify-schema-contract.js @@ -1,17 +1,28 @@ 'use strict'; /** - * DRIFT CHECK — compare the committed schema contract against a LIVE database. + * RELEASE-AUTHORIZED model_snapshots INSERT CONTRACT — drift verifier. * - * 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. + * The committed contract is derived from the MIGRATION CHAIN, and it is the + * RELEASE AUTHORITY: it says which columns this release is permitted to write. + * Production is NOT the authority. A column that exists in production but in no + * migration is drift, and drift does not silently become permission — that is + * how an accidental artifact turns into schema. + * + * Classification: + * RELEASE COLUMN MISSING IN PROD -> HARD FAILURE (exit 1). Retention will + * 400 on the whole batch. + * PROD-ONLY COLUMN -> DRIFT WARNING / RECORDED DEBT (exit 0). + * OUTBOUND INSERT KEY OUTSIDE + * THE RELEASE CONTRACT -> CONTRACT FAILURE. Enforced by + * tests/unit/retentionSchemaContract.test.js + * against the real .upsert() payload. + * + * This script is READ-ONLY. It never rewrites the contract from live schema. + * Regenerating is a deliberate act: scripts/generate-schema-contract.js against + * a disposable database with the migration chain applied. * * 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'); @@ -34,14 +45,25 @@ async function main() { 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(); + const missingInProd = contract.columns.filter((c) => !live.has(c)); + const prodOnly = [...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); + console.log(`RELEASE-AUTHORIZED INSERT CONTRACT (${TABLE})`); + console.log(` release columns (migration-derived) : ${contract.columns.length}`); + console.log(` live production columns : ${live.size}`); + console.log(` PROD-ONLY COLUMN (drift/debt) : ${prodOnly.length}${prodOnly.length ? ` -> ${prodOnly.join(', ')}` : ''}`); + console.log(` RELEASE COLUMN MISSING IN PROD : ${missingInProd.length}${missingInProd.length ? ` -> ${missingInProd.join(', ')}` : ''}`); + + if (prodOnly.length) { + console.log(' NOTE: prod-only columns are RECORDED DEBT. They are NOT release-authorized'); + console.log(' and must not be added to the contract from live schema.'); + } + if (missingInProd.length) { + console.error(' HARD FAILURE: a release-authorized column does not exist in production.'); + process.exit(1); + } + console.log(' RESULT: PASS (no release column missing in production)'); + process.exit(0); } if (require.main === module) { diff --git a/src/routes/internal.js b/src/routes/internal.js index 6492b9f..14f4582 100644 --- a/src/routes/internal.js +++ b/src/routes/internal.js @@ -189,7 +189,7 @@ router.get('/snapshot/status', async (req, res) => { const ticker = await cacheGet('ticker:items'); redis_keys['ticker:items'] = !!ticker; // Same helpers the pipeline itself uses — one build identity, one canary parser. - const { codeSha } = require('../services/retentionService'); + const { codeSha, lastRetention } = require('../services/retentionService'); 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; @@ -204,6 +204,12 @@ 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(), + // 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 — + // this is the only place that distinction is observable in production. + // Counts and status only: no row payloads, no credentials. + last_retention: lastRetention(), cron_armed: process.env.SNAPSHOT_CRON === '1', cron_hours_utc: HOURS_UTC, last_snapshot, diff --git a/src/services/retentionService.js b/src/services/retentionService.js index f2b7bca..d3844e6 100644 --- a/src/services/retentionService.js +++ b/src/services/retentionService.js @@ -727,8 +727,107 @@ function newSnapshotId() { return crypto.randomUUID(); } +/* ------------------------------------------------------------------ * + * TERMINAL RETENTION STATE + * + * `persist()` writes in CHUNKS and STOPS ON THE FIRST FAILED CHUNK. Chunks + * committed before the failure are already durable, so a failed cycle can + * leave real, valid-looking rows behind. Row presence under a snapshot_id is + * therefore NOT completion evidence, and neither is a uniform captured_at. + * + * The invariant: written < attempted with attempted > 0 is a FAILED cycle. + * A partial cycle is never degraded success. + * + * `written` counts rows in successfully COMMITTED CHUNKS — not rows inserted. + * The upsert uses ignoreDuplicates, so a re-run legitimately inserts far fewer + * database rows than it writes. Comparing `written` to count(*) for a + * snapshot_id will disagree by design; that is not a partial write. + * ------------------------------------------------------------------ */ + +const TERMINAL = Object.freeze({ + /** attempted === 0 — a refusal-only slate is still a legitimate cycle. */ + NOTHING_TO_PERSIST: 'NOTHING_TO_PERSIST', + /** No database configured (dev/test). Not a failure. */ + SKIPPED_NO_DATABASE: 'SKIPPED_NO_DATABASE', + /** attempted > 0, written === attempted, no unresolved error. */ + COMPLETE: 'COMPLETE', + /** attempted > 0, written === 0 — the first chunk failed. */ + FAILED_ZERO_WRITE: 'FAILED_ZERO_WRITE', + /** attempted > 0, 0 < written < attempted — a later chunk failed. */ + FAILED_PARTIAL: 'FAILED_PARTIAL', + /** + * Counts look complete but an error is unresolved. Unreachable through + * today's loop (it breaks before crediting a failed chunk), and kept + * because the alternative is reporting COMPLETE with an error in hand. + */ + FAILED_UNRESOLVED_ERROR: 'FAILED_UNRESOLVED_ERROR', +}); + +const FAILURE_STATUSES = Object.freeze([ + TERMINAL.FAILED_ZERO_WRITE, TERMINAL.FAILED_PARTIAL, TERMINAL.FAILED_UNRESOLVED_ERROR, +]); + +function isRetentionFailure(status) { return FAILURE_STATUSES.includes(status); } + +/** + * Classify the EXACT object `persist()` returned. It never recomputes + * attempted or written — a second calculation could disagree with the writer, + * and then the status would describe something that did not happen. + */ +function classifyPersist(result) { + if (!result || typeof result.attempted !== 'number' || typeof result.written !== 'number') { + throw new TypeError('classifyPersist requires the exact persist() result'); + } + const { attempted, written, skipped, error } = result; + if (attempted === 0) return TERMINAL.NOTHING_TO_PERSIST; + if (skipped) return TERMINAL.SKIPPED_NO_DATABASE; + if (written === 0) return TERMINAL.FAILED_ZERO_WRITE; + if (written < attempted) return TERMINAL.FAILED_PARTIAL; + // written >= attempted from here. + if (error) return TERMINAL.FAILED_UNRESOLVED_ERROR; + return written === attempted ? TERMINAL.COMPLETE : TERMINAL.FAILED_UNRESOLVED_ERROR; +} + +/** + * Latest terminal retention result per sport, for the internal status probe. + * In-memory and per-process on purpose: this is an observability surface, not + * a record. The record is model_snapshots. + */ +const lastTerminal = new Map(); + +function recordTerminal({ sport, snapshotId, result, completedAt }) { + const status = classifyPersist(result); + const entry = Object.freeze({ + sport: sport || null, + snapshot_id: snapshotId || null, + attempted: result.attempted, + written: result.written, + status, + completed_at: completedAt || new Date().toISOString(), + code_sha: codeSha(), + // Message only — never row payloads. + error_summary: result.error ? String(result.error).slice(0, 300) : null, + }); + if (sport) lastTerminal.set(sport, entry); + return entry; +} + +function lastRetention() { + const out = {}; + for (const [sport, entry] of lastTerminal) out[sport] = entry; + return out; +} + +function resetTerminal() { lastTerminal.clear(); } + module.exports = { MODEL_VERSION, + TERMINAL, + classifyPersist, + isRetentionFailure, + recordTerminal, + lastRetention, + resetTerminal, REPAIRED_CHAMPION_VERSION, codeSha, rowsFromSides, diff --git a/src/services/snapshotService.js b/src/services/snapshotService.js index 46215cd..1504332 100644 --- a/src/services/snapshotService.js +++ b/src/services/snapshotService.js @@ -658,23 +658,35 @@ async function runSnapshot(sport, opts = {}) { // commitPublication AFTER the slate write, so at this point no row has a // 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. + // TERMINAL RETENTION STATE — classified from the EXACT persist() result, + // never recomputed. Chunked writes stop at the first failed chunk, so a + // failed cycle can leave durable rows behind: presence of rows under a + // snapshot_id proves nothing about whether the cohort is complete. + const terminal = retention.recordTerminal({ + sport: sp, + snapshotId: retentionCtx.snapshotId, + result: r, + completedAt: deps.now(), + }); + console.log(`[snapshot] retention ${sp}: ${terminal.status} ${r.written}/${r.attempted} rows snapshot_id=${terminal.snapshot_id}${r.error ? ` ERROR: ${r.error}` : ''}`); + // RETENTION FAILURE MUST NOT BE SILENT — INCLUDING A PARTIAL ONE. // - // Retention stays best-effort — the product keeps publishing to Redis and - // that is deliberate. But a failed batch previously reached only this log + // Retention stays best-effort: the product keeps publishing to Redis and + // that is deliberate. But a failed batch previously reached only a 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. + // every batch for every sport with nothing surfacing anywhere. // - // `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)) { + // FAILED_PARTIAL is the dangerous successor to that bug: earlier chunks + // committed, the cohort is short, and the surviving rows make it look + // healthy. It alerts exactly as loudly as a total failure. + // + // NOTHING_TO_PERSIST and SKIPPED_NO_DATABASE are not failures. + if (retention.isRetentionFailure(terminal.status)) { 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.', + `Retention ${terminal.status} for ${sp.toUpperCase()} — ${terminal.written}/${terminal.attempted} rows persisted. ` + + `stage=model_snapshots snapshot_id=${terminal.snapshot_id} code_sha=${terminal.code_sha || 'unknown'} ` + + `at=${terminal.completed_at} error=${terminal.error_summary || 'none reported'}. ` + + 'The product is unaffected; this retention cohort is INCOMPLETE and must not be used as evidence.', { title: 'VYNDR pipeline', priority: 'high', tags: ['rotating_light'] }, ); } diff --git a/supabase/schema/model_snapshots.columns.json b/supabase/schema/model_snapshots.columns.json index 5119314..9263c38 100644 --- a/supabase/schema/model_snapshots.columns.json +++ b/supabase/schema/model_snapshots.columns.json @@ -9,7 +9,8 @@ "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." + "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.", + "classification": "PROD-ONLY COLUMN -> DRIFT WARNING / RECORDED DEBT (not release-authorized)" }, "column_count": 64, "columns": [ @@ -77,5 +78,7 @@ "team", "under_odds", "value" - ] + ], + "authority": "RELEASE-AUTHORIZED model_snapshots INSERT CONTRACT", + "authority_note": "Derived from the migration chain, which is the release authority. Production is not. A column present in production but in no migration is DRIFT / RECORDED DEBT and is NOT release-authorized; it must never be copied into this contract from live schema." } diff --git a/tests/unit/retentionCompletion.test.js b/tests/unit/retentionCompletion.test.js new file mode 100644 index 0000000..864bfa5 --- /dev/null +++ b/tests/unit/retentionCompletion.test.js @@ -0,0 +1,264 @@ +'use strict'; + +/** + * RETENTION COMPLETION CONTRACT. + * + * The previous bug made the recorder write NOTHING. The dangerous successor is + * a recorder that writes HALF and looks healthy: `persist()` writes in chunks + * of 250 and stops at the first failed chunk, so chunks committed before the + * failure are already durable. Rows exist, `captured_at` is uniform, Redis kept + * working — and the cohort is short. + * + * So completeness is proven by the WRITER (attempted N, persisted N), never by + * row presence. + */ + +const retention = require('../../src/services/retentionService'); +const { TERMINAL } = retention; + +/** Chunk-faithful fake: records every chunk and fails the nominated one. */ +function chunkClient({ failAtCall = null } = {}) { + const calls = []; + return { + calls, + client: { + from: () => ({ + upsert: async (rows) => { + calls.push(rows.length); + if (failAtCall !== null && calls.length === failAtCall) { + return { error: { message: 'PostgREST 400: chunk rejected', code: 'PGRST204' } }; + } + return { error: null }; + }, + }), + }, + }; +} + +const rowsN = (n) => Array.from({ length: n }, (_, i) => ({ snapshot_id: 's', player_key: `p${i}`, stat: 'hits', line: 0.5, side: 'over' })); + +describe('CHUNKED PERSISTENCE — traced contract', () => { + test('attempted is the input length; chunks are sequential and in array order', async () => { + const f = chunkClient(); + const out = await retention.persist(rowsN(600), { getClient: () => f.client }); + expect(out.attempted).toBe(600); + expect(f.calls).toEqual([250, 250, 100]); // chunk size 250, final short chunk + expect(out.written).toBe(600); + }); + + test('written counts COMMITTED CHUNK ROWS, not database inserts', async () => { + // The upsert uses ignoreDuplicates, so a re-run legitimately inserts far + // fewer rows than it writes. Comparing `written` to count(*) will disagree + // by design — that is not a partial write, and this is why row counting + // cannot be the completion test. + const src = require('fs').readFileSync(require.resolve('../../src/services/retentionService'), 'utf8'); + expect(src).toMatch(/ignoreDuplicates: true/); + expect(src).toMatch(/out\.written \+= chunk\.length/); + }); + + test('persist RETURNS an error, never throws it', async () => { + const boom = { from: () => ({ upsert: async () => { throw new Error('socket hang up'); } }) }; + const out = await retention.persist(rowsN(10), { getClient: () => boom }); + expect(out.error).toMatch(/socket hang up/); + expect(out.written).toBe(0); + }); +}); + +describe('TERMINAL STATE CLASSIFICATION', () => { + const classify = (r) => retention.classifyPersist(r); + + test('ALL CHUNKS SUCCEED -> COMPLETE', async () => { + const f = chunkClient(); + const out = await retention.persist(rowsN(600), { getClient: () => f.client }); + expect([out.attempted, out.written]).toEqual([600, 600]); + expect(classify(out)).toBe(TERMINAL.COMPLETE); + expect(retention.isRetentionFailure(classify(out))).toBe(false); + }); + + test('FIRST CHUNK FAILS -> FAILED_ZERO_WRITE', async () => { + const f = chunkClient({ failAtCall: 1 }); + const out = await retention.persist(rowsN(600), { getClient: () => f.client }); + expect(out.attempted).toBe(600); + expect(out.written).toBe(0); + expect(classify(out)).toBe(TERMINAL.FAILED_ZERO_WRITE); + expect(retention.isRetentionFailure(classify(out))).toBe(true); + expect(f.calls).toHaveLength(1); // stops immediately + }); + + test('MIDDLE CHUNK FAILS -> FAILED_PARTIAL, never COMPLETE', async () => { + const f = chunkClient({ failAtCall: 2 }); + const out = await retention.persist(rowsN(600), { getClient: () => f.client }); + expect(out.attempted).toBe(600); + expect(out.written).toBe(250); // chunk 1 is DURABLE + expect(out.written).toBeGreaterThan(0); + expect(out.written).toBeLessThan(out.attempted); + expect(classify(out)).toBe(TERMINAL.FAILED_PARTIAL); + expect(classify(out)).not.toBe(TERMINAL.COMPLETE); + expect(retention.isRetentionFailure(classify(out))).toBe(true); + }); + + test('FINAL CHUNK FAILS -> FAILED_PARTIAL', async () => { + const f = chunkClient({ failAtCall: 3 }); + const out = await retention.persist(rowsN(600), { getClient: () => f.client }); + expect([out.attempted, out.written]).toEqual([600, 500]); + expect(classify(out)).toBe(TERMINAL.FAILED_PARTIAL); + expect(retention.isRetentionFailure(classify(out))).toBe(true); + }); + + test('SKIPPED (no database) -> SKIPPED_NO_DATABASE, not a failure', async () => { + const out = await retention.persist(rowsN(600), { getClient: () => null }); + expect(out.skipped).toBe(true); + expect(classify(out)).toBe(TERMINAL.SKIPPED_NO_DATABASE); + expect(retention.isRetentionFailure(classify(out))).toBe(false); + }); + + test('EMPTY COHORT -> NOTHING_TO_PERSIST, not a failure', async () => { + const out = await retention.persist([], { getClient: () => null }); + expect(out.attempted).toBe(0); + expect(classify(out)).toBe(TERMINAL.NOTHING_TO_PERSIST); + expect(retention.isRetentionFailure(classify(out))).toBe(false); + }); + + test('THE INVARIANT: any written < attempted with attempted > 0 is FAILED', () => { + for (const [a, w] of [[600, 0], [600, 1], [600, 250], [600, 599], [1, 0]]) { + const st = classify({ attempted: a, written: w, skipped: false, error: 'x' }); + expect(retention.isRetentionFailure(st)).toBe(true); + expect(st).not.toBe(TERMINAL.COMPLETE); + } + }); + + test('counts that look complete but carry an error are NEVER COMPLETE', () => { + expect(classify({ attempted: 600, written: 600, skipped: false, error: 'late failure' })) + .toBe(TERMINAL.FAILED_UNRESOLVED_ERROR); + }); + + test('classification refuses anything that is not the real persist result', () => { + // It must never recompute or infer counts — a second calculation could + // disagree with the writer, and then the status describes a cycle that + // did not happen. + expect(() => classify(null)).toThrow(TypeError); + expect(() => classify({ attempted: '600', written: 600 })).toThrow(TypeError); + expect(() => classify({ written: 600 })).toThrow(TypeError); + }); +}); + +describe('COMPLETION OBSERVABILITY', () => { + beforeEach(() => retention.resetTerminal()); + + test('recordTerminal reports the EXACT persist counts, not a recalculation', async () => { + const f = chunkClient({ failAtCall: 2 }); + const out = await retention.persist(rowsN(600), { getClient: () => f.client }); + const e = retention.recordTerminal({ + sport: 'mlb', snapshotId: 'snap-1', result: out, completedAt: '2026-08-27T23:00:00Z', + }); + expect(e.attempted).toBe(out.attempted); + expect(e.written).toBe(out.written); + expect(e.status).toBe(TERMINAL.FAILED_PARTIAL); + expect(e.snapshot_id).toBe('snap-1'); + expect(e.sport).toBe('mlb'); + expect(e.completed_at).toBe('2026-08-27T23:00:00Z'); + expect(e).toHaveProperty('code_sha'); + expect(e.error_summary).toMatch(/chunk rejected/); + }); + + test('lastRetention exposes the latest terminal result per sport', async () => { + const ok = chunkClient(); + const bad = chunkClient({ failAtCall: 1 }); + retention.recordTerminal({ sport: 'mlb', snapshotId: 'a', result: await retention.persist(rowsN(10), { getClient: () => ok.client }) }); + retention.recordTerminal({ sport: 'wnba', snapshotId: 'b', result: await retention.persist(rowsN(10), { getClient: () => bad.client }) }); + const last = retention.lastRetention(); + expect(last.mlb.status).toBe(TERMINAL.COMPLETE); + expect(last.wnba.status).toBe(TERMINAL.FAILED_ZERO_WRITE); + expect(Object.keys(last).sort()).toEqual(['mlb', 'wnba']); + }); + + test('the observable result carries every required field', () => { + const e = retention.recordTerminal({ + sport: 'mlb', snapshotId: 's', result: { attempted: 5, written: 5, skipped: false, error: null }, + }); + for (const k of ['sport', 'snapshot_id', 'attempted', 'written', 'status', 'completed_at', 'code_sha', 'error_summary']) { + expect(Object.keys(e)).toContain(k); + } + }); + + test('it never exposes row payloads', async () => { + const bad = { from: () => ({ upsert: async () => ({ error: { message: 'x'.repeat(5000) } }) }) }; + const out = await retention.persist(rowsN(10), { getClient: () => bad }); + const e = retention.recordTerminal({ sport: 'mlb', snapshotId: 's', result: out }); + expect(e.error_summary.length).toBeLessThanOrEqual(300); + expect(JSON.stringify(e)).not.toMatch(/player_key/); + }); + + test('the internal status route exposes it behind existing internal auth', () => { + const src = require('fs').readFileSync(require.resolve('../../src/routes/internal.js'), 'utf8'); + expect(src).toMatch(/last_retention: lastRetention\(\)/); + // Read-only, and no second calculation in the route. + expect(src).not.toMatch(/last_retention[\s\S]{0,200}attempted:/); + expect(src).toMatch(/requireInternalAuth/); + }); +}); + +describe('PARTIAL FAILURE IS VISIBLE AT THE CALL SITE', () => { + const src = require('fs').readFileSync(require.resolve('../../src/services/snapshotService.js'), 'utf8'); + const block = src.slice(src.indexOf('const terminal = retention.recordTerminal'), + src.indexOf('const terminal = retention.recordTerminal') + 2200); + + test('the alert is driven by TERMINAL STATUS, not by a written===0 test', () => { + expect(block).toMatch(/retention\.isRetentionFailure\(terminal\.status\)/); + // The old condition could not see a partial as a distinct state. + expect(src).not.toMatch(/r\.attempted > 0 && r\.written === 0/); + }); + + test('a failed cohort alert carries every required field', () => { + for (const f of ['${terminal.status}', 'stage=model_snapshots', 'snapshot_id=${terminal.snapshot_id}', + 'code_sha=${terminal.code_sha', 'at=${terminal.completed_at}', 'error=${terminal.error_summary']) { + expect(block).toContain(f); + } + expect(block).toContain('${terminal.written}/${terminal.attempted}'); + expect(block).toMatch(/priority: 'high'/); + }); + + test('a partial cohort is labelled unusable as evidence, not degraded success', () => { + expect(block).toMatch(/INCOMPLETE and must not be used as evidence/); + expect(block).toMatch(/product is unaffected/); // best-effort preserved + }); + + test('the terminal result recorded is the one persist returned', () => { + expect(block).toMatch(/result: r,/); + expect(block).not.toMatch(/attempted:\s*rows\.length/); + }); +}); + +describe('ROW PRESENCE IS NEVER COMPLETION', () => { + const fs = require('fs'); + const strip = (t) => t.replace(/\/\*[\s\S]*?\*\//g, '').replace(/(^|[^:])\/\/.*$/gm, '$1'); + + test('COMPLETE is decided in exactly one place — classifyPersist', () => { + const src = strip(fs.readFileSync(require.resolve('../../src/services/retentionService'), 'utf8')); + // One definition in the TERMINAL map, one return in classifyPersist. + const assigns = src.match(/TERMINAL\.COMPLETE/g) || []; + expect(assigns.length).toBe(1); + const body = src.slice(src.indexOf('function classifyPersist')); + expect(body).toMatch(/TERMINAL\.COMPLETE/); + }); + + test('no consumer derives completion from row counts or timestamps', () => { + for (const f of ['../../src/services/snapshotService', '../../src/routes/internal.js']) { + const src = strip(fs.readFileSync(require.resolve(f), 'utf8')); + // A completion verdict must come from the writer, never from counting + // rows that survived, or from a uniform captured_at. + expect(src).not.toMatch(/status\s*[:=]\s*['"]COMPLETE['"]/); + expect(src).not.toMatch(/COMPLETE[\s\S]{0,80}(rows\.length|count\(|captured_at)/); + } + }); + + test('the call site passes the persist result through, never a substitute', () => { + const src = strip(fs.readFileSync(require.resolve('../../src/services/snapshotService'), 'utf8')); + const i = src.indexOf('retention.recordTerminal'); + const block = src.slice(i, i + 400); + expect(block).toMatch(/result: r,/); + expect(block).not.toMatch(/status:/); // no hand-set status + expect(block).not.toMatch(/attempted:/); // no recalculated counts + expect(block).not.toMatch(/written:/); + }); +}); diff --git a/tests/unit/retentionSchemaContract.test.js b/tests/unit/retentionSchemaContract.test.js index 0100b1e..9a26225 100644 --- a/tests/unit/retentionSchemaContract.test.js +++ b/tests/unit/retentionSchemaContract.test.js @@ -168,14 +168,19 @@ describe('RETENTION FAILURE IS NOT SILENT', () => { }); test('the pipeline emits a high-severity structured event on that failure', () => { + // MECHANISM CHANGED after 35da190: the alert condition was + // `r.error || (!r.skipped && r.attempted > 0 && r.written === 0)`, which + // could not distinguish a PARTIAL cohort from a complete one. It is now + // driven by the terminal status classified from the exact persist result, + // so FAILED_PARTIAL alerts as loudly as FAILED_ZERO_WRITE. 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); + const i = src.indexOf('const terminal = retention.recordTerminal'); + expect(i).toBeGreaterThan(-1); + const block = src.slice(i, i + 2200); 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/); }); @@ -183,13 +188,16 @@ describe('RETENTION FAILURE IS NOT SILENT', () => { 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/); + // The skip exemption now lives in the classifier, not in an inline + // `!r.skipped` test at the call site. + expect(retention.classifyPersist(out)).toBe(retention.TERMINAL.SKIPPED_NO_DATABASE); + expect(retention.isRetentionFailure(retention.classifyPersist(out))).toBe(false); }); 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/); + const st = retention.classifyPersist({ attempted: 600, written: 0, skipped: false, error: null }); + expect(st).toBe(retention.TERMINAL.FAILED_ZERO_WRITE); + expect(retention.isRetentionFailure(st)).toBe(true); }); }); diff --git a/tests/unit/runtimeObservability.test.js b/tests/unit/runtimeObservability.test.js index 59c21a1..6ae1e93 100644 --- a/tests/unit/runtimeObservability.test.js +++ b/tests/unit/runtimeObservability.test.js @@ -45,7 +45,9 @@ const loadConfig = (val) => { describe('RUNTIME SHA — one build identity', () => { test('the probe uses the production codeSha resolver, not git', () => { - expect(ROUTE_SRC).toMatch(/const \{ codeSha \} = require\('\.\.\/services\/retentionService'\)/); + // The destructure now also pulls `lastRetention` (terminal retention + // status), so match the resolver rather than the exact import list. + expect(ROUTE_SRC).toMatch(/const \{ codeSha[^}]*\} = require\('\.\.\/services\/retentionService'\)/); expect(ROUTE_SRC).toMatch(/code_sha: codeSha\(\)/); // Never repository state. expect(ROUTE_SRC).not.toMatch(/rev-parse|child_process|execSync|gitea|refs\/heads/);