slice 4b (omlx/ornith-35): CacheRepository adjustments kind + subscribe queues backfill keys

Local ornith-35 dispatch (~2min). +PriceAdjustment, +adjustmentsHandler
(price_adjustments table, permanent), subscribe now queues candles+adjustments on
first demand. Surgical, 99/99 tests, no regression.
This commit is contained in:
Investor Flow Build
2026-06-29 21:55:27 -04:00
parent 190e9500a9
commit 796d7b58b3
+19 -2
View File
@@ -19,6 +19,7 @@ export interface Provenance { fetchedAt: string; sourceKind: SourceKind; rawSour
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; }
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<void>; }
@@ -103,6 +104,21 @@ const candlesHandler: KindHandler = {
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<Record<string, unknown>>;
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) {
@@ -127,6 +143,7 @@ const HANDLERS = new Map<string, KindHandler>([
['quote', quoteHandler],
['candles', candlesHandler],
['symbol', symbolHandler],
['adjustments', adjustmentsHandler],
]);
export interface CacheRepository {
@@ -185,8 +202,8 @@ export class CacheRepositoryImpl implements CacheRepository {
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}`]) {
// 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 */ }
}
}