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.
This commit is contained in:
@@ -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<string, { value: unknown; stale: boolean }>();
|
||||
setValue(key: string, value: unknown, stale = false): this { this.store.set(key, { value, stale }); return this; }
|
||||
async get<T>(key: string): Promise<CacheEntry<T>> {
|
||||
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<void> { throw new Error('not used'); }
|
||||
stale(key: string): boolean { return !this.store.has(key) || this.store.get(key)!.stale; }
|
||||
async subscribe(): Promise<void> {}
|
||||
async unsubscribe(): Promise<void> {}
|
||||
async ensureInDemand(): Promise<void> {}
|
||||
async pinSystemSymbol(): Promise<void> {}
|
||||
async demandSet(): Promise<string[]> { return []; }
|
||||
async getMany<T>(): Promise<Array<{ key: string; value: T | null; isStale: boolean }>> { return []; }
|
||||
async del(): Promise<void> {}
|
||||
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}$/);
|
||||
});
|
||||
});
|
||||
@@ -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:<symbol>:<granularity>`) 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<CandleResolution>;
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// 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:<symbol>:<granularity>`
|
||||
* 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<CandleResolution> {
|
||||
const sym = symbol.toUpperCase();
|
||||
const entry = await this._cache.get<PriceCandle[]>(`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<Quote>(`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,
|
||||
};
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user