feat(alerts): add confluence_change picture-transition alert producer (M22 slice 5)
This commit is contained in:
@@ -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<AlertType, { maxPerHour: number }> = {
|
||||
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 },
|
||||
|
||||
@@ -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
|
||||
});
|
||||
});
|
||||
@@ -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<PictureQuality, string> = {
|
||||
'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<Record<string, unknown> & { 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<string, StoredPair>();
|
||||
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<string, unknown>): 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<Alert[]> {
|
||||
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;
|
||||
},
|
||||
};
|
||||
@@ -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 {
|
||||
|
||||
@@ -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.
|
||||
|
||||
Reference in New Issue
Block a user