// Investor Flow — backend entry (DESIGN.md §2.1: one process exposes /api/trpc/*). // Adapted from Bun.serve to node:http + @trpc/server fetch adapter (runtime glue only). import { createServer, type IncomingMessage, type ServerResponse } from 'node:http'; import { fetchRequestHandler } from '@trpc/server/adapters/fetch'; import { db } from './db/client.ts'; import { createCacheRepository, type SourceKind } from './cache/CacheRepository.ts'; import { YFinanceAdapter } from './adapters/YFinanceAdapter.ts'; import { OptionsAdapter } from './adapters/OptionsAdapter.ts'; import { NasdaqAdapter } from './adapters/NasdaqAdapter.ts'; import { FinraBulkAdapter } from './adapters/FinraBulkAdapter.ts'; import { FinraShortInterestAdapter } from './adapters/FinraShortInterestAdapter.ts'; import { SecFetchAdapter } from './adapters/SecFetchAdapter.ts'; import { SecCompanyTickersAdapter } from './adapters/SecCompanyTickersAdapter.ts'; import { SecLintAdapter } from './adapters/SecLintAdapter.ts'; import { XCookieAdapter } from './adapters/XCookieAdapter.ts'; import { FredAdapterImpl } from './macro/FredAdapter.ts'; import type { SourceFetch } from './adapters/SourceAdapter.ts'; import { composeYFinanceWithOptions } from './options/OptionsChainRouter.ts'; import cryptoMod from './lib/crypto.ts'; import { AdapterQueue } from './queue/AdapterQueue.ts'; import { seedCuratedCusips } from './services/cusipRegistry.ts'; import { seedAdminDefaultAlertSubscriptions } from './db/alertSubscriptionRepository.ts'; import { makeCreateContext } from './trpc/context.ts'; import { appRouter } from './trpc/router.ts'; import type { Alert } from './alerts/AlertEngine.ts'; const PORT = Number(process.env.PORT ?? 3001); const database = db(); const yfinanceAdapter = composeYFinanceWithOptions( new YFinanceAdapter({ db: database }) as unknown as SourceFetch, new OptionsAdapter() as unknown as SourceFetch, ); const adapters = new Map([ ['yfinance' as const, yfinanceAdapter], ['nasdaq' as const, new NasdaqAdapter() as unknown as SourceFetch], ['finra-bulk' as const, new FinraBulkAdapter(database) as unknown as SourceFetch], ['finra-si' as const, new FinraShortInterestAdapter(database) as unknown as SourceFetch], ['sec-fetch' as const, new SecFetchAdapter(database) as unknown as SourceFetch], ['sec-sc-fetch' as const, new SecFetchAdapter(database) as unknown as SourceFetch], ['sec-tickers' as const, new SecCompanyTickersAdapter(database) as unknown as SourceFetch], ['sec-lint-holders' as const, new SecLintAdapter(() => database, 'sec-lint-holders') as unknown as SourceFetch], ['sec-lint-insiders' as const, new SecLintAdapter(() => database, 'sec-lint-insiders') as unknown as SourceFetch], ]); // Load X credentials at startup and register XCookieAdapter if available. let xAdapter: XCookieAdapter | null = null; (function initXAdapter() { const row = database.prepare('SELECT ct0_enc, auth_token_enc FROM x_credentials WHERE id=?').get('singleton') as { ct0_enc?: string; auth_token_enc?: string } | undefined; if (!row?.ct0_enc || !row?.auth_token_enc) return; let creds: { ct0: string; auth_token: string }; try { creds = { ct0: cryptoMod.decrypt(row.ct0_enc), auth_token: cryptoMod.decrypt(row.auth_token_enc) }; } catch { return; } const updateHealth = (status: string, err?: string | null) => { database.prepare( `INSERT INTO x_credentials (id, healthy, last_error, updated_at) VALUES ('singleton', ?, ?, ?) ON CONFLICT(id) DO UPDATE SET healthy=excluded.healthy, last_error=excluded.last_error, updated_at=excluded.updated_at` ).run(status === 'healthy' ? 1 : 0, err ?? null, new Date().toISOString()); }; xAdapter = new XCookieAdapter(creds, (health) => updateHealth(health.sourceStatus, health.lastError), database); adapters.set('x' as const, xAdapter as unknown as SourceFetch); console.log('[investor-flow] X adapter registered (credentials configured)'); })(); // Register FredAdapter when a FRED API key is configured (encrypted in x_credentials). // FRED series are warmed by the queue schedule (fred_macro tier), never on the // market.condition request path. (function initFredAdapter() { const row = database.prepare('SELECT fred_api_key_enc FROM x_credentials WHERE id=?').get('singleton') as { fred_api_key_enc?: string | null } | undefined; if (!row?.fred_api_key_enc) return; let apiKey: string; try { apiKey = cryptoMod.decrypt(row.fred_api_key_enc); } catch { return; } if (!apiKey) return; adapters.set('fred' as const, new FredAdapterImpl(apiKey) as unknown as SourceFetch); console.log('[investor-flow] FRED adapter registered (series warm-up via queue schedule)'); })(); const queue = new AdapterQueue({ db: database, adapters }); const cache = createCacheRepository({ db: database, scheduler: queue }); queue.cache = cache; // break the cache<->scheduler cycle // Seed / migrate tiered schedules (safe every boot) queue.seedDefaultSchedules(); // Offline CUSIP registry → kv_cache so institutional alerts never depend solely on EFTS. try { const seeded = seedCuratedCusips(database); if (seeded > 0) console.log(`[investor-flow] seeded ${seeded} curated sec:cusip cache entries`); } catch (e) { console.error('[investor-flow] cusip seed failed', e); } // Admin default alert subscriptions: every catalog alert type ON for is_admin=1 // (idempotent; respects any toggled preference). Seeded now, not just on the // 10-min batch producer tick, so operators receive alerts immediately. try { const seeded = seedAdminDefaultAlertSubscriptions(database); if (seeded > 0) console.log(`[investor-flow] seeded ${seeded} admin default alert subscription(s)`); } catch (e) { console.error('[investor-flow] admin alert seed failed', e); } // Demand hygiene: junk test symbols + inflated refcounts from page-view subscribe spam try { const hygiene = queue.cleanupDemandHygiene(); if (hygiene.removedJunk || hygiene.cappedRefcounts) { console.log(`[investor-flow] demand hygiene: removedJunk=${hygiene.removedJunk} cappedRefcounts=${hygiene.cappedRefcounts}`); } } catch (e) { console.error('[investor-flow] demand hygiene failed', e); } // Pin rotation universe + SPY + VIX (no refcount inflation) queue.pinSystemUniverse().then(() => { console.log('[investor-flow] system pins (rotation universe) ready'); }).catch((e) => console.error('[investor-flow] pinSystemUniverse failed', e)); // Startup recovery: any job left 'in_flight' was interrupted by a restart/crash. // Reset to 'pending' so the drain loop reprocesses it. const recovered = database.prepare("UPDATE adapter_queue SET status='pending', error=NULL, retry_count=0 WHERE status='in_flight'").run(); if (Number(recovered.changes) > 0) console.log(`[investor-flow] recovered ${recovered.changes} interrupted in_flight jobs`); const createContext = makeCreateContext({ db: database, cache, queue, xAdapter }); // Background drain: stale-while-revalidate refreshes are queued by CacheRepository.get; // this loop drains them (fetch via adapter -> write to cache), deduped + backed off. const DRAIN_MS = Number(process.env.IFLOW_DRAIN_MS ?? 2000); const drainTimer = setInterval(() => { queue.drain().catch((e) => console.error('[drain error]', e)); }, DRAIN_MS); drainTimer.unref(); // Auto-scheduler: every 30s, enqueue refreshes for due schedules const SCHEDULE_MS = 30_000; const scheduleTimer = setInterval(() => { queue.enqueueDueSchedules().catch((e) => console.error('[schedule error]', e)); }, SCHEDULE_MS); scheduleTimer.unref(); // Realtime + per-fetch alert check: run alongside the schedule cycle for // low-latency alert types (e.g. VIX band crosses). Enqueues to the email // outbox only; SMTP is drained on its own timer so it never blocks this loop. const perFetchAlertTimer = setInterval(async () => { try { const alerts: Alert[] = [ ...(await runProducers(database, 'realtime')), ...(await runProducers(database, 'per-fetch')), ]; if (alerts.length > 0) { console.log(`[alert] ${alerts.length} realtime alert(s) created`); const { enqueueAlertEmail } = await import('./services/emailAlertService.ts'); alerts.forEach((a) => enqueueAlertEmail(database, a)); } } catch (e) { console.error('[alert] realtime check failed:', e); } }, SCHEDULE_MS); perFetchAlertTimer.unref(); // Register alert producers. import { registerProducer, runProducers, registerMirrorProducers } from './alerts/producers/index.ts'; import { informedBuyProducer, informedSellProducer } from './alerts/producers/insiderProducer.ts'; import { new13daProducer, new13fFilingProducer } from './alerts/producers/new13daProducer.ts'; import { vixLevelProducer } from './alerts/producers/vixLevelProducer.ts'; import { rotationIncipientProducer, regimeShiftProducer } from './alerts/producers/rotationProducer.ts'; import { convictionUnlockProducer } from './alerts/producers/unlockProducer.ts'; import { thesisBrokenProducer, thesisWeakeningProducer } from './alerts/producers/thesisProducer.ts'; import { clusterBreachProducer, drawdownHaltProducer, asymmetryWarningProducer } from './alerts/producers/portfolioRiskProducer.ts'; import { confluenceChangeProducer } from './alerts/producers/confluenceProducer.ts'; registerProducer(informedBuyProducer); registerProducer(informedSellProducer); registerProducer(new13daProducer); registerProducer(new13fFilingProducer); registerProducer(vixLevelProducer); registerProducer(rotationIncipientProducer); registerProducer(regimeShiftProducer); registerProducer(convictionUnlockProducer); registerProducer(thesisBrokenProducer); registerProducer(thesisWeakeningProducer); registerProducer(clusterBreachProducer); registerProducer(drawdownHaltProducer); registerProducer(asymmetryWarningProducer); registerProducer(confluenceChangeProducer); registerMirrorProducers(); // fund_capture / fund_13f / mirror_diff (lazy, offline-safe) // Batched alert check: run every 10 minutes for non-critical producers. 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`); // Enqueue email for each new alert (SMTP drained separately). const { enqueueAlertEmail } = await import('./services/emailAlertService.ts'); for (const alert of alerts) enqueueAlertEmail(database, alert); } catch (e) { console.error('[alert] batch check failed:', e); } }, ALERT_BATCH_MS); alertBatchTimer.unref(); // Email outbox drain: decouples SMTP latency/errors from the alert loops. const OUTBOX_DRAIN_MS = 60_000; const outboxDrainTimer = setInterval(async () => { try { const { drainEmailOutbox } = await import('./services/emailAlertService.ts'); const sent = await drainEmailOutbox(database); if (sent > 0) console.log(`[alert:email] drained ${sent} outbox email(s)`); } catch (e) { console.error('[alert:email] outbox drain failed:', e); } }, OUTBOX_DRAIN_MS); outboxDrainTimer.unref(); // Daily queue housekeeping: purge stale 'done' rows and age out queue_errors. // Both tables previously grew unbounded (live DB had 819k error rows). const DAILY_HOUSEKEEP_MS = 24 * 60 * 60 * 1000; const housekeepTimer = setInterval(() => { try { const cleared = queue.clearDone(); const pruned = queue.pruneQueueErrors(); if (cleared > 0 || pruned > 0) { console.log(`[queue] housekeeping: cleared ${cleared} done jobs, pruned ${pruned} old queue_errors`); } } catch (e) { console.error('[queue] daily housekeeping failed:', e); } }, DAILY_HOUSEKEEP_MS); housekeepTimer.unref(); function readBody(req: IncomingMessage): Promise { return new Promise((resolve, reject) => { let data = ''; req.on('data', (c: Buffer) => { data += c.toString(); }); req.on('end', () => resolve(data)); req.on('error', reject); }); } const server = createServer(async (req, res) => { const url = new URL(req.url ?? '/', `http://localhost:${PORT}`); if (url.pathname === '/health') { res.writeHead(200, { 'content-type': 'application/json' }); res.end(JSON.stringify({ ok: true, queue: queue.health() })); return; } if (url.pathname.startsWith('/api/trpc')) { const headers = new Headers(); for (const [k, v] of Object.entries(req.headers)) if (v != null) headers.set(k, Array.isArray(v) ? v.join(', ') : String(v)); const method = req.method ?? 'GET'; const body = method === 'GET' || method === 'HEAD' ? undefined : await readBody(req); const requestUrl = `http://localhost:${PORT}${url.pathname}${url.search}`; const init: RequestInit = { method, headers }; if (body) { init.body = new Blob([body], { type: 'application/json' }); } const request = new Request(requestUrl, init); try { const response = await fetchRequestHandler({ router: appRouter, createContext, endpoint: '/api/trpc', req: request }); const buf = Buffer.from(await response.arrayBuffer()); res.writeHead(response.status, Object.fromEntries(response.headers.entries())); res.end(buf); } catch (e) { console.error('[trpc error]', e); res.writeHead(500, { 'content-type': 'application/json' }); res.end(JSON.stringify({ error: 'tRPC handler error' })); } return; } res.writeHead(404, { 'content-type': 'application/json' }); res.end(JSON.stringify({ error: 'not found' })); }); const HOST = process.env.HOST ?? '0.0.0.0'; server.listen(PORT, HOST, () => { console.log(`[investor-flow] backend on http://${HOST}:${PORT} (tRPC /api/trpc, health /health, drain every ${DRAIN_MS}ms)`); });