From 9110d5023e1fa2d3ac826dc73fca8b1eb2462750 Mon Sep 17 00:00:00 2001 From: Investor Flow Build Date: Mon, 10 Aug 2026 16:02:20 -0400 Subject: [PATCH] feat(alerts): add confluence_change picture-transition alert producer (M22 slice 5) --- app/server/src/alerts/AlertEngine.ts | 9 +- .../__tests__/confluenceProducer.test.ts | 134 +++++++++++++ .../alerts/producers/confluenceProducer.ts | 177 ++++++++++++++++++ .../src/db/alertSubscriptionRepository.ts | 1 + app/server/src/index.ts | 2 + 5 files changed, 322 insertions(+), 1 deletion(-) create mode 100644 app/server/src/alerts/__tests__/confluenceProducer.test.ts create mode 100644 app/server/src/alerts/producers/confluenceProducer.ts diff --git a/app/server/src/alerts/AlertEngine.ts b/app/server/src/alerts/AlertEngine.ts index a4ca717..85f728c 100644 --- a/app/server/src/alerts/AlertEngine.ts +++ b/app/server/src/alerts/AlertEngine.ts @@ -24,7 +24,8 @@ export type AlertType = | 'fund_capture' | 'fund_13f' | 'mirror_diff' - | 'vix_level'; + | 'vix_level' + | 'confluence_change'; // ─── Alert Severity ────────────────────────────────────────────────────────── @@ -85,6 +86,7 @@ export function defaultSeverity(type: AlertType): AlertSeverity { case 'fund_13f': case 'mirror_diff': case 'vix_level': + case 'confluence_change': return 'info'; } } @@ -127,6 +129,8 @@ export function alertTitle(type: AlertType, symbol?: string): string { return `Mirror target changed${symbolTag(symbol)}`; case 'vix_level': return `Volatility index level changed`; + case 'confluence_change': + return `Confluence picture changed${symbolTag(symbol)}`; } } @@ -163,6 +167,8 @@ export function alertDescription(type: AlertType, details?: string, symbol?: str return `The mirror target changed${symbolTag(symbol)}. Showing the arithmetic delta; it is not advice.`; case 'vix_level': return `The VIX, a market-wide measure of expected near-term volatility, has moved into a new level that historically mattered to market participants.`; + case 'confluence_change': + return `The confluence picture for a symbol changed. This is an observation of the evidence picture, not advice.`; } })(); @@ -303,6 +309,7 @@ export const ALERT_THROTTLE: Record = { thesis_weakening: { maxPerHour: 3 }, cluster_breach: { maxPerHour: 2 }, drawdown_halt: { maxPerHour: 1 }, + confluence_change: { maxPerHour: 5 }, asymmetry_warning: { maxPerHour: 2 }, fund_capture: { maxPerHour: 3 }, fund_13f: { maxPerHour: 3 }, diff --git a/app/server/src/alerts/__tests__/confluenceProducer.test.ts b/app/server/src/alerts/__tests__/confluenceProducer.test.ts new file mode 100644 index 0000000..2bbd90c --- /dev/null +++ b/app/server/src/alerts/__tests__/confluenceProducer.test.ts @@ -0,0 +1,134 @@ +// 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 + }); +}); \ No newline at end of file diff --git a/app/server/src/alerts/producers/confluenceProducer.ts b/app/server/src/alerts/producers/confluenceProducer.ts new file mode 100644 index 0000000..df04a61 --- /dev/null +++ b/app/server/src/alerts/producers/confluenceProducer.ts @@ -0,0 +1,177 @@ +// Investor Flow — confluence_change alert producer (M22, slice 5) +// +// Fires when a symbol's confluence PICTURE changes meaningfully between two +// stored rack evaluations: quality tier shift (e.g. moderate-bullish → +// mixed) or a material net-evidence move with slot flips. +// +// Primary-rule safe (ADR-0007): the message describes the picture transition +// ("the confluence picture for PLTR deteriorated"), it never says buy or sell. +// Watermark = last evaluation as-of seen per (rack, symbol), so an unchanged +// evaluation never re-fires. + +import type { DatabaseSync } from 'node:sqlite'; +import { randomUUID } from 'node:crypto'; +import type { Alert } from '../AlertEngine.ts'; +import { createAlert, buildDedupKey } from '../AlertEngine.ts'; +import type { AlertProducer } from './types.ts'; +import { getSubscribedUsers, persistAlert, isDuplicate, readComparisonState, writeComparisonState } from './types.ts'; +import { detectPictureChange, type ConfluenceEvaluation, type PictureQuality } from '../../confluence/confluenceRack.ts'; +import { ConfluenceRepository } from '../../db/confluenceRepository.ts'; + +/** Human label used in the alert title (compass calibration, not a directive). */ +const QUALITY_LABEL: Record = { + 'strong-bullish': 'strongly constructive', + 'moderate-bullish': 'constructive', + 'weak-bullish': 'mildly constructive', + 'mixed': 'mixed', + 'weak-bearish': 'mildly cautious', + 'moderate-bearish': 'cautious', + 'strong-bearish': 'strongly cautious', + 'sparse': 'insufficient data', +}; + +interface EvaluationPair { + rackId: string; + symbol: string; + current: ConfluenceEvaluation; + previous: ConfluenceEvaluation | null; +} + +interface StoredPair { + current: ConfluenceEvaluation | null; + previous: ConfluenceEvaluation | null; +} + +/** Load the latest two evaluations per (symbol, rack) that have assessments. */ +function loadLatestPairs(db: DatabaseSync): EvaluationPair[] { + const repo = new ConfluenceRepository(db); + const rows = db.prepare( + `SELECT id, symbol, as_of, rack_id, assessments_json, bull_evidence, bear_evidence, + bull_count, bear_count, assessed_count, net_evidence, total_evidence, quality, created_at, + ROW_NUMBER() OVER (PARTITION BY symbol, rack_id ORDER BY as_of DESC, created_at DESC) AS rn + FROM confluence_evaluations`, + ).all() as Array & { rn: number }>; + + // repo has no bulk loader; reconstruct via its mapper by pulling each row id + // through the evaluation SELECT is wasteful — do a light local rebuild instead. + const byKey = new Map(); + const parseAssessments = (raw: string) => { + try { + const v = JSON.parse(raw) as unknown; + return Array.isArray(v) ? (v as ConfluenceEvaluation['assessments']) : []; + } catch { + return []; + } + }; + const toEval = (r: Record): ConfluenceEvaluation => ({ + symbol: r.symbol as string, + asOf: r.as_of as string, + assessments: parseAssessments(r.assessments_json as string), + bullEvidence: Number(r.bull_evidence), + bearEvidence: Number(r.bear_evidence), + bullCount: Number(r.bull_count), + bearCount: Number(r.bear_count), + assessedCount: Number(r.assessed_count), + netEvidence: Number(r.net_evidence), + totalEvidence: Number(r.total_evidence), + quality: r.quality as PictureQuality, + }); + + for (const row of rows) { + const key = `${row.rack_id}::${row.symbol}`; + const pair = byKey.get(key) ?? { current: null, previous: null } as StoredPair; + if (pair.current === null) { + pair.current = toEval(row); + } else if (pair.previous === null) { + pair.previous = toEval(row); + } + byKey.set(key, pair); + } + + const out: EvaluationPair[] = []; + for (const [key, pair] of byKey) { + if (pair.current === null) continue; + const [rackId, symbol] = key.split('::'); + out.push({ rackId, symbol, current: pair.current, previous: pair.previous }); + } + return out; +} + +export const confluenceChangeProducer: AlertProducer = { + alertType: 'confluence_change', + frequency: 'batched', + + async check(db: DatabaseSync): Promise { + const alerts: Alert[] = []; + + // Sanity: only run when there is data (guards empty installs). + const count = db.prepare('SELECT COUNT(*) AS n FROM confluence_evaluations').get() as { n: number }; + if (Number(count.n) === 0) return alerts; + + const pairs = loadLatestPairs(db); + for (const pair of pairs) { + const { rackId, symbol, current, previous } = pair; + if (current.assessedCount === 0) continue; // empty picture, no basis to alert + + const change = detectPictureChange(previous, current); + if (!change.changed) { + // Still watermark the seen evaluation so it won't refire after restart. + writeComparisonState(db, `${rackId}::${symbol}`, 'confluence_change', { + lastAsOf: current.asOf, + quality: current.quality, + }); + continue; + } + + const users = getSubscribedUsers(db, 'confluence_change', symbol); + for (const userId of users) { + const eventId = `${rackId}:${symbol}:${current.asOf}:${current.quality}`; + const dedupKey = buildDedupKey('confluence_change', symbol, eventId); + if (isDuplicate(db, dedupKey)) continue; + + const fromLabel = change.previous ? QUALITY_LABEL[change.previous] : 'a prior evaluation'; + const toLabel = QUALITY_LABEL[current.quality]; + const direction = change.changeType === 'improved' ? 'improved' : 'became more cautious'; + + const flipped = change.flippedSlotIds.length > 0 + ? `Slots that changed: ${change.flippedSlotIds.join(', ')}. ` + : ''; + const shift = current.quality === 'sparse' + ? '' + : `Net evidence shift: ${change.netEvidenceShift >= 0 ? '+' : ''}${change.netEvidenceShift.toFixed(2)}.`; + + const description = + `The confluence picture for ${symbol} ${direction}: ${fromLabel} → ${toLabel} ` + + `(as of ${current.asOf}). ${flipped}${shift}`; + + const alert = createAlert( + randomUUID(), + userId, + 'confluence_change', + symbol, + description, + eventId, + { + rackId, + symbol, + asOf: current.asOf, + from: change.previous ?? undefined, + to: change.current ?? undefined, + changeType: change.changeType, + netEvidenceShift: change.netEvidenceShift, + flippedSlots: change.flippedSlotIds, + }, + ); + persistAlert(db, alert); + alerts.push(alert); + } + + writeComparisonState(db, `${rackId}::${symbol}`, 'confluence_change', { + lastAsOf: current.asOf, + quality: current.quality, + }); + } + + return alerts; + }, +}; \ No newline at end of file diff --git a/app/server/src/db/alertSubscriptionRepository.ts b/app/server/src/db/alertSubscriptionRepository.ts index 5435766..540f39e 100644 --- a/app/server/src/db/alertSubscriptionRepository.ts +++ b/app/server/src/db/alertSubscriptionRepository.ts @@ -37,6 +37,7 @@ export const ALERT_TYPE_CATALOG: Array<{ type: string; label: string; descriptio { type: 'fund_13f', label: 'Tracked Fund 13F', description: 'A tracked fund filed a new 13F.' }, { type: 'mirror_diff', label: 'Mirror Diff', description: 'The mirror target book changed materially.' }, { type: 'vix_level', label: 'VIX Level', description: 'VIX moves into a new historical volatility band.' }, + { type: 'confluence_change', label: 'Confluence Picture Change', description: 'A symbol\'s confluence picture improved or became more cautious.' }, ]; export interface AlertTypeToggle { diff --git a/app/server/src/index.ts b/app/server/src/index.ts index 0c579bd..1199159 100644 --- a/app/server/src/index.ts +++ b/app/server/src/index.ts @@ -160,6 +160,7 @@ import { rotationIncipientProducer, regimeShiftProducer } from './alerts/produce import { convictionUnlockProducer } from './alerts/producers/unlockProducer.ts'; import { thesisBrokenProducer, thesisWeakeningProducer } from './alerts/producers/thesisProducer.ts'; import { clusterBreachProducer, drawdownHaltProducer, asymmetryWarningProducer } from './alerts/producers/portfolioRiskProducer.ts'; +import { confluenceChangeProducer } from './alerts/producers/confluenceProducer.ts'; registerProducer(informedBuyProducer); registerProducer(informedSellProducer); registerProducer(new13daProducer); @@ -173,6 +174,7 @@ registerProducer(thesisWeakeningProducer); registerProducer(clusterBreachProducer); registerProducer(drawdownHaltProducer); registerProducer(asymmetryWarningProducer); +registerProducer(confluenceChangeProducer); registerMirrorProducers(); // fund_capture / fund_13f / mirror_diff (lazy, offline-safe) // Batched alert check: run every 10 minutes for non-critical producers.