Files
vyndr/tests/unit/retentionCompletion.test.js
T
builtbykev 11277a1b99 Materialization truth: a completed write is not a complete cohort
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) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01CQJeAG8vcDoL5zkiaJyVb8
2026-08-27 19:44:56 -04:00

279 lines
13 KiB
JavaScript

'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');
// 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\)/);
// 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', () => {
// 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/);
});
});
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.
// 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', () => {
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:/);
});
});