From 11277a1b99359c8f6dc6ab10c5e0df90e6f9062a Mon Sep 17 00:00:00 2001 From: Kev Date: Thu, 27 Aug 2026 19:44:56 -0400 Subject: [PATCH] Materialization truth: a completed write is not a complete cohort MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Transport truth says every intended write succeeded. Materialization truth says every identity that should exist actually exists. The retention writer could only report the first, and the gap is not theoretical. THE CONFLICT IDENTITY, traced to the real index: model_snapshots_cycle_prop_uniq UNIQUE (snapshot_id, player_key, stat, line, side) written by `upsert(..., { onConflict: same five columns, ignoreDuplicates: true })`. Measured in production: no key column is ever NULL (0 of 328,262 rows), so NULLS DISTINCT never applies and the identity is plain column equality. `line` is an unconstrained numeric, so identity normalises it — 0.5 in and "0.50" out must not read as two identities for one stored row. PRIOR-CYCLE COLLISION IS IMPOSSIBLE. `snapshot_id` is in the identity and is a fresh UUID per cycle, so no row can be suppressed by an earlier cycle. Append-only chronology across cycles is safe, and `captured_at` is not in the identity, so a cohort cannot be split by timing. INTRA-CYCLE COLLISION IS REAL, AND WE CAUSED IT. `canonical_event_id` is NOT in the identity. A doubleheader — same hitter, same stat, same line, two genuinely different games — is ONE identity. Demonstrated through the real collector: 4 outbound rows, 2 distinct identities, 2 rows discarded by ignoreDuplicates with no error, `written` counting all 4 and the terminal status reading COMPLETE. Before event-aware dedupe the second game was dropped before grading, so the collision could not arise; that fix moved the loss downstream into retention. The conflict identity is NOT changed here — that is a separate decision with its own before/after. This makes the loss visible instead of silent. EXPECTED vs ACTUAL. `expectedMaterialization(rows)` derives the identity set from the FINAL outbound payload using the exact database identity — never from `attempted`, which counts rows sent, not identities that can exist. `reconcileMaterialization` compares SETS, not counts: two sets of equal size can still differ, and a cohort that swapped one identity for another passes every count test ever written. A collision passes set equality by construction (the discarded row was never in the expected set) while real rows were lost, so collision_count > 0 fails the cohort on its own. A cohort is evidence-complete only when transport is COMPLETE, missing = 0, extra = 0, and collisions = 0. OBSERVABILITY stayed minimal. `last_retention` was already PER SPORT (a Map keyed by sport), so no fix was needed there and the route is UNCHANGED — the new fields ride the existing entry: outbound_rows, expected_materialized_count, outbound_collision_count, expected_identity_digest. Counts and a digest only, never the identities, which carry player names. The expected set is the one materialization fact unrecoverable from the database afterwards, which is why it is the only thing recorded at runtime. A collision leaves transport COMPLETE, so the existing failure alert could never see it. It now has its own high-severity alert naming the counts, the cycle and the build, and says the cohort is not evidence-complete. Seven teeth, each injection verified present, against a GREEN baseline of 63: 1 attempted===written as evidence completeness -> 1 fail 2 COUNT(*) equality instead of set equality -> 1 fail 3 snapshot_id dropped from expected identity -> 4 fail 4 unexpected collision allowed to qualify -> 1 fail 5 single global last_retention slot -> 2 fail 6 partial chunk failure treated as usable -> 3 fail 7 collision loses its announcement -> 1 fail Restored byte-identically (retention 742f116473d97f49, snapshot 81129facbabeb280). Three brittle assertions repaired, with the reason recorded: two windowed on a byte count that a neighbouring block outgrew — a test failing because of its neighbour, not its subject — now windowed to syntactic landmarks; and one counted TERMINAL.COMPLETE occurrences, which a legitimate comparison incremented. It now asserts one DECISION and one READ. persist() and createCollector are BYTE-IDENTICAL. onConflict and ignoreDuplicates appear in the diff only as prose. Model, event, ledger, calibration, chain, lineage config, and the status route: UNCHANGED. Zero lineage/publication files, zero cacheSet changes, zero web paths. Lineage OFF. Schema contract unchanged: release 64, prod 67, prod-only 3 (debt, not authorized), missing in prod 0. 384 suites / 5,144 tests pass. web tsc exit 0. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_01CQJeAG8vcDoL5zkiaJyVb8 --- src/services/retentionService.js | 141 +++++++++- src/services/snapshotService.js | 24 ++ tests/unit/retentionCompletion.test.js | 30 +- tests/unit/retentionMaterialization.test.js | 289 ++++++++++++++++++++ tests/unit/retentionSchemaContract.test.js | 4 +- 5 files changed, 477 insertions(+), 11 deletions(-) create mode 100644 tests/unit/retentionMaterialization.test.js diff --git a/src/services/retentionService.js b/src/services/retentionService.js index d3844e6..31e6adf 100644 --- a/src/services/retentionService.js +++ b/src/services/retentionService.js @@ -727,6 +727,130 @@ function newSnapshotId() { return crypto.randomUUID(); } +/* ------------------------------------------------------------------ * + * MATERIALIZATION IDENTITY + * + * TRANSPORT COMPLETE and MATERIALIZATION COMPLETE are different facts. + * + * The write is `upsert(..., { onConflict: 'snapshot_id,player_key,stat,line,side', + * ignoreDuplicates: true })`, and the production index behind it is + * model_snapshots_cycle_prop_uniq UNIQUE (snapshot_id, player_key, stat, line, side) + * so two outbound rows sharing that tuple collapse to ONE stored row and the + * loser is discarded WITHOUT AN ERROR. `written` counts rows in committed + * chunks, so transport reports both as written and the status reads COMPLETE. + * + * `snapshot_id` IS in the identity and is a fresh UUID per cycle, so a row can + * never be suppressed by a PREVIOUS cycle — append-only chronology across + * cycles is safe. `canonical_event_id` is NOT in the identity, so two DISTINCT + * events inside ONE cycle (a doubleheader: same hitter, same stat, same line, + * both games) DO collide. + * + * No duplicate is intentional here: `rowsFromSides` emits one row per (base, + * side) and bases are already deduped by event::player::stat::line. Any + * duplicate identity in an outbound payload is therefore an UNEXPECTED + * COLLISION, and a cohort carrying one does not qualify as evidence. + * + * The conflict identity is NOT changed here. This only makes the loss visible. + * ------------------------------------------------------------------ */ + +/** The exact columns of model_snapshots_cycle_prop_uniq, in index order. */ +const CONFLICT_IDENTITY = Object.freeze(['snapshot_id', 'player_key', 'stat', 'line', 'side']); + +const IDENTITY_SEP = '\u001f'; +const IDENTITY_NULL = '\u0000NULL'; + +/** + * The database identity of one row. `line` is an unconstrained `numeric`, so it + * is normalised through Number() — otherwise 0.5 on the way in and "0.50" on + * the way out would read as two identities for one stored row. + * + * Every key column is NOT NULL in production (measured: 0 nulls across 328,262 + * rows), so PostgreSQL's NULLS DISTINCT behaviour never applies. A null still + * gets its own marker rather than collapsing to '' — an empty string is a real + * value and must not be confused with an absent one. + */ +function rowIdentity(row) { + return CONFLICT_IDENTITY.map((c) => { + const v = row ? row[c] : undefined; + if (v === null || v === undefined) return IDENTITY_NULL; + return c === 'line' ? String(Number(v)) : String(v); + }).join(IDENTITY_SEP); +} + +function identityDigest(identities) { + return crypto.createHash('sha256').update([...identities].sort().join('\n')).digest('hex').slice(0, 32); +} + +/** + * What SHOULD materialize, derived from the FINAL outbound payload using the + * exact database identity — never from `attempted`, which counts rows sent, not + * identities that can exist. + */ +function expectedMaterialization(rows) { + const list = Array.isArray(rows) ? rows : []; + const seen = new Set(); + let collisions = 0; + for (const r of list) { + const id = rowIdentity(r); + if (seen.has(id)) collisions += 1; + else seen.add(id); + } + return { + outbound_rows: list.length, + expected_identities: seen.size, + collision_count: collisions, + digest: identityDigest(seen), + identities: seen, + }; +} + +const MATERIALIZATION = Object.freeze({ + COMPLETE: 'MATERIALIZATION_COMPLETE', + MISSING: 'MATERIALIZATION_MISSING', + EXTRA: 'MATERIALIZATION_EXTRA', + COLLISION: 'MATERIALIZATION_UNEXPECTED_COLLISION', + TRANSPORT_FAILED: 'MATERIALIZATION_UNPROVEN_TRANSPORT_FAILED', + NOT_RECONCILED: 'MATERIALIZATION_NOT_RECONCILED', +}); + +/** + * Exact SET comparison, not a count comparison. Two sets of equal size can + * still differ, and a cohort that swapped one identity for another would pass + * every count test ever written. + * + * A cohort qualifies ONLY when transport completed, the sets are equal, AND the + * outbound payload held no colliding identity — a collision passes set equality + * by construction (the discarded row was never in the expected set) while real + * rows were lost. + */ +function reconcileMaterialization({ expected, actualIdentities, transportStatus }) { + if (!expected || typeof expected.expected_identities !== 'number') { + throw new TypeError('reconcileMaterialization requires an expectedMaterialization() result'); + } + const actual = actualIdentities instanceof Set ? actualIdentities : new Set(actualIdentities || []); + const exp = expected.identities instanceof Set ? expected.identities : new Set(); + const missing = [...exp].filter((i) => !actual.has(i)); + const extra = [...actual].filter((i) => !exp.has(i)); + const out = { + expected_materialized_count: expected.expected_identities, + actual_materialized_count: actual.size, + outbound_rows: expected.outbound_rows, + missing_identity_count: missing.length, + extra_identity_count: extra.length, + collision_count: expected.collision_count, + status: MATERIALIZATION.NOT_RECONCILED, + }; + if (transportStatus && transportStatus !== TERMINAL.COMPLETE) { + out.status = MATERIALIZATION.TRANSPORT_FAILED; + return out; + } + if (missing.length) out.status = MATERIALIZATION.MISSING; + else if (extra.length) out.status = MATERIALIZATION.EXTRA; + else if (expected.collision_count > 0) out.status = MATERIALIZATION.COLLISION; + else out.status = MATERIALIZATION.COMPLETE; + return out; +} + /* ------------------------------------------------------------------ * * TERMINAL RETENTION STATE * @@ -795,8 +919,12 @@ function classifyPersist(result) { */ const lastTerminal = new Map(); -function recordTerminal({ sport, snapshotId, result, completedAt }) { +function recordTerminal({ sport, snapshotId, result, completedAt, rows }) { const status = classifyPersist(result); + // Derived from the FINAL outbound payload, which does not survive the call. + // This is the one materialization fact unrecoverable from the database + // afterwards, which is why it is the only one recorded at runtime. + const expected = rows ? expectedMaterialization(rows) : null; const entry = Object.freeze({ sport: sport || null, snapshot_id: snapshotId || null, @@ -807,6 +935,11 @@ function recordTerminal({ sport, snapshotId, result, completedAt }) { code_sha: codeSha(), // Message only — never row payloads. error_summary: result.error ? String(result.error).slice(0, 300) : null, + // Counts and a digest only — never the identities, which carry player names. + outbound_rows: expected ? expected.outbound_rows : null, + expected_materialized_count: expected ? expected.expected_identities : null, + outbound_collision_count: expected ? expected.collision_count : null, + expected_identity_digest: expected ? expected.digest : null, }); if (sport) lastTerminal.set(sport, entry); return entry; @@ -823,6 +956,12 @@ function resetTerminal() { lastTerminal.clear(); } module.exports = { MODEL_VERSION, TERMINAL, + MATERIALIZATION, + CONFLICT_IDENTITY, + rowIdentity, + identityDigest, + expectedMaterialization, + reconcileMaterialization, classifyPersist, isRetentionFailure, recordTerminal, diff --git a/src/services/snapshotService.js b/src/services/snapshotService.js index 1504332..1133251 100644 --- a/src/services/snapshotService.js +++ b/src/services/snapshotService.js @@ -667,6 +667,9 @@ async function runSnapshot(sport, opts = {}) { snapshotId: retentionCtx.snapshotId, result: r, completedAt: deps.now(), + // The outbound payload, so the expected materialized identity set is + // derived from what was actually sent. It does not survive this call. + rows, }); 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. @@ -681,6 +684,27 @@ async function runSnapshot(sport, opts = {}) { // healthy. It alerts exactly as loudly as a total failure. // // NOTHING_TO_PERSIST and SKIPPED_NO_DATABASE are not failures. + // AN UNEXPECTED COLLISION IS SILENT DATA LOSS WITH A HEALTHY TRANSPORT. + // + // Two outbound rows sharing (snapshot_id, player_key, stat, line, side) + // collapse to one stored row under ignoreDuplicates, with no error. The + // transport reports both as written and the status reads COMPLETE, so the + // failure alert below can never see it. The known cause is a doubleheader: + // canonical_event_id is not in the conflict identity, so the same hitter's + // same line in two real games is one identity. + // + // The cohort is NOT evidence-complete. The conflict identity is not + // changed here; this only refuses to lose the rows quietly. + if (terminal.outbound_collision_count > 0) { + await deps.notify( + `Retention UNEXPECTED COLLISION for ${sp.toUpperCase()} — ${terminal.outbound_collision_count} of ` + + `${terminal.outbound_rows} outbound rows share a conflict identity and were discarded by ` + + `ignoreDuplicates. transport=${terminal.status} expected_identities=${terminal.expected_materialized_count} ` + + `snapshot_id=${terminal.snapshot_id} code_sha=${terminal.code_sha || 'unknown'} at=${terminal.completed_at}. ` + + 'This retention cohort is NOT evidence-complete.', + { title: 'VYNDR pipeline', priority: 'high', tags: ['rotating_light'] }, + ); + } if (retention.isRetentionFailure(terminal.status)) { await deps.notify( `Retention ${terminal.status} for ${sp.toUpperCase()} — ${terminal.written}/${terminal.attempted} rows persisted. ` diff --git a/tests/unit/retentionCompletion.test.js b/tests/unit/retentionCompletion.test.js index 864bfa5..09db1f5 100644 --- a/tests/unit/retentionCompletion.test.js +++ b/tests/unit/retentionCompletion.test.js @@ -200,8 +200,11 @@ describe('COMPLETION OBSERVABILITY', () => { 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); + // Windowed to a SYNTACTIC LANDMARK, not a byte count: a byte window breaks + // the moment a neighbouring block grows, which is a test failing because of + // its neighbour rather than its subject. + const block = src.slice(src.indexOf('if (retention.isRetentionFailure(terminal.status))'), + src.indexOf('} catch (e) {', src.indexOf('if (retention.isRetentionFailure(terminal.status))'))); test('the alert is driven by TERMINAL STATUS, not by a written===0 test', () => { expect(block).toMatch(/retention\.isRetentionFailure\(terminal\.status\)/); @@ -224,8 +227,11 @@ describe('PARTIAL FAILURE IS VISIBLE AT THE CALL SITE', () => { }); test('the terminal result recorded is the one persist returned', () => { - expect(block).toMatch(/result: r,/); - expect(block).not.toMatch(/attempted:\s*rows\.length/); + // Its own window: this asserts on the RECORDER call, not the alert. + const i = src.indexOf('const terminal = retention.recordTerminal'); + const rec = src.slice(i, src.indexOf('});', i) + 3); + expect(rec).toMatch(/result: r,/); + expect(rec).not.toMatch(/attempted:\s*rows\.length/); }); }); @@ -236,10 +242,18 @@ describe('ROW PRESENCE IS NEVER COMPLETION', () => { 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/); + // TERMINAL.COMPLETE is RETURNED in exactly one place. It is also READ once, + // by reconcileMaterialization, to refuse a cohort whose transport failed — + // a comparison, not a second decision. + // Exactly one DECISION, inside the classifier. + const classifier = src.slice(src.indexOf('function classifyPersist')); + expect((classifier.match(/TERMINAL\.COMPLETE/g) || []).length).toBe(1); + // And exactly one READ elsewhere — reconcileMaterialization refusing a + // cohort whose transport failed. A comparison, not a second decision. + const reconciler = src.slice(src.indexOf('function reconcileMaterialization'), + src.indexOf('function classifyPersist')); + expect(reconciler).toMatch(/transportStatus !== TERMINAL\.COMPLETE/); + expect((src.match(/TERMINAL\.COMPLETE/g) || []).length).toBe(2); }); test('no consumer derives completion from row counts or timestamps', () => { diff --git a/tests/unit/retentionMaterialization.test.js b/tests/unit/retentionMaterialization.test.js new file mode 100644 index 0000000..1a43dcf --- /dev/null +++ b/tests/unit/retentionMaterialization.test.js @@ -0,0 +1,289 @@ +'use strict'; + +/** + * TRANSPORT TRUTH vs MATERIALIZATION TRUTH. + * + * Transport truth says every intended write operation succeeded. Materialization + * truth says every identity that should exist actually exists. They are not the + * same fact, because the write is + * + * upsert(rows, { onConflict: 'snapshot_id,player_key,stat,line,side', + * ignoreDuplicates: true }) + * + * against the production index + * + * model_snapshots_cycle_prop_uniq UNIQUE (snapshot_id, player_key, stat, line, side) + * + * so two outbound rows sharing that tuple collapse to ONE stored row, silently, + * while `written` counts both and the terminal status reads COMPLETE. + * + * Measured against production 2026-08-27: no key column is ever NULL + * (0 of 328,262 rows), so NULLS DISTINCT never applies; `line` is an + * unconstrained numeric, which is why identity normalises it. + */ + +const retention = require('../../src/services/retentionService'); +const { TERMINAL, MATERIALIZATION } = retention; + +const SNAP = '11111111-1111-1111-1111-111111111111'; +const OTHER_SNAP = '22222222-2222-2222-2222-222222222222'; + +const row = (o = {}) => ({ + snapshot_id: SNAP, player_key: 'nolan arenado', stat: 'hits', line: 0.5, side: 'over', ...o, +}); + +/** The identity a real database row would present when read back. */ +const actualSetFrom = (rows) => new Set(rows.map(retention.rowIdentity)); + +describe('THE CONFLICT IDENTITY IS THE REAL DATABASE IDENTITY', () => { + test('it is exactly model_snapshots_cycle_prop_uniq', () => { + expect(retention.CONFLICT_IDENTITY).toEqual( + ['snapshot_id', 'player_key', 'stat', 'line', 'side']); + }); + + test('the production upsert names that same identity', () => { + const src = require('fs').readFileSync(require.resolve('../../src/services/retentionService'), 'utf8'); + const body = src.slice(src.indexOf('async function persist(')); + expect(body).toMatch(/onConflict: 'snapshot_id,player_key,stat,line,side'/); + expect(body).toMatch(/ignoreDuplicates: true/); + }); + + test('numeric line normalises — 0.5 and "0.50" are ONE identity', () => { + expect(retention.rowIdentity(row({ line: 0.5 }))) + .toBe(retention.rowIdentity(row({ line: '0.50' }))); + }); + + test('a null key value is marked, never collapsed into an empty string', () => { + expect(retention.rowIdentity(row({ side: null }))) + .not.toBe(retention.rowIdentity(row({ side: '' }))); + }); +}); + +describe('PRIOR-CYCLE COLLISION IS STRUCTURALLY IMPOSSIBLE', () => { + test('snapshot_id is part of the identity', () => { + expect(retention.CONFLICT_IDENTITY).toContain('snapshot_id'); + }); + + test('captured_at is NOT part of the identity', () => { + expect(retention.CONFLICT_IDENTITY).not.toContain('captured_at'); + }); + + test('the same prop in a different cycle is a DIFFERENT identity', () => { + expect(retention.rowIdentity(row())) + .not.toBe(retention.rowIdentity(row({ snapshot_id: OTHER_SNAP }))); + }); + + test('snapshot_id is freshly generated per cycle', () => { + expect(retention.newSnapshotId()).not.toBe(retention.newSnapshotId()); + const src = require('fs').readFileSync(require.resolve('../../src/services/snapshotService'), 'utf8'); + expect(src).toMatch(/retention\.newSnapshotId\(\)/); + }); + + test('so a new-cycle row can never be suppressed by an old-cycle row', () => { + // Both cycles' rows coexist: append-only chronology across cycles is safe. + const cycleA = [row()]; + const cycleB = [row({ snapshot_id: OTHER_SNAP })]; + const ids = actualSetFrom([...cycleA, ...cycleB]); + expect(ids.size).toBe(2); + }); +}); + +describe('EXPECTED MATERIALIZATION', () => { + test('NO DUPLICATES — every expected identity materializes', () => { + const rows = [row({ side: 'over' }), row({ side: 'under' }), row({ stat: 'total_bases' })]; + const exp = retention.expectedMaterialization(rows); + expect(exp.outbound_rows).toBe(3); + expect(exp.expected_identities).toBe(3); + expect(exp.collision_count).toBe(0); + const rec = retention.reconcileMaterialization({ + expected: exp, actualIdentities: actualSetFrom(rows), transportStatus: TERMINAL.COMPLETE, + }); + expect(rec.status).toBe(MATERIALIZATION.COMPLETE); + expect([rec.missing_identity_count, rec.extra_identity_count]).toEqual([0, 0]); + }); + + test('expected is derived from the payload, NOT from attempted', () => { + // `attempted` counts rows SENT. Using it as the denominator would make a + // collided cohort look short by exactly the rows the database discarded, + // and a deduped one look broken. + const rows = [row(), row()]; // identical -> one identity + const exp = retention.expectedMaterialization(rows); + expect(exp.outbound_rows).toBe(2); + expect(exp.expected_identities).toBe(1); + }); + + test('it refuses anything but an expectedMaterialization result', () => { + expect(() => retention.reconcileMaterialization({ expected: null })).toThrow(TypeError); + expect(() => retention.reconcileMaterialization({ expected: { count: 3 } })).toThrow(TypeError); + }); +}); + +describe('UNEXPECTED COLLISION — the doubleheader', () => { + // canonical_event_id is NOT in the conflict identity, so the same hitter's + // same line in two REAL games is one identity. Before event-aware dedupe the + // second game was dropped before grading; now both grade, and both collide. + const rows = [ + row({ canonical_event_id: 'mlb:gamepk:824514', side: 'over' }), + row({ canonical_event_id: 'mlb:gamepk:824514', side: 'under' }), + row({ canonical_event_id: 'mlb:gamepk:824478', side: 'over' }), + row({ canonical_event_id: 'mlb:gamepk:824478', side: 'under' }), + ]; + + test('two semantically distinct records collapse under the database identity', () => { + const exp = retention.expectedMaterialization(rows); + expect(exp.outbound_rows).toBe(4); + expect(exp.expected_identities).toBe(2); + expect(exp.collision_count).toBe(2); + }); + + test('the cohort does NOT silently qualify, even though set equality holds', () => { + const exp = retention.expectedMaterialization(rows); + // The discarded rows were never in the expected set, so a naive set + // comparison passes while two real records were lost. + const rec = retention.reconcileMaterialization({ + expected: exp, actualIdentities: exp.identities, transportStatus: TERMINAL.COMPLETE, + }); + expect(rec.missing_identity_count).toBe(0); + expect(rec.extra_identity_count).toBe(0); + expect(rec.status).toBe(MATERIALIZATION.COLLISION); + expect(rec.status).not.toBe(MATERIALIZATION.COMPLETE); + }); + + test('transport still reports COMPLETE — which is why transport alone is not evidence', () => { + const st = retention.classifyPersist({ attempted: 4, written: 4, skipped: false, error: null }); + expect(st).toBe(TERMINAL.COMPLETE); + expect(retention.isRetentionFailure(st)).toBe(false); + }); +}); + +describe('MATERIALIZATION FAILURES', () => { + const rows = [row({ side: 'over' }), row({ side: 'under' }), row({ stat: 'rbi' })]; + + test('COMPLETE TRANSPORT + MISSING EXPECTED IDENTITY -> fails', () => { + const exp = retention.expectedMaterialization(rows); + const actual = actualSetFrom(rows.slice(0, 2)); + const rec = retention.reconcileMaterialization({ + expected: exp, actualIdentities: actual, transportStatus: TERMINAL.COMPLETE, + }); + expect(rec.missing_identity_count).toBe(1); + expect(rec.status).toBe(MATERIALIZATION.MISSING); + }); + + test('COMPLETE TRANSPORT + EXTRA UNEXPECTED IDENTITY -> fails', () => { + const exp = retention.expectedMaterialization(rows); + const actual = actualSetFrom([...rows, row({ stat: 'runs' })]); + const rec = retention.reconcileMaterialization({ + expected: exp, actualIdentities: actual, transportStatus: TERMINAL.COMPLETE, + }); + expect(rec.extra_identity_count).toBe(1); + expect(rec.status).toBe(MATERIALIZATION.EXTRA); + }); + + test('EQUAL COUNTS, DIFFERENT SETS -> fails (count equality is not set equality)', () => { + const exp = retention.expectedMaterialization(rows); + const swapped = actualSetFrom([rows[0], rows[1], row({ stat: 'runs' })]); + const rec = retention.reconcileMaterialization({ + expected: exp, actualIdentities: swapped, transportStatus: TERMINAL.COMPLETE, + }); + expect(rec.expected_materialized_count).toBe(rec.actual_materialized_count); // counts agree + expect(rec.missing_identity_count).toBe(1); + expect(rec.extra_identity_count).toBe(1); + expect(rec.status).not.toBe(MATERIALIZATION.COMPLETE); + }); + + test('PARTIAL CHUNK FAILURE -> cohort cannot qualify however many rows persisted', () => { + const exp = retention.expectedMaterialization(rows); + const rec = retention.reconcileMaterialization({ + expected: exp, + actualIdentities: actualSetFrom(rows), // rows ARE all present + transportStatus: TERMINAL.FAILED_PARTIAL, + }); + expect(rec.status).toBe(MATERIALIZATION.TRANSPORT_FAILED); + expect(rec.status).not.toBe(MATERIALIZATION.COMPLETE); + }); + + test('row presence alone never qualifies a cohort', () => { + for (const st of [TERMINAL.FAILED_PARTIAL, TERMINAL.FAILED_ZERO_WRITE, TERMINAL.FAILED_UNRESOLVED_ERROR]) { + const rec = retention.reconcileMaterialization({ + expected: retention.expectedMaterialization(rows), + actualIdentities: actualSetFrom(rows), + transportStatus: st, + }); + expect(rec.status).toBe(MATERIALIZATION.TRANSPORT_FAILED); + } + }); +}); + +describe('PER-SPORT TERMINAL EVIDENCE', () => { + beforeEach(() => retention.resetTerminal()); + + test('last_retention is keyed per sport, not a single global slot', async () => { + const ok = { from: () => ({ upsert: async () => ({ error: null }) }) }; + const rowsA = [row({ side: 'over' })]; + const rowsB = [row({ snapshot_id: OTHER_SNAP, side: 'under' })]; + retention.recordTerminal({ + sport: 'mlb', snapshotId: SNAP, rows: rowsA, + result: await retention.persist(rowsA, { getClient: () => ok }), + }); + retention.recordTerminal({ + sport: 'wnba', snapshotId: OTHER_SNAP, rows: rowsB, + result: await retention.persist(rowsB, { getClient: () => ok }), + }); + const last = retention.lastRetention(); + // A later sport must not overwrite the target sport's evidence. + expect(last.mlb.snapshot_id).toBe(SNAP); + expect(last.wnba.snapshot_id).toBe(OTHER_SNAP); + expect(Object.keys(last).sort()).toEqual(['mlb', 'wnba']); + }); + + test('the terminal record carries the expected-materialization facts', () => { + const rows = [ + row({ canonical_event_id: 'g1' }), row({ canonical_event_id: 'g2' }), // collide + ]; + const e = retention.recordTerminal({ + sport: 'mlb', snapshotId: SNAP, rows, + result: { attempted: 2, written: 2, skipped: false, error: null }, + }); + expect(e.outbound_rows).toBe(2); + expect(e.expected_materialized_count).toBe(1); + expect(e.outbound_collision_count).toBe(1); + expect(e.expected_identity_digest).toMatch(/^[0-9a-f]{32}$/); + }); + + test('it exposes counts and a digest, never the identities themselves', () => { + const e = retention.recordTerminal({ + sport: 'mlb', snapshotId: SNAP, rows: [row()], + result: { attempted: 1, written: 1, skipped: false, error: null }, + }); + expect(JSON.stringify(e)).not.toMatch(/nolan arenado/); + expect(e).not.toHaveProperty('identities'); + }); + + test('the call site hands the outbound payload to the recorder', () => { + const src = require('fs').readFileSync(require.resolve('../../src/services/snapshotService'), 'utf8'); + const i = src.indexOf('retention.recordTerminal'); + const block = src.slice(i, i + 500); + expect(block).toMatch(/\brows,/); + expect(block).toMatch(/result: r,/); + }); +}); + +describe('AN UNEXPECTED COLLISION IS ANNOUNCED, NOT INFERRED', () => { + const src = require('fs').readFileSync(require.resolve('../../src/services/snapshotService'), 'utf8'); + const block = src.slice(src.indexOf('UNEXPECTED COLLISION IS SILENT DATA LOSS'), + src.indexOf('UNEXPECTED COLLISION IS SILENT DATA LOSS') + 1800); + + test('it alerts even though the transport status is COMPLETE', () => { + expect(block).toMatch(/terminal\.outbound_collision_count > 0/); + expect(block).toMatch(/priority: 'high'/); + }); + + test('the alert names the counts, the cycle and the build', () => { + for (const f of ['${terminal.outbound_collision_count}', '${terminal.outbound_rows}', + 'transport=${terminal.status}', 'expected_identities=${terminal.expected_materialized_count}', + 'snapshot_id=${terminal.snapshot_id}', 'code_sha=${terminal.code_sha']) { + expect(block).toContain(f); + } + expect(block).toMatch(/NOT evidence-complete/); + }); +}); diff --git a/tests/unit/retentionSchemaContract.test.js b/tests/unit/retentionSchemaContract.test.js index 9a26225..e11508b 100644 --- a/tests/unit/retentionSchemaContract.test.js +++ b/tests/unit/retentionSchemaContract.test.js @@ -174,9 +174,9 @@ describe('RETENTION FAILURE IS NOT SILENT', () => { // 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 i = src.indexOf('const terminal = retention.recordTerminal'); + const i = src.indexOf('if (retention.isRetentionFailure(terminal.status))'); expect(i).toBeGreaterThan(-1); - const block = src.slice(i, i + 2200); + const block = src.slice(i, src.indexOf('} catch (e) {', i)); for (const field of ['stage=model_snapshots', 'snapshot_id=', 'code_sha=', 'at=', 'error=']) { expect(block).toContain(field); }