// Investor Flow — confluenceProducer.test.ts // // Guards the confluence_change producer contract: // - First evaluation establishes a baseline (silent on first ever check // only when the picture can't yet be judged — but with two evaluations it // fires on the transition). // - Watermarks prevent re-firing on repeated checks of the same evaluation. // - Only subscribed users are notified. // - Description is a picture observation, never a directive. import { test, describe } from 'node:test'; import { strict as assert } from 'node:assert'; import { DatabaseSync } from 'node:sqlite'; import { readFileSync } from 'node:fs'; import { dirname, join } from 'node:path'; import { fileURLToPath } from 'node:url'; import { confluenceChangeProducer } from '../producers/confluenceProducer.ts'; import { ConfluenceRepository, rackFromSlots } from '../../db/confluenceRepository.ts'; import { evaluateRack, type SlotAssessment } from '../../confluence/confluenceRack.ts'; const __dirname = dirname(fileURLToPath(import.meta.url)); const SCHEMA_SQL = readFileSync(join(__dirname, '..', '..', 'db', 'schema.sql'), 'utf8'); function freshDb(): DatabaseSync { const db = new DatabaseSync(':memory:', { enableForeignKeyConstraints: true }); db.exec(SCHEMA_SQL); return db; } function addUser(db: DatabaseSync, id: string) { db.prepare( `INSERT OR IGNORE INTO users (id, email, pw_hash, complexity, drawdown_tolerance, created_at) VALUES (?, ?, '', 'beginner', 20, ?)`, ).run(id, `${id}@t.local`, new Date().toISOString()); } function subscribeGlobal(db: DatabaseSync, userId: string, alertType: string) { db.prepare( `INSERT INTO alerts (id, owner_id, symbol, alert_type, enabled, params, created_at) VALUES (?, ?, NULL, ?, 1, '{}', ?)`, ).run(`sub-${userId}-${alertType}`, userId, alertType, new Date().toISOString()); } function alertCount(db: DatabaseSync, type: string): number { const r = db.prepare('SELECT COUNT(*) AS n FROM alert_events WHERE type = ?').get(type) as { n: number }; return r.n; } const fired = (id: string): SlotAssessment => ({ id, state: 'fired' }); const notFired = (id: string): SlotAssessment => ({ id, state: 'not-fired' }); function seedRack(db: DatabaseSync): string { const repo = new ConfluenceRepository(db); repo.saveRack(rackFromSlots('rt', 'Rack T', ['goldenCross', 'relVolume', 'cotPositioning'], { isSystem: true })); return 'rt'; } describe('confluence_change producer', () => { test('fires once when the picture deteriorates between two evaluations', async () => { const db = freshDb(); addUser(db, 'u1'); subscribeGlobal(db, 'u1', 'confluence_change'); const rackId = seedRack(db); const repo = new ConfluenceRepository(db); // Baseline: constructive. repo.saveEvaluation(evaluateRack('PLTR', '2026-01-05', [ fired('goldenCross'), fired('relVolume'), notFired('cotPositioning'), ]), rackId, 'ev1'); // Deterioration: death-cross + cautious signals. repo.saveEvaluation(evaluateRack('PLTR', '2026-02-02', [ notFired('goldenCross'), fired('cotPositioning'), { id: 'deathCross', state: 'fired' }, ]), rackId, 'ev2'); const alerts = await confluenceChangeProducer.check(db); assert.equal(alerts.length, 1); assert.equal(alerts[0].symbol, 'PLTR'); assert.equal(alerts[0].type, 'confluence_change'); assert.match(alerts[0].description, /PLTR became more cautious/); assert.doesNotMatch(alerts[0].description, /\b(buy|sell|long|short)\b/i); assert.equal(alertCount(db, 'confluence_change'), 1); // Re-run: watermark prevents re-fire. const again = await confluenceChangeProducer.check(db); assert.equal(again.length, 0); assert.equal(alertCount(db, 'confluence_change'), 1); }); test('fires on improvement too', async () => { const db = freshDb(); addUser(db, 'u1'); subscribeGlobal(db, 'u1', 'confluence_change'); const rackId = seedRack(db); const repo = new ConfluenceRepository(db); repo.saveEvaluation(evaluateRack('PLTR', '2026-01-05', [ { id: 'deathCross', state: 'fired' }, notFired('goldenCross'), ]), rackId, 'ev1'); repo.saveEvaluation(evaluateRack('PLTR', '2026-02-02', [ fired('goldenCross'), fired('relVolume'), notFired('cotPositioning'), ]), rackId, 'ev2'); const alerts = await confluenceChangeProducer.check(db); assert.equal(alerts.length, 1); assert.match(alerts[0].description, /PLTR improved/); }); test('does not alert unsubscribed or empty-data installs', async () => { const db = freshDb(); addUser(db, 'u1'); const rackId = seedRack(db); const repo = new ConfluenceRepository(db); repo.saveEvaluation(evaluateRack('PLTR', '2026-01-05', [fired('goldenCross')]), rackId, 'ev1'); repo.saveEvaluation(evaluateRack('PLTR', '2026-02-02', [notFired('goldenCross')]), rackId, 'ev2'); const alerts = await confluenceChangeProducer.check(db); // no subscription assert.equal(alerts.length, 0); }); test('silent when the picture is unchanged', async () => { const db = freshDb(); addUser(db, 'u1'); subscribeGlobal(db, 'u1', 'confluence_change'); const rackId = seedRack(db); const repo = new ConfluenceRepository(db); repo.saveEvaluation(evaluateRack('PLTR', '2026-01-05', [fired('goldenCross'), fired('relVolume')]), rackId, 'ev1'); repo.saveEvaluation(evaluateRack('PLTR', '2026-02-02', [fired('goldenCross'), fired('relVolume')]), rackId, 'ev2'); const alerts = await confluenceChangeProducer.check(db); assert.equal(alerts.length, 0); // mixed→mixed at same tier, no drift }); });