|
|
|
@@ -0,0 +1,302 @@
|
|
|
|
|
/**
|
|
|
|
|
* LINEAGE PRODUCTIZATION CLOSEOUT — automatic coverage monitoring and the
|
|
|
|
|
* additive ancestry product contract.
|
|
|
|
|
*
|
|
|
|
|
* Two properties: the monitor speaks without being asked and never lies in the
|
|
|
|
|
* reassuring direction; the product contract owns ancestry and nothing else.
|
|
|
|
|
*/
|
|
|
|
|
|
|
|
|
|
const coverage = require('../../src/services/lineageCoverage');
|
|
|
|
|
const ancestry = require('../../src/services/read/readAncestry');
|
|
|
|
|
const fs = require('fs');
|
|
|
|
|
const path = require('path');
|
|
|
|
|
const ROOT = path.resolve(__dirname, '..', '..');
|
|
|
|
|
|
|
|
|
|
/* ───────────────── AUTOMATIC MONITOR ───────────────── */
|
|
|
|
|
describe('the monitor runs without being asked', () => {
|
|
|
|
|
test('the scheduler invokes the durable observer on its own cadence', () => {
|
|
|
|
|
const src = fs.readFileSync(path.join(ROOT, 'src/snapshotScheduler.js'), 'utf8');
|
|
|
|
|
expect(src).toMatch(/coverageTick/);
|
|
|
|
|
expect(src).toMatch(/lineageCoverage\.auditLatestCohort/);
|
|
|
|
|
// Called from the per-minute tick, not from a slot branch: a monitor that
|
|
|
|
|
// only runs when a snapshot fires cannot report that no snapshot fired.
|
|
|
|
|
const tick = src.slice(src.indexOf('const tick = async () =>'));
|
|
|
|
|
expect(tick.slice(0, 300)).toMatch(/await coverageTick\(\)/);
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
test('the monitor consumes the observer, never the writer', () => {
|
|
|
|
|
const src = fs.readFileSync(path.join(ROOT, 'src/snapshotScheduler.js'), 'utf8');
|
|
|
|
|
const fn = src.slice(src.indexOf('const coverageTick'), src.indexOf('const tick = async () =>'));
|
|
|
|
|
expect(fn).not.toMatch(/attachLineage|lineage\.origins|pub\.lineage/);
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
test('cadence is a pure predicate', () => {
|
|
|
|
|
expect(coverage.coverageDue(null, 1_000)).toBe(true);
|
|
|
|
|
expect(coverage.coverageDue(1_000, 1_000 + 29 * 60_000)).toBe(false);
|
|
|
|
|
expect(coverage.coverageDue(1_000, 1_000 + 31 * 60_000)).toBe(true);
|
|
|
|
|
expect(coverage.coverageDue(1_000, Number.NaN)).toBe(false);
|
|
|
|
|
});
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
describe('alert policy — silence only when healthy', () => {
|
|
|
|
|
const A = (h, extra = {}) => coverage.coverageAlarm(null, { health: h, sport: 'mlb', snapshot_id: 'abc12345', ...extra });
|
|
|
|
|
|
|
|
|
|
test('HEALTHY and NO_ELIGIBLE_CLAIMS are silent', () => {
|
|
|
|
|
expect(A(coverage.HEALTH.HEALTHY).alert).toBe(false);
|
|
|
|
|
expect(A(coverage.HEALTH.NO_ELIGIBLE_CLAIMS).alert).toBe(false);
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
test('every unhealthy state alerts', () => {
|
|
|
|
|
for (const h of [coverage.HEALTH.MISSING_COVERAGE, coverage.HEALTH.PARTIAL_COVERAGE,
|
|
|
|
|
coverage.HEALTH.INVALID_GRAPH, coverage.HEALTH.STALE, coverage.HEALTH.AUDIT_UNAVAILABLE]) {
|
|
|
|
|
expect(A(h).alert).toBe(true);
|
|
|
|
|
expect(typeof A(h).message).toBe('string');
|
|
|
|
|
}
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
test('AUDIT_UNAVAILABLE says it could not measure — it is NOT a coverage result', () => {
|
|
|
|
|
const m = A(coverage.HEALTH.AUDIT_UNAVAILABLE, { reason: 'statement timeout' }).message;
|
|
|
|
|
expect(m).toMatch(/COULD NOT BE MEASURED/);
|
|
|
|
|
expect(m).toMatch(/did not run/);
|
|
|
|
|
expect(m).not.toMatch(/0 recorded|MISSING/);
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
test('a standing condition states itself once', () => {
|
|
|
|
|
const first = coverage.coverageAlarm(null, { health: coverage.HEALTH.MISSING_COVERAGE, snapshot_id: 's1' });
|
|
|
|
|
expect(first.alert).toBe(true);
|
|
|
|
|
const second = coverage.coverageAlarm(first.key, { health: coverage.HEALTH.MISSING_COVERAGE, snapshot_id: 's1' });
|
|
|
|
|
expect(second.alert).toBe(false);
|
|
|
|
|
// A NEW cohort with the same fault speaks again.
|
|
|
|
|
const third = coverage.coverageAlarm(first.key, { health: coverage.HEALTH.MISSING_COVERAGE, snapshot_id: 's2' });
|
|
|
|
|
expect(third.alert).toBe(true);
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
test('an unreadable audit never reports HEALTHY', () => {
|
|
|
|
|
expect(coverage.classify({ audit_available: false })).toBe(coverage.HEALTH.AUDIT_UNAVAILABLE);
|
|
|
|
|
expect(A(coverage.HEALTH.AUDIT_UNAVAILABLE).alert).toBe(true);
|
|
|
|
|
});
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
describe('STALE is provable without an outage', () => {
|
|
|
|
|
test('a healthy cohort older than the window becomes STALE by injected time', async () => {
|
|
|
|
|
const rows = [{
|
|
|
|
|
id: 1, sport: 'mlb', game_date: '2026-08-31', player_key: 'a', stat: 'hits', side: 'over',
|
|
|
|
|
line: 0.5, published: true, captured_at: '2026-08-31T03:00:00Z',
|
|
|
|
|
read_id: 'r1', read_natural_key: 'mlb|2026-08-31|a|hits|over|0.500000',
|
|
|
|
|
lineage_action: 'ORIGIN', claim_digest: 'd', revision_ordinal: 0, lineage_state: 's',
|
|
|
|
|
lineage_version: 'lin@1', claim_schema_version: 'claim@1', digest_algorithm_version: 'x',
|
|
|
|
|
}];
|
|
|
|
|
const client = () => ({
|
|
|
|
|
from: () => ({
|
|
|
|
|
select: () => ({
|
|
|
|
|
eq() { return this; }, gte() { return this; }, lte() { return this; },
|
|
|
|
|
order() { return this; },
|
|
|
|
|
limit: async () => ({ data: [{ snapshot_id: 's1', game_date: '2026-08-31', captured_at: '2026-08-31T03:00:00Z', code_sha: 'sha' }], error: null }),
|
|
|
|
|
range: async (from) => ({ data: from === 0 ? rows : [], error: null }),
|
|
|
|
|
}),
|
|
|
|
|
}),
|
|
|
|
|
});
|
|
|
|
|
const fresh = await coverage.auditLatestCohort('mlb', { getClient: client, now: () => '2026-08-31T04:00:00Z' });
|
|
|
|
|
expect(fresh.health).toBe(coverage.HEALTH.HEALTHY);
|
|
|
|
|
const stale = await coverage.auditLatestCohort('mlb', { getClient: client, now: () => '2026-09-01T12:00:00Z' });
|
|
|
|
|
expect(stale.health).toBe(coverage.HEALTH.STALE);
|
|
|
|
|
});
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
/** A test double must HONOUR a filter, not merely accept one: a pass-through
|
|
|
|
|
* fake would return every row regardless of the query and make a narrowed
|
|
|
|
|
* cohort walk look correct. */
|
|
|
|
|
function fakeClient(rows, headRow) {
|
|
|
|
|
const build = () => {
|
|
|
|
|
const preds = [];
|
|
|
|
|
const q = {
|
|
|
|
|
eq(col, val) { preds.push((r) => String(r[col]) === String(val)); return this; },
|
|
|
|
|
gte(col, val) { preds.push((r) => String(r[col]) >= String(val)); return this; },
|
|
|
|
|
lte(col, val) { preds.push((r) => String(r[col]) <= String(val)); return this; },
|
|
|
|
|
order() { return this; },
|
|
|
|
|
limit: async () => ({ data: [headRow], error: null }),
|
|
|
|
|
range: async (from) => ({
|
|
|
|
|
data: from === 0 ? rows.filter((r) => preds.every((p) => p(r))) : [],
|
|
|
|
|
error: null,
|
|
|
|
|
}),
|
|
|
|
|
};
|
|
|
|
|
return q;
|
|
|
|
|
};
|
|
|
|
|
return () => ({ from: () => ({ select: () => build() }) });
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
describe('a cohort spanning two game_dates is walked whole (behavioural)', () => {
|
|
|
|
|
const act = (d, p) => ({
|
|
|
|
|
id: p === 'a' ? 1 : 2, sport: 'mlb', game_date: d, player_key: p, stat: 'hits',
|
|
|
|
|
side: 'over', line: 0.5, published: true, captured_at: '2026-08-31T03:00:00Z',
|
|
|
|
|
snapshot_id: 's1', code_sha: 'sha',
|
|
|
|
|
read_id: `r-${p}`, read_natural_key: `mlb|${d}|${p}|hits|over|0.500000`,
|
|
|
|
|
lineage_action: 'ORIGIN', claim_digest: 'd', revision_ordinal: 0, lineage_state: 's',
|
|
|
|
|
lineage_version: 'lin@1', claim_schema_version: 'claim@1', digest_algorithm_version: 'x',
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
test('both date slices are counted, not just the head row\'s date', async () => {
|
|
|
|
|
const rows = [act('2026-08-31', 'a'), act('2026-08-30', 'b')];
|
|
|
|
|
const head = { snapshot_id: 's1', game_date: '2026-08-31', captured_at: '2026-08-31T03:00:00Z', code_sha: 'sha' };
|
|
|
|
|
const out = await coverage.auditLatestCohort('mlb', {
|
|
|
|
|
getClient: fakeClient(rows, head), now: () => '2026-08-31T03:30:00Z',
|
|
|
|
|
});
|
|
|
|
|
expect(out.rows).toBe(2);
|
|
|
|
|
expect(out.expected_keys).toBe(2);
|
|
|
|
|
expect(out.covered_keys).toBe(2);
|
|
|
|
|
expect(out.game_dates_audited).toEqual(['2026-08-30', '2026-08-31']);
|
|
|
|
|
expect(out.health).toBe(coverage.HEALTH.HEALTHY);
|
|
|
|
|
});
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
describe('the monitor cannot break the scheduler (behavioural)', () => {
|
|
|
|
|
const OLD = process.env.SNAPSHOT_CRON;
|
|
|
|
|
afterEach(() => { if (OLD === undefined) delete process.env.SNAPSHOT_CRON; else process.env.SNAPSHOT_CRON = OLD; });
|
|
|
|
|
|
|
|
|
|
test('an observer that throws is swallowed, and no alert is attempted', async () => {
|
|
|
|
|
process.env.SNAPSHOT_CRON = '1';
|
|
|
|
|
const { startSnapshotScheduler } = require('../../src/snapshotScheduler');
|
|
|
|
|
const notifies = [];
|
|
|
|
|
const sched = startSnapshotScheduler({
|
|
|
|
|
runAllSnapshots: async () => [],
|
|
|
|
|
notify: async (m) => { notifies.push(m); },
|
|
|
|
|
lineageCoverage: {
|
|
|
|
|
coverageDue: () => true,
|
|
|
|
|
auditLatestCohort: async () => { throw new Error('observer exploded'); },
|
|
|
|
|
coverageAlarm: () => { throw new Error('should not be reached'); },
|
|
|
|
|
},
|
|
|
|
|
});
|
|
|
|
|
expect(sched).toBeTruthy();
|
|
|
|
|
if (sched.interval) clearInterval(sched.interval);
|
|
|
|
|
// Must resolve, not reject. A monitor that can take down the pipeline it
|
|
|
|
|
// watches is worse than no monitor.
|
|
|
|
|
await expect(sched.coverageTick()).resolves.toBeUndefined();
|
|
|
|
|
expect(notifies).toEqual([]);
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
test('a healthy audit alerts nothing; an unhealthy one alerts once', async () => {
|
|
|
|
|
process.env.SNAPSHOT_CRON = '1';
|
|
|
|
|
const { startSnapshotScheduler } = require('../../src/snapshotScheduler');
|
|
|
|
|
const notifies = [];
|
|
|
|
|
const real = require('../../src/services/lineageCoverage');
|
|
|
|
|
let health = coverage.HEALTH.HEALTHY;
|
|
|
|
|
const sched = startSnapshotScheduler({
|
|
|
|
|
runAllSnapshots: async () => [],
|
|
|
|
|
notify: async (m) => { notifies.push(m); },
|
|
|
|
|
lineageCoverage: {
|
|
|
|
|
coverageDue: () => true,
|
|
|
|
|
auditLatestCohort: async () => ({ health, sport: 'mlb', snapshot_id: 's1', expected_keys: 801, covered_keys: 0 }),
|
|
|
|
|
coverageAlarm: real.coverageAlarm,
|
|
|
|
|
HEALTH: real.HEALTH,
|
|
|
|
|
},
|
|
|
|
|
});
|
|
|
|
|
if (sched.interval) clearInterval(sched.interval);
|
|
|
|
|
await sched.coverageTick();
|
|
|
|
|
expect(notifies).toEqual([]);
|
|
|
|
|
health = coverage.HEALTH.MISSING_COVERAGE;
|
|
|
|
|
await sched.coverageTick();
|
|
|
|
|
expect(notifies).toHaveLength(1);
|
|
|
|
|
expect(notifies[0]).toMatch(/MISSING/);
|
|
|
|
|
await sched.coverageTick();
|
|
|
|
|
expect(notifies).toHaveLength(1); // deduped
|
|
|
|
|
});
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
describe('index discipline is preserved', () => {
|
|
|
|
|
test('every cohort read is bounded on game_date', () => {
|
|
|
|
|
const src = fs.readFileSync(path.join(ROOT, 'src/services/lineageCoverage.js'), 'utf8');
|
|
|
|
|
for (const fn of ['auditLatestCohort', 'auditCohort']) {
|
|
|
|
|
const body = src.slice(src.indexOf(`async function ${fn}`));
|
|
|
|
|
const end = body.indexOf('\n}\n');
|
|
|
|
|
const scoped = body.slice(0, end);
|
|
|
|
|
expect(scoped).toMatch(/\.gte\('game_date'/);
|
|
|
|
|
expect(scoped).toMatch(/\.lte\('game_date'|\.gte\('game_date', from\)/);
|
|
|
|
|
}
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
test('a named cohort audit refuses without an id rather than auditing something else', async () => {
|
|
|
|
|
const out = await coverage.auditCohort('mlb', null, {});
|
|
|
|
|
expect(out.health).toBe(coverage.HEALTH.AUDIT_UNAVAILABLE);
|
|
|
|
|
expect(out.reason).toBe('no_snapshot_id');
|
|
|
|
|
});
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
/* ───────────────── PRODUCT ANCESTRY CONTRACT ───────────────── */
|
|
|
|
|
describe('the ancestry product contract', () => {
|
|
|
|
|
const routeSrc = fs.readFileSync(path.join(ROOT, 'src/routes/ancestry.js'), 'utf8');
|
|
|
|
|
|
|
|
|
|
test('it is authenticated', () => {
|
|
|
|
|
expect(routeSrc).toMatch(/requireAuth/);
|
|
|
|
|
expect(routeSrc).toMatch(/router\.get\('\/ledger\/:id', requireAuth/);
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
test('it is keyed on a stable product identifier, not a raw natural key', () => {
|
|
|
|
|
expect(routeSrc).toMatch(/ledger_entries/);
|
|
|
|
|
expect(routeSrc).not.toMatch(/read_natural_key/);
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
test('it declares LOCAL authority scoped to ancestry', () => {
|
|
|
|
|
expect(routeSrc).toMatch(/authority: 'LOCAL'/);
|
|
|
|
|
expect(routeSrc).toMatch(/authority_scope: 'ANCESTRY_ONLY'/);
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
test('it is reachable from the browser (the S25 proxy rule)', () => {
|
|
|
|
|
const proxy = path.join(ROOT, 'web/src/app/api/ancestry/ledger/[id]/route.ts');
|
|
|
|
|
expect(fs.existsSync(proxy)).toBe(true);
|
|
|
|
|
expect(fs.readFileSync(proxy, 'utf8')).toMatch(/api\/ancestry\/ledger/);
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
test('it is mounted', () => {
|
|
|
|
|
expect(fs.readFileSync(path.join(ROOT, 'src/app.js'), 'utf8'))
|
|
|
|
|
.toMatch(/app\.use\('\/api\/ancestry', require\('\.\/routes\/ancestry'\)\)/);
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
test('the ledger and profile routers still contain no lineage at all', () => {
|
|
|
|
|
for (const f of ['src/routes/ledger.js', 'src/routes/profiles.js']) {
|
|
|
|
|
const s = fs.readFileSync(path.join(ROOT, f), 'utf8');
|
|
|
|
|
expect(s).toMatch(/revised_from_grade/);
|
|
|
|
|
// CASE-INSENSITIVE on purpose: `LINEAGE_ANCESTRY` and `readAncestry`
|
|
|
|
|
// both slipped past a case-sensitive guard when this was injected.
|
|
|
|
|
expect(s).not.toMatch(/lineage|ancestry|read_id/i);
|
|
|
|
|
}
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
test('the route writes nothing', () => {
|
|
|
|
|
const stripped = routeSrc.replace(/\/\*[\s\S]*?\*\//g, '').replace(/^\s*\/\/.*$/gm, '');
|
|
|
|
|
for (const verb of ['.insert(', '.update(', '.upsert(', '.delete(']) {
|
|
|
|
|
expect(stripped).not.toContain(verb);
|
|
|
|
|
}
|
|
|
|
|
});
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
describe('ancestry response states, end to end', () => {
|
|
|
|
|
const A = ancestry.ANCESTRY_STATE;
|
|
|
|
|
const act = (o) => ({
|
|
|
|
|
id: 1, read_id: 'r1', read_natural_key: 'k', lineage_action: 'ORIGIN', claim_digest: 'd0',
|
|
|
|
|
revision_ordinal: 0, lineage_state: 'PUBLISHED', lineage_version: 'lin@1',
|
|
|
|
|
claim_schema_version: 'claim@1', digest_algorithm_version: 'x', publication_id: 'p1', ...o,
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
test('ORIGIN / REVISION / RECAPTURE / legacy / failed / incomplete / invalid', () => {
|
|
|
|
|
expect(ancestry.classifyAncestry([act({})]).state).toBe(A.LINEAGE_AVAILABLE);
|
|
|
|
|
expect(ancestry.classifyAncestry([
|
|
|
|
|
act({ id: 1 }), act({ id: 2, lineage_action: 'REVISION', revision_ordinal: 1, supersedes_id: 1 }),
|
|
|
|
|
]).revision_count).toBe(1);
|
|
|
|
|
expect(ancestry.classifyAncestry([
|
|
|
|
|
act({ id: 1 }), act({ id: 2, lineage_action: 'RECAPTURE', recaptures_id: 1 }),
|
|
|
|
|
]).recapture_count).toBe(1);
|
|
|
|
|
expect(ancestry.classifyAncestry([{ id: 1, publication_id: null }]).state).toBe(A.LEGACY_UNVERIFIED);
|
|
|
|
|
expect(ancestry.classifyAncestry([{ id: 1, publication_id: 'p' }]).state).toBe(A.LINEAGE_UNAVAILABLE);
|
|
|
|
|
expect(ancestry.classifyAncestry([act({ id: 1 }), { id: 2, publication_id: 'p2' }]).chronology_complete).toBe(false);
|
|
|
|
|
expect(ancestry.classifyAncestry([
|
|
|
|
|
act({ id: 1 }), act({ id: 2, lineage_action: 'REVISION', revision_ordinal: 1, supersedes_id: 999 }),
|
|
|
|
|
]).state).toBe(A.INVALID_LINEAGE);
|
|
|
|
|
expect(ancestry.classifyAncestry([]).state).toBe(A.NOT_FOUND);
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
test('other sports are unaffected — nothing here is sport-specific product logic', () => {
|
|
|
|
|
const out = ancestry.classifyAncestry([act({ id: 1 })]);
|
|
|
|
|
expect(out.state).toBe(A.LINEAGE_AVAILABLE);
|
|
|
|
|
expect(out.authority).toBe('NON_AUTHORITATIVE');
|
|
|
|
|
});
|
|
|
|
|
});
|