// Investor Flow — CacheRepository deep module (DESIGN.md §3a Module 1). // The ONLY interface the SPA (via tRPC) touches for cached data. Owns staleness windows, // refcount/demand-set, and stale-while-revalidate. Does NOT talk to external sources — // that is SourceAdapter's job; CacheRepository only schedules background refreshes via // the injected scheduler (SourceAdapter/AdapterQueue satisfy `CacheScheduler`). import { DatabaseSync } from 'node:sqlite'; import { db as defaultDb } from '../db/client.ts'; export type SourceKind = 'yfinance' | 'sec' | 'sec-fetch' | 'reddit' | 'x' | 'macro' | 'llm' | 'sec-lint-holders' | 'sec-lint-insiders' | 'fred'; export type TickerKind = 'equity' | 'crypto' | 'etf' | 'index'; export type CacheKey = string; // `${SourceKind}:${kind}:${id}` e.g. 'yfinance:quote:NVDA', 'yfinance:candles:NVDA:1d' export type TtlClass = | 'live_quote' | 'intraday' | 'daily_permanent' | 'options_snapshot' | 'filing_immutable' | 'quarterly_immutable' | 'thread_7d' | 'macro_event' | 'regime_classification' | 'llm_summary' | 'symbol_meta' | 'short_interest'; export interface Provenance { fetchedAt: string; sourceKind: SourceKind; rawSourceId?: string; } export interface Quote { symbol: string; price: number; bid?: number | null; ask?: number | null; change?: number | null; changePercent?: number | null; iv?: number | null; } export interface PriceCandle { ts: string; o: number; h: number; l: number; c: number; v: number; adjClose?: number | null; } export interface SymbolMeta { symbol: string; name?: string | null; sector?: string | null; industry?: string | null; exchange?: string | null; tickerKind: TickerKind; peers?: string[] | null; description?: string | null; website?: string | null; marketCap?: number | null; currency?: string | null; employees?: number | null; country?: string | null; city?: string | null; } export interface PriceAdjustment { symbol: string; exDate: string; type: "split" | "dividend"; ratio: number } /** Port CacheRepository depends on to schedule background refreshes. SourceAdapter/AdapterQueue satisfy this. */ export interface CacheScheduler { queue(key: CacheKey): Promise; } export interface CacheEntry { value: T | null; provenance: Provenance | null; isStale: boolean; } export const TTL_MS: Record = { live_quote: 60_000, // 1min (mkt hrs); after-hours 15min refined in adapter slice intraday: 5 * 60_000, daily_permanent: Number.POSITIVE_INFINITY, // immutable once fetched; stale only when absent options_snapshot: 15 * 60_000, filing_immutable: Number.POSITIVE_INFINITY, quarterly_immutable: Number.POSITIVE_INFINITY, thread_7d: 7 * 24 * 60 * 60_000, macro_event: Number.POSITIVE_INFINITY, regime_classification: 24 * 60 * 60_000, llm_summary: Number.POSITIVE_INFINITY, // cached forever by prompt-hash symbol_meta: 7 * 24 * 60 * 60_000, // weekly (slow-changing sector/industry) short_interest: 24 * 60 * 60_000, // refreshed twice/month per source }; /** Parse 'source:kind:id...' into { source, kind, id } (id may contain colons). */ export function parseCacheKey(key: CacheKey): { source: SourceKind; kind: string; id: string } { const parts = key.split(':'); if (parts.length < 3) throw new Error(`invalid cache key: ${key}`); const source = parts[0] as SourceKind; const kind = parts[1]; const id = parts.slice(2).join(':'); return { source, kind, id }; } interface KindHandler { ttlClass: TtlClass; /** Read cached value + the timestamp to judge staleness against. null if not cached. */ read(d: DatabaseSync, id: string): { value: unknown; stalenessTs: string } | null; /** Write value to cache, stamping provenance. */ write(d: DatabaseSync, id: string, value: unknown, provenance: Provenance): void; /** Staleness verdict given the staleness timestamp (null = never cached) and now. */ isStale(stalenessTs: string | null, now: number): boolean; } function tsAgeMs(ts: string | null, now: number): number { if (!ts) return Number.POSITIVE_INFINITY; const t = Date.parse(ts); if (Number.isNaN(t)) return Number.POSITIVE_INFINITY; return now - t; } // ----- Kind handlers (slice 1: quote, candles, symbol). Later slices register more. ----- const quoteHandler: KindHandler = { ttlClass: 'live_quote', read(d, symbol) { const r = d.prepare('SELECT symbol,price,bid,ask,change,change_pct,iv,observed_at FROM quotes WHERE symbol=?').get(symbol) as Record | undefined; if (!r) return null; return { value: { symbol: r.symbol, price: r.price, bid: r.bid, ask: r.ask, change: r.change, changePercent: r.change_pct, iv: r.iv } as Quote, stalenessTs: r.observed_at as string, }; }, write(d, symbol, value, provenance) { const q = value as Quote; d.prepare('INSERT OR REPLACE INTO quotes (symbol,price,bid,ask,change,change_pct,iv,observed_at) VALUES (?,?,?,?,?,?,?,?)') .run(symbol, q.price, q.bid ?? null, q.ask ?? null, q.change ?? null, q.changePercent ?? null, q.iv ?? null, provenance.fetchedAt); }, isStale(ts, now) { return tsAgeMs(ts, now) > TTL_MS.live_quote; }, }; const candlesHandler: KindHandler = { ttlClass: 'daily_permanent', read(d, id) { const [symbol, timeframe] = id.split(':'); if (!timeframe) return null; const rows = d.prepare('SELECT ts,o,h,l,c,v,adj_close,observed_at FROM price_candles WHERE symbol=? AND timeframe=? ORDER BY ts ASC').all(symbol, timeframe) as Array>; if (!rows.length) return null; const value: PriceCandle[] = rows.map((r) => ({ ts: r.ts as string, o: r.o as number, h: r.h as number, l: r.l as number, c: r.c as number, v: r.v as number, adjClose: r.adj_close as number | null })); return { value, stalenessTs: rows[rows.length - 1].observed_at as string }; }, write(d, id, value, provenance) { const [symbol, timeframe] = id.split(':'); const ins = d.prepare('INSERT OR REPLACE INTO price_candles (symbol,timeframe,ts,o,h,l,c,v,adj_close,observed_at) VALUES (?,?,?,?,?,?,?,?,?,?)'); for (const c of value as PriceCandle[]) ins.run(symbol, timeframe, c.ts, c.o, c.h, c.l, c.c, c.v, c.adjClose ?? null, provenance.fetchedAt); }, isStale(ts) { return ts === null; }, // permanent: stale only when absent }; const adjustmentsHandler: KindHandler = { ttlClass: 'daily_permanent', read(d, symbol) { const rows = d.prepare('SELECT ex_date, type, ratio FROM price_adjustments WHERE symbol=? ORDER BY ex_date ASC').all(symbol) as Array>; if (!rows.length) return null; const value: PriceAdjustment[] = rows.map((r) => ({ symbol, exDate: r.ex_date as string, type: r.type as "split" | "dividend", ratio: r.ratio as number })); return { value, stalenessTs: rows[rows.length - 1].ex_date as string }; }, write(d, _symbol, value, provenance) { const ins = d.prepare('INSERT OR REPLACE INTO price_adjustments (symbol, ex_date, type, ratio) VALUES (?,?,?,?)'); for (const a of value as PriceAdjustment[]) ins.run(a.symbol, a.exDate, a.type, a.ratio); }, isStale(ts) { return ts === null; }, // permanent: stale only when absent }; const symbolHandler: KindHandler = { ttlClass: 'symbol_meta', read(d, symbol) { const r = d.prepare('SELECT symbol,name,sector,industry,exchange,ticker_kind,peers,updated_at FROM symbols WHERE symbol=?').get(symbol) as Record | undefined; if (!r) return null; let peers: string[] | null = null; if (typeof r.peers === 'string') { try { peers = JSON.parse(r.peers); } catch { peers = null; } } return { value: { symbol: r.symbol, name: r.name, sector: r.sector, industry: r.industry, exchange: r.exchange, tickerKind: r.ticker_kind, peers } as SymbolMeta, stalenessTs: r.updated_at as string, }; }, write(d, symbol, value, provenance) { const s = value as SymbolMeta; d.prepare('INSERT OR REPLACE INTO symbols (symbol,name,sector,industry,exchange,ticker_kind,peers,updated_at) VALUES (?,?,?,?,?,?,?,?)') .run(symbol, s.name ?? null, s.sector ?? null, s.industry ?? null, s.exchange ?? null, s.tickerKind, s.peers ? JSON.stringify(s.peers) : null, provenance.fetchedAt); }, isStale(ts, now) { return tsAgeMs(ts, now) > TTL_MS.symbol_meta; }, }; // ----- Options handlers (slice 15) ----- const optionsChainHandler: KindHandler = { ttlClass: 'options_snapshot', read(d, id) { const [symbol, expiry] = id.split(':'); if (!expiry) return null; const rows = d.prepare( 'SELECT symbol,expiry,strike,type,bid,ask,iv,delta,gamma,theta,vega,open_interest,volume,ts FROM options_chains WHERE symbol=? AND expiry=? ORDER BY strike ASC, type ASC' ).all(symbol, expiry) as Array>; if (!rows.length) return null; const value = rows.map((r) => ({ contractSymbol: `${r.symbol}_${r.expiry}_${r.strike}_${r.type}`, strike: r.strike as number, right: r.type as 'call' | 'put', expiration: r.expiry as string, bid: r.bid as number | null, ask: r.ask as number | null, impliedVolatility: r.iv as number | null, delta: r.delta as number | null, gamma: r.gamma as number | null, theta: r.theta as number | null, vega: r.vega as number | null, openInterest: r.open_interest as number | null, volume: r.volume as number | null, })); return { value, stalenessTs: rows[rows.length - 1].ts as string }; }, write(d, id, value, provenance) { const [symbol, expiry] = id.split(':'); if (!expiry) return; const rows = (value as Array>); const ins = d.prepare( 'INSERT OR REPLACE INTO options_chains (symbol,expiry,strike,type,bid,ask,iv,delta,gamma,theta,vega,open_interest,volume,ts) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?)' ); for (const r of rows) { const strike = typeof r.strike === 'number' ? r.strike : 0; const right = r.right === 'put' ? 'put' : 'call'; ins.run( symbol, expiry, strike, right, numOrNull(r.bid), numOrNull(r.ask), numOrNull(r.impliedVolatility), numOrNull(r.delta), numOrNull(r.gamma), numOrNull(r.theta), numOrNull(r.vega), numOrNull(r.openInterest), numOrNull(r.volume), provenance.fetchedAt ); } }, isStale(ts, now) { return tsAgeMs(ts, now) > TTL_MS.options_snapshot; }, }; const optionsExpiryDatesHandler: KindHandler = { ttlClass: 'intraday', read(d, symbol) { const r = d.prepare('SELECT value, observed_at FROM kv_cache WHERE key=?').get(`options_expiry:${symbol}`) as Record | undefined; if (!r) return null; try { const value = JSON.parse(r.value as string); return { value, stalenessTs: r.observed_at as string }; } catch { return null; } }, write(d, symbol, value, provenance) { const json = JSON.stringify(value); d.prepare('INSERT OR REPLACE INTO kv_cache (key, value, observed_at) VALUES (?,?,?)') .run(`options_expiry:${symbol}`, json, provenance.fetchedAt); }, isStale(ts, now) { return tsAgeMs(ts, now) > TTL_MS.intraday; }, }; /** Coerce a possibly-undefined/unknown value to number | null for SQL binding. */ function numOrNull(v: unknown): number | null { return typeof v === 'number' ? v : null; } const greeksHandler: KindHandler = { ttlClass: 'options_snapshot', read(d, id) { const { symbol, expiry, strike } = parseGreeksId(id); const r = d.prepare( 'SELECT delta, gamma, theta, vega, strike, open_interest, iv, ts, type FROM options_chains WHERE symbol=? AND expiry=? AND strike=? LIMIT 1' ).get(symbol, expiry, parseFloat(strike ?? '0')) as Record | undefined; if (!r) return null; return { value: { delta: numOrNull(r.delta), gamma: numOrNull(r.gamma), theta: numOrNull(r.theta), vega: numOrNull(r.vega), strike: typeof r.strike === 'number' ? r.strike : 0, openInterest: numOrNull(r.open_interest), impliedVolatility: numOrNull(r.iv), right: r.type as 'call' | 'put', }, stalenessTs: r.ts as string, }; }, write(d, id, value, provenance) { const { symbol, expiry, strike } = parseGreeksId(id); const v = value as Record; const right = v.right === 'put' ? 'put' : 'call'; d.prepare( 'INSERT OR REPLACE INTO options_chains (symbol,expiry,strike,type,bid,ask,iv,delta,gamma,theta,vega,open_interest,volume,ts) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?)' ).run( symbol, expiry, parseFloat(strike ?? '0'), right, numOrNull(v.bid), numOrNull(v.ask), numOrNull(v.impliedVolatility), numOrNull(v.delta), numOrNull(v.gamma), numOrNull(v.theta), numOrNull(v.vega), numOrNull(v.openInterest), numOrNull(v.volume), provenance.fetchedAt ); }, isStale(ts, now) { return tsAgeMs(ts, now) > TTL_MS.options_snapshot; }, }; /** Parse 'symbol:expiry:strike' from greeks cache key id. */ function parseGreeksId(id: string): { symbol: string; expiry: string; strike: string } { const parts = id.split(':'); return { symbol: parts[0] ?? '', expiry: parts[1] ?? '', strike: parts[2] ?? '0' }; } const fetchHandler: KindHandler = { ttlClass: 'daily_permanent', read(d, id) { const r = d.prepare('SELECT value, observed_at FROM kv_cache WHERE key=?').get(`sec-fetch:${id}`) as { value: string; observed_at: string } | undefined; if (!r) return null; try { return { value: JSON.parse(r.value), stalenessTs: r.observed_at }; } catch { return null; } }, write(d, id, value, provenance) { d.prepare('INSERT OR REPLACE INTO kv_cache (key, value, observed_at) VALUES (?,?,?)').run(`sec-fetch:${id}`, JSON.stringify(value), provenance.fetchedAt); }, isStale(ts) { return ts === null; }, // never stale once written }; const lintHoldersHandler: KindHandler = { ttlClass: 'daily_permanent', read() { return null; }, // never read — work happens in DB tables directly write(d, key, _value, provenance) { d.prepare('INSERT OR REPLACE INTO kv_cache (key,value,observed_at) VALUES (?,?,?)').run(`sec-lint-holders:${key}`, JSON.stringify({ ok: true }), provenance.fetchedAt); }, isStale() { return false; }, // never stale once written (lint writes are permanent) }; const lintInsidersHandler: KindHandler = { ttlClass: 'daily_permanent', read() { return null; }, write(d, key, _value, provenance) { d.prepare('INSERT OR REPLACE INTO kv_cache (key,value,observed_at) VALUES (?,?,?)').run(`sec-lint-insiders:${key}`, JSON.stringify({ ok: true }), provenance.fetchedAt); }, isStale() { return false; }, }; const shortInterestHandler: KindHandler = { ttlClass: 'short_interest', read(d, id) { const r = d.prepare('SELECT value, observed_at FROM kv_cache WHERE key=?').get(`yfinance:shortinterest:${id}`) as { value: string; observed_at: string } | undefined; if (!r) return null; try { return { value: JSON.parse(r.value), stalenessTs: r.observed_at }; } catch { return null; } }, write(d, id, value, provenance) { d.prepare('INSERT OR REPLACE INTO kv_cache (key, value, observed_at) VALUES (?,?,?)').run(`yfinance:shortinterest:${id}`, JSON.stringify(value), provenance.fetchedAt); }, isStale(ts, now) { return tsAgeMs(ts, now) > TTL_MS.short_interest; }, }; const HANDLERS = new Map([ ['quote', quoteHandler], ['candles', candlesHandler], ['symbol', symbolHandler], ['adjustments', adjustmentsHandler], ['chain', optionsChainHandler], ['expiry_dates', optionsExpiryDatesHandler], ['greeks', greeksHandler], ['fetch', fetchHandler], ['holders', lintHoldersHandler], ['insiders', lintInsidersHandler], ['shortinterest', shortInterestHandler], ]); export interface CacheRepository { get(key: CacheKey): Promise>; set(key: CacheKey, value: T, ttlClass: TtlClass, provenance: Provenance): Promise; stale(key: CacheKey): boolean; subscribe(symbol: string, tickerKind: TickerKind): Promise; unsubscribe(symbol: string): Promise; demandSet(): Promise; getMany(keys: CacheKey[]): Promise>; /** Delete a cache entry by key (or, for wildcard keys ending in `:*`, all matching entries). */ del(key: CacheKey): Promise; } export class CacheRepositoryImpl implements CacheRepository { private readonly _db: DatabaseSync; private readonly _scheduler: CacheScheduler; constructor(opts: { db: DatabaseSync; scheduler: CacheScheduler }) { this._db = opts.db; this._scheduler = opts.scheduler; } private handler(kind: string): KindHandler { const h = HANDLERS.get(kind); if (!h) throw new Error(`unknown cache kind: ${kind}`); return h; } async get(key: CacheKey): Promise> { const { source, kind, id } = parseCacheKey(key); const h = this.handler(kind); const row = h.read(this._db, id); const now = Date.now(); const stale = h.isStale(row ? row.stalenessTs : null, now); if (stale) { try { await this._scheduler.queue(key); } catch { /* background refresh; never block readers */ } } return { value: (row ? row.value : null) as T | null, provenance: row ? { fetchedAt: row.stalenessTs, sourceKind: source } : null, isStale: stale, }; } async set(key: CacheKey, value: T, ttlClass: TtlClass, provenance: Provenance): Promise { const { kind, id } = parseCacheKey(key); const h = this.handler(kind); if (h.ttlClass !== ttlClass) throw new Error(`ttlClass mismatch for kind '${kind}': expected ${h.ttlClass}, got ${ttlClass}`); h.write(this._db, id, value, provenance); } stale(key: CacheKey): boolean { const { kind, id } = parseCacheKey(key); const h = this.handler(kind); const row = h.read(this._db, id); return h.isStale(row ? row.stalenessTs : null, Date.now()); } async subscribe(symbol: string, tickerKind: TickerKind): Promise { const d = this._db; d.prepare('INSERT OR IGNORE INTO symbol_demand (symbol,refcount,ticker_kind,in_demand,last_refreshed_at) VALUES (?,?,?,?,?)').run(symbol, 0, tickerKind, 1, null); const prev = d.prepare('SELECT refcount FROM symbol_demand WHERE symbol=?').get(symbol) as { refcount: number } | undefined; const before = prev?.refcount ?? 0; d.prepare('UPDATE symbol_demand SET refcount = refcount + 1, in_demand = 1 WHERE symbol=?').run(symbol); if (before === 0) { // First demand: schedule initial cache population (slice 1: yfinance quote + symbol meta + candles + adjustments) for (const k of [`yfinance:quote:${symbol}`, `yfinance:symbol:${symbol}`, `yfinance:candles:${symbol}:1d`, `yfinance:adjustments:${symbol}`]) { try { await this._scheduler.queue(k); } catch { /* ignore */ } } } } async unsubscribe(symbol: string): Promise { const d = this._db; d.prepare('UPDATE symbol_demand SET refcount = MAX(refcount - 1, 0) WHERE symbol=?').run(symbol); d.prepare('UPDATE symbol_demand SET in_demand = 0 WHERE symbol=? AND refcount = 0').run(symbol); } async demandSet(): Promise { return (this._db.prepare('SELECT symbol FROM symbol_demand WHERE refcount > 0 ORDER BY symbol').all() as Array<{ symbol: string }>).map((r) => r.symbol); } async getMany(keys: CacheKey[]): Promise> { return Promise.all(keys.map(async (key) => { const e = await this.get(key); return { key, value: e.value, isStale: e.isStale }; })); } async del(key: CacheKey): Promise { const { source, kind, id } = parseCacheKey(key); const d = this._db; switch (kind) { case 'quote': d.prepare('DELETE FROM quotes WHERE symbol=?').run(id); break; case 'candles': { const [symbol, tf] = id.split(':'); d.prepare('DELETE FROM price_candles WHERE symbol=? AND timeframe=?').run(symbol, tf); break; } case 'symbol': d.prepare('DELETE FROM symbols WHERE symbol=?').run(id); break; case 'adjustments': d.prepare('DELETE FROM price_adjustments WHERE symbol=?').run(id); break; case 'chain': case 'greeks': { const [symbol, expiry] = id.split(':'); d.prepare('DELETE FROM options_chains WHERE symbol=? AND expiry=?').run(symbol, expiry); break; } case 'expiry_dates': d.prepare('DELETE FROM kv_cache WHERE key=?').run(`options_expiry:${id}`); break; case 'shortinterest': d.prepare('DELETE FROM kv_cache WHERE key=?').run(`yfinance:shortinterest:${id}`); break; default: { // Unknown/wildcard kind (e.g. `x:cashtag:*`): best-effort delete from kv_cache via LIKE. const like = key.endsWith(':*') ? `${key.slice(0, -1)}%` : key; d.prepare('DELETE FROM kv_cache WHERE key LIKE ?').run(like); } } void source; } } export function createCacheRepository(opts: { db: DatabaseSync; scheduler: CacheScheduler }): CacheRepository { return new CacheRepositoryImpl(opts); } // Prod singleton — wired in slice 1f once AdapterQueue (the scheduler) exists. let _cache: CacheRepository | null = null; export function cacheRepository(scheduler: CacheScheduler): CacheRepository { if (!_cache) _cache = createCacheRepository({ db: defaultDb(), scheduler }); return _cache; }