Confluence startup seed: pins the15-symbol research universe + SPY benchmark into the permanent demand set, and creates three system rack presets (Full Confluence, Technical Momentum, Macro+Flows+Sentiment) if absent. Called once from index.ts; idempotent on every restart.
285 lines
14 KiB
TypeScript
285 lines
14 KiB
TypeScript
// 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 { CotAdapter } from './adapters/CotAdapter.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 { seedConfluence } from './confluence/confluenceSeed.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<SourceKind, SourceFetch>([
|
|
['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],
|
|
['cot' as const, new CotAdapter() 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);
|
|
}
|
|
|
|
// Confluence starter seed: pin15-symbol universe + 3 system rack presets
|
|
try {
|
|
const { symbolsPinned, racksCreated } = await seedConfluence(database, cache);
|
|
if (symbolsPinned > 0 || racksCreated > 0) {
|
|
console.log(`[investor-flow] confluence seed: ${symbolsPinned} symbols pinned, ${racksCreated} system racks created`);
|
|
}
|
|
} catch (e) {
|
|
console.error('[investor-flow] confluence 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<string> {
|
|
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)`);
|
|
});
|