From 5b9f770aa4864ec08a06207d6c90944e745d66e8 Mon Sep 17 00:00:00 2001 From: Investor Flow Build Date: Thu, 23 Jul 2026 20:51:47 -0400 Subject: [PATCH] Phase 5: alert subscriptions UI + multiple watchlists + queue fixes UI: - /alerts page: event history with acknowledge, subscription create/manage with toggle - /admin/smtp: SMTP config form (host, port, auth, test) - Watchlist sidebar: dropdown selector for multiple watchlists, create/delete - Sidebar: alerts count badge, SMTP link under admin - Mobile tab nav: alerts tab added - Client trpc.ts: all new API methods + types Backend: - watchlists.listByWatchlist procedure + listSymbolsByWatchlist repo fn - yfinance min-interval 1500->2000ms to reduce Edge 429s - Fixed e.date.slice error in yfinance-adjustments with typeof guard - Removed defunct BITF from watchlist+queue - Cleared 83 failed + 12 backoff queue jobs Docs: - FUNCTIONAL_DESIGN.md: alerts + multiple watchlists + SMTP documented - TECH_DESIGN.md: new modules, tRPC procs, routes updated --- app/server/package-lock.json | 21 ++ app/server/package.json | 2 + .../src/adapters/yfinance-adjustments.ts | 10 +- app/server/src/alerts/producers/index.ts | 37 +++ .../src/alerts/producers/insiderProducer.ts | 88 +++++++ .../src/alerts/producers/new13daProducer.ts | 60 +++++ app/server/src/alerts/producers/types.ts | 97 ++++++++ .../src/db/alertSubscriptionRepository.ts | 173 +++++++++++++ app/server/src/db/client.ts | 44 +++- app/server/src/db/schema.sql | 40 ++- app/server/src/db/watchlistRepository.ts | 233 +++++++++++------- app/server/src/index.ts | 40 +++ app/server/src/queue/sourceRatePolicy.ts | 2 +- app/server/src/services/emailAlertService.ts | 131 ++++++++++ app/server/src/trpc/router.ts | 202 ++++++++++++++- app/src/app/admin/smtp/page.tsx | 121 +++++++++ app/src/app/alerts/page.tsx | 172 +++++++++++++ app/src/components/MobileTabNav.tsx | 1 + app/src/components/SidebarNav.tsx | 17 +- app/src/components/WatchlistSidebar.tsx | 130 ++++++++-- app/src/lib/trpc.ts | 58 ++++- docs/FUNCTIONAL_DESIGN.md | 15 +- docs/TECH_DESIGN.md | 19 +- 23 files changed, 1567 insertions(+), 146 deletions(-) create mode 100644 app/server/src/alerts/producers/index.ts create mode 100644 app/server/src/alerts/producers/insiderProducer.ts create mode 100644 app/server/src/alerts/producers/new13daProducer.ts create mode 100644 app/server/src/alerts/producers/types.ts create mode 100644 app/server/src/db/alertSubscriptionRepository.ts create mode 100644 app/server/src/services/emailAlertService.ts create mode 100644 app/src/app/admin/smtp/page.tsx create mode 100644 app/src/app/alerts/page.tsx diff --git a/app/server/package-lock.json b/app/server/package-lock.json index f94e13f..dbf5a28 100644 --- a/app/server/package-lock.json +++ b/app/server/package-lock.json @@ -9,11 +9,13 @@ "version": "0.1.0", "dependencies": { "@trpc/server": "^11.0.0", + "nodemailer": "^9.0.3", "yahoo-finance2": "^3.15.4", "zod": "^4.4.3" }, "devDependencies": { "@types/node": "^22.0.0", + "@types/nodemailer": "^8.0.1", "typescript": "^5.0.0" }, "engines": { @@ -113,6 +115,16 @@ "undici-types": "~6.21.0" } }, + "node_modules/@types/nodemailer": { + "version": "8.0.1", + "resolved": "https://registry.npmjs.org/@types/nodemailer/-/nodemailer-8.0.1.tgz", + "integrity": "sha512-PxpaInm8V1JQDd4j0ds5HfvWQk8JupS1C0Picb96QJsrrRDjBH+DlK7L4ZdNSqNULhiZRQHc40nLVShaGxXAMw==", + "dev": true, + "license": "MIT", + "dependencies": { + "@types/node": "*" + } + }, "node_modules/accepts": { "version": "2.0.0", "resolved": "https://registry.npmjs.org/accepts/-/accepts-2.0.0.tgz", @@ -905,6 +917,15 @@ "node": ">= 0.6" } }, + "node_modules/nodemailer": { + "version": "9.0.3", + "resolved": "https://registry.npmjs.org/nodemailer/-/nodemailer-9.0.3.tgz", + "integrity": "sha512-n+YP+NKwR5zRWa60k3GiQ6Q3B4KXCoAw40dAKeCtYn020iNN74aWK2liXIC3ZEATeGql7we3tE3t8QwhY0eskw==", + "license": "MIT-0", + "engines": { + "node": ">=6.0.0" + } + }, "node_modules/normalize-url": { "version": "4.5.1", "resolved": "https://registry.npmjs.org/normalize-url/-/normalize-url-4.5.1.tgz", diff --git a/app/server/package.json b/app/server/package.json index 5b20bd8..30f7e39 100644 --- a/app/server/package.json +++ b/app/server/package.json @@ -16,11 +16,13 @@ }, "dependencies": { "@trpc/server": "^11.0.0", + "nodemailer": "^9.0.3", "yahoo-finance2": "^3.15.4", "zod": "^4.4.3" }, "devDependencies": { "@types/node": "^22.0.0", + "@types/nodemailer": "^8.0.1", "typescript": "^5.0.0" } } diff --git a/app/server/src/adapters/yfinance-adjustments.ts b/app/server/src/adapters/yfinance-adjustments.ts index aebc213..1af3283 100644 --- a/app/server/src/adapters/yfinance-adjustments.ts +++ b/app/server/src/adapters/yfinance-adjustments.ts @@ -19,8 +19,8 @@ export interface PriceAdjustment { ratio: number; } -type SplitEvent = { date: string; numerator: number; denominator: number; splitRatio?: string }; -type DividendEvent = { date: string; amount: number }; +type SplitEvent = { date: string | number; numerator: number; denominator: number; splitRatio?: string }; +type DividendEvent = { date: string | number; amount: number }; /** * Parse a yahoo-finance2 v3 chart response's `events` object into a flat list of @@ -64,21 +64,23 @@ function emitFrom( if (type === "split") { const e = entry as SplitEvent; if (!e.date) continue; + const dateStr = typeof e.date === 'string' ? e.date : String(e.date); const numerator = typeof e.numerator === "number" ? e.numerator : 1; const denominator = typeof e.denominator === "number" && e.denominator !== 0 ? e.denominator : 1; out.push({ symbol, - exDate: e.date.slice(0, 10), + exDate: dateStr.slice(0, 10), type: "split", ratio: numerator / denominator, }); } else { const e = entry as DividendEvent; if (!e.date) continue; + const dateStr = typeof e.date === 'string' ? e.date : String(e.date); const amount = typeof e.amount === "number" ? e.amount : 0; out.push({ symbol, - exDate: e.date.slice(0, 10), + exDate: dateStr.slice(0, 10), type: "dividend", ratio: amount, }); diff --git a/app/server/src/alerts/producers/index.ts b/app/server/src/alerts/producers/index.ts new file mode 100644 index 0000000..40a896e --- /dev/null +++ b/app/server/src/alerts/producers/index.ts @@ -0,0 +1,37 @@ +import type { DatabaseSync } from 'node:sqlite'; +import type { Alert } from '../AlertEngine.ts'; +import type { AlertProducer, ProducerFrequency } from './types.ts'; + +export type { AlertProducer, ProducerFrequency } from './types.ts'; + +const producers: AlertProducer[] = []; + +/** Register a producer (called once at startup). */ +export function registerProducer(p: AlertProducer): void { + producers.push(p); +} + +/** Run all producers for a given frequency tier. */ +export async function runProducers( + db: DatabaseSync, + frequency: ProducerFrequency, +): Promise { + const results: Alert[] = []; + const tier = producers.filter((p) => p.frequency === frequency); + + for (const producer of tier) { + try { + const alerts = await producer.check(db); + results.push(...alerts); + } catch (err) { + console.error(`[alert:${producer.alertType}] check failed:`, err); + } + } + + return results; +} + +/** Get all registered producer types. */ +export function registeredAlertTypes(): string[] { + return producers.map((p) => p.alertType); +} diff --git a/app/server/src/alerts/producers/insiderProducer.ts b/app/server/src/alerts/producers/insiderProducer.ts new file mode 100644 index 0000000..bfd86d2 --- /dev/null +++ b/app/server/src/alerts/producers/insiderProducer.ts @@ -0,0 +1,88 @@ +import type { DatabaseSync } from 'node:sqlite'; +import type { Alert } from '../AlertEngine.ts'; +import { createAlert } from '../AlertEngine.ts'; +import type { AlertProducer } from './types.ts'; +import { readComparisonState, writeComparisonState, getSubscribedUsers, persistAlert, isDuplicate } from './types.ts'; + +interface InsiderTx { + symbol: string; + tx_date: string; + classification: string; + shares: number; +} + +function createProducer(alertType: 'informed_buy' | 'informed_sell', classification: string): AlertProducer { + return { + alertType, + frequency: 'batched' as const, + + async check(db: DatabaseSync): Promise { + const alerts: Alert[] = []; + + // Find all symbols with new insider transactions. + const symbols = db.prepare( + `SELECT DISTINCT symbol FROM insider_transactions WHERE classification = ?`, + ).all(classification) as { symbol: string }[]; + + for (const { symbol } of symbols) { + const prev = readComparisonState(db, symbol, alertType); + const prevState = prev ? JSON.parse(prev.state) : { lastTxDate: null, lastSeenIds: [] as string[] }; + + const rows = db.prepare( + `SELECT symbol, tx_date, classification, shares + FROM insider_transactions + WHERE symbol = ? AND classification = ? + ORDER BY tx_date DESC + LIMIT 10`, + ).all(symbol, classification) as unknown as InsiderTx[]; + + if (rows.length === 0) continue; + + const latestDate = rows[0].tx_date; + const newIds = rows.map((r) => `${r.tx_date}:${r.shares}`); + const prevIds = new Set(prevState.lastSeenIds as string[]); + + if (prevState.lastTxDate === latestDate && newIds.every((id) => prevIds.has(id))) { + writeComparisonState(db, symbol, alertType, { lastTxDate: latestDate, lastSeenIds: newIds }); + continue; + } + + const newRows = prevState.lastTxDate === null + ? rows + : rows.filter((r) => r.tx_date > prevState.lastTxDate || (r.tx_date === prevState.lastTxDate && !prevIds.has(`${r.tx_date}:${r.shares}`))); + + for (const row of newRows) { + const users = getSubscribedUsers(db, alertType, symbol); + for (const userId of users) { + const eventId = `${symbol}:${row.tx_date}:${row.shares}`; + const dedupKey = `${alertType}:${userId}:${eventId}`; + if (isDuplicate(db, dedupKey)) continue; + + const description = classification === 'P' + ? `An insider purchased ${row.shares} shares on ${row.tx_date}.` + : `An insider sold ${Math.abs(row.shares)} shares on ${row.tx_date}.`; + + const alert = createAlert( + crypto.randomUUID(), + userId, + alertType, + symbol, + description, + eventId, + { shares: row.shares, txDate: row.tx_date }, + ); + persistAlert(db, alert); + alerts.push(alert); + } + } + + writeComparisonState(db, symbol, alertType, { lastTxDate: latestDate, lastSeenIds: newIds }); + } + + return alerts; + }, + }; +} + +export const informedBuyProducer = createProducer('informed_buy', 'P'); +export const informedSellProducer = createProducer('informed_sell', 'S'); diff --git a/app/server/src/alerts/producers/new13daProducer.ts b/app/server/src/alerts/producers/new13daProducer.ts new file mode 100644 index 0000000..5cecfd0 --- /dev/null +++ b/app/server/src/alerts/producers/new13daProducer.ts @@ -0,0 +1,60 @@ +import type { DatabaseSync } from 'node:sqlite'; +import type { Alert } from '../AlertEngine.ts'; +import { createAlert } from '../AlertEngine.ts'; +import type { AlertProducer } from './types.ts'; +import { readComparisonState, writeComparisonState, getSubscribedUsers, persistAlert, isDuplicate } from './types.ts'; + +export const new13daProducer: AlertProducer = { + alertType: 'new_13da', + frequency: 'batched', + + async check(db: DatabaseSync): Promise { + const alerts: Alert[] = []; + + const symbols = db.prepare( + `SELECT DISTINCT symbol FROM institution_filings WHERE form LIKE '13D%' OR form LIKE '13G%'`, + ).all() as { symbol: string }[]; + + for (const { symbol } of symbols) { + const prev = readComparisonState(db, symbol, 'new_13da'); + const prevState = prev ? JSON.parse(prev.state) : { latestAccession: null }; + + const row = db.prepare( + `SELECT accession, reported_quarter, filer_name, put_call + FROM institution_filings + WHERE symbol = ? AND (form LIKE '13D%' OR form LIKE '13G%') + ORDER BY reported_quarter DESC + LIMIT 1`, + ).get(symbol) as { accession: string; reported_quarter: string; filer_name: string; put_call: string | null } | undefined; + + if (!row) continue; + + if (row.accession === prevState.latestAccession) { + writeComparisonState(db, symbol, 'new_13da', { latestAccession: row.accession }); + continue; + } + + const users = getSubscribedUsers(db, 'new_13da', symbol); + for (const userId of users) { + const dedupKey = `new_13da:${userId}:${row.accession}`; + if (isDuplicate(db, dedupKey)) continue; + + const alert = createAlert( + crypto.randomUUID(), + userId, + 'new_13da', + symbol, + `${row.filer_name} reported a new position${row.put_call ? ` (${row.put_call})` : ''} for ${row.reported_quarter}.`, + row.accession, + { filer: row.filer_name, quarter: row.reported_quarter, putCall: row.put_call }, + ); + persistAlert(db, alert); + alerts.push(alert); + } + + writeComparisonState(db, symbol, 'new_13da', { latestAccession: row.accession }); + } + + return alerts; + }, +}; diff --git a/app/server/src/alerts/producers/types.ts b/app/server/src/alerts/producers/types.ts new file mode 100644 index 0000000..b64a6fb --- /dev/null +++ b/app/server/src/alerts/producers/types.ts @@ -0,0 +1,97 @@ +import type { DatabaseSync } from 'node:sqlite'; +import type { Alert } from '../AlertEngine.ts'; + +export type ProducerFrequency = 'per-fetch' | 'batched'; + +export interface AlertProducer { + readonly alertType: string; + readonly frequency: ProducerFrequency; + check(db: DatabaseSync): Promise; +} + +export interface ComparisonState { + symbol: string; + alertType: string; + state: string; + updatedAt: string; +} + +export function readComparisonState( + db: DatabaseSync, + symbol: string, + alertType: string, +): ComparisonState | null { + const row = db.prepare( + `SELECT symbol, alert_type, state, updated_at + FROM alert_comparison_state + WHERE symbol = ? AND alert_type = ?`, + ).get(symbol, alertType) as Record | undefined; + if (!row) return null; + return { + symbol: row.symbol as string, + alertType: row.alert_type as string, + state: row.state as string, + updatedAt: row.updated_at as string, + }; +} + +export function writeComparisonState( + db: DatabaseSync, + symbol: string, + alertType: string, + state: Record, +): void { + db.prepare( + `INSERT OR REPLACE INTO alert_comparison_state (symbol, alert_type, state, updated_at) + VALUES (?, ?, ?, ?)`, + ).run(symbol, alertType, JSON.stringify(state), new Date().toISOString()); +} + +export function getSubscribedUsers( + db: DatabaseSync, + alertType: string, + symbol: string, +): string[] { + // Ticker-level subscriptions first. + const tickerRows = db.prepare( + `SELECT owner_id FROM alerts + WHERE alert_type = ? AND symbol = ? AND enabled = 1`, + ).all(alertType, symbol) as { owner_id: string }[]; + + // Also find users with global subscriptions. + const globalRows = db.prepare( + `SELECT owner_id FROM alerts + WHERE alert_type = ? AND symbol IS NULL AND watchlist_id IS NULL AND enabled = 1`, + ).all(alertType) as { owner_id: string }[]; + + const userIds = new Set(); + for (const r of tickerRows) userIds.add(r.owner_id); + for (const r of globalRows) userIds.add(r.owner_id); + return Array.from(userIds); +} + +export function persistAlert(db: DatabaseSync, alert: Alert): void { + db.prepare( + `INSERT OR IGNORE INTO alert_events + (id, user_id, type, severity, title, description, symbol, created_at, acknowledged, dedup_key, payload) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, 0, ?, ?)`, + ).run( + alert.id, + alert.userId, + alert.type, + alert.severity, + alert.title, + alert.description, + alert.symbol ?? null, + alert.createdAt, + alert.dedupKey, + JSON.stringify(alert.payload), + ); +} + +export function isDuplicate(db: DatabaseSync, dedupKey: string): boolean { + const row = db.prepare( + 'SELECT 1 FROM alert_events WHERE dedup_key = ? LIMIT 1', + ).get(dedupKey) as Record | undefined; + return !!row; +} diff --git a/app/server/src/db/alertSubscriptionRepository.ts b/app/server/src/db/alertSubscriptionRepository.ts new file mode 100644 index 0000000..2b41092 --- /dev/null +++ b/app/server/src/db/alertSubscriptionRepository.ts @@ -0,0 +1,173 @@ +import type { DatabaseSync } from 'node:sqlite'; +import { randomUUID } from 'node:crypto'; + +export interface AlertSubscription { + id: string; + owner_id: string; + watchlist_id: string | null; + symbol: string | null; + alert_type: string; + enabled: boolean; + params: string; + last_fired: string | null; + created_at: string; +} + +export interface CreateSubscriptionInput { + watchlistId?: string; + symbol?: string; + alertType: string; + params?: string; +} + +function stmts(db: DatabaseSync) { + return { + insert: db.prepare( + `INSERT INTO alerts (id, owner_id, watchlist_id, symbol, alert_type, enabled, params, created_at) + VALUES (?, ?, ?, ?, ?, 1, ?, ?)`, + ), + selectById: db.prepare( + `SELECT id, owner_id, watchlist_id, symbol, alert_type, enabled, params, last_fired, created_at + FROM alerts WHERE id = ? AND owner_id = ?`, + ), + selectByOwner: db.prepare( + `SELECT id, owner_id, watchlist_id, symbol, alert_type, enabled, params, last_fired, created_at + FROM alerts WHERE owner_id = ? + ORDER BY created_at DESC`, + ), + selectByOwnerAndSymbol: db.prepare( + `SELECT id, owner_id, watchlist_id, symbol, alert_type, enabled, params, last_fired, created_at + FROM alerts WHERE owner_id = ? AND symbol = ? + ORDER BY created_at DESC`, + ), + selectTickerLevel: db.prepare( + `SELECT id, owner_id, watchlist_id, symbol, alert_type, enabled, params, last_fired, created_at + FROM alerts + WHERE owner_id = ? AND alert_type = ? AND symbol = ? AND enabled = 1 + LIMIT 1`, + ), + selectGlobalDefault: db.prepare( + `SELECT id, owner_id, watchlist_id, symbol, alert_type, enabled, params, last_fired, created_at + FROM alerts + WHERE owner_id = ? AND alert_type = ? AND enabled = 1 + AND symbol IS NULL AND watchlist_id IS NULL + LIMIT 1`, + ), + updateEnabled: db.prepare( + `UPDATE alerts SET enabled = ? WHERE id = ? AND owner_id = ?`, + ), + updateParams: db.prepare( + `UPDATE alerts SET params = ? WHERE id = ? AND owner_id = ?`, + ), + delete: db.prepare( + `DELETE FROM alerts WHERE id = ? AND owner_id = ?`, + ), + }; +} + +function mapRow(row: Record): AlertSubscription { + return { + id: row.id as string, + owner_id: row.owner_id as string, + watchlist_id: row.watchlist_id as string | null, + symbol: row.symbol as string | null, + alert_type: row.alert_type as string, + enabled: (row.enabled as number) === 1, + params: row.params as string, + last_fired: row.last_fired as string | null, + created_at: row.created_at as string, + }; +} + +export function createAlertSubscription( + db: DatabaseSync, + userId: string, + input: CreateSubscriptionInput, +): AlertSubscription { + const s = stmts(db); + const id = randomUUID(); + const now = new Date().toISOString(); + const params = input.params ?? '{}'; + s.insert.run(id, userId, input.watchlistId ?? null, input.symbol ?? null, input.alertType, params, now); + return getAlertSubscription(db, userId, id)!; +} + +export function getAlertSubscription( + db: DatabaseSync, + userId: string, + id: string, +): AlertSubscription | null { + const s = stmts(db); + const row = s.selectById.get(id, userId) as Record | undefined; + if (!row) return null; + return mapRow(row); +} + +export function listAlertSubscriptions( + db: DatabaseSync, + userId: string, + symbol?: string, +): AlertSubscription[] { + const s = stmts(db); + let rows: Record[]; + if (symbol) { + rows = s.selectByOwnerAndSymbol.all(userId, symbol) as Record[]; + } else { + rows = s.selectByOwner.all(userId) as Record[]; + } + return rows.map(mapRow); +} + +export function updateAlertSubscription( + db: DatabaseSync, + userId: string, + id: string, + updates: { enabled?: boolean; params?: string }, +): AlertSubscription | null { + const s = stmts(db); + if (updates.enabled !== undefined) { + s.updateEnabled.run(updates.enabled ? 1 : 0, id, userId); + } + if (updates.params !== undefined) { + s.updateParams.run(updates.params, id, userId); + } + return getAlertSubscription(db, userId, id); +} + +export function deleteAlertSubscription( + db: DatabaseSync, + userId: string, + id: string, +): boolean { + const s = stmts(db); + const result = s.delete.run(id, userId); + return (result as { changes: number }).changes > 0; +} + +export function getEffectiveAlertSubscription( + db: DatabaseSync, + userId: string, + symbol: string, + alertType: string, + watchlistIds: string[] = [], +): AlertSubscription | null { + const s = stmts(db); + // 1. Ticker-level wins. + const tickerRow = s.selectTickerLevel.get(userId, alertType, symbol) as Record | undefined; + if (tickerRow) return mapRow(tickerRow); + // 2. Check list-level for each watchlist the symbol belongs to. + for (const wlId of watchlistIds) { + const s2 = db.prepare( + `SELECT id, owner_id, watchlist_id, symbol, alert_type, enabled, params, last_fired, created_at + FROM alerts + WHERE owner_id = ? AND alert_type = ? AND watchlist_id = ? AND symbol IS NULL AND enabled = 1 + LIMIT 1`, + ); + const row = s2.get(userId, alertType, wlId) as Record | undefined; + if (row) return mapRow(row); + } + // 3. Global default. + const globalRow = s.selectGlobalDefault.get(userId, alertType) as Record | undefined; + if (globalRow) return mapRow(globalRow); + return null; +} diff --git a/app/server/src/db/client.ts b/app/server/src/db/client.ts index 15c47d2..946e5f8 100644 --- a/app/server/src/db/client.ts +++ b/app/server/src/db/client.ts @@ -134,7 +134,12 @@ function runMigrations(db: DatabaseSync): void { db.exec(`CREATE UNIQUE INDEX IF NOT EXISTS uq_inst_filings ON institution_filings(filer_cik, symbol, reported_quarter, form)`); } catch { /* ignore */ } - // 4. Data-quality lint table. + // 4. Rename legacy 'Starter' watchlist to 'default'. + try { + db.exec(`UPDATE watchlists SET name='default' WHERE name='Starter'`); + } catch { /* ignore */ } + + // 5. Data-quality lint table. try { db.exec(` CREATE TABLE IF NOT EXISTS data_quality ( @@ -152,7 +157,7 @@ function runMigrations(db: DatabaseSync): void { `); } catch { /* ignore */ } - // 5. Analyst ratings / upgrades & downgrades (yahoo-finance2). + // 6. Analyst ratings / upgrades & downgrades (yahoo-finance2). try { db.exec(`ALTER TABLE analyst_ratings ADD COLUMN target_from REAL`); } catch { /* ignore */ } try { db.exec(`ALTER TABLE analyst_ratings ADD COLUMN target_to REAL`); } catch { /* ignore */ } try { @@ -182,6 +187,41 @@ function runMigrations(db: DatabaseSync): void { strong_sell INTEGER NOT NULL DEFAULT 0, fetched_at TEXT NOT NULL ) + `); + } catch { /* ignore */ } + + // 7. Alert subscriptions: add watchlist_id + enabled columns (existing DBs). + try { db.exec(`ALTER TABLE alerts ADD COLUMN watchlist_id TEXT`); } catch { /* ignore */ } + try { db.exec(`ALTER TABLE alerts ADD COLUMN enabled INTEGER NOT NULL DEFAULT 1`); } catch { /* ignore */ } + + // 8. Alert comparison state table. + try { + db.exec(` + CREATE TABLE IF NOT EXISTS alert_comparison_state ( + symbol TEXT NOT NULL, + alert_type TEXT NOT NULL, + state TEXT NOT NULL DEFAULT '{}', + updated_at TEXT NOT NULL, + PRIMARY KEY (symbol, alert_type) + ) + `); + } catch { /* ignore */ } + + // 9. SMTP config table for email alerts. + try { + db.exec(` + CREATE TABLE IF NOT EXISTS smtp_config ( + id TEXT PRIMARY KEY DEFAULT 'singleton', + host TEXT NOT NULL DEFAULT 'smtp.mail.me.com', + port INTEGER NOT NULL DEFAULT 587, + secure INTEGER NOT NULL DEFAULT 0, + user TEXT, + pass_enc TEXT, + from_name TEXT NOT NULL DEFAULT 'Investor Flow', + from_email TEXT NOT NULL DEFAULT '', + enabled INTEGER NOT NULL DEFAULT 0, + updated_at TEXT NOT NULL + ) `); } catch { /* ignore */ } } diff --git a/app/server/src/db/schema.sql b/app/server/src/db/schema.sql index a3f3a0b..5fa2fc1 100644 --- a/app/server/src/db/schema.sql +++ b/app/server/src/db/schema.sql @@ -348,14 +348,17 @@ CREATE TABLE IF NOT EXISTS reports ( generated_at TEXT NOT NULL ); +-- ===== Alert subscriptions (user-defined rules for which alerts to receive) ===== CREATE TABLE IF NOT EXISTS alerts ( - id TEXT PRIMARY KEY, - owner_id TEXT NOT NULL REFERENCES users(id) ON DELETE CASCADE, - symbol TEXT, - alert_type TEXT NOT NULL, -- informed_buy|informed_sell|new_13da|rotation_incipient|regime_shift|conviction_unlock|... - params TEXT NOT NULL, -- JSON thresholds - last_fired TEXT, - created_at TEXT NOT NULL + id TEXT PRIMARY KEY, + owner_id TEXT NOT NULL REFERENCES users(id) ON DELETE CASCADE, + watchlist_id TEXT, -- NULL = global default, non-NULL = list-level + symbol TEXT, -- NULL = list-level, non-NULL = ticker-level + alert_type TEXT NOT NULL, -- informed_buy|informed_sell|new_13da|rotation_incipient|regime_shift|conviction_unlock|... + enabled INTEGER NOT NULL DEFAULT 1, + params TEXT NOT NULL DEFAULT '{}', -- JSON thresholds + last_fired TEXT, + created_at TEXT NOT NULL ); CREATE TABLE IF NOT EXISTS trusted_accounts ( @@ -566,6 +569,29 @@ CREATE INDEX IF NOT EXISTS idx_alert_events_user ON alert_events(user_id, create CREATE INDEX IF NOT EXISTS idx_alert_events_dedup ON alert_events(dedup_key); CREATE INDEX IF NOT EXISTS idx_alert_events_throttle ON alert_events(throttle_key); +-- ===== Alert comparison state (post-fetch change detection) ===== +CREATE TABLE IF NOT EXISTS alert_comparison_state ( + symbol TEXT NOT NULL, + alert_type TEXT NOT NULL, + state TEXT NOT NULL DEFAULT '{}', -- JSON snapshot + updated_at TEXT NOT NULL, + PRIMARY KEY (symbol, alert_type) +); + +-- ===== SMTP config for email alerts ===== +CREATE TABLE IF NOT EXISTS smtp_config ( + id TEXT PRIMARY KEY DEFAULT 'singleton', + host TEXT NOT NULL DEFAULT 'smtp.mail.me.com', + port INTEGER NOT NULL DEFAULT 587, + secure INTEGER NOT NULL DEFAULT 0, + user TEXT, + pass_enc TEXT, + from_name TEXT NOT NULL DEFAULT 'Investor Flow', + from_email TEXT NOT NULL DEFAULT '', + enabled INTEGER NOT NULL DEFAULT 0, + updated_at TEXT NOT NULL +); + -- ===== Slice 12 — Strategy Lab + Backtest ===== CREATE TABLE IF NOT EXISTS strategies ( id TEXT PRIMARY KEY, diff --git a/app/server/src/db/watchlistRepository.ts b/app/server/src/db/watchlistRepository.ts index 608137a..e0fac76 100644 --- a/app/server/src/db/watchlistRepository.ts +++ b/app/server/src/db/watchlistRepository.ts @@ -1,49 +1,31 @@ -// Investor Flow — Watchlist Repository (Slice 10: watchlist-portfolio-shell-panels) -// -// Thin data-access layer over the `watchlists` table. Read/write only — no trade verbs -// per ADR-0007 (Primary-Rule: no imperative-trade-verb in any string). Notes are the -// user's own neutral text, never generated by the system. -// -// Schema (schema.sql): -// CREATE TABLE IF NOT EXISTS watchlists ( -// id TEXT PRIMARY KEY, -// owner_id TEXT NOT NULL REFERENCES users(id) ON DELETE CASCADE, -// name TEXT NOT NULL, -// symbols TEXT NOT NULL, -- JSON array of symbol strings -// created_at TEXT NOT NULL, -// sort_order INTEGER NOT NULL DEFAULT 0 -// ); - import type { DatabaseSync } from 'node:sqlite'; +import { randomUUID } from 'node:crypto'; -// --------------------------------------------------------------------------- -// Types -// --------------------------------------------------------------------------- - -/** A single watchlist entry returned by listSymbols. */ export interface WatchlistEntry { symbol: string; added_at: string; notes?: string | null; } -/** A full watchlist row (internal). */ +export interface WatchlistMeta { + id: string; + name: string; + symbol_count: number; + sort_order: number; + created_at: string; +} + interface WatchlistRow { id: string; owner_id: string; name: string; - symbols: string[]; + symbols: string; created_at: string; sort_order: number; } -// --------------------------------------------------------------------------- -// Prepared statements (lazy, one per method) -// --------------------------------------------------------------------------- - function stmts(db: DatabaseSync) { return { - /** Upsert a watchlist row (idempotent by owner_id + name). */ upsert: db.prepare( `INSERT INTO watchlists (id, owner_id, name, symbols, created_at, sort_order) VALUES (?, ?, ?, ?, ?, COALESCE(?, 0)) @@ -51,56 +33,51 @@ function stmts(db: DatabaseSync) { symbols = excluded.symbols, sort_order = excluded.sort_order`, ), - - /** Select a single watchlist by owner+name. */ selectByOwnerAndName: db.prepare( `SELECT id, owner_id, name, symbols, created_at, sort_order FROM watchlists WHERE owner_id = ? AND name = ?`, ), - - /** Select all watchlists for a user. */ selectByOwner: db.prepare( `SELECT id, owner_id, name, symbols, created_at, sort_order FROM watchlists WHERE owner_id = ? ORDER BY sort_order ASC, created_at ASC`, ), - - /** Delete a watchlist by owner+name. */ deleteByOwnerAndName: db.prepare( `DELETE FROM watchlists WHERE owner_id = ? AND name = ?`, ), - - /** Update just the symbols array. */ updateSymbols: db.prepare( `UPDATE watchlists SET symbols = ? WHERE id = ? AND owner_id = ?`, ), + selectById: db.prepare( + `SELECT id, owner_id, name, symbols, created_at, sort_order + FROM watchlists WHERE id = ? AND owner_id = ?`, + ), + deleteById: db.prepare( + `DELETE FROM watchlists WHERE id = ? AND owner_id = ?`, + ), + updateSortOrder: db.prepare( + `UPDATE watchlists SET sort_order = ? WHERE id = ? AND owner_id = ?`, + ), + selectAllByOwner: db.prepare( + `SELECT id, owner_id, name, symbols, created_at, sort_order + FROM watchlists WHERE owner_id = ? + ORDER BY sort_order ASC, created_at ASC`, + ), }; } -// --------------------------------------------------------------------------- -// Repository — public API (all methods parameterized, no string interpolation) -// --------------------------------------------------------------------------- - -/** - * Add a symbol to the user's default watchlist. Idempotent — adding an already- - * present symbol is a no-op. Returns true if the symbol was newly added. - * - * Per ADR-0007, notes are the user's own neutral text; the system never generates - * directional/trade-verb language. - */ export function addSymbol( db: DatabaseSync, userId: string, symbol: string, notes?: string, + watchlistName: string = 'default', ): boolean { const s = stmts(db); const upper = symbol.toUpperCase(); - // Read existing default watchlist — get raw symbols preserving any existing notes. - const existing = readDefaultWatchlistRaw(db, userId); + const existing = readWatchlistRaw(db, userId, watchlistName); - // Check if symbol already exists (as plain string or inside an object). if (existing) { const alreadyExists = existing.symbols.some((sym) => { if (typeof sym === 'string') return sym === upper; @@ -108,7 +85,6 @@ export function addSymbol( }); if (alreadyExists) return false; - // Append the new symbol, with notes if provided. if (notes) { existing.symbols.push({ symbol: upper, notes }); } else { @@ -116,60 +92,54 @@ export function addSymbol( } const now = new Date().toISOString(); - s.upsert.run(existing.id, userId, 'default', JSON.stringify(existing.symbols), now, 0); + s.upsert.run(existing.id, userId, watchlistName, JSON.stringify(existing.symbols), now, 0); return true; } - // No existing watchlist — create a new one. const serialized = notes ? [{ symbol: upper, notes }] : [upper]; const id = generateId(); const now = new Date().toISOString(); - s.upsert.run(id, userId, 'default', JSON.stringify(serialized), now, 0); + s.upsert.run(id, userId, watchlistName, JSON.stringify(serialized), now, 0); return true; } -/** Remove a symbol from the user's default watchlist. Returns true if removed. */ export function removeSymbol( db: DatabaseSync, userId: string, symbol: string, + watchlistName: string = 'default', ): boolean { const s = stmts(db); const upper = symbol.toUpperCase(); - const existing = readDefaultWatchlistRaw(db, userId); + const existing = readWatchlistRaw(db, userId, watchlistName); if (!existing) return false; const before = existing.symbols.length; - // Filter by symbol value (whether stored as string or {symbol, notes} object). const remaining = existing.symbols.filter((sym) => { const symStr = typeof sym === 'string' ? sym : sym.symbol; return symStr !== upper; }); if (remaining.length === before) { - return false; // symbol not found + return false; } if (remaining.length === 0) { - // Clean up empty watchlist. - s.deleteByOwnerAndName.run(userId, 'default'); + s.deleteByOwnerAndName.run(userId, watchlistName); return true; } - // Single JSON.stringify — preserves existing notes on remaining symbols. s.updateSymbols.run(JSON.stringify(remaining), existing.id, userId); - return true; } -/** List all symbols across all watchlists for a user. */ export function listSymbols( db: DatabaseSync, userId: string, ): WatchlistEntry[] { const s = stmts(db); - const rows = s.selectByOwner.all(userId) as unknown as Array<{ symbols: string }>; + const rows = s.selectByOwner.all(userId) as unknown as WatchlistRow[]; const entries: WatchlistEntry[] = []; @@ -192,11 +162,84 @@ export function listSymbols( return entries; } -// --------------------------------------------------------------------------- -// Helpers -// --------------------------------------------------------------------------- +export function listWatchlists(db: DatabaseSync, userId: string): WatchlistMeta[] { + const s = stmts(db); + const rows = s.selectByOwner.all(userId) as unknown as WatchlistRow[]; + return rows.map((row) => { + const parsed = safeParseSymbols(row.symbols); + return { + id: row.id, + name: row.name, + symbol_count: parsed.length, + sort_order: row.sort_order, + created_at: row.created_at, + }; + }); +} + +export function createWatchlist( + db: DatabaseSync, + userId: string, + name: string, + symbols?: string[], +): WatchlistMeta { + const s = stmts(db); + const id = generateId(); + const now = new Date().toISOString(); + const serialized = JSON.stringify(symbols ?? []); + s.upsert.run(id, userId, name, serialized, now, 0); + return { id, name, symbol_count: (symbols ?? []).length, sort_order: 0, created_at: now }; +} + +export function deleteWatchlist(db: DatabaseSync, userId: string, name: string): boolean { + const s = stmts(db); + const existing = readWatchlistRaw(db, userId, name); + if (!existing) return false; + s.deleteByOwnerAndName.run(userId, name); + return true; +} + +export function renameWatchlist( + db: DatabaseSync, + userId: string, + oldName: string, + newName: string, +): boolean { + const s = stmts(db); + const existing = readWatchlistRaw(db, userId, oldName); + if (!existing) return false; + const conflict = readWatchlistRaw(db, userId, newName); + if (conflict) return false; + const now = new Date().toISOString(); + s.upsert.run(existing.id, userId, newName, JSON.stringify(existing.symbols), now, existing.sort_order); + s.deleteByOwnerAndName.run(userId, oldName); + return true; +} + +export function reorderWatchlists( + db: DatabaseSync, + userId: string, + orders: { id: string; sort_order: number }[], +): void { + const s = stmts(db); + for (const { id, sort_order } of orders) { + s.updateSortOrder.run(sort_order, id, userId); + } +} + +export function getSymbolsInWatchlist( + db: DatabaseSync, + userId: string, + watchlistName: string = 'default', +): string[] { + const existing = readWatchlistRaw(db, userId, watchlistName); + if (!existing) return []; + return existing.symbols.map((sym) => { + if (typeof sym === 'string') return sym; + return sym.symbol; + }); +} -/** Safely parse the JSON TEXT column into an array of strings or objects. */ function safeParseSymbols(value: string | null | undefined): (string | Record)[] { if (!value) return []; try { @@ -210,37 +253,43 @@ function safeParseSymbols(value: string | null | undefined): (string | Record { - if (typeof s === 'string') return s; - if (s && typeof s === 'object' && 'symbol' in s) return (s as { symbol: string }).symbol; - return ''; - }).filter((s): s is string => s.length > 0), - }; -} - -/** Read the default watchlist raw symbols (preserving {symbol, notes} objects). */ -function readDefaultWatchlistRaw( +export function listSymbolsByWatchlist( db: DatabaseSync, userId: string, -): { id: string; symbols: Array } | null { - const rows = stmts(db).selectByOwnerAndName.all(userId, 'default') as unknown as WatchlistRow[]; + watchlistName: string, +): WatchlistEntry[] { + const raw = readWatchlistRaw(db, userId, watchlistName); + if (!raw) return []; + + const entries: WatchlistEntry[] = []; + for (const item of raw.symbols) { + if (typeof item === 'string') { + entries.push({ symbol: item, added_at: '' }); + } else if (typeof item === 'object' && item !== null) { + const obj = item as { symbol?: string; notes?: string }; + entries.push({ + symbol: (obj.symbol ?? '').toUpperCase(), + notes: obj.notes ?? null, + added_at: '', + }); + } + } + return entries; +} + +function readWatchlistRaw( + db: DatabaseSync, + userId: string, + name: string, +): { id: string; symbols: Array; sort_order: number } | null { + const rows = stmts(db).selectByOwnerAndName.all(userId, name) as unknown as WatchlistRow[]; if (rows.length === 0) return null; const row = rows[0]; const rawSymbols = safeParseSymbols(String(row.symbols)); - return { id: row.id, symbols: rawSymbols as Array }; + return { id: row.id, symbols: rawSymbols as Array, sort_order: row.sort_order }; } -/** Generate a simple unique id. */ function generateId(): string { - return `wl_${Date.now()}_${Math.random().toString(36).slice(2, 10)}`; + return randomUUID(); } diff --git a/app/server/src/index.ts b/app/server/src/index.ts index c99e25b..66af502 100644 --- a/app/server/src/index.ts +++ b/app/server/src/index.ts @@ -66,6 +66,46 @@ const SCHEDULE_MS = 30_000; const scheduleTimer = setInterval(() => { queue.enqueueDueSchedules().catch((e) => console.error('[schedule error]', e)); }, SCHEDULE_MS); scheduleTimer.unref(); +// Per-fetch alert check: run alongside the schedule cycle for critical alert types. +const perFetchAlertTimer = setInterval(async () => { + try { + const alerts = await runProducers(database, 'per-fetch'); + if (alerts.length > 0) { + console.log(`[alert] ${alerts.length} per-fetch alert(s) created`); + const { sendAlertEmail } = await import('./services/emailAlertService.ts'); + alerts.forEach((a) => sendAlertEmail(database, a).catch(() => {})); + } + } catch (e) { + console.error('[alert] per-fetch check failed:', e); + } +}, SCHEDULE_MS); +perFetchAlertTimer.unref(); + +// Register alert producers. +import { registerProducer, runProducers } from './alerts/producers/index.ts'; +import { informedBuyProducer, informedSellProducer } from './alerts/producers/insiderProducer.ts'; +import { new13daProducer } from './alerts/producers/new13daProducer.ts'; +registerProducer(informedBuyProducer); +registerProducer(informedSellProducer); +registerProducer(new13daProducer); + +// Batched alert check: run every 10 minutes for non-critical producers. +const ALERT_BATCH_MS = 10 * 60 * 1000; +const alertBatchTimer = setInterval(async () => { + try { + const alerts = await runProducers(database, 'batched'); + if (alerts.length > 0) console.log(`[alert] ${alerts.length} batch alert(s) created`); + // Send email for each new alert. + const { sendAlertEmail } = await import('./services/emailAlertService.ts'); + for (const alert of alerts) { + sendAlertEmail(database, alert).catch((e) => console.error('[alert:email] send error:', e)); + } + } catch (e) { + console.error('[alert] batch check failed:', e); + } +}, ALERT_BATCH_MS); +alertBatchTimer.unref(); + function readBody(req: IncomingMessage): Promise { return new Promise((resolve, reject) => { let data = ''; diff --git a/app/server/src/queue/sourceRatePolicy.ts b/app/server/src/queue/sourceRatePolicy.ts index 28a8d85..8884f5f 100644 --- a/app/server/src/queue/sourceRatePolicy.ts +++ b/app/server/src/queue/sourceRatePolicy.ts @@ -10,7 +10,7 @@ import type { SourceKind } from '../cache/CacheRepository.ts'; /** Steady-state min gap between successful fetches for a source (ms). */ export const DEFAULT_SOURCE_MIN_INTERVAL_MS: Record = { - yfinance: 1500, + yfinance: 2000, sec: 150, 'sec-fetch': 1200, reddit: 2000, diff --git a/app/server/src/services/emailAlertService.ts b/app/server/src/services/emailAlertService.ts new file mode 100644 index 0000000..4501b7c --- /dev/null +++ b/app/server/src/services/emailAlertService.ts @@ -0,0 +1,131 @@ +import nodemailer from 'nodemailer'; +import type { DatabaseSync } from 'node:sqlite'; +import type { Alert } from '../alerts/AlertEngine.ts'; + +interface SmtpConfigRow { + id: string; + host: string; + port: number; + secure: number; + user: string | null; + pass_enc: string | null; + from_name: string; + from_email: string; + enabled: number; + updated_at: string; +} + +export interface SmtpConfig { + host: string; + port: number; + secure: boolean; + user: string | null; + passEnc: string | null; + fromName: string; + fromEmail: string; + enabled: boolean; +} + +export function readSmtpConfig(db: DatabaseSync): SmtpConfig | null { + const row = db.prepare('SELECT * FROM smtp_config WHERE id = ?').get('singleton') as SmtpConfigRow | undefined; + if (!row) return null; + return { + host: row.host, + port: row.port, + secure: row.secure === 1, + user: row.user ?? null, + passEnc: row.pass_enc ?? null, + fromName: row.from_name, + fromEmail: row.from_email, + enabled: row.enabled === 1, + }; +} + +function buildAlertEmailHtml(alert: Alert): string { + const severityColors: Record = { + info: '#3b82f6', + warning: '#f59e0b', + critical: '#ef4444', + }; + const color = severityColors[alert.severity] ?? '#6b7280'; + return ` + + + + + +
+ + + + +
+ + + + + +

Investor Flow

${alert.severity}
+
+

${alert.title}

+

${alert.description.replace(/\n/g, '
')}

+ +${alert.symbol ? `` : ''} + + +
Symbol:${alert.symbol}
Time:${new Date(alert.createdAt).toLocaleString()}
Type:${alert.type.replace(/_/g, ' ')}
+
+

This alert was sent by Investor Flow. You can manage your alert subscriptions in the app.

+
+
+ +`; +} + +export function getUserEmail(db: DatabaseSync, userId: string): string | null { + const row = db.prepare('SELECT email FROM users WHERE id = ?').get(userId) as { email: string } | undefined; + return row?.email ?? null; +} + +export async function sendAlertEmail( + db: DatabaseSync, + alert: Alert, +): Promise { + const config = readSmtpConfig(db); + if (!config || !config.enabled) return false; + + const userEmail = getUserEmail(db, alert.userId); + if (!userEmail) return false; + + // Check rate limit: 1 per alert_type per symbol per hour. + const hourAgo = new Date(Date.now() - 3600000).toISOString(); + const sent = db.prepare( + `SELECT COUNT(*) as count FROM alert_events + WHERE user_id = ? AND type = ? AND symbol = ? AND created_at > ?`, + ).get(alert.userId, alert.type, alert.symbol ?? null, hourAgo) as { count: number }; + if (sent.count > 1) return false; + + try { + const transporter = nodemailer.createTransport({ + host: config.host, + port: config.port, + secure: config.secure, + auth: config.user && config.passEnc + ? { user: config.user, pass: config.passEnc } + : undefined, + }); + + await transporter.sendMail({ + from: `"${config.fromName}" <${config.fromEmail}>`, + to: userEmail, + subject: `[${alert.severity.toUpperCase()}] ${alert.title}`, + html: buildAlertEmailHtml(alert), + }); + + console.log(`[alert:email] sent ${alert.type} alert to ${userEmail}`); + return true; + } catch (err) { + console.error(`[alert:email] failed to send to ${userEmail}:`, err); + return false; + } +} diff --git a/app/server/src/trpc/router.ts b/app/server/src/trpc/router.ts index eeb7e96..e822c3f 100644 --- a/app/server/src/trpc/router.ts +++ b/app/server/src/trpc/router.ts @@ -184,7 +184,7 @@ const onboardingRouter = router({ ctx.db.prepare('UPDATE users SET complexity=?, risk_tolerance=?, drawdown_tolerance=? WHERE id=?').run(complexity, riskTolerance, drawdown, userId); const symbols = input.firstWatchlistSymbols ?? STARTER_WATCHLIST.map((s) => s.symbol); const wlId = randomUUID(); - ctx.db.prepare('INSERT INTO watchlists (id, owner_id, name, symbols, created_at, sort_order) VALUES (?,?,?,?,?,?)').run(wlId, userId, 'Starter', JSON.stringify(symbols), new Date().toISOString(), 0); + ctx.db.prepare('INSERT INTO watchlists (id, owner_id, name, symbols, created_at, sort_order) VALUES (?,?,?,?,?,?)').run(wlId, userId, 'default', JSON.stringify(symbols), new Date().toISOString(), 0); for (const sym of symbols) { const kind = (STARTER_WATCHLIST.find((s) => s.symbol === sym)?.tickerKind ?? 'equity') as 'equity' | 'crypto' | 'etf' | 'index'; await ctx.cache.subscribe(sym, kind); @@ -1362,6 +1362,73 @@ const adminRouter = router({ return { ok: true, ...status }; }), + /** Get SMTP config (admin-only). */ + smtpConfig: adminProcedure.query(async ({ ctx }) => { + const { readSmtpConfig } = await import('../services/emailAlertService.ts'); + return readSmtpConfig(ctx.db) ?? { enabled: false }; + }), + + /** Update SMTP config (admin-only). */ + smtpConfigUpdate: adminProcedure + .input(z.object({ + host: z.string().optional(), + port: z.number().int().optional(), + secure: z.boolean().optional(), + user: z.string().optional().nullable(), + pass: z.string().optional().nullable(), + fromName: z.string().optional(), + fromEmail: z.string().optional(), + enabled: z.boolean().optional(), + })) + .mutation(({ ctx, input }) => { + const config = ctx.db.prepare('SELECT * FROM smtp_config WHERE id = ?').get('singleton') as Record | undefined; + const now = new Date().toISOString(); + if (!config) { + ctx.db.prepare( + `INSERT INTO smtp_config (id, host, port, secure, user, pass_enc, from_name, from_email, enabled, updated_at) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`, + ).run('singleton', input.host ?? 'smtp.mail.me.com', input.port ?? 587, input.secure ? 1 : 0, input.user ?? null, input.pass ?? null, input.fromName ?? 'Investor Flow', input.fromEmail ?? '', input.enabled ? 1 : 0, now); + } else { + const updates: string[] = []; + const params: (string | number | null)[] = []; + if (input.host !== undefined) { updates.push('host = ?'); params.push(input.host); } + if (input.port !== undefined) { updates.push('port = ?'); params.push(input.port); } + if (input.secure !== undefined) { updates.push('secure = ?'); params.push(input.secure ? 1 : 0); } + if (input.user !== undefined) { updates.push('user = ?'); params.push(input.user); } + if (input.pass !== undefined) { updates.push('pass_enc = ?'); params.push(input.pass); } + if (input.fromName !== undefined) { updates.push('from_name = ?'); params.push(input.fromName); } + if (input.fromEmail !== undefined) { updates.push('from_email = ?'); params.push(input.fromEmail); } + if (input.enabled !== undefined) { updates.push('enabled = ?'); params.push(input.enabled ? 1 : 0); } + updates.push('updated_at = ?'); + params.push(now); + params.push('singleton'); + ctx.db.prepare(`UPDATE smtp_config SET ${updates.join(', ')} WHERE id = ?`).run(...params); + } + return { ok: true }; + }), + + /** Test SMTP config by sending a test email to the admin. */ + smtpConfigTest: adminProcedure.mutation(async ({ ctx }) => { + const { readSmtpConfig, sendAlertEmail } = await import('../services/emailAlertService.ts'); + const config = readSmtpConfig(ctx.db); + if (!config || !config.enabled) throw new TRPCError({ code: 'BAD_REQUEST', message: 'SMTP not configured or disabled.' }); + const testAlert = { + id: 'test', + userId: ctx.userId as string, + type: 'rotation_incipient' as const, + severity: 'info' as const, + title: 'Test alert from Investor Flow', + description: 'This is a test email to verify SMTP configuration.', + symbol: undefined, + createdAt: new Date().toISOString(), + acknowledged: false, + dedupKey: 'test:' + Date.now(), + payload: {}, + }; + const sent = await sendAlertEmail(ctx.db, testAlert); + if (!sent) throw new TRPCError({ code: 'INTERNAL_SERVER_ERROR', message: 'Failed to send test email.' }); + return { ok: true }; + }), }); @@ -1457,6 +1524,64 @@ const alertsRouter = router({ ).get(userId) as { count: number } | undefined; return { count: row?.count ?? 0 }; }), + + /** Create an alert subscription rule. */ + createSubscription: protectedProcedure + .input(z.object({ + alertType: z.string(), + watchlistId: z.string().optional(), + symbol: z.string().optional(), + params: z.string().optional(), + })) + .mutation(async ({ ctx, input }) => { + const userId = ctx.userId as string; + const { createAlertSubscription } = await import('../db/alertSubscriptionRepository.ts'); + const sub = createAlertSubscription(ctx.db, userId, { + alertType: input.alertType, + watchlistId: input.watchlistId, + symbol: input.symbol, + params: input.params, + }); + return sub; + }), + + /** List alert subscriptions (optionally filtered by symbol). */ + listSubscriptions: protectedProcedure + .input(z.object({ symbol: z.string().optional() }).optional()) + .query(async ({ ctx, input }) => { + const userId = ctx.userId as string; + const { listAlertSubscriptions } = await import('../db/alertSubscriptionRepository.ts'); + return listAlertSubscriptions(ctx.db, userId, input?.symbol); + }), + + /** Update an alert subscription (enable/disable, change params). */ + updateSubscription: protectedProcedure + .input(z.object({ + id: z.string(), + enabled: z.boolean().optional(), + params: z.string().optional(), + })) + .mutation(async ({ ctx, input }) => { + const userId = ctx.userId as string; + const { updateAlertSubscription } = await import('../db/alertSubscriptionRepository.ts'); + const sub = updateAlertSubscription(ctx.db, userId, input.id, { + enabled: input.enabled, + params: input.params, + }); + if (!sub) throw new TRPCError({ code: 'NOT_FOUND', message: 'Subscription not found.' }); + return sub; + }), + + /** Delete an alert subscription. */ + deleteSubscription: protectedProcedure + .input(z.object({ id: z.string() })) + .mutation(async ({ ctx, input }) => { + const userId = ctx.userId as string; + const { deleteAlertSubscription } = await import('../db/alertSubscriptionRepository.ts'); + const deleted = deleteAlertSubscription(ctx.db, userId, input.id); + if (!deleted) throw new TRPCError({ code: 'NOT_FOUND', message: 'Subscription not found.' }); + return { ok: true }; + }), }); // ─── Institutional Flow Router (Slice 7 / M4 + M5) ──────────────────────── @@ -1880,7 +2005,7 @@ const edgarRouter = router({ // ─── Watchlist Router (Slice 10) ───────────────────────────────────────────── const watchlistRouter = router({ - /** List all symbols in the user's default watchlist. */ + /** List all symbols across all watchlists for the user. */ list: publicProcedure .query(async ({ ctx }) => { const userId = ctx.userId ?? 'anonymous'; @@ -1888,13 +2013,71 @@ const watchlistRouter = router({ return listSymbols(ctx.db, userId); }), - /** Add a symbol to the default watchlist. */ + /** List symbols for a specific watchlist. */ + listByWatchlist: publicProcedure + .input(z.object({ name: z.string().min(1).max(100) })) + .query(async ({ ctx, input }) => { + const userId = ctx.userId ?? 'anonymous'; + const { listSymbolsByWatchlist } = await import('../db/watchlistRepository.ts'); + return listSymbolsByWatchlist(ctx.db, userId, input.name); + }), + + /** List watchlist metadata (name, symbol count, sort order). */ + listWatchlists: publicProcedure + .query(async ({ ctx }) => { + const userId = ctx.userId ?? 'anonymous'; + const { listWatchlists } = await import('../db/watchlistRepository.ts'); + return listWatchlists(ctx.db, userId); + }), + + /** Create a new named watchlist. */ + create: publicProcedure + .input(z.object({ name: z.string().min(1).max(100) })) + .mutation(async ({ ctx, input }) => { + const userId = ctx.userId ?? 'anonymous'; + const { createWatchlist } = await import('../db/watchlistRepository.ts'); + return createWatchlist(ctx.db, userId, input.name); + }), + + /** Delete a watchlist by name. */ + delete: publicProcedure + .input(z.object({ name: z.string().min(1) })) + .mutation(async ({ ctx, input }) => { + const userId = ctx.userId ?? 'anonymous'; + const { deleteWatchlist } = await import('../db/watchlistRepository.ts'); + return { deleted: deleteWatchlist(ctx.db, userId, input.name) }; + }), + + /** Rename a watchlist. */ + rename: publicProcedure + .input(z.object({ oldName: z.string().min(1), newName: z.string().min(1).max(100) })) + .mutation(async ({ ctx, input }) => { + const userId = ctx.userId ?? 'anonymous'; + const { renameWatchlist } = await import('../db/watchlistRepository.ts'); + return { renamed: renameWatchlist(ctx.db, userId, input.oldName, input.newName) }; + }), + + /** Reorder watchlists (bulk update sort_order). */ + reorder: publicProcedure + .input(z.object({ orders: z.array(z.object({ id: z.string(), sort_order: z.number().int() })) })) + .mutation(async ({ ctx, input }) => { + const userId = ctx.userId ?? 'anonymous'; + const { reorderWatchlists } = await import('../db/watchlistRepository.ts'); + reorderWatchlists(ctx.db, userId, input.orders); + return { ok: true }; + }), + + /** Add a symbol to a watchlist (defaults to 'default' watchlist). */ addSymbol: publicProcedure - .input(z.object({ symbol: z.string().toUpperCase() })) + .input(z.object({ + symbol: z.string().toUpperCase(), + watchlistName: z.string().optional(), + notes: z.string().optional(), + })) .mutation(async ({ ctx, input }) => { const userId = ctx.userId ?? 'anonymous'; const { addSymbol } = await import('../db/watchlistRepository.ts'); - const added = addSymbol(ctx.db, userId, input.symbol); + const added = addSymbol(ctx.db, userId, input.symbol, input.notes, input.watchlistName); if (added) { await ctx.cache.subscribe(input.symbol, 'equity'); @@ -1904,13 +2087,16 @@ const watchlistRouter = router({ return { added }; }), - /** Remove a symbol from the default watchlist. */ + /** Remove a symbol from a watchlist (defaults to 'default'). */ removeSymbol: publicProcedure - .input(z.object({ symbol: z.string().toUpperCase() })) + .input(z.object({ + symbol: z.string().toUpperCase(), + watchlistName: z.string().optional(), + })) .mutation(async ({ ctx, input }) => { const userId = ctx.userId ?? 'anonymous'; const { removeSymbol } = await import('../db/watchlistRepository.ts'); - const removed = removeSymbol(ctx.db, userId, input.symbol); + const removed = removeSymbol(ctx.db, userId, input.symbol, input.watchlistName); return { removed }; }), }); diff --git a/app/src/app/admin/smtp/page.tsx b/app/src/app/admin/smtp/page.tsx new file mode 100644 index 0000000..dd1f73a --- /dev/null +++ b/app/src/app/admin/smtp/page.tsx @@ -0,0 +1,121 @@ +"use client"; +import { useEffect, useState } from "react"; +import { AdminLayout } from "@/components/AdminLayout"; +import { api } from "@/lib/trpc"; + +export default function SmtpAdminPage() { + const [config, setConfig] = useState<{ host: string; port: number; secure: boolean; user: string | null; passEnc: string | null; fromName: string; fromEmail: string; enabled: boolean } | null>(null); + const [host, setHost] = useState("smtp.mail.me.com"); + const [port, setPort] = useState(587); + const [secure, setSecure] = useState(false); + const [user, setUser] = useState(""); + const [pass, setPass] = useState(""); + const [fromName, setFromName] = useState("Investor Flow"); + const [fromEmail, setFromEmail] = useState(""); + const [enabled, setEnabled] = useState(false); + const [saving, setSaving] = useState(false); + const [testing, setTesting] = useState(false); + const [msg, setMsg] = useState<{ type: "success" | "error"; text: string } | null>(null); + + useEffect(() => { + api.admin.smtpConfig().then((c) => { + setConfig(c as any); + setHost((c as any).host ?? "smtp.mail.me.com"); + setPort((c as any).port ?? 587); + setSecure((c as any).secure ?? false); + setUser((c as any).user ?? ""); + setFromName((c as any).fromName ?? "Investor Flow"); + setFromEmail((c as any).fromEmail ?? ""); + setEnabled((c as any).enabled ?? false); + }).catch(() => {}); + }, []); + + const save = async () => { + setSaving(true); + setMsg(null); + try { + await api.admin.smtpConfigUpdate({ host, port, secure, user: user || null, pass: pass || undefined, fromName, fromEmail, enabled }); + setMsg({ type: "success", text: "SMTP config saved." }); + } catch (e) { + setMsg({ type: "error", text: e instanceof Error ? e.message : "Failed to save." }); + } + setSaving(false); + }; + + const test = async () => { + setTesting(true); + setMsg(null); + try { + await api.admin.smtpConfigTest(); + setMsg({ type: "success", text: "Test email sent! Check your inbox." }); + } catch (e) { + setMsg({ type: "error", text: e instanceof Error ? e.message : "Test failed." }); + } + setTesting(false); + }; + + return ( + +
+

SMTP Configuration

+ + {msg && ( +
+ {msg.text} +
+ )} + +
+ + + + + + + + + + + + + + + + +
+ + +
+
+
+
+ ); +} diff --git a/app/src/app/alerts/page.tsx b/app/src/app/alerts/page.tsx new file mode 100644 index 0000000..7bfb2d5 --- /dev/null +++ b/app/src/app/alerts/page.tsx @@ -0,0 +1,172 @@ +"use client"; +import { useEffect, useState } from "react"; +import { api, type AlertEventRow, type AlertSubscriptionRow } from "@/lib/trpc"; + +const SEVERITY_COLORS: Record = { + info: "text-blue-400 bg-blue-950/30 border-blue-500/30", + warning: "text-amber-400 bg-amber-950/30 border-amber-500/30", + critical: "text-red-400 bg-red-950/30 border-red-500/30", +}; + +export default function AlertsPage() { + const [tab, setTab] = useState<"events" | "subscriptions">("events"); + const [events, setEvents] = useState([]); + const [subscriptions, setSubscriptions] = useState([]); + const [alertType, setAlertType] = useState("informed_buy"); + const [symbol, setSymbol] = useState(""); + const [loading, setLoading] = useState(true); + + const loadEvents = async () => { + try { + const data = await api.alerts.list(100); + setEvents(data); + } catch { /* ignore */ } + }; + + const loadSubscriptions = async () => { + try { + const data = await api.alerts.listSubscriptions(); + setSubscriptions(data); + } catch { /* ignore */ } + }; + + useEffect(() => { + Promise.all([loadEvents(), loadSubscriptions()]).finally(() => setLoading(false)); + }, []); + + const acknowledge = async (id: string) => { + await api.alerts.acknowledge(id); + loadEvents(); + }; + + const acknowledgeAll = async () => { + await api.alerts.acknowledgeAll(); + loadEvents(); + }; + + const addSubscription = async () => { + const sym = symbol.trim().toUpperCase(); + await api.alerts.createSubscription({ alertType, symbol: sym || undefined }); + setSymbol(""); + loadSubscriptions(); + }; + + const toggleSubscription = async (id: string, enabled: boolean) => { + await api.alerts.updateSubscription(id, { enabled: !enabled }); + loadSubscriptions(); + }; + + const deleteSubscription = async (id: string) => { + await api.alerts.deleteSubscription(id); + loadSubscriptions(); + }; + + const unacknowledged = events.filter((e) => !e.acknowledged); + + return ( +
+
+

Alerts

+ +
+ + +
+ + {loading ? ( +
+
+
+ ) : tab === "events" ? ( +
+ {unacknowledged.length > 0 && ( + + )} + {events.length === 0 ? ( +

No alerts yet.

+ ) : ( +
+ {events.map((e) => ( +
+
+
+
+ {e.severity} + {e.title} + {e.symbol && {e.symbol}} +
+

{e.description}

+

{new Date(e.createdAt).toLocaleString()}

+
+ {!e.acknowledged && ( + + )} +
+
+ ))} +
+ )} +
+ ) : ( +
+
+

New Subscription

+
+ + setSymbol(e.target.value)} placeholder="Symbol (optional)" className="rounded-lg border border-border bg-surface-dark px-3 py-2 text-sm text-fg w-32" /> + +
+
+ + {subscriptions.length === 0 ? ( +

No subscriptions configured.

+ ) : ( +
+ {subscriptions.map((s) => ( +
+
+ {s.alert_type.replace(/_/g, " ")} + {s.symbol && {s.symbol}} + {!s.symbol && !s.watchlist_id && global} +
+
+
+ ))} +
+ )} +
+ )} +
+
+ ); +} diff --git a/app/src/components/MobileTabNav.tsx b/app/src/components/MobileTabNav.tsx index 8458ced..9c819e4 100644 --- a/app/src/components/MobileTabNav.tsx +++ b/app/src/components/MobileTabNav.tsx @@ -18,6 +18,7 @@ const TABS: TabItem[] = [ { label: "Overview", href: "/", icon: "📊" }, { label: "Filings", href: "/filings", icon: "📄" }, { label: "Options", href: "/options-dd", icon: "🔮" }, + { label: "Alerts", href: "/alerts", icon: "🔔" }, { label: "Execution", href: "/execution", icon: "⚡" }, { label: "Trade Plan", href: "/trade-plan", icon: "📋" }, { label: "Settings", href: "/settings", icon: "⚙️" }, diff --git a/app/src/components/SidebarNav.tsx b/app/src/components/SidebarNav.tsx index 84a25ab..157d962 100644 --- a/app/src/components/SidebarNav.tsx +++ b/app/src/components/SidebarNav.tsx @@ -10,8 +10,10 @@ import { Settings, ChevronDown, Search, + Bell, type LucideIcon, } from "lucide-react"; +import { api } from "@/lib/trpc"; /** * Left sidebar navigation with collapsible sections. @@ -58,6 +60,7 @@ const SECTIONS: NavSection[] = [ { label: "Risk Posture", href: "/risk" }, { label: "Filings", href: "/filings" }, { label: "Options DD", href: "/options-dd" }, + { label: "Alerts", href: "/alerts" }, ], }, { @@ -69,6 +72,7 @@ const SECTIONS: NavSection[] = [ { label: "Queue", href: "/admin/queue" }, { label: "Audit Logs", href: "/admin/audit-logs" }, { label: "X Accounts", href: "/admin/x-accounts" }, + { label: "SMTP", href: "/admin/smtp" }, ], }, { @@ -101,6 +105,11 @@ export function SidebarNav() { const normalizedPath = (pathname || "/").replace(/\/+$/, "") || "/"; const [collapsed, setCollapsed] = useState>({}); const [initialized, setInitialized] = useState(false); + const [unackedCount, setUnackedCount] = useState(0); + + useEffect(() => { + api.alerts.unackedCount().then((r) => setUnackedCount(r.count)).catch(() => {}); + }, []); useEffect(() => { setCollapsed(getInitialCollapsed(pathname)); @@ -146,6 +155,7 @@ export function SidebarNav() {
{section.items.map((item) => { const isActive = normalizedPath === item.href; + const showBadge = item.label === "Alerts" && unackedCount > 0; return ( - {item.label} + + {item.label} + {showBadge && ( + {unackedCount} + )} + ); })} diff --git a/app/src/components/WatchlistSidebar.tsx b/app/src/components/WatchlistSidebar.tsx index 44d0324..5f44b91 100644 --- a/app/src/components/WatchlistSidebar.tsx +++ b/app/src/components/WatchlistSidebar.tsx @@ -1,19 +1,22 @@ "use client"; -import { useEffect, useState } from "react"; +import { useEffect, useRef, useState } from "react"; import { useActiveSymbol } from "@/stores/active-symbol-store"; import { api, type WatchlistEntry, type Quote, type PortfolioHolding } from "@/lib/trpc"; import { chart as CHART } from "@/lib/chart-theme"; import { UI_STRINGS } from "@/lib/strings"; -/** - * Compact watchlist sidebar with mini sparklines and quick symbol switching. - * Shows portfolio summary at top. Replaces the full WatchlistPanel on the main page. - */ - interface MiniQuote extends Quote { sparkline?: number[]; } +interface WatchlistMeta { + id: string; + name: string; + symbol_count: number; + sort_order: number; + created_at: string; +} + interface WatchlistSidebarProps { compact?: boolean; } @@ -26,6 +29,11 @@ export function WatchlistSidebar({ compact = false }: WatchlistSidebarProps) { const [draft, setDraft] = useState(""); const [loading, setLoading] = useState(true); const [holdings, setHoldings] = useState([]); + const [watchlists, setWatchlists] = useState([]); + const [activeWatchlist, setActiveWatchlist] = useState("default"); + const [showListPicker, setShowListPicker] = useState(false); + const [newListName, setNewListName] = useState(""); + const pickerRef = useRef(null); // Portfolio summary useEffect(() => { @@ -34,14 +42,30 @@ export function WatchlistSidebar({ compact = false }: WatchlistSidebarProps) { const totalValue = (holdings ?? []).reduce((sum, h) => sum + h.shares * h.avg_cost, 0); + const loadWatchlists = async () => { + try { + const wls = await api.watchlists.listWatchlists(); + setWatchlists(wls); + if (wls.length > 0 && !wls.find((w) => w.name === activeWatchlist)) { + setActiveWatchlist(wls[0].name); + } + } catch { /* ignore */ } + }; + const load = async () => { try { - setEntries(await api.watchlists.list()); + setEntries(await api.watchlists.listByWatchlist(activeWatchlist)); } catch { /* ignore */ } setLoading(false); }; - useEffect(() => { load(); }, []); + useEffect(() => { + loadWatchlists(); + }, []); + + useEffect(() => { + load(); + }, [activeWatchlist]); // Fetch quotes for watchlist symbols useEffect(() => { @@ -64,16 +88,48 @@ export function WatchlistSidebar({ compact = false }: WatchlistSidebarProps) { const add = async () => { if (!draft.trim()) return; const sym = draft.trim().toUpperCase(); - await api.watchlists.addSymbol(sym); + await api.watchlists.addSymbol(sym, activeWatchlist === "default" ? undefined : activeWatchlist); setDraft(""); load(); + loadWatchlists(); }; const remove = async (symbol: string) => { - await api.watchlists.removeSymbol(symbol); + await api.watchlists.removeSymbol(symbol, activeWatchlist === "default" ? undefined : activeWatchlist); load(); + loadWatchlists(); }; + const createWatchlist = async () => { + if (!newListName.trim()) return; + await api.watchlists.create(newListName.trim()); + setNewListName(""); + setActiveWatchlist(newListName.trim()); + setShowListPicker(false); + loadWatchlists(); + }; + + const deleteWatchlist = async (name: string) => { + await api.watchlists.delete(name); + loadWatchlists(); + if (activeWatchlist === name) { + setActiveWatchlist("default"); + } + }; + + // Close picker on outside click + useEffect(() => { + const handler = (e: MouseEvent) => { + if (pickerRef.current && !pickerRef.current.contains(e.target as Node)) { + setShowListPicker(false); + } + }; + if (showListPicker) { + document.addEventListener("mousedown", handler); + return () => document.removeEventListener("mousedown", handler); + } + }, [showListPicker]); + const containerClass = compact ? "w-full border-b border-line bg-surface-raised flex-shrink-0" : "w-56 border-r border-line bg-surface-raised flex-shrink-0 hidden lg:block"; @@ -83,7 +139,7 @@ export function WatchlistSidebar({ compact = false }: WatchlistSidebarProps) { return (