fix: backfill symbol_demand for sidebar-added symbols + analyst ratings schema fix

- Add await ctx.cache.subscribe() to addSymbol mutation so symbols
  added via the sidebar get registered in symbol_demand and yfinance
  jobs are queued immediately
- Backfill PEP, WYNN, STZ, CELH into symbol_demand + adapter_queue
- Upgrade yahoo-finance2 3.15.3 -> 3.15.4 and pass validateResult:false
  to quoteSummary() to handle Yahoo schema drift
- Add error detail logging for analyst ratings schema failures
- Update .gitignore with common ignores
This commit is contained in:
Investor Flow Build
2026-07-23 18:02:24 -04:00
parent 5ef2b2f060
commit e262187c3c
204 changed files with 25014 additions and 2934 deletions
+283
View File
@@ -0,0 +1,283 @@
// Investor Flow — AlertEngine (Slice 17): hybrid event-driven + polling alert system.
//
// ADR-0007: Alert text says "something changed" not "action needed". Never
// imperative (no "buy" / "sell" / "cut" / "trim"). Every alert is a
// notification of a change in state, not a recommendation to act.
//
// Pure/cache-deterministic core: no I/O in the pure functions below.
// Database reads happen in the tRPC layer, not here.
// ─── Alert Types ─────────────────────────────────────────────────────────────
export type AlertType =
| 'informed_buy'
| 'informed_sell'
| 'new_13da'
| 'rotation_incipient'
| 'regime_shift'
| 'conviction_unlock'
| 'thesis_broken'
| 'thesis_weakening'
| 'cluster_breach'
| 'drawdown_halt'
| 'asymmetry_warning';
// ─── Alert Severity ──────────────────────────────────────────────────────────
export type AlertSeverity = 'info' | 'warning' | 'critical';
// ─── Alert Payload ───────────────────────────────────────────────────────────
export interface Alert {
id: string;
userId: string;
type: AlertType;
severity: AlertSeverity;
title: string;
description: string;
symbol?: string;
createdAt: string;
acknowledged: boolean;
dedupKey: string;
payload: Record<string, unknown>;
}
// ─── Alert Dedup Store ───────────────────────────────────────────────────────
export class DedupStore {
private seen = new Set<string>();
isDuplicate(key: string): boolean {
return this.seen.has(key);
}
mark(key: string): void {
this.seen.add(key);
}
reset(): void {
this.seen.clear();
}
}
// ─── Alert Severity Mapping ──────────────────────────────────────────────────
export function defaultSeverity(type: AlertType): AlertSeverity {
switch (type) {
case 'drawdown_halt':
return 'critical';
case 'thesis_broken':
case 'cluster_breach':
case 'asymmetry_warning':
return 'warning';
case 'informed_buy':
case 'informed_sell':
case 'new_13da':
case 'rotation_incipient':
case 'regime_shift':
case 'conviction_unlock':
case 'thesis_weakening':
return 'info';
}
}
// ─── Alert Title / Description Builders ──────────────────────────────────────
function symbolTag(symbol?: string): string {
return symbol ? ` for ${symbol}` : '';
}
export function alertTitle(type: AlertType, symbol?: string): string {
switch (type) {
case 'informed_buy':
return `Insider bought${symbolTag(symbol)}`;
case 'informed_sell':
return `Insider sold${symbolTag(symbol)}`;
case 'new_13da':
return `New institutional position${symbolTag(symbol)}`;
case 'rotation_incipient':
return `Sector rotation signal detected`;
case 'regime_shift':
return `Market regime changed`;
case 'conviction_unlock':
return `Conviction tier unlocked`;
case 'thesis_broken':
return `Thesis invalidation criteria met${symbolTag(symbol)}`;
case 'thesis_weakening':
return `Thesis showing signs of weakening${symbolTag(symbol)}`;
case 'cluster_breach':
return `Cluster exposure limit reached${symbolTag(symbol)}`;
case 'drawdown_halt':
return `Drawdown tolerance breached`;
case 'asymmetry_warning':
return `Portfolio asymmetry below threshold`;
}
}
export function alertDescription(type: AlertType, details?: string, symbol?: string): string {
const base = (() => {
switch (type) {
case 'informed_buy':
return `A company insider purchased shares${symbolTag(symbol)}. This filing was not part of a 10b5-1 trading plan.`;
case 'informed_sell':
return `A company insider sold shares${symbolTag(symbol)}. This filing was not part of a 10b5-1 trading plan.`;
case 'new_13da':
return `An institutional investor reported a new position${symbolTag(symbol)} in a 13F filing.`;
case 'rotation_incipient':
return `The sector rotation detector identified an incipient rotation signal. Capital may be moving between sectors.`;
case 'regime_shift':
return `The market regime has changed. This affects portfolio-level risk assessments.`;
case 'conviction_unlock':
return `A new conviction tier is now available based on your trading history.`;
case 'thesis_broken':
return `The invalidation criteria for your thesis${symbolTag(symbol)} have been met. Consider reviewing your thesis.`;
case 'thesis_weakening':
return `Some signals suggest your thesis${symbolTag(symbol)} may be weakening, but invalidation criteria are not yet met.`;
case 'cluster_breach':
return `Your exposure in this cluster has exceeded the recommended cap${symbolTag(symbol)}.`;
case 'drawdown_halt':
return `Your portfolio drawdown has exceeded the tolerance threshold. The circuit breaker has paused new entries for 24 hours. Existing positions continue unaffected.`;
case 'asymmetry_warning':
return `Your portfolio's reward-to-risk ratio has fallen below 1.0, meaning risk outweighs expected reward across your positions.`;
}
})();
if (details) {
return `${base}\n\n${details}`;
}
return base;
}
// ─── Dedup Key Builder ───────────────────────────────────────────────────────
export function buildDedupKey(
type: AlertType,
symbol?: string,
eventId?: string,
): string {
const parts: string[] = [type];
if (symbol) parts.push(symbol);
if (eventId) parts.push(eventId);
return parts.join(':');
}
// ─── Alert Factory ───────────────────────────────────────────────────────────
export function createAlert(
id: string,
userId: string,
type: AlertType,
symbol?: string,
details?: string,
eventId?: string,
extraPayload?: Record<string, unknown>,
): Alert {
return {
id,
userId,
type,
severity: defaultSeverity(type),
title: alertTitle(type, symbol),
description: alertDescription(type, details, symbol),
symbol,
createdAt: new Date().toISOString(),
acknowledged: false,
dedupKey: buildDedupKey(type, symbol, eventId),
payload: {
...extraPayload,
...(eventId ? { eventId } : {}),
},
};
}
// ─── AlertEngine ─────────────────────────────────────────────────────────────
export interface AlertEngineDeps {
dedup: DedupStore;
poll: () => Promise<Alert[]>;
persist: (alert: Alert) => Promise<Alert>;
listAlerts: (userId: string, limit?: number) => Promise<Alert[]>;
acknowledge: (alertId: string, userId: string) => Promise<boolean>;
}
export class AlertEngine {
private deps: AlertEngineDeps;
private pollingIntervalMs: number;
private pollTimer: ReturnType<typeof setInterval> | null = null;
constructor(deps: AlertEngineDeps, pollingIntervalMs = 5 * 60 * 1000) {
this.deps = deps;
this.pollingIntervalMs = pollingIntervalMs;
}
async fireAndForget(
type: AlertType,
userId: string,
symbol?: string,
details?: string,
eventId?: string,
extraPayload?: Record<string, unknown>,
): Promise<Alert | null> {
const dedupKey = buildDedupKey(type, symbol, eventId);
if (this.deps.dedup.isDuplicate(dedupKey)) return null;
const id = crypto.randomUUID();
const alert = createAlert(id, userId, type, symbol, details, eventId, extraPayload);
this.deps.dedup.mark(dedupKey);
await this.deps.persist(alert);
return alert;
}
async pollCycle(): Promise<Alert[]> {
const candidates = await this.deps.poll();
const created: Alert[] = [];
for (const candidate of candidates) {
if (!this.deps.dedup.isDuplicate(candidate.dedupKey)) {
this.deps.dedup.mark(candidate.dedupKey);
await this.deps.persist(candidate);
created.push(candidate);
}
}
return created;
}
start(): void {
if (this.pollTimer) return;
this.pollTimer = setInterval(() => {
this.pollCycle().catch((err) => {
console.error('[AlertEngine] poll cycle failed:', err);
});
}, this.pollingIntervalMs);
}
stop(): void {
if (this.pollTimer) {
clearInterval(this.pollTimer);
this.pollTimer = null;
}
}
async listAlerts(userId: string, limit?: number): Promise<Alert[]> {
return this.deps.listAlerts(userId, limit);
}
async acknowledge(alertId: string, userId: string): Promise<boolean> {
return this.deps.acknowledge(alertId, userId);
}
}
// ─── Throttle Configuration ──────────────────────────────────────────────────
export const ALERT_THROTTLE: Record<AlertType, { maxPerHour: number }> = {
informed_buy: { maxPerHour: 5 },
informed_sell: { maxPerHour: 5 },
new_13da: { maxPerHour: 3 },
rotation_incipient: { maxPerHour: 2 },
regime_shift: { maxPerHour: 1 },
conviction_unlock: { maxPerHour: 1 },
thesis_broken: { maxPerHour: 3 },
thesis_weakening: { maxPerHour: 3 },
cluster_breach: { maxPerHour: 2 },
drawdown_halt: { maxPerHour: 1 },
asymmetry_warning: { maxPerHour: 2 },
};
@@ -0,0 +1,245 @@
// Tests — AlertEngine (Slice 17). Pure, no network.
import { test } from 'node:test';
import { strict as assert } from 'node:assert';
import {
AlertEngine,
DedupStore,
createAlert,
alertTitle,
alertDescription,
defaultSeverity,
buildDedupKey,
ALERT_THROTTLE,
type Alert,
type AlertType,
} from '../AlertEngine.ts';
// ─── Pure utility tests ──────────────────────────────────────────────────────
test('defaultSeverity returns expected severity for each alert type', () => {
assert.equal(defaultSeverity('drawdown_halt'), 'critical');
assert.equal(defaultSeverity('thesis_broken'), 'warning');
assert.equal(defaultSeverity('cluster_breach'), 'warning');
assert.equal(defaultSeverity('asymmetry_warning'), 'warning');
assert.equal(defaultSeverity('informed_buy'), 'info');
assert.equal(defaultSeverity('informed_sell'), 'info');
assert.equal(defaultSeverity('new_13da'), 'info');
assert.equal(defaultSeverity('rotation_incipient'), 'info');
assert.equal(defaultSeverity('regime_shift'), 'info');
assert.equal(defaultSeverity('conviction_unlock'), 'info');
assert.equal(defaultSeverity('thesis_weakening'), 'info');
});
test('alertTitle returns expected title for each alert type', () => {
assert.equal(alertTitle('informed_buy'), 'Insider bought');
assert.equal(alertTitle('informed_sell', 'NVDA'), 'Insider sold for NVDA');
assert.equal(alertTitle('regime_shift'), 'Market regime changed');
assert.equal(alertTitle('drawdown_halt'), 'Drawdown tolerance breached');
assert.equal(alertTitle('informed_buy', 'AAPL'), 'Insider bought for AAPL');
});
test('alertDescription never contains trade verbs', () => {
const types: AlertType[] = [
'informed_buy', 'informed_sell', 'new_13da', 'rotation_incipient',
'regime_shift', 'conviction_unlock', 'thesis_broken', 'thesis_weakening',
'cluster_breach', 'drawdown_halt', 'asymmetry_warning',
];
const forbidden = ['buy ', 'sell ', 'cut ', 'trim ', 'buy.', 'sell.', 'cut.', 'trim.'];
for (const type of types) {
const desc = alertDescription(type);
const lower = desc.toLowerCase();
for (const word of forbidden) {
assert.ok(!lower.includes(word),
`"${type}" description should not contain "${word}": "${desc.substring(0, 60)}..."`);
}
}
});
test('alertDescription includes ADR-0007 compliant language', () => {
const desc = alertDescription('drawdown_halt');
assert.ok(desc.includes('circuit breaker'), 'Should mention circuit breaker');
assert.ok(desc.includes('Existing positions continue unaffected'),
'Should mention existing positions unaffected');
assert.ok(!desc.includes('cut'), 'Should not contain "cut"');
});
test('buildDedupKey produces consistent keys', () => {
assert.equal(buildDedupKey('informed_buy'), 'informed_buy');
assert.equal(buildDedupKey('informed_buy', 'NVDA'), 'informed_buy:NVDA');
assert.equal(buildDedupKey('informed_buy', 'NVDA', 'filing-123'),
'informed_buy:NVDA:filing-123');
assert.equal(buildDedupKey('drawdown_halt'), 'drawdown_halt');
});
// ─── Alert factory tests ─────────────────────────────────────────────────────
test('createAlert produces a complete Alert object', () => {
const alert = createAlert('id-1', 'user-1', 'informed_buy', 'NVDA');
assert.equal(alert.id, 'id-1');
assert.equal(alert.userId, 'user-1');
assert.equal(alert.type, 'informed_buy');
assert.equal(alert.symbol, 'NVDA');
assert.equal(alert.severity, 'info');
assert.ok(alert.title.startsWith('Insider bought'));
assert.ok(alert.description.length > 20);
assert.equal(alert.acknowledged, false);
assert.ok(alert.dedupKey.startsWith('informed_buy:NVDA'));
assert.ok(alert.createdAt.length > 0);
assert.ok(alert.payload);
});
test('createAlert includes eventId in dedupKey and payload', () => {
const alert = createAlert('id-2', 'user-1', 'informed_sell', 'AAPL',
'details', 'evt-456');
assert.ok(alert.dedupKey.endsWith(':evt-456'));
assert.equal(alert.payload.eventId, 'evt-456');
});
test('createAlert with details appends to description', () => {
const alert = createAlert('id-3', 'user-1', 'new_13da', 'TSLA',
'Additional context about the filing.');
assert.ok(alert.description.includes('Additional context about the filing.'));
});
// ─── DedupStore tests ────────────────────────────────────────────────────────
test('DedupStore tracks seen keys', () => {
const store = new DedupStore();
assert.ok(!store.isDuplicate('test-key'));
store.mark('test-key');
assert.ok(store.isDuplicate('test-key'));
});
test('DedupStore can be reset', () => {
const store = new DedupStore();
store.mark('key-1');
assert.ok(store.isDuplicate('key-1'));
store.reset();
assert.ok(!store.isDuplicate('key-1'));
});
// ─── AlertEngine tests ───────────────────────────────────────────────────────
test('AlertEngine.fireAndForget creates alert on first call, dedups on second',
async () => {
const persisted: Alert[] = [];
const engine = new AlertEngine({
dedup: new DedupStore(),
poll: async () => [],
persist: async (a) => { persisted.push(a); return a; },
listAlerts: async () => [],
acknowledge: async () => true,
});
const first = await engine.fireAndForget('informed_buy', 'user-1', 'NVDA',
undefined, 'filing-1');
assert.ok(first !== null, 'First call should create alert');
assert.equal(persisted.length, 1);
const second = await engine.fireAndForget('informed_buy', 'user-1', 'NVDA',
undefined, 'filing-1');
assert.equal(second, null, 'Duplicate should return null');
assert.equal(persisted.length, 1, 'Should not persist duplicate');
});
test('AlertEngine.fireAndForget different eventIds are not duplicates',
async () => {
const persisted: Alert[] = [];
const engine = new AlertEngine({
dedup: new DedupStore(),
poll: async () => [],
persist: async (a) => { persisted.push(a); return a; },
listAlerts: async () => [],
acknowledge: async () => true,
});
const first = await engine.fireAndForget('informed_buy', 'user-1', 'NVDA',
undefined, 'filing-1');
assert.ok(first !== null);
const second = await engine.fireAndForget('informed_buy', 'user-1', 'NVDA',
undefined, 'filing-2');
assert.ok(second !== null, 'Different eventId should not be duplicate');
assert.equal(persisted.length, 2);
});
test('AlertEngine.pollCycle dedupes within same cycle', async () => {
// If the poll function returns two alerts with the same dedup key,
// only the first should be persisted.
const persisted: Alert[] = [];
const engine = new AlertEngine({
dedup: new DedupStore(),
poll: async () => {
const a = createAlert('dup-a', 'user-1', 'asymmetry_warning', 'NVDA');
const b = createAlert('dup-b', 'user-1', 'asymmetry_warning', 'NVDA');
return [a, b];
},
persist: async (a) => { persisted.push(a); return a; },
listAlerts: async () => [],
acknowledge: async () => true,
});
const created = await engine.pollCycle();
assert.equal(created.length, 1, 'Only one of the duplicates should be created');
assert.equal(persisted.length, 1, 'Only one should be persisted');
});
test('AlertEngine.pollCycle skips already-known dedup keys', async () => {
// If a dedup key was already marked (e.g. from a prior event-driven fire),
// the poll cycle should skip it.
const persisted: Alert[] = [];
const dedupStore = new DedupStore();
dedupStore.mark('asymmetry_warning:NVDA'); // already seen
const engine = new AlertEngine({
dedup: dedupStore,
poll: async () => [
createAlert('poll-1', 'user-1', 'asymmetry_warning', 'NVDA'),
],
persist: async (a) => { persisted.push(a); return a; },
listAlerts: async () => [],
acknowledge: async () => true,
});
const created = await engine.pollCycle();
assert.equal(created.length, 0, 'Known dedup keys should be skipped');
assert.equal(persisted.length, 0, 'Nothing should be persisted');
});
test('AlertEngine start/stop polling', () => {
const engine = new AlertEngine({
dedup: new DedupStore(),
poll: async () => [],
persist: async (a) => a,
listAlerts: async () => [],
acknowledge: async () => true,
}, 1000);
engine.start();
// Should be idempotent
engine.start();
engine.stop();
// Should be safe to stop twice
engine.stop();
});
// ─── Throttle tests ──────────────────────────────────────────────────────────
test('ALERT_THROTTLE defines limits for all alert types', () => {
const types: AlertType[] = [
'informed_buy', 'informed_sell', 'new_13da', 'rotation_incipient',
'regime_shift', 'conviction_unlock', 'thesis_broken', 'thesis_weakening',
'cluster_breach', 'drawdown_halt', 'asymmetry_warning',
];
for (const type of types) {
assert.ok(ALERT_THROTTLE[type], `Throttle config exists for ${type}`);
assert.ok(ALERT_THROTTLE[type].maxPerHour > 0,
`maxPerHour > 0 for ${type}`);
}
});
test('Critical alerts have stricter throttle', () => {
assert.equal(ALERT_THROTTLE['drawdown_halt'].maxPerHour, 1);
assert.equal(ALERT_THROTTLE['regime_shift'].maxPerHour, 1);
assert.equal(ALERT_THROTTLE['conviction_unlock'].maxPerHour, 1);
});