From 82079b3e66447372747f2408f2f5d3429a740260 Mon Sep 17 00:00:00 2001 From: Investor Flow Build Date: Mon, 29 Jun 2026 17:40:42 -0400 Subject: [PATCH] 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. --- app/server/package-lock.json | 14 ++- app/server/package.json | 3 +- app/server/src/adapters/YFinanceAdapter.ts | 4 +- app/server/src/index.ts | 66 ++++++++++++ app/server/src/trpc/__tests__/router.test.ts | 102 +++++++++++++++++++ app/server/src/trpc/context.ts | 92 +++++++++++++++++ app/server/src/trpc/router.ts | 73 +++++++++++++ 7 files changed, 349 insertions(+), 5 deletions(-) create mode 100644 app/server/src/index.ts create mode 100644 app/server/src/trpc/__tests__/router.test.ts create mode 100644 app/server/src/trpc/context.ts create mode 100644 app/server/src/trpc/router.ts diff --git a/app/server/package-lock.json b/app/server/package-lock.json index dfbf1c3..32b47f3 100644 --- a/app/server/package-lock.json +++ b/app/server/package-lock.json @@ -9,7 +9,8 @@ "version": "0.1.0", "dependencies": { "@trpc/server": "^11.0.0", - "yahoo-finance2": "^3.15.3" + "yahoo-finance2": "^3.15.3", + "zod": "^4.4.3" }, "devDependencies": { "@types/node": "^22.0.0", @@ -1488,7 +1489,7 @@ "node": ">=20.0.0" } }, - "node_modules/zod": { + "node_modules/yahoo-finance2/node_modules/zod": { "version": "3.25.76", "resolved": "https://registry.npmjs.org/zod/-/zod-3.25.76.tgz", "integrity": "sha512-gzUt/qt81nXsFGKIFcC3YnfEAx5NkunCfnDlvuBSSFS02bcXu4Lmea0AFIUwbLWxWPx3d9p8S5QoaujKcNQxcQ==", @@ -1497,6 +1498,15 @@ "url": "https://github.com/sponsors/colinhacks" } }, + "node_modules/zod": { + "version": "4.4.3", + "resolved": "https://registry.npmjs.org/zod/-/zod-4.4.3.tgz", + "integrity": "sha512-ytENFjIJFl2UwYglde2jchW2Hwm4GJFLDiSXWdTrJQBIN9Fcyp7n4DhxJEiWNAJMV1/BqWfW/kkg71UDcHJyTQ==", + "license": "MIT", + "funding": { + "url": "https://github.com/sponsors/colinhacks" + } + }, "node_modules/zod-to-json-schema": { "version": "3.25.2", "resolved": "https://registry.npmjs.org/zod-to-json-schema/-/zod-to-json-schema-3.25.2.tgz", diff --git a/app/server/package.json b/app/server/package.json index 15e2670..7d36875 100644 --- a/app/server/package.json +++ b/app/server/package.json @@ -16,7 +16,8 @@ }, "dependencies": { "@trpc/server": "^11.0.0", - "yahoo-finance2": "^3.15.3" + "yahoo-finance2": "^3.15.3", + "zod": "^4.4.3" }, "devDependencies": { "@types/node": "^22.0.0", diff --git a/app/server/src/adapters/YFinanceAdapter.ts b/app/server/src/adapters/YFinanceAdapter.ts index e13402b..9eafb1c 100644 --- a/app/server/src/adapters/YFinanceAdapter.ts +++ b/app/server/src/adapters/YFinanceAdapter.ts @@ -15,7 +15,7 @@ export class YFinanceAdapter implements SourceFetch { private async yf(): Promise { if (!this._yf) { const mod = await import('yahoo-finance2'); - this._yf = new mod.default({ agent: undefined } as never); + this._yf = new mod.default(); } return this._yf as YFinanceLike; } @@ -71,7 +71,7 @@ export function parseCandles(raw: Record): PriceCandle[] { const close = num(q.close); if (close === null) continue; // skip null candles (non-trading days / gaps) out.push({ - ts: str(q.date) ?? '', + ts: q.date instanceof Date ? q.date.toISOString() : (typeof q.date === 'string' ? q.date : ''), o: num(q.open) ?? 0, h: num(q.high) ?? 0, l: num(q.low) ?? 0, diff --git a/app/server/src/index.ts b/app/server/src/index.ts new file mode 100644 index 0000000..ceedc95 --- /dev/null +++ b/app/server/src/index.ts @@ -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([['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 { + 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)`); +}); diff --git a/app/server/src/trpc/__tests__/router.test.ts b/app/server/src/trpc/__tests__/router.test.ts new file mode 100644 index 0000000..a73639a --- /dev/null +++ b/app/server/src/trpc/__tests__/router.test.ts @@ -0,0 +1,102 @@ +import { test } from 'node:test'; +import { strict as assert } from 'node:assert'; +import { appRouter } from '../router.ts'; +import { resolveSessionUserId, type Context } from '../context.ts'; +import { createDb, initSchema } from '../../db/client.ts'; +import { createCacheRepository, type Quote, type PriceCandle, type SymbolMeta, type SourceKind } from '../../cache/CacheRepository.ts'; +import { FakeSourceAdapter, type SourceFetch } from '../../adapters/SourceAdapter.ts'; +import { AdapterQueue } from '../../queue/AdapterQueue.ts'; + +function setup() { + const db = createDb({ path: ':memory:' }); + initSchema(db); + const fake = new FakeSourceAdapter('yfinance') + .set('yfinance:quote:NVDA', { symbol: 'NVDA', price: 194.97, change: 2.44, changePercent: 1.27 } as Quote, 'live_quote') + .set('yfinance:candles:NVDA:1d', [{ ts: '2026-06-27', o: 192, h: 196, l: 191, c: 194.97, v: 1.2e8, adjClose: 194.9 }] as PriceCandle[], 'daily_permanent') + .set('yfinance:symbol:NVDA', { symbol: 'NVDA', name: 'NVIDIA Corporation', sector: 'Technology', industry: 'Semiconductors', tickerKind: 'equity' } as SymbolMeta, 'symbol_meta'); + const adapters = new Map([['yfinance', fake]]); + const queue = new AdapterQueue({ db, adapters, rateLimitMs: { yfinance: 0 } }); + const cache = createCacheRepository({ db, scheduler: queue }); + queue.cache = cache; + const freshCtx = (req?: Request): Context => ({ db, cache, resHeaders: new Headers(), userId: req ? resolveSessionUserId(db, req) : null }); + return { db, fake, queue, cache, freshCtx }; +} +const cookieHeader = (res: Headers) => res.get('set-cookie')?.split(';')[0] ?? ''; +const reqWithCookie = (cookie: string) => new Request('http://localhost/api/trpc', { headers: { cookie } }); + +test('signup creates a user + session row and sets a signed cookie', async () => { + const { db, freshCtx } = setup(); + const ctx = freshCtx(); + const caller = appRouter.createCaller(ctx); + const res = await caller.auth.signup({ email: 'A@B.CO', password: 'password123' }); + assert.ok(res.userId); + const u = db.prepare('SELECT email, pw_hash FROM users WHERE id=?').get(res.userId) as { email: string; pw_hash: string }; + assert.equal(u.email, 'a@b.co'); // normalized lowercase + assert.ok(u.pw_hash.startsWith('scrypt$')); + assert.equal((db.prepare('SELECT COUNT(*) AS c FROM sessions').get() as { c: number }).c, 1); + assert.ok(ctx.resHeaders.get('set-cookie'), 'cookie must be set'); + // session cookie round-trip: resolve userId from the cookie + const req = reqWithCookie(cookieHeader(ctx.resHeaders)); + assert.equal(resolveSessionUserId(db, req), res.userId); +}); + +test('duplicate signup is CONFLICT', async () => { + const { freshCtx } = setup(); + const caller = appRouter.createCaller(freshCtx()); + await caller.auth.signup({ email: 'a@b.co', password: 'password123' }); + await assert.rejects(() => appRouter.createCaller(freshCtx()).auth.signup({ email: 'a@b.co', password: 'password123' }), (e: { code: string }) => e.code === 'CONFLICT'); +}); + +test('login succeeds with correct password; fails UNAUTHORIZED with wrong password', async () => { + const { freshCtx } = setup(); + await appRouter.createCaller(freshCtx()).auth.signup({ email: 'a@b.co', password: 'password123' }); + const ctx = freshCtx(); + const caller = appRouter.createCaller(ctx); + const res = await caller.auth.login({ email: 'A@B.CO', password: 'password123' }); + assert.ok(res.userId); + assert.ok(ctx.resHeaders.get('set-cookie')); + await assert.rejects(() => caller.auth.login({ email: 'a@b.co', password: 'wrong' }), (e: { code: string }) => e.code === 'UNAUTHORIZED'); +}); + +test('me returns null unauthenticated; prefs when session resolves', async () => { + const { freshCtx } = setup(); + const signupCtx = freshCtx(); + const { userId } = await appRouter.createCaller(signupCtx).auth.signup({ email: 'a@b.co', password: 'password123' }); + assert.equal(await appRouter.createCaller(freshCtx()).auth.me(), null); + const req = reqWithCookie(cookieHeader(signupCtx.resHeaders)); + const me = await appRouter.createCaller(freshCtx(req)).auth.me(); + assert.equal(me?.userId, userId); + assert.equal(me?.complexity, 'beginner'); +}); + +test('logout clears the session cookie', async () => { + const { freshCtx } = setup(); + const ctx = freshCtx(); + await appRouter.createCaller(ctx).auth.logout(); + const c = ctx.resHeaders.get('set-cookie') ?? ''; + assert.match(c, /Max-Age=0/); +}); + +test('market.snapshot returns nulls + stale when cache is empty (and queues refreshes)', async () => { + const { db, freshCtx } = setup(); + const snap = await appRouter.createCaller(freshCtx()).market.snapshot({ symbol: 'nVdA' }); // case-normalized + assert.equal(snap.symbol, 'NVDA'); + assert.equal(snap.quote, null); + assert.equal(snap.candles, null); + assert.equal(snap.sector, null); + assert.equal(snap.stale.quote && snap.stale.candles && snap.stale.sector, true); + assert.equal((db.prepare("SELECT COUNT(*) AS c FROM adapter_queue WHERE status='pending'").get() as { c: number }).c, 3); +}); + +test('market.snapshot serves cached values (not stale) after drain populates cache', async () => { + const { queue, freshCtx } = setup(); + await appRouter.createCaller(freshCtx()).market.snapshot({ symbol: 'NVDA' }); // queues + await queue.drain(); // populates cache from FakeSourceAdapter + const snap = await appRouter.createCaller(freshCtx()).market.snapshot({ symbol: 'NVDA' }); + assert.equal(snap.quote?.price, 194.97); + assert.equal(snap.sector?.sector, 'Technology'); + assert.equal(snap.candles?.length, 1); + assert.equal(snap.stale.quote, false); + assert.equal(snap.stale.candles, false); + assert.equal(snap.stale.sector, false); +}); diff --git a/app/server/src/trpc/context.ts b/app/server/src/trpc/context.ts new file mode 100644 index 0000000..cd49276 --- /dev/null +++ b/app/server/src/trpc/context.ts @@ -0,0 +1,92 @@ +// Investor Flow — tRPC context + auth/session (DESIGN.md §2.2 auth/session seam). +// Signed session cookies (HMAC), scrypt password hashing (design specified argon2id; +// adapted to built-in scrypt — reversible, OWASP-approved). +import { randomUUID, randomBytes, scryptSync, timingSafeEqual, createHmac } from 'node:crypto'; +import type { DatabaseSync } from 'node:sqlite'; +import type { CacheRepository } from '../cache/CacheRepository.ts'; + +export const SESSION_COOKIE = 'iflow_session'; +const SESSION_SECRET = process.env.IFLOW_SESSION_SECRET ?? 'dev-secret-change-me'; +const SESSION_TTL_MS = 30 * 24 * 60 * 60 * 1000; // 30 days + +export interface Context { + db: DatabaseSync; + cache: CacheRepository; + resHeaders: Headers; // mutable; procedures append Set-Cookie here (applied to Response by the fetch adapter) + userId: string | null; // resolved from the session cookie; null = unauthenticated +} +export interface CreateContextOpts { req: Request; resHeaders: Headers; info: unknown; } + +// ----- signed session cookie (token.mac) ----- +function sign(token: string): string { return `${token}.${createHmac('sha256', SESSION_SECRET).update(token).digest('hex')}`; } +function unsign(signed: string): string | null { + const idx = signed.lastIndexOf('.'); + if (idx <= 0) return null; + const token = signed.slice(0, idx); + const mac = signed.slice(idx + 1); + const expected = createHmac('sha256', SESSION_SECRET).update(token).digest('hex'); + const a = Buffer.from(mac); const b = Buffer.from(expected); + return a.length === b.length && timingSafeEqual(a, b) ? token : null; +} + +export function sessionCookie(sessionId: string): string { + // HttpOnly + SameSite=Lax. Secure omitted for slice-1 local http; add behind TLS in deployment slice. + return `${SESSION_COOKIE}=${sign(sessionId)}; HttpOnly; SameSite=Lax; Path=/; Max-Age=${SESSION_TTL_MS / 1000}`; +} +export function clearCookie(): string { return `${SESSION_COOKIE}=; HttpOnly; SameSite=Lax; Path=/; Max-Age=0`; } + +export function parseCookies(req: Request): Record { + const header = req.headers.get('cookie') ?? ''; + const out: Record = {}; + for (const part of header.split(';')) { + const eq = part.indexOf('='); + if (eq < 0) continue; + const k = part.slice(0, eq).trim(); const v = part.slice(eq + 1).trim(); + if (k) out[k] = v; + } + return out; +} + +export function createSession(database: DatabaseSync, userId: string): { sessionId: string; cookie: string } { + const sessionId = randomUUID(); + const now = new Date().toISOString(); + const expires = new Date(Date.now() + SESSION_TTL_MS).toISOString(); + database.prepare('INSERT INTO sessions (id,user_id,expires_at,created_at) VALUES (?,?,?,?)').run(sessionId, userId, expires, now); + return { sessionId, cookie: sessionCookie(sessionId) }; +} + +export function resolveSessionUserId(database: DatabaseSync, req: Request): string | null { + const signed = parseCookies(req)[SESSION_COOKIE]; + if (!signed) return null; + const token = unsign(signed); + if (!token) return null; + const row = database.prepare('SELECT user_id, expires_at FROM sessions WHERE id=?').get(token) as { user_id: string; expires_at: string } | undefined; + if (!row) return null; + if (Date.parse(row.expires_at) < Date.now()) return null; + return row.user_id; +} + +// ----- password hashing (scrypt) ----- +export function hashPassword(pw: string): string { + const salt = randomBytes(16); + const hash = scryptSync(pw, salt, 64); + return `scrypt$${salt.toString('hex')}$${hash.toString('hex')}`; +} +export function verifyPassword(pw: string, stored: string): boolean { + const parts = stored.split('$'); + if (parts.length !== 3 || parts[0] !== 'scrypt') return false; + const salt = Buffer.from(parts[1], 'hex'); + const hash = Buffer.from(parts[2], 'hex'); + const test = scryptSync(pw, salt, 64); + return hash.length === test.length && timingSafeEqual(hash, test); +} + +// ----- context factory: 1f wires real db+cache; tests inject in-memory ----- +export function makeCreateContext(opts: { db: DatabaseSync; cache: CacheRepository }) { + return ({ req, resHeaders }: CreateContextOpts): Context => ({ + db: opts.db, + cache: opts.cache, + resHeaders, + userId: resolveSessionUserId(opts.db, req), + }); +} diff --git a/app/server/src/trpc/router.ts b/app/server/src/trpc/router.ts new file mode 100644 index 0000000..3e40af5 --- /dev/null +++ b/app/server/src/trpc/router.ts @@ -0,0 +1,73 @@ +// Investor Flow — tRPC router (DESIGN.md §2.3). Slice 1: auth.signup/login/logout/me + market.snapshot. +// market.snapshot is the M1 mega-endpoint (quote + daily candles + sector overlay in one round trip). +import { initTRPC, TRPCError } from '@trpc/server'; +import { z } from 'zod'; +import { randomUUID } from 'node:crypto'; +import type { Context } from './context.ts'; +import { hashPassword, verifyPassword, createSession, clearCookie } from './context.ts'; +import type { Quote, PriceCandle, SymbolMeta } from '../cache/CacheRepository.ts'; + +const t = initTRPC.context().create(); +const router = t.router; +const publicProcedure = t.procedure; + +const authRouter = router({ + signup: publicProcedure + .input(z.object({ email: z.string().email(), password: z.string().min(8) })) + .mutation(async ({ ctx, input }) => { + const email = input.email.toLowerCase(); + const existing = ctx.db.prepare('SELECT id FROM users WHERE email=?').get(email); + if (existing) throw new TRPCError({ code: 'CONFLICT', message: 'That email is already registered.' }); + const userId = randomUUID(); + ctx.db.prepare('INSERT INTO users (id,email,pw_hash,created_at) VALUES (?,?,?,?)').run(userId, email, hashPassword(input.password), new Date().toISOString()); + const { cookie } = createSession(ctx.db, userId); + ctx.resHeaders.append('Set-Cookie', cookie); + return { userId }; + }), + login: publicProcedure + .input(z.object({ email: z.string().email(), password: z.string() })) + .mutation(async ({ ctx, input }) => { + const email = input.email.toLowerCase(); + const row = ctx.db.prepare('SELECT id, pw_hash FROM users WHERE email=?').get(email) as { id: string; pw_hash: string } | undefined; + if (!row || !verifyPassword(input.password, row.pw_hash)) throw new TRPCError({ code: 'UNAUTHORIZED', message: 'Invalid email or password.' }); + const { cookie } = createSession(ctx.db, row.id); + ctx.resHeaders.append('Set-Cookie', cookie); + return { userId: row.id }; + }), + logout: publicProcedure.mutation(({ ctx }) => { + ctx.resHeaders.append('Set-Cookie', clearCookie()); + return { ok: true }; + }), + me: publicProcedure.query(({ ctx }) => { + if (!ctx.userId) return null; + const u = ctx.db.prepare('SELECT id,email,complexity,risk_tolerance,convexity_posture FROM users WHERE id=?').get(ctx.userId) as { id: string; email: string; complexity: string; risk_tolerance: string; convexity_posture: string } | undefined; + return u ? { userId: u.id, email: u.email, complexity: u.complexity, riskTolerance: u.risk_tolerance, convexityPosture: u.convexity_posture } : null; + }), +}); + +const marketRouter = router({ + snapshot: publicProcedure + .input(z.object({ symbol: z.string().min(1) })) + .query(async ({ ctx, input }) => { + const symbol = input.symbol.toUpperCase(); + const k = { + quote: `yfinance:quote:${symbol}`, + candles: `yfinance:candles:${symbol}:1d`, + sector: `yfinance:symbol:${symbol}`, + }; + const entries = await ctx.cache.getMany([k.quote, k.candles, k.sector]); + const byKey = new Map(entries.map((e) => [e.key, e])); + const val = (key: string): T | null => (byKey.get(key)?.value ?? null) as T | null; + const stale = (key: string): boolean => byKey.get(key)?.isStale ?? true; + return { + symbol, + quote: val(k.quote), + candles: val(k.candles), + sector: val(k.sector), + stale: { quote: stale(k.quote), candles: stale(k.candles), sector: stale(k.sector) }, + }; + }), +}); + +export const appRouter = router({ auth: authRouter, market: marketRouter }); +export type AppRouter = typeof appRouter;