Files
investor-flow/app/server/src/index.ts
T

158 lines
7.7 KiB
TypeScript
Raw Normal View History

// 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 { NasdaqAdapter } from './adapters/NasdaqAdapter.ts';
import { SecFetchAdapter } from './adapters/SecFetchAdapter.ts';
import { SecLintAdapter } from './adapters/SecLintAdapter.ts';
import { XCookieAdapter } from './adapters/XCookieAdapter.ts';
import type { SourceFetch } from './adapters/SourceAdapter.ts';
import cryptoMod from './lib/crypto.ts';
import { AdapterQueue } from './queue/AdapterQueue.ts';
import { makeCreateContext } from './trpc/context.ts';
import { appRouter } from './trpc/router.ts';
const PORT = Number(process.env.PORT ?? 3001);
const database = db();
const adapters = new Map<SourceKind, SourceFetch>([
['yfinance' as const, new YFinanceAdapter() as unknown as SourceFetch],
['nasdaq' as const, new NasdaqAdapter() as unknown as SourceFetch],
['sec-fetch' as const, new SecFetchAdapter(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)');
})();
const queue = new AdapterQueue({ db: database, adapters });
const cache = createCacheRepository({ db: database, scheduler: queue });
queue.cache = cache; // break the cache<->scheduler cycle
// Seed default schedules (noop if already seeded)
queue.seedDefaultSchedules();
// 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();
// 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 = '';
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)`);
});