Retention completion: a cohort is complete only when the writer says N of N

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) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01CQJeAG8vcDoL5zkiaJyVb8
This commit is contained in:
Kev
2026-08-27 19:12:59 -04:00
parent 35da190f2c
commit 9809626c99
8 changed files with 455 additions and 39 deletions
+264
View File
@@ -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:/);
});
});
+15 -7
View File
@@ -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);
});
});
+3 -1
View File
@@ -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/);