slice 1e-1f: tRPC router (auth + market.snapshot) + node:http server

auth.signup/login/logout/me with signed HMAC session cookies + scrypt hashing;
market.snapshot mega-endpoint (quote+candles+sector). node:http server mounts
tRPC at /api/trpc + /health + background drain loop. Live-verified: real NVDA
194.97/65 candles/Technology served after stale-while-revalidate drain. 36 tests green.
This commit is contained in:
Investor Flow Build
2026-06-29 17:40:42 -04:00
parent 1962ecc740
commit 82079b3e66
7 changed files with 349 additions and 5 deletions
+66
View File
@@ -0,0 +1,66 @@
// 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 { 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, YFinanceAdapter>([['yfinance', new YFinanceAdapter()]]);
const queue = new AdapterQueue({ db: database, adapters });
const cache = createCacheRepository({ db: database, scheduler: queue });
queue.cache = cache; // break the cache<->scheduler cycle
const createContext = makeCreateContext({ db: database, cache });
// 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();
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 request = new Request(`http://localhost:${PORT}${url.pathname}${url.search}`, { method, headers, body });
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' }));
});
server.listen(PORT, () => {
console.log(`[investor-flow] backend on http://localhost:${PORT} (tRPC /api/trpc, health /health, drain every ${DRAIN_MS}ms)`);
});