220 lines
11 KiB
TypeScript
220 lines
11 KiB
TypeScript
// 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' | 'reddit' | 'x' | 'macro' | 'llm';
|
||
|
|
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';
|
||
|
|
|
||
|
|
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; }
|
||
|
|
|
||
|
|
/** Port CacheRepository depends on to schedule background refreshes. SourceAdapter/AdapterQueue satisfy this. */
|
||
|
|
export interface CacheScheduler { queue(key: CacheKey): Promise<void>; }
|
||
|
|
|
||
|
|
export interface CacheEntry<T> { value: T | null; provenance: Provenance | null; isStale: boolean; }
|
||
|
|
|
||
|
|
export const TTL_MS: Record<TtlClass, number> = {
|
||
|
|
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)
|
||
|
|
};
|
||
|
|
|
||
|
|
/** 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<string, unknown> | 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<Record<string, unknown>>;
|
||
|
|
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 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<string, unknown> | 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; },
|
||
|
|
};
|
||
|
|
|
||
|
|
const HANDLERS = new Map<string, KindHandler>([
|
||
|
|
['quote', quoteHandler],
|
||
|
|
['candles', candlesHandler],
|
||
|
|
['symbol', symbolHandler],
|
||
|
|
]);
|
||
|
|
|
||
|
|
export interface CacheRepository {
|
||
|
|
get<T>(key: CacheKey): Promise<CacheEntry<T>>;
|
||
|
|
set<T>(key: CacheKey, value: T, ttlClass: TtlClass, provenance: Provenance): Promise<void>;
|
||
|
|
stale(key: CacheKey): boolean;
|
||
|
|
subscribe(symbol: string, tickerKind: TickerKind): Promise<void>;
|
||
|
|
unsubscribe(symbol: string): Promise<void>;
|
||
|
|
demandSet(): Promise<string[]>;
|
||
|
|
getMany<T>(keys: CacheKey[]): Promise<Array<{ key: CacheKey; value: T | null; isStale: boolean }>>;
|
||
|
|
}
|
||
|
|
|
||
|
|
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<T>(key: CacheKey): Promise<CacheEntry<T>> {
|
||
|
|
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<T>(key: CacheKey, value: T, ttlClass: TtlClass, provenance: Provenance): Promise<void> {
|
||
|
|
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<void> {
|
||
|
|
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)
|
||
|
|
for (const k of [`yfinance:quote:${symbol}`, `yfinance:symbol:${symbol}`]) {
|
||
|
|
try { await this._scheduler.queue(k); } catch { /* ignore */ }
|
||
|
|
}
|
||
|
|
}
|
||
|
|
}
|
||
|
|
async unsubscribe(symbol: string): Promise<void> {
|
||
|
|
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<string[]> {
|
||
|
|
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<T>(keys: CacheKey[]): Promise<Array<{ key: CacheKey; value: T | null; isStale: boolean }>> {
|
||
|
|
return Promise.all(keys.map(async (key) => {
|
||
|
|
const e = await this.get<T>(key);
|
||
|
|
return { key, value: e.value, isStale: e.isStale };
|
||
|
|
}));
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
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;
|
||
|
|
}
|