Phase 5: alert subscriptions UI + multiple watchlists + queue fixes
CI / Test & Type-Check (push) Canceled after 0s
CI / Test & Type-Check (push) Canceled after 0s
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
This commit is contained in:
@@ -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,
|
||||
});
|
||||
|
||||
@@ -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<Alert[]> {
|
||||
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);
|
||||
}
|
||||
@@ -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<Alert[]> {
|
||||
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<string>(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');
|
||||
@@ -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<Alert[]> {
|
||||
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;
|
||||
},
|
||||
};
|
||||
@@ -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<Alert[]>;
|
||||
}
|
||||
|
||||
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<string, unknown> | 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<string, unknown>,
|
||||
): 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<string>();
|
||||
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<string, unknown> | undefined;
|
||||
return !!row;
|
||||
}
|
||||
@@ -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<string, unknown>): 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<string, unknown> | 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<string, unknown>[];
|
||||
if (symbol) {
|
||||
rows = s.selectByOwnerAndSymbol.all(userId, symbol) as Record<string, unknown>[];
|
||||
} else {
|
||||
rows = s.selectByOwner.all(userId) as Record<string, unknown>[];
|
||||
}
|
||||
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<string, unknown> | 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<string, unknown> | undefined;
|
||||
if (row) return mapRow(row);
|
||||
}
|
||||
// 3. Global default.
|
||||
const globalRow = s.selectGlobalDefault.get(userId, alertType) as Record<string, unknown> | undefined;
|
||||
if (globalRow) return mapRow(globalRow);
|
||||
return null;
|
||||
}
|
||||
@@ -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 */ }
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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<string, unknown>)[] {
|
||||
if (!value) return [];
|
||||
try {
|
||||
@@ -210,37 +253,43 @@ function safeParseSymbols(value: string | null | undefined): (string | Record<st
|
||||
return [];
|
||||
}
|
||||
|
||||
/** Read the default watchlist for a user. Returns null if not found. */
|
||||
function readDefaultWatchlist(db: DatabaseSync, userId: string): WatchlistRow | null {
|
||||
const rows = stmts(db).selectByOwnerAndName.all(userId, 'default') as unknown as WatchlistRow[];
|
||||
if (rows.length === 0) return null;
|
||||
|
||||
const row = rows[0];
|
||||
return {
|
||||
...row,
|
||||
// Extract plain symbol strings from the mixed array (handles both legacy strings and {symbol,notes} objects).
|
||||
symbols: safeParseSymbols(String(row.symbols)).map((s) => {
|
||||
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<string | { symbol: string; notes?: string }> } | 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<string | { symbol: string; notes?: string }>; 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<string | { symbol: string; notes?: string }> };
|
||||
return { id: row.id, symbols: rawSymbols as Array<string | { symbol: string; notes?: string }>, 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();
|
||||
}
|
||||
|
||||
@@ -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<string> {
|
||||
return new Promise((resolve, reject) => {
|
||||
let data = '';
|
||||
|
||||
@@ -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<SourceKind, number> = {
|
||||
yfinance: 1500,
|
||||
yfinance: 2000,
|
||||
sec: 150,
|
||||
'sec-fetch': 1200,
|
||||
reddit: 2000,
|
||||
|
||||
@@ -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<string, string> = {
|
||||
info: '#3b82f6',
|
||||
warning: '#f59e0b',
|
||||
critical: '#ef4444',
|
||||
};
|
||||
const color = severityColors[alert.severity] ?? '#6b7280';
|
||||
return `<!DOCTYPE html>
|
||||
<html>
|
||||
<head><meta charset="utf-8"><meta name="viewport" content="width=device-width,initial-scale=1"></head>
|
||||
<body style="font-family:-apple-system,BlinkMacSystemFont,'Segoe UI',Roboto,sans-serif;margin:0;padding:0;background:#f9fafb">
|
||||
<table role="presentation" width="100%" cellpadding="0" cellspacing="0" style="background:#f9fafb;padding:20px">
|
||||
<tr><td align="center">
|
||||
<table role="presentation" width="560" cellpadding="0" cellspacing="0" style="background:#fff;border-radius:8px;overflow:hidden;box-shadow:0 1px 3px rgba(0,0,0,0.1)">
|
||||
<tr><td style="padding:24px;border-bottom:3px solid ${color}">
|
||||
<table width="100%" cellpadding="0" cellspacing="0">
|
||||
<tr>
|
||||
<td><h1 style="margin:0;font-size:20px;color:#111827">Investor Flow</h1></td>
|
||||
<td align="right"><span style="display:inline-block;padding:4px 12px;border-radius:12px;font-size:12px;font-weight:600;color:#fff;background:${color};text-transform:uppercase">${alert.severity}</span></td>
|
||||
</tr>
|
||||
</table>
|
||||
</td></tr>
|
||||
<tr><td style="padding:24px">
|
||||
<h2 style="margin:0 0 8px;font-size:18px;color:#111827">${alert.title}</h2>
|
||||
<p style="margin:0 0 16px;font-size:14px;color:#6b7280;line-height:1.5">${alert.description.replace(/\n/g, '<br>')}</p>
|
||||
<table width="100%" cellpadding="0" cellspacing="0" style="font-size:13px;color:#6b7280">
|
||||
${alert.symbol ? `<tr><td style="padding:4px 0">Symbol:</td><td style="padding:4px 0;font-weight:600;color:#111827">${alert.symbol}</td></tr>` : ''}
|
||||
<tr><td style="padding:4px 0">Time:</td><td style="padding:4px 0;font-weight:600;color:#111827">${new Date(alert.createdAt).toLocaleString()}</td></tr>
|
||||
<tr><td style="padding:4px 0">Type:</td><td style="padding:4px 0;font-weight:600;color:#111827">${alert.type.replace(/_/g, ' ')}</td></tr>
|
||||
</table>
|
||||
</td></tr>
|
||||
<tr><td style="padding:16px 24px;background:#f9fafb;border-top:1px solid #e5e7eb;font-size:12px;color:#9ca3af;text-align:center">
|
||||
<p style="margin:0">This alert was sent by Investor Flow. You can manage your alert subscriptions in the app.</p>
|
||||
</td></tr>
|
||||
</table>
|
||||
</td></tr>
|
||||
</table>
|
||||
</body>
|
||||
</html>`;
|
||||
}
|
||||
|
||||
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<boolean> {
|
||||
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;
|
||||
}
|
||||
}
|
||||
@@ -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<string, unknown> | 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 };
|
||||
}),
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user