From 07d95ba6016e69bd4355f66c8ca9c1ff4f49cdf9 Mon Sep 17 00:00:00 2001 From: Investor Flow Build Date: Mon, 10 Aug 2026 17:08:33 -0400 Subject: [PATCH] feat(confluence): add candle resolution seam for confluence (M22 slice 8) CANDLE_PROVIDER resolves a symbol's candles per-granularity (1d/1wk) from the shared cache, with a reversible fold-in of the freshest live quote so a mid-session evaluation sees the current price. Pure fold logic + todayIso are unit-tested; CacheCandleProvider is a thin cache-backed shim. Evaluators use this seam instead of cache.get inline, so a future realtime/replay source can slot in without touching slot logic. --- .../__tests__/candleProvider.test.ts | 141 ++++++++++++++ app/server/src/confluence/candleProvider.ts | 176 ++++++++++++++++++ 2 files changed, 317 insertions(+) create mode 100644 app/server/src/confluence/__tests__/candleProvider.test.ts create mode 100644 app/server/src/confluence/candleProvider.ts diff --git a/app/server/src/confluence/__tests__/candleProvider.test.ts b/app/server/src/confluence/__tests__/candleProvider.test.ts new file mode 100644 index 0000000..94737a9 --- /dev/null +++ b/app/server/src/confluence/__tests__/candleProvider.test.ts @@ -0,0 +1,141 @@ +// Investor Flow — candleProvider.test.ts (M22 slice 8) +// Tests for the candle-resolution seam: pure fold-in logic + todayIso + the +// cache-backed provider against a fake CacheRepository. + +import { describe, it } from 'node:test'; +import assert from 'node:assert/strict'; + +import type { CacheRepository, CacheEntry, PriceCandle, Quote } from '../../cache/CacheRepository.ts'; +import { + CacheCandleProvider, + type CandleProvider, + foldRealtimeBar, + todayIso, +} from '../candleProvider.ts'; + +const bar = (ts: string, c: number, o = c, h = Math.max(o, c), l = Math.min(o, c), v = 1000): PriceCandle => + ({ ts, o, h, l, c, v, adjClose: c }); + +class FakeCache implements CacheRepository { + private readonly store = new Map(); + setValue(key: string, value: unknown, stale = false): this { this.store.set(key, { value, stale }); return this; } + async get(key: string): Promise> { + const e = this.store.get(key); + return { + value: (e ? e.value : null) as T | null, + provenance: null, + isStale: e ? e.stale : true, + }; + } + async set(): Promise { throw new Error('not used'); } + stale(key: string): boolean { return !this.store.has(key) || this.store.get(key)!.stale; } + async subscribe(): Promise {} + async unsubscribe(): Promise {} + async ensureInDemand(): Promise {} + async pinSystemSymbol(): Promise {} + async demandSet(): Promise { return []; } + async getMany(): Promise> { return []; } + async del(): Promise {} + readonly db: never = undefined as never; +} + +describe('foldRealtimeBar', () => { + const candles = [bar('2026-08-07', 100), bar('2026-08-10', 105)]; + + it('returns unchanged when there is no quote or no candles', () => { + const r1 = foldRealtimeBar(candles, null, '2026-08-11'); + assert.equal(r1.folded, false); + assert.equal(r1.candles, candles); + const r2 = foldRealtimeBar([], { symbol: 'X', price: 100 }, '2026-08-11'); + assert.equal(r2.folded, false); + assert.deepEqual(r2.candles, []); + }); + + it('appends a realtime bar when the quote is on a later day', () => { + const { candles: out, folded } = foldRealtimeBar( + candles, + { symbol: 'X', price: 108 }, + '2026-08-11', + ); + assert.equal(folded, true); + assert.equal(out.length, 3); + const last = out[out.length - 1]; + assert.equal(last.ts, '2026-08-11'); + assert.equal(last.c, 108); + assert.equal(last.o, 105); // prior close + assert.ok(last.h >= 108 && last.h >= 105); + assert.ok(last.l <= 108 && last.l <= 105); + assert.equal(last.v, 0); + }); + + it('replaces today’s bar close with the live price when the series already has today', () => { + const withToday = [bar('2026-08-07', 100), bar('2026-08-11', 110, 108, 112, 107)]; + const { candles: out, folded } = foldRealtimeBar(withToday, { symbol: 'X', price: 113 }, '2026-08-11'); + assert.equal(folded, true); + assert.equal(out.length, 2); + const last = out[out.length - 1]; + assert.equal(last.ts, '2026-08-11'); + assert.equal(last.c, 113); // live price wins + assert.equal(last.h, 113); // expanded to contain the print + assert.equal(last.o, 108); + }); + + it('ignores a quote not newer than the last bar (no duplicate bar)', () => { + const { candles: out, folded } = foldRealtimeBar(candles, { symbol: 'X', price: 90 }, '2026-08-09'); + assert.equal(folded, false); + assert.equal(out.length, 2); + }); +}); + +describe('todayIso', () => { + it('formats YYYY-MM-DD', () => { + const s = todayIso(new Date('2026-08-10T12:00:00Z')); + assert.match(s, /^\d{4}-\d{2}-\d{2}$/); + }); +}); + +describe('CacheCandleProvider', () => { + it('resolves stored candles with asOf = last bar ts and yfinance provenance', async () => { + const cache = new FakeCache() + .setValue('yfinance:candles:SPY:1d', [bar('2026-08-07', 100), bar('2026-08-10', 105)], false); + const provider: CandleProvider = new CacheCandleProvider(cache); + const r = await provider.resolve('spy', '1d'); + assert.equal(r.symbol, 'SPY'); + assert.equal(r.granularity, '1d'); + assert.equal(r.candles.length, 2); + assert.equal(r.asOf, '2026-08-10'); + assert.equal(r.lastBar, 'yfinance'); + assert.equal(r.realtimeFolded, false); + assert.equal(r.isStale, false); + }); + + it('folds a newer live quote into a 1d series', async () => { + const quote: Quote = { symbol: 'SPY', price: 110 }; + const cache = new FakeCache() + .setValue('yfinance:candles:SPY:1d', [bar('2026-08-07', 100), bar('2026-08-10', 105)], false) + .setValue('yfinance:quote:SPY', quote, false); + const r = await new CacheCandleProvider(cache).resolve('SPY', '1d'); + assert.equal(r.realtimeFolded, true); + assert.equal(r.lastBar, 'realtime'); + assert.equal(r.asOf, todayIso()); + assert.equal(r.candles[r.candles.length - 1].c, 110); + }); + + it('does not fold a quote for weekly granularity', async () => { + const cache = new FakeCache() + .setValue('yfinance:candles:SPY:1wk', [bar('2026-08-07', 500)], false) + .setValue('yfinance:quote:SPY', { symbol: 'SPY', price: 520 }, false); + const r = await new CacheCandleProvider(cache).resolve('SPY', '1wk'); + assert.equal(r.realtimeFolded, false); + assert.equal(r.lastBar, 'yfinance'); + assert.equal(r.candles.length, 1); + }); + + it('flags staleness when candles are absent', async () => { + const cache = new FakeCache(); // nothing stored + const r = await new CacheCandleProvider(cache).resolve('NVDA', '1d'); + assert.equal(r.isStale, true); + assert.deepEqual(r.candles, []); + assert.match(r.asOf, /^\d{4}-\d{2}-\d{2}$/); + }); +}); \ No newline at end of file diff --git a/app/server/src/confluence/candleProvider.ts b/app/server/src/confluence/candleProvider.ts new file mode 100644 index 0000000..a4e5ab3 --- /dev/null +++ b/app/server/src/confluence/candleProvider.ts @@ -0,0 +1,176 @@ +// Investor Flow — Candle Resolution Seam for confluence (M22, slice 8) +// +// CANDLE_PROVIDER: the single way confluence slot evaluators obtain a symbol's +// price history. It resolves daily or weekly candles from the shared cache +// (`yfinance:candles::`) and optionally folds the freshest +// live quote into the series so a mid-session evaluation sees the current price +// instead of only the last EOD close. +// +// Why a seam instead of calling `cache.get` inline: +// • evaluators stay testable against fake candle streams, +// • one place owns "what does confluence mean by candles" (sorted ascending, +// quote fold-in rules, staleness), so a future realtime/replay source can +// slot in without touching any slot logic. +// +// ADR-0007: this is a data seam. It resolves price history; it never emits a +// directive. `asOf` on the resolution is the effective evaluation date. +// +// Pure where possible: `foldRealtimeBar` is a pure function; the cache-backed +// provider is a thin shim over CacheRepository. + +import type { CacheRepository, PriceCandle, Quote } from '../cache/CacheRepository.ts'; +import type { SlotGranularity } from './confluenceSlots.ts'; + +// --------------------------------------------------------------------------- +// Types +// --------------------------------------------------------------------------- + +/** Where the last bar of the resolved series came from. */ +export type BarProvenance = 'yfinance' | 'realtime'; + +/** A symbol's resolved candle series for one granularity. */ +export interface CandleResolution { + symbol: string; + granularity: SlotGranularity; + /** Candles sorted ascending by ts. May include a folded-in realtime bar. */ + candles: PriceCandle[]; + /** Effective evaluation date (YYYY-MM-DD) = last bar ts, or quote date when folded. */ + asOf: string; + /** Last-bar provenance: folded live quote vs stored EOD bar. */ + lastBar: BarProvenance; + /** True when the folded realtime bar was appended/updated (not a stored bar). */ + realtimeFolded: boolean; + /** True when the underlying cached series is absent or past its freshness window. */ + isStale: boolean; +} + +/** + * Candle-resolution seam for confluence evaluators and the slot engine. + * + * `resolve` must return candles sorted ascending by ts. Implementations may be + * cache-backed (CacheCandleProvider), precomputed fixtures (tests), or a future + * realtime source — evaluators must not care which. + */ +export interface CandleProvider { + resolve(symbol: string, granularity: SlotGranularity): Promise; +} + +// --------------------------------------------------------------------------- +// Pure helpers +// --------------------------------------------------------------------------- + +/** + * Fold the freshest live quote into a daily series. + * + * Rules: + * • no quote / non-finite price ⇒ unchanged + * • last bar is already `today` ⇒ its close is replaced with the live price + * (O/L/H expanded to contain the print); keeps bar count stable + * • last bar is before `today` ⇒ a new bar for `today` is appended with the + * quote price (O = last close, L/H bracketing it, V = 0) + * + * Returns the new array plus whether anything was folded. Pure. + */ +export function foldRealtimeBar( + candles: PriceCandle[], + quote: Quote | null | undefined, + todayIso: string, +): { candles: PriceCandle[]; folded: boolean } { + if (!quote?.price || !Number.isFinite(quote.price) || candles.length === 0) { + return { candles, folded: false }; + } + + const last = candles[candles.length - 1]; + const lastTs = (last.ts ?? '').slice(0, 10); + const price = quote.price; + + if (lastTs === todayIso) { + const updated: PriceCandle = { + ts: last.ts, + o: last.o, + h: Math.max(last.h, price), + l: Math.min(last.l, price), + c: price, + v: last.v, + adjClose: last.adjClose, + }; + return { candles: [...candles.slice(0, -1), updated], folded: true }; + } + + if (lastTs < todayIso) { + const o = last.c; + return { + candles: [ + ...candles, + { ts: todayIso, o, h: Math.max(o, price), l: Math.min(o, price), c: price, v: 0 }, + ], + folded: true, + }; + } + + return { candles, folded: false }; +} + +/** Current date as YYYY-MM-DD in US/Eastern (the market session's clock). */ +export function todayIso(now: Date = new Date()): string { + const parts = new Intl.DateTimeFormat('en-US', { + timeZone: 'America/New_York', + year: 'numeric', + month: '2-digit', + day: '2-digit', + }).formatToParts(now); + const get = (t: string) => parts.find((p) => p.type === t)?.value ?? ''; + return `${get('year')}-${get('month')}-${get('day')}`; +} + +// --------------------------------------------------------------------------- +// Cache-backed provider +// --------------------------------------------------------------------------- + +/** + * Cache-backed CandleProvider. Reads `yfinance:candles::` + * from the shared cache; for daily granularity it folds the freshest quote in + * when the quote is newer than the last stored bar. + */ +export class CacheCandleProvider implements CandleProvider { + private readonly _cache: CacheRepository; + + constructor(cache: CacheRepository) { + this._cache = cache; + } + + async resolve(symbol: string, granularity: SlotGranularity): Promise { + const sym = symbol.toUpperCase(); + const entry = await this._cache.get(`yfinance:candles:${sym}:${granularity}`); + const stored = (entry?.value ?? []).slice(); + const isStale = entry?.isStale ?? true; + + let candles = stored; + let realtimeFolded = false; + let lastBar: BarProvenance = 'yfinance'; + + if (granularity === '1d') { + const quoteEntry = await this._cache.get(`yfinance:quote:${sym}`); + const quote = quoteEntry?.value; + const fold = foldRealtimeBar(candles, quote, todayIso()); + if (fold.folded && fold.candles.length > 0) { + candles = fold.candles; + realtimeFolded = true; + lastBar = 'realtime'; + } + } + + const lastTs = candles.length > 0 ? (candles[candles.length - 1].ts ?? '').slice(0, 10) : ''; + const asOf = lastTs || todayIso(); + + return { + symbol: sym, + granularity, + candles, + asOf, + lastBar, + realtimeFolded, + isStale, + }; + } +} \ No newline at end of file