fix: IREN/ASTS/IRE stuck-pending + corridor feature
Fix stuck adapter_queue jobs (ASTS/IRE/IREN stuck pending forever): 1. fetchSpec early-return paths (cooldown checks) now update job status to 'backoff' with last_attempt set and 30s backoff_until, instead of returning without any status change. Prevents jobs from being re-processed every drain cycle indefinitely. 2. Wrap adapter.fetchOne() in 30s Promise.race timeout. A hung HTTP request no longer blocks the entire per-source promise chain forever, preventing all subsequent jobs for that source. Also includes corridor/confluence feature, tiered quote schedules, cache improvements, and HoldingsBookView refinements.
This commit is contained in:
@@ -0,0 +1,239 @@
|
||||
// Investor Flow — Corridor Method backtest grader (M24, slice 4)
|
||||
//
|
||||
// Grades @alojoh's weekly "Corridor Method" entry rankings against ACTUAL
|
||||
// forward price returns — the user asked to judge the methodology itself, not
|
||||
// merely replicate it. Each week we hold his top-3 entry-attractive names
|
||||
// versus his bottom-3 least-attractive names and measure whether the top set
|
||||
// outperformed over a forward window.
|
||||
//
|
||||
// This is a ONE-TIME / periodic offline analysis ledger, not live tracking:
|
||||
// the ranking corpus is baked in (extracted from the articles the user pulled),
|
||||
// and the grader replays it against candle data. Results land in
|
||||
// corridor_backtest for the scorecard UI.
|
||||
//
|
||||
// Pure: the grader consumes candle streams via an injected provider so it is
|
||||
// testable; the DB shim only persists/resolves the ledger.
|
||||
|
||||
import type { DatabaseSync } from 'node:sqlite';
|
||||
import type { PriceCandle } from '../cache/CacheRepository.ts';
|
||||
import { CorridorRepository, type CorridorBacktestRow } from '../db/corridorRepository.ts';
|
||||
import { BACKTEST_HORIZON_DAYS } from '../confluence/corridorData.ts';
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Types
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/** One weekly ranking (top/bottom 3) for one ranking type, from an article. */
|
||||
export interface AlojohRanking {
|
||||
/** The weekly report's date (as close as the article states). YYYY-MM-DD. */
|
||||
articleDate: string;
|
||||
/** X post id of the article the ranking came from (for attribution). */
|
||||
articleId: string;
|
||||
rankingType: 'entry_1y' | 'entry_90d';
|
||||
top: string[]; // most entry-attractive names (max 3)
|
||||
bottom: string[]; // least entry-attractive names (max 3)
|
||||
}
|
||||
|
||||
/** Grader output for one ranking. */
|
||||
export interface CorridorGrade {
|
||||
articleDate: string;
|
||||
articleId: string;
|
||||
rankingType: 'entry_1y' | 'entry_90d';
|
||||
topSymbols: string[];
|
||||
bottomSymbols: string[];
|
||||
topAvgReturn: number | null;
|
||||
bottomAvgReturn: number | null;
|
||||
spread: number | null;
|
||||
isWin: boolean | null; // null when not gradeable (missing data)
|
||||
horizonDays: number;
|
||||
resolvableCount: number; // symbols with enough forward data
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Corpus (extracted from @alojoh's weekly "U.S. Tech Coverage" reports)
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/**
|
||||
* The ranked weeks with complete top/bottom 3 for both observable windows.
|
||||
* Extracted by hand from the 7 fetched subscriber articles (June 21 – Aug 9,
|
||||
* 2026). Only weeks with BOTH top and bottom lists are included so the spread
|
||||
* is always computable when price data exists.
|
||||
*/
|
||||
export const ALOJOH_WEEKLY_RANKINGS: AlojohRanking[] = [
|
||||
{ articleDate: '2026-06-21', articleId: '2068666606801285446', rankingType: 'entry_1y', top: ['PLTR', 'MSFT', 'META'], bottom: ['AMD', 'ORCL', 'TSM'] },
|
||||
{ articleDate: '2026-06-21', articleId: '2068666606801285446', rankingType: 'entry_90d', top: ['MSFT', 'PLTR', 'META'], bottom: ['AMD', 'ORCL', 'TSM'] },
|
||||
{ articleDate: '2026-06-27', articleId: '2070871205624893886', rankingType: 'entry_1y', top: ['PLTR', 'MSFT', 'AVGO'], bottom: ['AMD', 'TSM', 'ORCL'] },
|
||||
{ articleDate: '2026-06-27', articleId: '2070871205624893886', rankingType: 'entry_90d', top: ['PLTR', 'AVGO', 'MSFT'], bottom: ['AMD', 'TSM', 'AAPL'] },
|
||||
{ articleDate: '2026-07-03', articleId: '2073087491511644461', rankingType: 'entry_1y', top: ['PLTR', 'MSFT', 'AVGO'], bottom: ['AMD', 'TSM', 'GOOG'] },
|
||||
{ articleDate: '2026-07-03', articleId: '2073087491511644461', rankingType: 'entry_90d', top: ['AVGO', 'ORCL', 'NVDA'], bottom: ['AMD', 'AAPL', 'TSM'] },
|
||||
{ articleDate: '2026-07-12', articleId: '2076169632717877370', rankingType: 'entry_1y', top: ['PLTR', 'MSFT', 'AMZN'], bottom: ['AMD', 'TSM', 'ORCL'] },
|
||||
{ articleDate: '2026-07-12', articleId: '2076169632717877370', rankingType: 'entry_90d', top: ['ORCL', 'MSFT', 'PLTR'], bottom: ['AMD', 'META', 'AAPL'] },
|
||||
{ articleDate: '2026-07-25', articleId: '2081022231732523491', rankingType: 'entry_1y', top: ['AMZN', 'MSFT', 'PLTR'], bottom: ['AMD', 'AAPL', 'TSM'] },
|
||||
{ articleDate: '2026-07-25', articleId: '2081022231732523491', rankingType: 'entry_90d', top: ['ORCL', 'GOOG', 'TSM'], bottom: ['AAPL', 'AMD', 'AVGO'] },
|
||||
{ articleDate: '2026-08-01', articleId: '2083570089702596888', rankingType: 'entry_1y', top: ['META', 'NVDA', 'AVGO'], bottom: ['AMD', 'MSFT', 'ORCL'] },
|
||||
{ articleDate: '2026-08-01', articleId: '2083570089702596888', rankingType: 'entry_90d', top: ['META', 'ORCL', 'TSM'], bottom: ['MSFT', 'AMZN', 'AVGO'] },
|
||||
{ articleDate: '2026-08-09', articleId: '2086315964740809181', rankingType: 'entry_1y', top: ['TSLA', 'META', 'AMZN'], bottom: ['AMD', 'PLTR', 'ORCL'] },
|
||||
{ articleDate: '2026-08-09', articleId: '2086315964740809181', rankingType: 'entry_90d', top: ['TSLA', 'AMD', 'TSM'], bottom: ['PLTR', 'ORCL', 'MSFT'] },
|
||||
];
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Pure grader
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/**
|
||||
* 7-day forward return for a symbol from a cached daily candle stream.
|
||||
* Returns null when the close on/after `fromDate` or 7 calendar days later is
|
||||
* unavailable. Pure.
|
||||
*/
|
||||
export function forwardReturnFor(
|
||||
candles: PriceCandle[],
|
||||
fromDate: string,
|
||||
horizonDays = BACKTEST_HORIZON_DAYS,
|
||||
): number | null {
|
||||
const sorted = candles.slice().sort((a, b) => (a.ts < b.ts ? -1 : a.ts > b.ts ? 1 : 0));
|
||||
if (sorted.length === 0) return null;
|
||||
|
||||
const targetMs = Date.parse(fromDate);
|
||||
if (Number.isNaN(targetMs)) return null;
|
||||
|
||||
// entry bar: the first bar at or after fromDate (or the last bar before it).
|
||||
let entryIdx = -1;
|
||||
for (let i = 0; i < sorted.length; i++) {
|
||||
const t = Date.parse(sorted[i].ts);
|
||||
if (t >= targetMs) { entryIdx = i; break; }
|
||||
entryIdx = i; // keep the last bar strictly before fromDate as fallback
|
||||
}
|
||||
if (entryIdx < 0) return null;
|
||||
|
||||
const entry = sorted[entryIdx];
|
||||
const horizonMs = horizonDays * 86_400_000;
|
||||
|
||||
let exitIdx = -1;
|
||||
for (let i = entryIdx; i < sorted.length; i++) {
|
||||
const t = Date.parse(sorted[i].ts);
|
||||
if (t >= targetMs + horizonMs) { exitIdx = i; break; }
|
||||
}
|
||||
if (exitIdx < 0) return null; // not enough forward data
|
||||
|
||||
const exit = sorted[exitIdx];
|
||||
if (!Number.isFinite(entry.c) || !Number.isFinite(exit.c) || entry.c <= 0) return null;
|
||||
return exit.c / entry.c - 1;
|
||||
}
|
||||
|
||||
/**
|
||||
* Grade one ranking: mean forward return of the top set minus the bottom set.
|
||||
* Pure — candleProvider is injected so tests can stub streams.
|
||||
*/
|
||||
export async function gradeRanking(
|
||||
ranking: AlojohRanking,
|
||||
candleProvider: (symbol: string) => Promise<PriceCandle[]>,
|
||||
horizonDays = BACKTEST_HORIZON_DAYS,
|
||||
): Promise<CorridorGrade> {
|
||||
const topReturns: number[] = [];
|
||||
const botReturns: number[] = [];
|
||||
|
||||
for (const sym of ranking.top) {
|
||||
const r = forwardReturnFor(await candleProvider(sym), ranking.articleDate, horizonDays);
|
||||
if (r !== null) topReturns.push(r);
|
||||
}
|
||||
for (const sym of ranking.bottom) {
|
||||
const r = forwardReturnFor(await candleProvider(sym), ranking.articleDate, horizonDays);
|
||||
if (r !== null) botReturns.push(r);
|
||||
}
|
||||
|
||||
const topAvg = topReturns.length > 0 ? topReturns.reduce((a, b) => a + b, 0) / topReturns.length : null;
|
||||
const botAvg = botReturns.length > 0 ? botReturns.reduce((a, b) => a + b, 0) / botReturns.length : null;
|
||||
const spread = topAvg !== null && botAvg !== null ? topAvg - botAvg : null;
|
||||
|
||||
return {
|
||||
articleDate: ranking.articleDate,
|
||||
articleId: ranking.articleId,
|
||||
rankingType: ranking.rankingType,
|
||||
topSymbols: ranking.top,
|
||||
bottomSymbols: ranking.bottom,
|
||||
topAvgReturn: topAvg,
|
||||
bottomAvgReturn: botAvg,
|
||||
spread,
|
||||
isWin: spread === null ? null : spread > 0,
|
||||
horizonDays,
|
||||
resolvableCount: topReturns.length + botReturns.length,
|
||||
};
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Aggregates
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/** Aggregate scorecard from a set of grades (pure). */
|
||||
export interface CorridorBacktestAggregate {
|
||||
entries: number;
|
||||
resolved: number; // grades with a spread
|
||||
resolved1y: number;
|
||||
resolved90d: number;
|
||||
hitRate: number | null; // frac of resolved grades where top outperformed
|
||||
avgSpread: number | null; // mean spread (returns, e.g. 0.012 = +1.2%)
|
||||
bestSpread: number | null;
|
||||
worstSpread: number | null;
|
||||
avgTopReturn: number | null;
|
||||
avgBottomReturn: number | null;
|
||||
}
|
||||
|
||||
export function aggregateBacktest(grades: CorridorGrade[]): CorridorBacktestAggregate {
|
||||
const resolved = grades.filter((g) => g.spread !== null);
|
||||
const resolved1y = resolved.filter((g) => g.rankingType === 'entry_1y');
|
||||
const resolved90d = resolved.filter((g) => g.rankingType === 'entry_90d');
|
||||
const mean = (arr: number[]): number | null => (arr.length > 0 ? arr.reduce((a, b) => a + b, 0) / arr.length : null);
|
||||
|
||||
return {
|
||||
entries: grades.length,
|
||||
resolved: resolved.length,
|
||||
resolved1y: resolved1y.length,
|
||||
resolved90d: resolved90d.length,
|
||||
hitRate: resolved.length > 0 ? resolved.filter((g) => g.isWin === true).length / resolved.length : null,
|
||||
avgSpread: mean(resolved.map((g) => g.spread as number)),
|
||||
bestSpread: resolved.length > 0 ? Math.max(...resolved.map((g) => g.spread as number)) : null,
|
||||
worstSpread: resolved.length > 0 ? Math.min(...resolved.map((g) => g.spread as number)) : null,
|
||||
avgTopReturn: mean(resolved.map((g) => g.topAvgReturn as number)),
|
||||
avgBottomReturn: mean(resolved.map((g) => g.bottomAvgReturn as number)),
|
||||
};
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Ledger shim (offline replay; persists grades to corridor_backtest)
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/**
|
||||
* Replay the whole corpus against a candle provider and persist the ledger.
|
||||
* Returns the freshly graded rows + aggregate. Callers control cadence (this is
|
||||
* an offline analysis run, not a request-path fetch).
|
||||
*/
|
||||
export async function runCorridorBacktest(
|
||||
db: DatabaseSync,
|
||||
candleProvider: (symbol: string) => Promise<PriceCandle[]>,
|
||||
opts: { horizonDays?: number } = {},
|
||||
): Promise<{ grades: CorridorGrade[]; aggregate: CorridorBacktestAggregate }> {
|
||||
const repo = new CorridorRepository(db);
|
||||
const grades: CorridorGrade[] = [];
|
||||
for (const ranking of ALOJOH_WEEKLY_RANKINGS) {
|
||||
const grade = await gradeRanking(ranking, candleProvider, opts.horizonDays ?? BACKTEST_HORIZON_DAYS);
|
||||
grades.push(grade);
|
||||
if (grade.spread !== null) {
|
||||
repo.saveBacktest({
|
||||
articleDate: grade.articleDate,
|
||||
articleId: grade.articleId,
|
||||
rankingType: grade.rankingType,
|
||||
topSymbols: grade.topSymbols,
|
||||
bottomSymbols: grade.bottomSymbols,
|
||||
topAvgReturn: grade.topAvgReturn,
|
||||
bottomAvgReturn: grade.bottomAvgReturn,
|
||||
spread: grade.spread,
|
||||
isWin: grade.isWin === true,
|
||||
horizonDays: grade.horizonDays,
|
||||
gradedAt: new Date().toISOString(),
|
||||
});
|
||||
}
|
||||
}
|
||||
return { grades, aggregate: aggregateBacktest(grades) };
|
||||
}
|
||||
|
||||
export type { CorridorBacktestRow };
|
||||
+22
@@ -8,6 +8,7 @@ import { db as defaultDb } from '../db/client.ts';
|
||||
import {
|
||||
CANDLE_FRESH_MS,
|
||||
quoteTtlMs,
|
||||
tieredQuoteTtlMs,
|
||||
SYMBOL_META_INCOMPLETE_TTL_MS,
|
||||
} from '../queue/sourceRatePolicy.ts';
|
||||
import { KvReadCache } from './LruCache.ts';
|
||||
@@ -190,6 +191,13 @@ export function needsQuoteRefresh(d: DatabaseSync, symbol: string, now = Date.no
|
||||
return tsAgeMs(row.observed_at, now) > quoteTtlMs(new Date(now));
|
||||
}
|
||||
|
||||
/** Tier-aware quote freshness check. Portfolio (T0) gets the tightest TTL. */
|
||||
export function needsTieredQuoteRefresh(d: DatabaseSync, symbol: string, tier: number, now = Date.now()): boolean {
|
||||
const row = d.prepare('SELECT observed_at FROM quotes WHERE symbol=?').get(symbol) as { observed_at: string } | undefined;
|
||||
if (!row?.observed_at) return true;
|
||||
return tsAgeMs(row.observed_at, now) > tieredQuoteTtlMs(tier, new Date(now));
|
||||
}
|
||||
|
||||
/** True when symbol meta missing, incomplete (no name), or past weekly TTL. */
|
||||
export function needsSymbolMetaRefresh(d: DatabaseSync, symbol: string, now = Date.now()): boolean {
|
||||
const row = d.prepare('SELECT name, sector, updated_at FROM symbols WHERE symbol=?').get(symbol) as
|
||||
@@ -617,6 +625,8 @@ export interface CacheRepository {
|
||||
* Safe to call on every Market Outlook / ticker context load.
|
||||
*/
|
||||
ensureInDemand(symbol: string, tickerKind: TickerKind): Promise<void>;
|
||||
/** Bump symbol to watched tier (2) on page view. Decays back after 10 min. */
|
||||
bumpToWatched(symbol: string, tickerKind: TickerKind): Promise<void>;
|
||||
/** Permanent system pin (rotation universe, SPY, VIX) — survives unsubscribe. */
|
||||
pinSystemSymbol(symbol: string, tickerKind: TickerKind): Promise<void>;
|
||||
demandSet(): Promise<string[]>;
|
||||
@@ -735,6 +745,18 @@ export class CacheRepositoryImpl implements CacheRepository {
|
||||
await this.queueIfNeeded(symbol);
|
||||
}
|
||||
|
||||
/** Bump a symbol to watched tier (2) on page view. The periodic tier
|
||||
* recompute decays it back to background after ~10 min of inactivity. */
|
||||
async bumpToWatched(symbol: string, tickerKind: TickerKind): Promise<void> {
|
||||
this.ensureDemandRow(symbol, tickerKind);
|
||||
const now = new Date().toISOString();
|
||||
// Only lower tier (raise priority) — never raise tier above current.
|
||||
this._db.prepare(
|
||||
"UPDATE symbol_demand SET tier = MIN(tier, 2), last_viewed_at = ?, in_demand = 1 WHERE symbol=?",
|
||||
).run(now, symbol);
|
||||
await this.queueIfNeeded(symbol);
|
||||
}
|
||||
|
||||
async pinSystemSymbol(symbol: string, tickerKind: TickerKind): Promise<void> {
|
||||
this.ensureDemandRow(symbol, tickerKind);
|
||||
this._db.prepare(
|
||||
|
||||
@@ -0,0 +1,129 @@
|
||||
// Investor Flow — confluenceEngine.test.ts (M24 slice 5)
|
||||
// Integration test for the evaluation engine: runs the wired slot families over
|
||||
// fake cache candle streams, verifies rack evaluations + corridor snapshots +
|
||||
// signal history are persisted and idempotent per (symbol, asOf, rack).
|
||||
|
||||
import { describe, it, beforeEach } from 'node:test';
|
||||
import assert from 'node:assert/strict';
|
||||
|
||||
import { createDb, initSchema } from '../../db/client.ts';
|
||||
import { ConfluenceRepository } from '../../db/confluenceRepository.ts';
|
||||
import { CorridorRepository } from '../../db/corridorRepository.ts';
|
||||
import { createCacheRepository, type CacheRepository, type CacheEntry, type PriceCandle, type Quote } from '../../cache/CacheRepository.ts';
|
||||
import { FakeSourceAdapter } from '../../adapters/SourceAdapter.ts';
|
||||
import { AdapterQueue } from '../../queue/AdapterQueue.ts';
|
||||
import { seedConfluence, CONFLUENCE_UNIVERSE, BENCHMARK_SYMBOL } from '../confluenceSeed.ts';
|
||||
import { runConfluenceEvaluationCycle } from '../confluenceEngine.ts';
|
||||
|
||||
// ---- fake cache that returns whatever we seeded ---------------------------------
|
||||
|
||||
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;
|
||||
}
|
||||
|
||||
// ---- fixtures -------------------------------------------------------------
|
||||
|
||||
/** Monotonic ramp up (bullish technicals) over `days` trading days. */
|
||||
function rampUp(days: number, start = 100, dailyPct = 0.0015): PriceCandle[] {
|
||||
const candles: PriceCandle[] = [];
|
||||
const base = Date.UTC(2020, 0, 1);
|
||||
for (let i = 0; i < days; i++) {
|
||||
const c = start * Math.pow(1 + dailyPct, i);
|
||||
candles.push({
|
||||
ts: new Date(base + i * 86400000).toISOString().slice(0, 10),
|
||||
o: c * (1 - dailyPct / 2),
|
||||
h: c * 1.003,
|
||||
l: c * 0.997,
|
||||
c,
|
||||
v: 2_000_000,
|
||||
adjClose: c,
|
||||
});
|
||||
}
|
||||
return candles;
|
||||
}
|
||||
|
||||
let db: ReturnType<typeof createDb>;
|
||||
let cache: FakeCache;
|
||||
let confluenceRepo: ConfluenceRepository;
|
||||
let corridorRepo: CorridorRepository;
|
||||
|
||||
beforeEach(async () => {
|
||||
db = createDb({ path: ':memory:' });
|
||||
initSchema(db);
|
||||
confluenceRepo = new ConfluenceRepository(db);
|
||||
corridorRepo = new CorridorRepository(db);
|
||||
cache = new FakeCache();
|
||||
await seedConfluence(db, cache as unknown as CacheRepository);
|
||||
|
||||
// Seed candle streams for the whole universe + SPY benchmark.
|
||||
for (const { symbol } of CONFLUENCE_UNIVERSE) {
|
||||
cache.setValue(`yfinance:candles:${symbol}:1d`, rampUp(300));
|
||||
cache.setValue(`yfinance:candles:${symbol}:1wk`, rampUp(80, 100, 0.01));
|
||||
cache.setValue(`yfinance:quote:${symbol}`, { price: 130 } as Quote);
|
||||
cache.setValue(`yfinance:dividendFundamentals:${symbol}`, { forwardPE: 25, trailingPE: 24, forwardEPS: 5.2, trailingEPS: 5.0, currentPrice: 130 });
|
||||
}
|
||||
cache.setValue(`yfinance:candles:SPY:1d`, rampUp(300, 400, 0.001));
|
||||
cache.setValue(`yfinance:candles:SPY:1wk`, rampUp(80, 400, 0.01));
|
||||
cache.setValue(`yfinance:quote:SPY`, { price: 440 } as Quote);
|
||||
cache.setValue(`yfinance:dividendFundamentals:SPY`, { forwardPE: 22, trailingPE: 21.5, forwardEPS: 20, trailingEPS: 19.3, currentPrice: 440 });
|
||||
});
|
||||
|
||||
describe('runConfluenceEvaluationCycle', () => {
|
||||
it('persists an evaluation + corridor snapshots per symbol, for every system rack', async () => {
|
||||
const summary = await runConfluenceEvaluationCycle(db, cache as unknown as CacheRepository);
|
||||
|
||||
// All 15 universe symbols evaluated (SPY benchmark not part of universe).
|
||||
assert.equal(summary.symbolsEvaluated.length, CONFLUENCE_UNIVERSE.length);
|
||||
assert.equal(summary.evaluationsStored, CONFLUENCE_UNIVERSE.length * 3, 'one eval per symbol per system rack');
|
||||
assert.ok(summary.corridorSnapshots >= CONFLUENCE_UNIVERSE.length, 'symbol + SPY corridor snapshots');
|
||||
|
||||
// spot-check one symbol persisted its "full" rack evaluation
|
||||
const pltr = CONFLUENCE_UNIVERSE.find((s) => s.symbol === 'PLTR')!;
|
||||
const evals = confluenceRepo.listEvaluationsForSymbol(pltr.symbol);
|
||||
assert.equal(evals.length, 3);
|
||||
const validQuality = ['strong-bullish', 'moderate-bullish', 'weak-bullish', 'mixed', 'weak-bearish', 'moderate-bearish', 'strong-bearish', 'sparse'];
|
||||
for (const ev of evals) {
|
||||
assert.equal(ev.symbol, pltr.symbol);
|
||||
assert.ok(validQuality.includes(ev.quality), `unexpected quality ${ev.quality}`);
|
||||
assert.ok(ev.assessedCount >= 10, 'rack has meaningful coverage');
|
||||
}
|
||||
|
||||
// corridor snapshot table populated for at least the spot-checked symbol
|
||||
const snap = corridorRepo.latestSnapshot(pltr.symbol);
|
||||
assert.ok(snap, 'corridor snapshot stored');
|
||||
assert.ok(snap.corridor1yMedian !== null, 'corridor median computed');
|
||||
});
|
||||
|
||||
it('is idempotent per (symbol, asOf, rack): a second run reuses and stores nothing new', async () => {
|
||||
const first = await runConfluenceEvaluationCycle(db, cache as unknown as CacheRepository);
|
||||
const second = await runConfluenceEvaluationCycle(db, cache as unknown as CacheRepository);
|
||||
|
||||
assert.equal(second.evaluationsStored, 0);
|
||||
assert.equal(second.evaluationsReused, first.evaluationsStored);
|
||||
const all = confluenceRepo.listEvaluationsForSymbol(CONFLUENCE_UNIVERSE[0].symbol);
|
||||
assert.equal(all.length, 3);
|
||||
});
|
||||
|
||||
it('skips symbols with no cached candles without failing the run', async () => {
|
||||
cache.setValue(`yfinance:candles:${CONFLUENCE_UNIVERSE[0].symbol}:1d`, [], false);
|
||||
const summary = await runConfluenceEvaluationCycle(db, cache as unknown as CacheRepository);
|
||||
assert.equal(summary.symbolsEvaluated.length, CONFLUENCE_UNIVERSE.length - 1);
|
||||
assert.equal(summary.symbolsSkipped.length, 1);
|
||||
assert.equal(summary.symbolsSkipped[0].symbol, CONFLUENCE_UNIVERSE[0].symbol);
|
||||
});
|
||||
});
|
||||
@@ -19,8 +19,8 @@ const fired = (id: string): SlotAssessment => ({ id, state: 'fired' });
|
||||
const notFired = (id: string): SlotAssessment => ({ id, state: 'not-fired' });
|
||||
|
||||
describe('confluence catalog', () => {
|
||||
it('defines exactly 34 slots across all six families', () => {
|
||||
assert.equal(CONFLUENCE_SLOTS.length, 34);
|
||||
it('defines slots across all six families in catalog order', () => {
|
||||
assert.ok(CONFLUENCE_SLOTS.length >= 34);
|
||||
const families = new Set(CONFLUENCE_SLOTS.map((s) => s.family));
|
||||
assert.deepEqual([...families].sort(), ['flows', 'institutional', 'macro', 'seasonal', 'sentiment', 'technical']);
|
||||
});
|
||||
|
||||
@@ -15,7 +15,7 @@ import {
|
||||
BENCHMARK_SYMBOL,
|
||||
defineSystemRackPresets,
|
||||
} from '../confluenceSeed.ts';
|
||||
import { CONFLUENCE_SLOT_IDS } from '../confluenceSlots.ts';
|
||||
import { CONFLUENCE_SLOTS, CONFLUENCE_SLOT_IDS } from '../confluenceSlots.ts';
|
||||
|
||||
let db: ReturnType<typeof createDb>;
|
||||
let cache: CacheRepository;
|
||||
@@ -63,19 +63,22 @@ describe('defineSystemRackPresets', () => {
|
||||
}
|
||||
});
|
||||
|
||||
it('"Full Confluence" uses all34 slots', () => {
|
||||
it('"Full Confluence" uses every slot in the catalog', () => {
|
||||
const full = defineSystemRackPresets().find((p) => p.id === 'confluence-full')!;
|
||||
assert.equal(full.slotIds.length, 34);
|
||||
assert.equal(full.slotIds.length, CONFLUENCE_SLOT_IDS.length);
|
||||
});
|
||||
|
||||
it('"Technical Momentum" has 15 slots', () => {
|
||||
it('"Technical Momentum" has only the technical slots', () => {
|
||||
const tech = defineSystemRackPresets().find((p) => p.id === 'confluence-technical')!;
|
||||
assert.equal(tech.slotIds.length, 15);
|
||||
assert.equal(tech.slotIds.length, CONFLUENCE_SLOTS.filter((s) => s.family === 'technical').length);
|
||||
});
|
||||
|
||||
it('"Macro + Flows + Sentiment" has 14 slots', () => {
|
||||
it('"Macro + Flows + Sentiment" has macro + seasonal + flows + sentiment slots', () => {
|
||||
const macro = defineSystemRackPresets().find((p) => p.id === 'confluence-macro-flows')!;
|
||||
assert.equal(macro.slotIds.length, 14);
|
||||
const expected = CONFLUENCE_SLOTS.filter((s) =>
|
||||
['macro', 'seasonal', 'flows', 'sentiment'].includes(s.family),
|
||||
).length;
|
||||
assert.equal(macro.slotIds.length, expected);
|
||||
});
|
||||
});
|
||||
|
||||
|
||||
@@ -0,0 +1,188 @@
|
||||
// Investor Flow — Confluence Evaluation Engine (M24, slice 5)
|
||||
//
|
||||
// The missing backbone of the Confluence Signal Engine: the daily cycle that
|
||||
// runs every wired slot family's evaluator over a symbol's cached data and
|
||||
// persists the resulting rack evaluations + slot-fire history.
|
||||
//
|
||||
// Each symbol in the confluence universe is resolved once (daily candles,
|
||||
// weekly candles, SPY benchmark daily, seasonality snapshot, corridor
|
||||
// snapshots for the symbol and SPY), then every system rack's slot subset is
|
||||
// sliced out, run through `evaluateRack` (redundancy-aware), and stored via
|
||||
// `ConfluenceRepository`. Fired slots are logged to signal history for the
|
||||
// reliability scorecard, and pending fires are resolved against forward prices
|
||||
//
|
||||
// ADR-0007: this engine computes description ("the picture is moderate-bullish"
|
||||
// because X evidence) — it never emits buy/sell directives.
|
||||
//
|
||||
// ADR-0009: the engine ONLY reads the shared cache (via CacheCandleProvider and
|
||||
// corridorData.resolveCorridorSnapshot). Any refreshing is already queued by
|
||||
// the cache; the engine never touches a vendor directly.
|
||||
|
||||
import type { DatabaseSync } from 'node:sqlite';
|
||||
import { randomUUID } from 'node:crypto';
|
||||
import type { CacheRepository } from '../cache/CacheRepository.ts';
|
||||
import { CONFLUENCE_SLOTS } from './confluenceSlots.ts';
|
||||
import { CacheCandleProvider, type CandleProvider } from './candleProvider.ts';
|
||||
import { evaluateTechnicalSlots } from './technicalEvaluator.ts';
|
||||
import { evaluateSeasonalSlots } from './seasonalEvaluator.ts';
|
||||
import { evaluateCorridorSlots } from './corridorEvaluator.ts';
|
||||
import { buildSeasonalitySnapshot } from '../analysis/seasonality.ts';
|
||||
import { resolveCorridorSnapshot } from './corridorData.ts';
|
||||
import { evaluateRack, type ConfluenceEvaluation, type SlotAssessment } from './confluenceRack.ts';
|
||||
import { ConfluenceRepository, type ConfluenceRack } from '../db/confluenceRepository.ts';
|
||||
import { resolveSignalHistory } from './confluenceBacktest.ts';
|
||||
import { CONFLUENCE_UNIVERSE, BENCHMARK_SYMBOL } from './confluenceSeed.ts';
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Unwired slots → honest fallback assessments
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/**
|
||||
* Emit a `not-fired` fallback for every catalog slot the wired evaluators did
|
||||
* not already cover (currently: macro non-corridor, institutional, flows,
|
||||
* sentiment). Keeps each rack's slot set fully covered so assessedCount stays
|
||||
* meaningful, but the fallbacks contribute no evidence to the picture.
|
||||
*/
|
||||
function unwiredFallbacks(emitted: SlotAssessment[]): SlotAssessment[] {
|
||||
const emittedIds = new Set(emitted.map((a) => a.id));
|
||||
return CONFLUENCE_SLOTS
|
||||
.filter((s) => !emittedIds.has(s.id))
|
||||
.map((s) => ({ id: s.id, state: 'not-fired' as const, note: 'No evaluator wired for this slot yet — contributes no evidence to the picture.' }));
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Types
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
export interface ConfluenceEngineRunSummary {
|
||||
symbolsEvaluated: string[];
|
||||
symbolsSkipped: Array<{ symbol: string; reason: string }>;
|
||||
evaluationsStored: number;
|
||||
evaluationsReused: number;
|
||||
signalsLogged: number;
|
||||
signalsResolved: number;
|
||||
corridorSnapshots: number;
|
||||
}
|
||||
|
||||
export interface RunOptions {
|
||||
/** Symbols to evaluate (default: the confluence universe). */
|
||||
symbols?: string[];
|
||||
/** Racks to evaluate per symbol (default: system racks). */
|
||||
racks?: ConfluenceRack[];
|
||||
/** Force re-evaluation even when a same-asOf evaluation already exists. */
|
||||
force?: boolean;
|
||||
/** Evaluation date override (default: the daily resolution asOf). */
|
||||
asOf?: string;
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Engine
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/**
|
||||
* Run one full confluence evaluation cycle over the universe.
|
||||
*
|
||||
* Per symbol: resolve candles once, run every wired evaluator across the slot
|
||||
* families, slice each rack's slots, evaluate the rack, persist the evaluation
|
||||
* (+ signal fires when new), then resolve the pending signal history. Returns a
|
||||
* summary for logging/tests. Idempotent per (symbol, asOf, rack).
|
||||
*/
|
||||
export async function runConfluenceEvaluationCycle(
|
||||
db: DatabaseSync,
|
||||
cache: CacheRepository,
|
||||
opts: RunOptions = {},
|
||||
): Promise<ConfluenceEngineRunSummary> {
|
||||
const repo = new ConfluenceRepository(db);
|
||||
const provider: CandleProvider = new CacheCandleProvider(cache);
|
||||
|
||||
const symbols = (opts.symbols?.map((s) => s.toUpperCase()) ?? CONFLUENCE_UNIVERSE.map((s) => s.symbol));
|
||||
const racks = opts.racks ?? repo.listSystemRacks();
|
||||
if (racks.length === 0) throw new Error('confluence engine: no system racks seeded');
|
||||
|
||||
const summary: ConfluenceEngineRunSummary = {
|
||||
symbolsEvaluated: [],
|
||||
symbolsSkipped: [],
|
||||
evaluationsStored: 0,
|
||||
evaluationsReused: 0,
|
||||
signalsLogged: 0,
|
||||
signalsResolved: 0,
|
||||
corridorSnapshots: 0,
|
||||
};
|
||||
|
||||
// SPY corridor snapshot is shared across every symbol (market-level slots).
|
||||
let spyCorridor = await resolveCorridorSnapshot(db, cache, BENCHMARK_SYMBOL, { date: opts.asOf });
|
||||
if (spyCorridor) summary.corridorSnapshots += 1;
|
||||
|
||||
for (const symbol of symbols) {
|
||||
const daily = await provider.resolve(symbol, '1d');
|
||||
if (daily.candles.length === 0) {
|
||||
summary.symbolsSkipped.push({ symbol, reason: 'no daily candles in cache yet' });
|
||||
continue;
|
||||
}
|
||||
|
||||
const weekly = await provider.resolve(symbol, '1wk');
|
||||
const benchmarkDaily = await provider.resolve(BENCHMARK_SYMBOL, '1d');
|
||||
const asOf = opts.asOf ?? daily.asOf;
|
||||
|
||||
// Corridor snapshot for the symbol (valuation corridor, entry/upside slots).
|
||||
let corridor = await resolveCorridorSnapshot(db, cache, symbol, { date: asOf });
|
||||
if (corridor) summary.corridorSnapshots += 1;
|
||||
if (spyCorridor === null) spyCorridor = await resolveCorridorSnapshot(db, cache, BENCHMARK_SYMBOL, { date: asOf });
|
||||
|
||||
// ----- run every wired family evaluator -----
|
||||
const wired: SlotAssessment[] = [
|
||||
...evaluateTechnicalSlots(symbol, daily.candles, {
|
||||
weekly: weekly.candles,
|
||||
benchmarkDaily: benchmarkDaily.candles,
|
||||
}),
|
||||
...evaluateSeasonalSlots(buildSeasonalitySnapshot(symbol, daily.candles), asOf),
|
||||
...(corridor ? evaluateCorridorSlots(corridor, spyCorridor) : corridorUnavailableAssessments()),
|
||||
];
|
||||
// Cover every catalog slot the wired evaluators left out (honest fallbacks).
|
||||
const assessments = [...wired, ...unwiredFallbacks(wired)];
|
||||
|
||||
// ----- slice per rack, evaluate, persist -----
|
||||
let symbolStored = 0;
|
||||
let symbolReused = 0;
|
||||
for (const rack of racks) {
|
||||
const rackSlots = new Set(rack.slotIds);
|
||||
const sliced = assessments.filter((a) => rackSlots.has(a.id));
|
||||
if (sliced.length === 0) continue;
|
||||
|
||||
const existing = repo.getEvaluation(symbol, asOf, rack.id);
|
||||
if (existing && !opts.force) {
|
||||
symbolReused += 1;
|
||||
continue;
|
||||
}
|
||||
|
||||
const evaluation: ConfluenceEvaluation = evaluateRack(symbol, asOf, sliced);
|
||||
const firesLogged = evaluation.assessments.filter((a) => a.state === 'fired').length;
|
||||
repo.saveEvaluation(evaluation, rack.id, randomUUID());
|
||||
repo.logSignalFires(evaluation, rack.id);
|
||||
symbolStored += 1;
|
||||
summary.signalsLogged += firesLogged;
|
||||
}
|
||||
|
||||
summary.evaluationsStored += symbolStored;
|
||||
summary.evaluationsReused += symbolReused;
|
||||
summary.symbolsEvaluated.push(symbol);
|
||||
}
|
||||
|
||||
// Resolve pending signal history against forward prices.
|
||||
const resolution = await resolveSignalHistory(db, async (sym) => {
|
||||
const res = await provider.resolve(sym.toUpperCase(), '1d');
|
||||
return res.candles;
|
||||
});
|
||||
summary.signalsResolved = resolution.resolved.filter((r) => r.verdict !== 'deferred').length;
|
||||
|
||||
return summary;
|
||||
}
|
||||
|
||||
/** Six corridor slots, all not-fired with a shared note when no snapshot exists. */
|
||||
function corridorUnavailableAssessments(): SlotAssessment[] {
|
||||
return [
|
||||
'corridorEntryCheap', 'corridorEntryStretched',
|
||||
'corridorUpsideHigh', 'corridorUpsideLow',
|
||||
'spyCorridorCheap', 'spyCorridorStretched',
|
||||
].map((id) => ({ id, state: 'not-fired', note: 'No valuation-corridor snapshot cached for this symbol yet.' }));
|
||||
}
|
||||
@@ -35,6 +35,8 @@ export const REDUNDANCY_GROUPS: RedundancyGroup[] = [
|
||||
{ family: 'institutional', slots: ['instNetActivePositive', 'insiderInformedBuy30d', 'new13da'] },
|
||||
{ family: 'macro', slots: ['ratesRegime', 'macroRegimeUp'] },
|
||||
{ family: 'macro', slots: ['consumerSentimentLow', 'breadthThrust'] },
|
||||
{ family: 'macro', slots: ['corridorEntryCheap', 'corridorUpsideHigh'] },
|
||||
{ family: 'macro', slots: ['corridorEntryStretched', 'corridorUpsideLow'] },
|
||||
{ family: 'seasonal', slots: ['seasonalFavorableMonth', 'winterHalfOn', 'electionCycleFavorableYear'] },
|
||||
{ family: 'flows', slots: ['etfFlowPositive', 'cotPositioning'] },
|
||||
];
|
||||
|
||||
@@ -66,31 +66,32 @@ function familySlots(...families: string[]): string[] {
|
||||
|
||||
/**
|
||||
* Three curated system rack presets. Each is a different lens on the same
|
||||
* symbol data, expressed as a subset of the 34-slot catalog:
|
||||
* symbol data, expressed as a subset of the catalog:
|
||||
*
|
||||
* 1. "Full Confluence" — all 34 slots (the default every-picture view)
|
||||
* 2. "Technical Momentum" — the15 technical slots only (price-action focus)
|
||||
* 3. "Macro + Flows + Sentiment" — macro 5 + seasonal 5 + flows 3 +
|
||||
* sentiment 1 = 14 slots (the macro/structural lens)
|
||||
* 1. "Full Confluence" — all slots (the default every-picture view)
|
||||
* 2. "Technical Momentum" — the 15 technical slots only (price-action focus)
|
||||
* 3. "Macro + Flows + Sentiment" — macro + seasonal + flows + sentiment
|
||||
* slots (the macro/structural lens, now including the valuation corridor)
|
||||
*/
|
||||
export function defineSystemRackPresets(): RackPreset[] {
|
||||
const techIds = familySlots('technical');
|
||||
return [
|
||||
{
|
||||
id: 'confluence-full',
|
||||
name: 'Full Confluence',
|
||||
description: 'All 34 slots. The broadest evidence view of a symbol\'s picture.',
|
||||
description: `All ${ALL_IDS.length} slots. The broadest evidence view of a symbol\'s picture.`,
|
||||
slotIds: [...ALL_IDS],
|
||||
},
|
||||
{
|
||||
id: 'confluence-technical',
|
||||
name: 'Technical Momentum',
|
||||
description: 'The 15 technical slots: trend, momentum, mean-reversion, and volume.',
|
||||
slotIds: familySlots('technical'),
|
||||
description: `The ${techIds.length} technical slots: trend, momentum, mean-reversion, and volume.`,
|
||||
slotIds: techIds,
|
||||
},
|
||||
{
|
||||
id: 'confluence-macro-flows',
|
||||
name: 'Macro + Flows + Sentiment',
|
||||
description: 'Macro regime, seasonal calendar, ETF/COT flows, and informed-commentator sentiment (14 slots).',
|
||||
description: 'Macro regime, valuation corridor, seasonal calendar, ETF/COT flows, and informed-commentator sentiment.',
|
||||
slotIds: familySlots('macro', 'seasonal', 'flows', 'sentiment'),
|
||||
},
|
||||
];
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
// Investor Flow — Confluence Slot Catalog (M22, slice 2)
|
||||
//
|
||||
// The 34-slot confluence inventory for the Confluence Signal Engine. Each slot is
|
||||
// The 40-slot confluence inventory for the Confluence Signal Engine. Each slot is
|
||||
// a named, independently-evaluable check whose *firing* state contributes bullish
|
||||
// or bearish evidence about a symbol's entry/exit quality.
|
||||
//
|
||||
@@ -46,7 +46,7 @@ export interface ConfluenceSlot {
|
||||
explain: string;
|
||||
}
|
||||
|
||||
/** Complete 34-slot confluence catalog in evaluation order. */
|
||||
/** Complete confluence catalog in evaluation order. */
|
||||
export const CONFLUENCE_SLOTS: ConfluenceSlot[] = [
|
||||
// ---------------------------------------------------------------- technical
|
||||
{ id: 'goldenCross', name: 'Golden Cross', family: 'technical', body: 'bull', granularity: '1wk', explain: 'The 50-window average has crossed above the 200-window average, a widely-watched trend-quality marker.' },
|
||||
@@ -93,6 +93,14 @@ export const CONFLUENCE_SLOTS: ConfluenceSlot[] = [
|
||||
|
||||
// ---------------------------------------------------------------- sentiment
|
||||
{ id: 'commentatorSentiment', name: 'Informed Commentator Sentiment', family: 'sentiment', body: 'bull', granularity: '1d', explain: 'Informed commentators tracked via the configured sentiment source are net-positive on the symbol in the measurement window.' },
|
||||
|
||||
// ----------------------------------------------------------------- corridor
|
||||
{ id: 'corridorEntryCheap', name: 'Corridor: Entry Cheap', family: 'macro', body: 'bull', granularity: '1d', explain: 'The current P/E sits in the lower third of the 1-year observable valuation corridor, indicating a relatively attractive entry point versus the symbol\'s own history.' },
|
||||
{ id: 'corridorEntryStretched', name: 'Corridor: Entry Stretched', family: 'macro', body: 'exit', granularity: '1d', explain: 'The current P/E sits in the upper third of the 1-year observable valuation corridor, indicating a stretched valuation versus the symbol\'s own history.' },
|
||||
{ id: 'corridorUpsideHigh', name: 'Corridor: Upside High', family: 'macro', body: 'bull', granularity: '1d', explain: 'Applying the 1-year median observable multiple to forward earnings implies meaningful upside from the current price.' },
|
||||
{ id: 'corridorUpsideLow', name: 'Corridor: Upside Low', family: 'macro', body: 'exit', granularity: '1d', explain: 'Applying the 1-year median observable multiple to forward earnings implies meaningful downside from the current price.' },
|
||||
{ id: 'spyCorridorCheap', name: 'SPY Corridor: Cheap', family: 'macro', body: 'bull', granularity: '1wk', explain: 'The S&P 500 forward P/E sits below its 3-year median, a market-level valuation tailwind that improves the odds for broad equity exposure.' },
|
||||
{ id: 'spyCorridorStretched', name: 'SPY Corridor: Stretched', family: 'macro', body: 'exit', granularity: '1wk', explain: 'The S&P 500 forward P/E sits at or above its 3-year median, a market-level valuation headwind that tempers the broad-equity picture.' },
|
||||
];
|
||||
|
||||
/** Indexed by slot id for O(1) lookup. */
|
||||
|
||||
@@ -0,0 +1,279 @@
|
||||
// Investor Flow — Price Corridor data pipeline (M24, slice 2)
|
||||
//
|
||||
// The Corridor Method (as reverse-engineered from @alojoh's weekly "U.S. Tech
|
||||
// Coverage / Market Valuation" reports): a symbol's *observable multiple* range
|
||||
// over a lookback window defines a valuation corridor. The current P/E position
|
||||
// within that corridor signals entry timing (cheap near the low band, stretched
|
||||
// near the high band), and applying the corridor's median multiple to forward
|
||||
// EPS derives an implied fair value / upside.
|
||||
//
|
||||
// ADR-0007: this is a valuation-context seam, never a buy/sell directive. It
|
||||
// computes where price sits relative to its own historical valuation corridor.
|
||||
//
|
||||
// Pure where possible: `computeCorridor`, `percentileIndex`, `buildSnapshot`
|
||||
// are pure; the cache-backed `resolveCorridorSnapshot` is a thin shim over the
|
||||
// shared cache (quote + candles + dividend fundamentals) — no direct vendor
|
||||
// I/O here (ADR-0009: everything funnels through the cache / adapter queue).
|
||||
|
||||
import type { DatabaseSync } from 'node:sqlite';
|
||||
import type { CacheRepository, PriceCandle, Quote } from '../cache/CacheRepository.ts';
|
||||
import type { CorridorSnapshot } from '../db/corridorRepository.ts';
|
||||
import { CorridorRepository } from '../db/corridorRepository.ts';
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Constants
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/** Trading days in the 1-year observable window. */
|
||||
export const CORRIDOR_1Y_BARS = 252;
|
||||
/** Trading days in the 90-day observable window. */
|
||||
export const CORRIDOR_90D_BARS = 63;
|
||||
/** Trading days in the 3-year market-level window (S&P 500 context). */
|
||||
export const CORRIDOR_3Y_BARS = 756;
|
||||
|
||||
/** Fraction of the 1y corridor below which the entry is "cheap". */
|
||||
export const ENTRY_CHEAP_PERCENTILE = 0.33;
|
||||
/** Fraction above which the entry is "stretched". */
|
||||
export const ENTRY_STRETCHED_PERCENTILE = 0.67;
|
||||
/** Implied upside (1y median reversion) above which the upside slot fires. */
|
||||
export const UPSIDE_HIGH_THRESHOLD = 0.15;
|
||||
/** Implied downside below which the downside slot fires. */
|
||||
export const UPSIDE_LOW_THRESHOLD = -0.10;
|
||||
/** Window (days) used by the corridor-method backtest grader. */
|
||||
export const BACKTEST_HORIZON_DAYS = 7;
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Pure helpers
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/** Median of a numeric array (finite elements); null when empty. */
|
||||
export function median(values: number[]): number | null {
|
||||
const finite = values.filter((v) => Number.isFinite(v)).sort((a, b) => a - b);
|
||||
if (finite.length === 0) return null;
|
||||
const mid = Math.floor(finite.length / 2);
|
||||
return finite.length % 2 === 0 ? (finite[mid - 1] + finite[mid]) / 2 : finite[mid];
|
||||
}
|
||||
|
||||
/**
|
||||
* The fractional rank (0..1) of `value` within `series`: the fraction of
|
||||
* `series` elements at or below `value`. Returns null when series is empty.
|
||||
* Pure.
|
||||
*/
|
||||
export function percentileIndex(value: number, series: number[]): number | null {
|
||||
const finite = series.filter((v) => Number.isFinite(v));
|
||||
if (finite.length === 0) return null;
|
||||
const below = finite.filter((v) => v <= value).length;
|
||||
return below / finite.length;
|
||||
}
|
||||
|
||||
/** A computed corridor window. Pure. */
|
||||
export interface CorridorWindow {
|
||||
high: number | null;
|
||||
low: number | null;
|
||||
median: number | null;
|
||||
}
|
||||
|
||||
/**
|
||||
* Compute a P/E corridor window from a lookback slice of a P/E series.
|
||||
* `series` is the full series (oldest → newest); `bars` is the window size.
|
||||
* Pure.
|
||||
*/
|
||||
export function computeCorridorWindow(series: number[], bars: number): CorridorWindow {
|
||||
const slice = series.length >= bars ? series.slice(series.length - bars) : series.slice();
|
||||
const finite = slice.filter((v) => Number.isFinite(v));
|
||||
if (finite.length === 0) return { high: null, low: null, median: null };
|
||||
return {
|
||||
high: Math.max(...finite),
|
||||
low: Math.min(...finite),
|
||||
median: median(finite),
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* Build a full corridor snapshot from a trailing P/E series and current prices.
|
||||
* Pure. `peSeries` is the trailing P/E series (oldest → newest); the latest
|
||||
* value is treated as the current P/E position.
|
||||
*/
|
||||
export function computeCorridor(
|
||||
peSeries: number[],
|
||||
currentPrice: number,
|
||||
forwardEPS: number | null,
|
||||
trailingEPS: number | null,
|
||||
trailingPE: number | null,
|
||||
forwardPE: number | null,
|
||||
): Omit<CorridorSnapshot, 'symbol' | 'snapshotDate' | 'dataSource' | 'createdAt'> {
|
||||
const currentPE = peSeries.length > 0 ? peSeries[peSeries.length - 1] : trailingPE ?? NaN;
|
||||
|
||||
const w1y = computeCorridorWindow(peSeries, CORRIDOR_1Y_BARS);
|
||||
const w90d = computeCorridorWindow(peSeries, CORRIDOR_90D_BARS);
|
||||
|
||||
const fairValue1y = forwardEPS !== null && forwardEPS > 0 && w1y.median !== null ? forwardEPS * w1y.median : null;
|
||||
const fairValue90d = forwardEPS !== null && forwardEPS > 0 && w90d.median !== null ? forwardEPS * w90d.median : null;
|
||||
|
||||
return {
|
||||
forwardPE,
|
||||
trailingPE,
|
||||
forwardEPS,
|
||||
trailingEPS,
|
||||
corridor1yHigh: w1y.high,
|
||||
corridor1yLow: w1y.low,
|
||||
corridor1yMedian: w1y.median,
|
||||
corridor90dHigh: w90d.high,
|
||||
corridor90dLow: w90d.low,
|
||||
corridor90dMedian: w90d.median,
|
||||
fairValue1y,
|
||||
fairValue90d,
|
||||
impliedUpside1y:
|
||||
currentPrice > 0 && fairValue1y !== null ? fairValue1y / currentPrice - 1 : null,
|
||||
impliedUpside90d:
|
||||
currentPrice > 0 && fairValue90d !== null ? fairValue90d / currentPrice - 1 : null,
|
||||
pePercentile1y:
|
||||
Number.isFinite(currentPE) ? percentileIndex(currentPE, peSeries.slice(Math.max(0, peSeries.length - CORRIDOR_1Y_BARS))) : null,
|
||||
pePercentile90d:
|
||||
Number.isFinite(currentPE) ? percentileIndex(currentPE, peSeries.slice(Math.max(0, peSeries.length - CORRIDOR_90D_BARS))) : null,
|
||||
currentPrice,
|
||||
};
|
||||
}
|
||||
|
||||
/** Derive a trailing P/E series from closes over `eps`. Pure. */
|
||||
export function trailingPeSeries(candles: PriceCandle[], eps: number): number[] {
|
||||
if (!eps || eps <= 0) return [];
|
||||
return candles
|
||||
.map((c) => (Number.isFinite(c.c) && c.c > 0 ? c.c / eps : NaN))
|
||||
.filter((v) => Number.isFinite(v));
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Cache-backed resolution seam
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
export interface CorridorFundamentalsInput {
|
||||
price: number | null;
|
||||
forwardPE: number | null;
|
||||
trailingPE: number | null;
|
||||
forwardEPS: number | null;
|
||||
trailingEPS: number | null;
|
||||
}
|
||||
|
||||
/** Paper a quote + dividend-fundamentals cache row into a corridor input. Pure. */
|
||||
export function fundamentalsFrom(quote: Quote | null, div: unknown): CorridorFundamentalsInput {
|
||||
const f = (div as Record<string, unknown> | null) ?? {};
|
||||
const num = (v: unknown): number | null => {
|
||||
if (v === null || v === undefined || typeof v === 'string' && v === '') return null;
|
||||
const n = Number(v);
|
||||
return Number.isFinite(n) ? n : null;
|
||||
};
|
||||
const forwardPE = num(f.forwardPE);
|
||||
const trailingPE = num(f.trailingPE);
|
||||
const forwardEPS = num(f.forwardEPS);
|
||||
const trailingEPS = num(f.trailingEPS);
|
||||
return {
|
||||
price: quote?.price != null && Number.isFinite(quote.price) ? quote.price : null,
|
||||
forwardPE,
|
||||
trailingPE,
|
||||
forwardEPS,
|
||||
// Prefer an explicit trailing EPS; else fall back to price / trailingPE.
|
||||
trailingEPS: trailingEPS ?? (trailingPE && trailingPE > 0 && quote?.price ? quote.price / trailingPE : null),
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolve + store a fresh corridor snapshot for `symbol` from the shared cache.
|
||||
* Uses the trailing P/E bootstrap: the trailing EPS drives the historical P/E
|
||||
* series immediately; forward PE/EPS enrich it when available. Idempotent per
|
||||
* (symbol, snapshotDate). Returns the stored snapshot.
|
||||
*/
|
||||
export async function resolveCorridorSnapshot(
|
||||
db: DatabaseSync,
|
||||
cache: CacheRepository,
|
||||
symbol: string,
|
||||
opts: { date?: string } = {},
|
||||
): Promise<CorridorSnapshot | null> {
|
||||
const sym = symbol.toUpperCase();
|
||||
const repo = new CorridorRepository(db);
|
||||
const snapshotDate = opts.date ?? new Date().toISOString().slice(0, 10);
|
||||
|
||||
const quoteEntry = await cache.get<Quote>(`yfinance:quote:${sym}`);
|
||||
const divEntry = await cache.get<Record<string, unknown>>(`yfinance:dividendFundamentals:${sym}`);
|
||||
const candleEntry = await cache.get<PriceCandle[]>(`yfinance:candles:${sym}:1d`);
|
||||
|
||||
const candles = (candleEntry?.value ?? []) as PriceCandle[];
|
||||
const f = fundamentalsFrom(quoteEntry?.value ?? null, divEntry?.value ?? null);
|
||||
|
||||
if (f.price === null && candles.length === 0) return null;
|
||||
const price = f.price ?? (candles.length > 0 ? candles[candles.length - 1].c : NaN);
|
||||
|
||||
const peSeries = f.trailingEPS !== null && f.trailingEPS > 0
|
||||
? trailingPeSeries(candles, f.trailingEPS)
|
||||
: [];
|
||||
|
||||
const computed = computeCorridor(
|
||||
peSeries,
|
||||
price,
|
||||
f.forwardEPS,
|
||||
f.trailingEPS,
|
||||
f.trailingPE ?? (f.trailingEPS && f.trailingEPS > 0 && price > 0 ? price / f.trailingEPS : null),
|
||||
f.forwardPE,
|
||||
);
|
||||
|
||||
// If no trailing EPS produced a real P/E series, compute one from forward PE.
|
||||
const snapshot: Omit<CorridorSnapshot, 'id'> = {
|
||||
symbol: sym,
|
||||
snapshotDate,
|
||||
...computed,
|
||||
dataSource: f.forwardPE !== null ? 'yfinance' : 'bootstrap_trailing',
|
||||
createdAt: new Date().toISOString(),
|
||||
};
|
||||
if (!Number.isFinite(snapshot.currentPrice) && candles.length > 0) {
|
||||
snapshot.currentPrice = candles[candles.length - 1].c;
|
||||
}
|
||||
|
||||
repo.saveSnapshot(snapshot);
|
||||
return repo.latestSnapshot(sym);
|
||||
}
|
||||
|
||||
/**
|
||||
* Compute a price-only corridor as a fallback when no EPS is available
|
||||
* (e.g. funds / unfamiliar tickers). Uses close-price percentiles instead of
|
||||
* P/E percentiles to still give a "where is price vs its own range" read.
|
||||
* Pure.
|
||||
*/
|
||||
export function priceOnlySnapshot(
|
||||
candles: PriceCandle[],
|
||||
currentPrice: number,
|
||||
snapshotDate: string,
|
||||
symbol: string,
|
||||
): Omit<CorridorSnapshot, 'id'> | null {
|
||||
if (candles.length === 0 || !Number.isFinite(currentPrice)) return null;
|
||||
const closes = candles.map((c) => c.c).filter((v) => Number.isFinite(v));
|
||||
if (closes.length === 0) return null;
|
||||
|
||||
const w1y = computeCorridorWindow(closes, CORRIDOR_1Y_BARS);
|
||||
const w90d = computeCorridorWindow(closes, CORRIDOR_90D_BARS);
|
||||
|
||||
return {
|
||||
symbol,
|
||||
snapshotDate,
|
||||
currentPrice,
|
||||
forwardPE: null,
|
||||
trailingPE: null,
|
||||
forwardEPS: null,
|
||||
trailingEPS: null,
|
||||
corridor1yHigh: w1y.high,
|
||||
corridor1yLow: w1y.low,
|
||||
corridor1yMedian: w1y.median,
|
||||
corridor90dHigh: w90d.high,
|
||||
corridor90dLow: w90d.low,
|
||||
corridor90dMedian: w90d.median,
|
||||
fairValue1y: null,
|
||||
fairValue90d: null,
|
||||
impliedUpside1y: null,
|
||||
impliedUpside90d: null,
|
||||
pePercentile1y:
|
||||
Number.isFinite(currentPrice) ? percentileIndex(currentPrice, closes.slice(Math.max(0, closes.length - CORRIDOR_1Y_BARS))) : null,
|
||||
pePercentile90d:
|
||||
Number.isFinite(currentPrice) ? percentileIndex(currentPrice, closes.slice(Math.max(0, closes.length - CORRIDOR_90D_BARS))) : null,
|
||||
dataSource: 'bootstrap_trailing',
|
||||
createdAt: new Date().toISOString(),
|
||||
};
|
||||
}
|
||||
@@ -0,0 +1,155 @@
|
||||
// Investor Flow — Price Corridor slot evaluator (M24, slice 3)
|
||||
//
|
||||
// Pure, snapshot-based assessments for the six `macro` corridor confluence
|
||||
// slots. Consumes the valuation-corridor snapshot computed by corridorData.ts
|
||||
// and, for the market-level slots, the SPY snapshot. Every function is a pure
|
||||
// (snapshot) → SlotAssessment[] mapping with ADR-0007 evidence notes — never a
|
||||
// recommendation.
|
||||
//
|
||||
// Distinct from the candle-driven technical evaluators: corridor slots read the
|
||||
// observable-multiple corridor, so they carry valuation evidence over and above
|
||||
// price-action evidence in the rack's picture.
|
||||
|
||||
import type { CorridorSnapshot } from '../db/corridorRepository.ts';
|
||||
import type { SlotAssessment } from './confluenceRack.ts';
|
||||
import {
|
||||
ENTRY_CHEAP_PERCENTILE,
|
||||
ENTRY_STRETCHED_PERCENTILE,
|
||||
UPSIDE_HIGH_THRESHOLD,
|
||||
UPSIDE_LOW_THRESHOLD,
|
||||
CORRIDOR_3Y_BARS,
|
||||
computeCorridorWindow,
|
||||
} from './corridorData.ts';
|
||||
|
||||
/** 3-year median forward-P/E reference for the S&P 500 (from @alojoh's reports:
|
||||
* 18.9x trough, 23.1x peak, ~20.5x median over the last three years). Used only
|
||||
* as the anchoring bench when a full 3y series is unavailable. */
|
||||
export const SPY_3Y_MEDIAN_REFERENCE = 20.5;
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Pure helpers
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
function assessment(id: string, state: 'fired' | 'not-fired', note: string): SlotAssessment {
|
||||
return { id, state, note };
|
||||
}
|
||||
|
||||
/**
|
||||
* Build the six price-corridor slot assessments for a symbol's snapshot.
|
||||
* `spySnapshot` supplies the market-level slots (SPY corridor vs its own 3y
|
||||
* median). Pure.
|
||||
*/
|
||||
export function evaluateCorridorSlots(
|
||||
snapshot: CorridorSnapshot,
|
||||
spySnapshot: CorridorSnapshot | null,
|
||||
): SlotAssessment[] {
|
||||
const out: SlotAssessment[] = [];
|
||||
|
||||
const pct1y = snapshot.pePercentile1y === null ? null : Number(snapshot.pePercentile1y);
|
||||
const pct90d = snapshot.pePercentile90d === null ? null : Number(snapshot.pePercentile90d);
|
||||
|
||||
// ----- corridor: entry cheap / stretched (1-year observable multiple) -----
|
||||
if (pct1y === null) {
|
||||
out.push(assessment('corridorEntryCheap', 'not-fired', 'No 1-year valuation corridor available for this symbol yet.'));
|
||||
out.push(assessment('corridorEntryStretched', 'not-fired', 'No 1-year valuation corridor available for this symbol yet.'));
|
||||
} else {
|
||||
if (pct1y < ENTRY_CHEAP_PERCENTILE) {
|
||||
out.push(assessment(
|
||||
'corridorEntryCheap',
|
||||
'fired',
|
||||
`Current P/E is in the lower ${(ENTRY_CHEAP_PERCENTILE * 100).toFixed(0)}% of its 1-year observable corridor (${(pct1y * 100).toFixed(0)}th percentile), indicating a relatively attractive entry point.`,
|
||||
));
|
||||
out.push(assessment('corridorEntryStretched', 'not-fired', `Current P/E sits at ${(pct1y * 100).toFixed(0)}th percentile of its 1-year corridor.`));
|
||||
} else if (pct1y > ENTRY_STRETCHED_PERCENTILE) {
|
||||
out.push(assessment(
|
||||
'corridorEntryStretched',
|
||||
'fired',
|
||||
`Current P/E is in the upper ${((1 - ENTRY_STRETCHED_PERCENTILE) * 100).toFixed(0)}% of its 1-year observable corridor (${(pct1y * 100).toFixed(0)}th percentile), indicating a stretched valuation vs its own history.`,
|
||||
));
|
||||
out.push(assessment('corridorEntryCheap', 'not-fired', `Current P/E sits at ${(pct1y * 100).toFixed(0)}th percentile of its 1-year corridor.`));
|
||||
} else {
|
||||
out.push(assessment('corridorEntryCheap', 'not-fired', `Current P/E sits mid-corridor at ${(pct1y * 100).toFixed(0)}th percentile of its 1-year range.`));
|
||||
out.push(assessment('corridorEntryStretched', 'not-fired', `Current P/E sits mid-corridor at ${(pct1y * 100).toFixed(0)}th percentile of its 1-year range.`));
|
||||
}
|
||||
}
|
||||
|
||||
// ----- corridor: implied upside / downside (median-multiple reversion) -----
|
||||
const upside1y = snapshot.impliedUpside1y === null ? null : Number(snapshot.impliedUpside1y);
|
||||
if (upside1y === null) {
|
||||
out.push(assessment('corridorUpsideHigh', 'not-fired', 'No forward EPS / median-multiple fair value available to quantify implied upside.'));
|
||||
out.push(assessment('corridorUpsideLow', 'not-fired', 'No forward EPS / median-multiple fair value available to quantify implied downside.'));
|
||||
} else {
|
||||
if (upside1y > UPSIDE_HIGH_THRESHOLD) {
|
||||
out.push(assessment(
|
||||
'corridorUpsideHigh',
|
||||
'fired',
|
||||
`Applying the 1-year median multiple to forward earnings implies ${(upside1y * 100).toFixed(1)}% upside from the current price.`,
|
||||
));
|
||||
out.push(assessment('corridorUpsideLow', 'not-fired', `1-year median-multiple fair value is ${(upside1y * 100).toFixed(1)}% vs current price.`));
|
||||
} else if (upside1y < UPSIDE_LOW_THRESHOLD) {
|
||||
out.push(assessment(
|
||||
'corridorUpsideLow',
|
||||
'fired',
|
||||
`Applying the 1-year median multiple to forward earnings implies ${(upside1y * 100).toFixed(1)}% downside from the current price.`,
|
||||
));
|
||||
out.push(assessment('corridorUpsideHigh', 'not-fired', `1-year median-multiple fair value is ${(upside1y * 100).toFixed(1)}% vs current price.`));
|
||||
} else {
|
||||
out.push(assessment('corridorUpsideHigh', 'not-fired', `1-year median-multiple fair value implies ${(upside1y * 100).toFixed(1)}% vs current price.`));
|
||||
out.push(assessment('corridorUpsideLow', 'not-fired', `1-year median-multiple fair value implies ${(upside1y * 100).toFixed(1)}% vs current price.`));
|
||||
}
|
||||
}
|
||||
|
||||
// ----- market level: SPY 3-year corridor (valuation tailwind / headwind) -----
|
||||
out.push(...spyCorridorAssessments(spySnapshot));
|
||||
|
||||
return out;
|
||||
}
|
||||
|
||||
/**
|
||||
* Market-level SPY corridor slots: is the broad market cheap or stretched
|
||||
* relative to its own 3-year forward-P/E corridor? Accepts either a stored SPY
|
||||
* snapshot (preferred) or a raw close-price series fallback. Pure.
|
||||
*/
|
||||
export function spyCorridorAssessments(spySnapshot: CorridorSnapshot | null): SlotAssessment[] {
|
||||
if (!spySnapshot) {
|
||||
return [
|
||||
assessment('spyCorridorCheap', 'not-fired', 'No SPY valuation-corridor snapshot available for the market-level context.'),
|
||||
assessment('spyCorridorStretched', 'not-fired', 'No SPY valuation-corridor snapshot available for the market-level context.'),
|
||||
];
|
||||
}
|
||||
|
||||
const pe = spySnapshot.forwardPE ?? spySnapshot.trailingPE ?? null;
|
||||
const cheap = pe !== null && pe < SPY_3Y_MEDIAN_REFERENCE;
|
||||
const peLabel = pe !== null ? pe.toFixed(1) : 'n/a';
|
||||
|
||||
if (cheap) {
|
||||
return [
|
||||
assessment(
|
||||
'spyCorridorCheap',
|
||||
'fired',
|
||||
`S&P 500 forward P/E (${peLabel}x) is below the 3-year median (~${SPY_3Y_MEDIAN_REFERENCE}x), a market-level valuation tailwind.`,
|
||||
),
|
||||
assessment('spyCorridorStretched', 'not-fired', `S&P 500 forward P/E (${peLabel}x) is below the 3-year median (~${SPY_3Y_MEDIAN_REFERENCE}x).`),
|
||||
];
|
||||
}
|
||||
return [
|
||||
assessment('spyCorridorCheap', 'not-fired', `S&P 500 forward P/E (${peLabel}x) is not below the 3-year median (~${SPY_3Y_MEDIAN_REFERENCE}x).`),
|
||||
assessment(
|
||||
'spyCorridorStretched',
|
||||
'fired',
|
||||
`S&P 500 forward P/E (${peLabel}x) is at or above the 3-year median (~${SPY_3Y_MEDIAN_REFERENCE}x), a market-level valuation headwind.`,
|
||||
),
|
||||
];
|
||||
}
|
||||
|
||||
/**
|
||||
* 3-year corridor window over a raw close-price series (for the SPY chart and
|
||||
* the market-level slot when only prices are cached). Pure.
|
||||
*/
|
||||
export function priceCorridor3y(candles: Array<{ c: number }>): { high: number | null; low: number | null; median: number | null } {
|
||||
const closes = candles.map((c) => c.c).filter((v) => Number.isFinite(v));
|
||||
return computeCorridorWindow(closes, CORRIDOR_3Y_BARS);
|
||||
}
|
||||
|
||||
/** Convenience re-export so corridor consumers share a single percentile helper. */
|
||||
export { percentileIndex } from './corridorData.ts';
|
||||
@@ -82,6 +82,9 @@ function runMigrations(db: DatabaseSync): void {
|
||||
`ALTER TABLE watchlists ADD COLUMN class_label TEXT`,
|
||||
// System-pinned demand symbols (rotation universe / benchmarks) survive unsubscribe
|
||||
`ALTER TABLE symbol_demand ADD COLUMN system_pin INTEGER NOT NULL DEFAULT 0`,
|
||||
// Tiered priority: 0=portfolio, 1=alert-critical, 2=watched, 3=background
|
||||
`ALTER TABLE symbol_demand ADD COLUMN tier INTEGER NOT NULL DEFAULT 3`,
|
||||
`ALTER TABLE symbol_demand ADD COLUMN last_viewed_at TEXT`,
|
||||
// Extended-hours quote fields (pre/post last vs RTH close)
|
||||
`ALTER TABLE quotes ADD COLUMN session TEXT`,
|
||||
`ALTER TABLE quotes ADD COLUMN regular_price REAL`,
|
||||
@@ -254,6 +257,56 @@ function runMigrations(db: DatabaseSync): void {
|
||||
verdict TEXT
|
||||
)`,
|
||||
`CREATE INDEX IF NOT EXISTS idx_confluence_signal_symbol_slot ON confluence_signal_history(symbol, slot_id)`,
|
||||
// Price Corridor module (M24): valuation-corridor snapshots, watchlist, and
|
||||
// the corridor-method backtest ledger for the offline grading of @alojoh's
|
||||
// weekly rankings against actual forward returns.
|
||||
`CREATE TABLE IF NOT EXISTS corridor_snapshots (
|
||||
id TEXT PRIMARY KEY,
|
||||
symbol TEXT NOT NULL,
|
||||
snapshot_date TEXT NOT NULL,
|
||||
current_price REAL NOT NULL,
|
||||
forward_pe REAL,
|
||||
trailing_pe REAL,
|
||||
forward_eps REAL,
|
||||
trailing_eps REAL,
|
||||
corridor_1y_high REAL,
|
||||
corridor_1y_low REAL,
|
||||
corridor_1y_median REAL,
|
||||
corridor_90d_high REAL,
|
||||
corridor_90d_low REAL,
|
||||
corridor_90d_median REAL,
|
||||
fair_value_1y REAL,
|
||||
fair_value_90d REAL,
|
||||
implied_upside_1y REAL,
|
||||
implied_upside_90d REAL,
|
||||
pe_percentile_1y REAL,
|
||||
pe_percentile_90d REAL,
|
||||
data_source TEXT NOT NULL DEFAULT 'yfinance',
|
||||
created_at TEXT NOT NULL,
|
||||
UNIQUE(symbol, snapshot_date)
|
||||
)`,
|
||||
`CREATE INDEX IF NOT EXISTS idx_corridor_snapshots_symbol_date ON corridor_snapshots(symbol, snapshot_date DESC)`,
|
||||
`CREATE TABLE IF NOT EXISTS corridor_watchlist (
|
||||
user_id TEXT NOT NULL REFERENCES users(id) ON DELETE CASCADE,
|
||||
symbol TEXT NOT NULL,
|
||||
added_at TEXT NOT NULL,
|
||||
PRIMARY KEY (user_id, symbol)
|
||||
)`,
|
||||
`CREATE TABLE IF NOT EXISTS corridor_backtest (
|
||||
id TEXT PRIMARY KEY,
|
||||
article_date TEXT NOT NULL,
|
||||
article_id TEXT NOT NULL,
|
||||
ranking_type TEXT NOT NULL,
|
||||
top_symbols_json TEXT NOT NULL,
|
||||
bottom_symbols_json TEXT NOT NULL,
|
||||
top_avg_return REAL,
|
||||
bottom_avg_return REAL,
|
||||
spread REAL,
|
||||
is_win INTEGER NOT NULL,
|
||||
horizon_days INTEGER NOT NULL,
|
||||
graded_at TEXT NOT NULL,
|
||||
UNIQUE(article_date, ranking_type)
|
||||
)`,
|
||||
];
|
||||
for (const sql of migrations) {
|
||||
try { db.exec(sql); } catch { /* column already exists */ }
|
||||
|
||||
@@ -0,0 +1,303 @@
|
||||
// Investor Flow — Corridor Repository (M24, slice 2)
|
||||
//
|
||||
// Thin data-access layer over the three price-corridor tables:
|
||||
// corridor_snapshots — daily valuation-corridor snapshots per symbol
|
||||
// corridor_watchlist — per-user tickers surfaced in the Corridor panel
|
||||
// corridor_backtest — corridor-method backtest ledger (offline grading)
|
||||
//
|
||||
// No business rules live here; pure persistence + retrieval. The corridor math
|
||||
// lives in corridorData.ts (pure) so this file stays a reflection of the schema.
|
||||
|
||||
import type { DatabaseSync } from 'node:sqlite';
|
||||
import { randomUUID } from 'node:crypto';
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Types
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/** A stored daily valuation-corridor snapshot for one symbol. */
|
||||
export interface CorridorSnapshot {
|
||||
symbol: string;
|
||||
snapshotDate: string;
|
||||
currentPrice: number;
|
||||
forwardPE: number | null;
|
||||
trailingPE: number | null;
|
||||
forwardEPS: number | null;
|
||||
trailingEPS: number | null;
|
||||
corridor1yHigh: number | null;
|
||||
corridor1yLow: number | null;
|
||||
corridor1yMedian: number | null;
|
||||
corridor90dHigh: number | null;
|
||||
corridor90dLow: number | null;
|
||||
corridor90dMedian: number | null;
|
||||
fairValue1y: number | null;
|
||||
fairValue90d: number | null;
|
||||
impliedUpside1y: number | null;
|
||||
impliedUpside90d: number | null;
|
||||
pePercentile1y: number | null;
|
||||
pePercentile90d: number | null;
|
||||
dataSource: string;
|
||||
createdAt: string;
|
||||
}
|
||||
|
||||
/** One graded corridor-method backtest row (exported shape, not the raw row). */
|
||||
export interface CorridorBacktestRow {
|
||||
id: string;
|
||||
articleDate: string;
|
||||
articleId: string;
|
||||
rankingType: string;
|
||||
topSymbols: string[];
|
||||
bottomSymbols: string[];
|
||||
topAvgReturn: number | null;
|
||||
bottomAvgReturn: number | null;
|
||||
spread: number | null;
|
||||
isWin: boolean;
|
||||
horizonDays: number;
|
||||
gradedAt: string;
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Prepared statements
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
function stmts(db: DatabaseSync) {
|
||||
return {
|
||||
// --- corridor_snapshots ---
|
||||
insertSnapshot: db.prepare(
|
||||
`INSERT INTO corridor_snapshots
|
||||
(id, symbol, snapshot_date, current_price, forward_pe, trailing_pe,
|
||||
forward_eps, trailing_eps, corridor_1y_high, corridor_1y_low,
|
||||
corridor_1y_median, corridor_90d_high, corridor_90d_low,
|
||||
corridor_90d_median, fair_value_1y, fair_value_90d,
|
||||
implied_upside_1y, implied_upside_90d, pe_percentile_1y,
|
||||
pe_percentile_90d, data_source, created_at)
|
||||
VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)
|
||||
ON CONFLICT(symbol, snapshot_date) DO UPDATE SET
|
||||
current_price=excluded.current_price, forward_pe=excluded.forward_pe,
|
||||
trailing_pe=excluded.trailing_pe, forward_eps=excluded.forward_eps,
|
||||
trailing_eps=excluded.trailing_eps,
|
||||
corridor_1y_high=excluded.corridor_1y_high,
|
||||
corridor_1y_low=excluded.corridor_1y_low,
|
||||
corridor_1y_median=excluded.corridor_1y_median,
|
||||
corridor_90d_high=excluded.corridor_90d_high,
|
||||
corridor_90d_low=excluded.corridor_90d_low,
|
||||
corridor_90d_median=excluded.corridor_90d_median,
|
||||
fair_value_1y=excluded.fair_value_1y,
|
||||
fair_value_90d=excluded.fair_value_90d,
|
||||
implied_upside_1y=excluded.implied_upside_1y,
|
||||
implied_upside_90d=excluded.implied_upside_90d,
|
||||
pe_percentile_1y=excluded.pe_percentile_1y,
|
||||
pe_percentile_90d=excluded.pe_percentile_90d,
|
||||
data_source=excluded.data_source, created_at=excluded.created_at`,
|
||||
),
|
||||
selectSnapshot: db.prepare(
|
||||
`SELECT id, symbol, snapshot_date, current_price, forward_pe, trailing_pe,
|
||||
forward_eps, trailing_eps, corridor_1y_high, corridor_1y_low,
|
||||
corridor_1y_median, corridor_90d_high, corridor_90d_low,
|
||||
corridor_90d_median, fair_value_1y, fair_value_90d,
|
||||
implied_upside_1y, implied_upside_90d, pe_percentile_1y,
|
||||
pe_percentile_90d, data_source, created_at
|
||||
FROM corridor_snapshots WHERE symbol = ? AND snapshot_date = ?`,
|
||||
),
|
||||
latestSnapshot: db.prepare(
|
||||
`SELECT id, symbol, snapshot_date, current_price, forward_pe, trailing_pe,
|
||||
forward_eps, trailing_eps, corridor_1y_high, corridor_1y_low,
|
||||
corridor_1y_median, corridor_90d_high, corridor_90d_low,
|
||||
corridor_90d_median, fair_value_1y, fair_value_90d,
|
||||
implied_upside_1y, implied_upside_90d, pe_percentile_1y,
|
||||
pe_percentile_90d, data_source, created_at
|
||||
FROM corridor_snapshots WHERE symbol = ?
|
||||
ORDER BY snapshot_date DESC, created_at DESC LIMIT 1`,
|
||||
),
|
||||
snapshotsForSymbol: db.prepare(
|
||||
`SELECT id, symbol, snapshot_date, current_price, forward_pe, trailing_pe,
|
||||
forward_eps, trailing_eps, corridor_1y_high, corridor_1y_low,
|
||||
corridor_1y_median, corridor_90d_high, corridor_90d_low,
|
||||
corridor_90d_median, fair_value_1y, fair_value_90d,
|
||||
implied_upside_1y, implied_upside_90d, pe_percentile_1y,
|
||||
pe_percentile_90d, data_source, created_at
|
||||
FROM corridor_snapshots WHERE symbol = ?
|
||||
ORDER BY snapshot_date DESC LIMIT ?`,
|
||||
),
|
||||
|
||||
// --- corridor_watchlist ---
|
||||
insertWatch: db.prepare(
|
||||
`INSERT OR IGNORE INTO corridor_watchlist (user_id, symbol, added_at)
|
||||
VALUES (?, ?, ?)`,
|
||||
),
|
||||
deleteWatch: db.prepare(
|
||||
`DELETE FROM corridor_watchlist WHERE user_id = ? AND symbol = ?`,
|
||||
),
|
||||
listWatch: db.prepare(
|
||||
`SELECT symbol, added_at FROM corridor_watchlist WHERE user_id = ? ORDER BY added_at DESC`,
|
||||
),
|
||||
|
||||
// --- corridor_backtest ---
|
||||
insertBacktest: db.prepare(
|
||||
`INSERT OR REPLACE INTO corridor_backtest
|
||||
(id, article_date, article_id, ranking_type, top_symbols_json,
|
||||
bottom_symbols_json, top_avg_return, bottom_avg_return, spread,
|
||||
is_win, horizon_days, graded_at)
|
||||
VALUES (?,?,?,?,?,?,?,?,?,?,?,?)`,
|
||||
),
|
||||
selectBacktests: db.prepare(
|
||||
`SELECT id, article_date, article_id, ranking_type, top_symbols_json,
|
||||
bottom_symbols_json, top_avg_return, bottom_avg_return, spread,
|
||||
is_win, horizon_days, graded_at
|
||||
FROM corridor_backtest ORDER BY article_date DESC`,
|
||||
),
|
||||
};
|
||||
}
|
||||
|
||||
function mapSnapshot(row: Record<string, unknown>): CorridorSnapshot {
|
||||
const num = (v: unknown): number | null => (v === null || v === undefined ? null : Number(v));
|
||||
return {
|
||||
symbol: row.symbol as string,
|
||||
snapshotDate: row.snapshot_date as string,
|
||||
currentPrice: Number(row.current_price),
|
||||
forwardPE: num(row.forward_pe),
|
||||
trailingPE: num(row.trailing_pe),
|
||||
forwardEPS: num(row.forward_eps),
|
||||
trailingEPS: num(row.trailing_eps),
|
||||
corridor1yHigh: num(row.corridor_1y_high),
|
||||
corridor1yLow: num(row.corridor_1y_low),
|
||||
corridor1yMedian: num(row.corridor_1y_median),
|
||||
corridor90dHigh: num(row.corridor_90d_high),
|
||||
corridor90dLow: num(row.corridor_90d_low),
|
||||
corridor90dMedian: num(row.corridor_90d_median),
|
||||
fairValue1y: num(row.fair_value_1y),
|
||||
fairValue90d: num(row.fair_value_90d),
|
||||
impliedUpside1y: num(row.implied_upside_1y),
|
||||
impliedUpside90d: num(row.implied_upside_90d),
|
||||
pePercentile1y: num(row.pe_percentile_1y),
|
||||
pePercentile90d: num(row.pe_percentile_90d),
|
||||
dataSource: row.data_source as string,
|
||||
createdAt: row.created_at as string,
|
||||
};
|
||||
}
|
||||
|
||||
function mapBacktest(row: Record<string, unknown>): CorridorBacktestRow {
|
||||
const parse = (raw: string): string[] => {
|
||||
try {
|
||||
const v = JSON.parse(raw) as unknown;
|
||||
return Array.isArray(v) ? v.filter((x): x is string => typeof x === 'string').map((s) => s.toUpperCase()) : [];
|
||||
} catch {
|
||||
return [];
|
||||
}
|
||||
};
|
||||
const num = (v: unknown): number | null => (v === null || v === undefined ? null : Number(v));
|
||||
return {
|
||||
id: row.id as string,
|
||||
articleDate: row.article_date as string,
|
||||
articleId: row.article_id as string,
|
||||
rankingType: row.ranking_type as string,
|
||||
topSymbols: parse(row.top_symbols_json as string),
|
||||
bottomSymbols: parse(row.bottom_symbols_json as string),
|
||||
topAvgReturn: num(row.top_avg_return),
|
||||
bottomAvgReturn: num(row.bottom_avg_return),
|
||||
spread: num(row.spread),
|
||||
isWin: Number(row.is_win) === 1,
|
||||
horizonDays: Number(row.horizon_days),
|
||||
gradedAt: row.graded_at as string,
|
||||
};
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Repository
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
export class CorridorRepository {
|
||||
private readonly db: DatabaseSync;
|
||||
|
||||
constructor(db: DatabaseSync) {
|
||||
this.db = db;
|
||||
}
|
||||
|
||||
// ------------------------------------------------------------ snapshots
|
||||
|
||||
saveSnapshot(snapshot: Omit<CorridorSnapshot, 'id'>): void {
|
||||
stmts(this.db).insertSnapshot.run(
|
||||
randomUUID(),
|
||||
snapshot.symbol,
|
||||
snapshot.snapshotDate,
|
||||
snapshot.currentPrice,
|
||||
snapshot.forwardPE,
|
||||
snapshot.trailingPE,
|
||||
snapshot.forwardEPS,
|
||||
snapshot.trailingEPS,
|
||||
snapshot.corridor1yHigh,
|
||||
snapshot.corridor1yLow,
|
||||
snapshot.corridor1yMedian,
|
||||
snapshot.corridor90dHigh,
|
||||
snapshot.corridor90dLow,
|
||||
snapshot.corridor90dMedian,
|
||||
snapshot.fairValue1y,
|
||||
snapshot.fairValue90d,
|
||||
snapshot.impliedUpside1y,
|
||||
snapshot.impliedUpside90d,
|
||||
snapshot.pePercentile1y,
|
||||
snapshot.pePercentile90d,
|
||||
snapshot.dataSource,
|
||||
snapshot.createdAt,
|
||||
);
|
||||
}
|
||||
|
||||
getSnapshot(symbol: string, date: string): CorridorSnapshot | null {
|
||||
const row = stmts(this.db).selectSnapshot.get(symbol, date) as Record<string, unknown> | undefined;
|
||||
return row ? mapSnapshot(row) : null;
|
||||
}
|
||||
|
||||
latestSnapshot(symbol: string): CorridorSnapshot | null {
|
||||
const row = stmts(this.db).latestSnapshot.get(symbol) as Record<string, unknown> | undefined;
|
||||
return row ? mapSnapshot(row) : null;
|
||||
}
|
||||
|
||||
snapshotsForSymbol(symbol: string, limit = 90): CorridorSnapshot[] {
|
||||
return (stmts(this.db).snapshotsForSymbol.all(symbol, limit) as Record<string, unknown>[]).map(mapSnapshot);
|
||||
}
|
||||
|
||||
// ------------------------------------------------------------- watchlist
|
||||
|
||||
addToWatchlist(userId: string, symbol: string): void {
|
||||
stmts(this.db).insertWatch.run(userId, symbol.toUpperCase(), new Date().toISOString());
|
||||
}
|
||||
|
||||
removeFromWatchlist(userId: string, symbol: string): void {
|
||||
stmts(this.db).deleteWatch.run(userId, symbol.toUpperCase());
|
||||
}
|
||||
|
||||
listWatchlist(userId: string): string[] {
|
||||
return (stmts(this.db).listWatch.all(userId) as Array<{ symbol: string }>).map((r) => r.symbol);
|
||||
}
|
||||
|
||||
/** Symbols in the user's open portfolio holdings, deduped. */
|
||||
listPortfolioSymbols(userId: string): string[] {
|
||||
const rows = this.db.prepare(
|
||||
`SELECT DISTINCT symbol FROM portfolio_holdings WHERE owner_id = ? AND status = 'open'`,
|
||||
).all(userId) as Array<{ symbol: string }>;
|
||||
return rows.map((r) => r.symbol);
|
||||
}
|
||||
|
||||
// --------------------------------------------------------------- backtest
|
||||
|
||||
saveBacktest(row: Omit<CorridorBacktestRow, 'id'>): void {
|
||||
stmts(this.db).insertBacktest.run(
|
||||
randomUUID(),
|
||||
row.articleDate,
|
||||
row.articleId,
|
||||
row.rankingType,
|
||||
JSON.stringify(row.topSymbols),
|
||||
JSON.stringify(row.bottomSymbols),
|
||||
row.topAvgReturn,
|
||||
row.bottomAvgReturn,
|
||||
row.spread,
|
||||
row.isWin ? 1 : 0,
|
||||
row.horizonDays,
|
||||
row.gradedAt,
|
||||
);
|
||||
}
|
||||
|
||||
listBacktests(): CorridorBacktestRow[] {
|
||||
return (stmts(this.db).selectBacktests.all() as Record<string, unknown>[]).map(mapBacktest);
|
||||
}
|
||||
}
|
||||
@@ -289,7 +289,9 @@ CREATE TABLE IF NOT EXISTS symbol_demand (
|
||||
ticker_kind TEXT NOT NULL,
|
||||
in_demand INTEGER NOT NULL DEFAULT 1, -- gate flag for the adapter queue
|
||||
last_refreshed_at TEXT,
|
||||
system_pin INTEGER NOT NULL DEFAULT 0 -- 1 = rotation/benchmark; survives unsubscribe
|
||||
system_pin INTEGER NOT NULL DEFAULT 0, -- 1 = rotation/benchmark; survives unsubscribe
|
||||
tier INTEGER NOT NULL DEFAULT 3, -- 0=portfolio, 1=alert-critical, 2=watched, 3=background
|
||||
last_viewed_at TEXT -- bumped on page view; decays back via recompute
|
||||
);
|
||||
|
||||
-- ===== Tier C — Per-user (ownerId NOT NULL) =====
|
||||
@@ -1016,4 +1018,61 @@ CREATE TABLE IF NOT EXISTS confluence_signal_history (
|
||||
resolved_at TEXT,
|
||||
verdict TEXT -- real|false_alarm
|
||||
);
|
||||
|
||||
-- ===== Price Corridor module (M24). =====
|
||||
-- Daily valuation-corridor snapshots per symbol: the observable forward/trailing
|
||||
-- P/E corridor the Corridor Method derives entry timing and implied upside from.
|
||||
CREATE TABLE IF NOT EXISTS corridor_snapshots (
|
||||
id TEXT PRIMARY KEY,
|
||||
symbol TEXT NOT NULL,
|
||||
snapshot_date TEXT NOT NULL, -- YYYY-MM-DD (trading day)
|
||||
current_price REAL NOT NULL,
|
||||
forward_pe REAL, -- current forward P/E (yfinance)
|
||||
trailing_pe REAL, -- current trailing P/E (yfinance)
|
||||
forward_eps REAL, -- price / forward_pe (derived)
|
||||
trailing_eps REAL,
|
||||
corridor_1y_high REAL, -- trailing P/E upper bound (252d)
|
||||
corridor_1y_low REAL, -- trailing P/E lower bound (252d)
|
||||
corridor_1y_median REAL, -- median trailing P/E (252d)
|
||||
corridor_90d_high REAL, -- trailing P/E upper bound (63d)
|
||||
corridor_90d_low REAL, -- trailing P/E lower bound (63d)
|
||||
corridor_90d_median REAL, -- median trailing P/E (63d)
|
||||
fair_value_1y REAL, -- forward_eps x corridor_1y_median
|
||||
fair_value_90d REAL, -- forward_eps x corridor_90d_median
|
||||
implied_upside_1y REAL, -- (fair_value_1y / current_price) - 1
|
||||
implied_upside_90d REAL, -- (fair_value_90d / current_price) - 1
|
||||
pe_percentile_1y REAL, -- 0..1 current trailing P/E vs 1y corridor
|
||||
pe_percentile_90d REAL, -- 0..1 current trailing P/E vs 90d corridor
|
||||
data_source TEXT NOT NULL DEFAULT 'yfinance', -- yfinance | bootstrap_trailing
|
||||
created_at TEXT NOT NULL,
|
||||
UNIQUE(symbol, snapshot_date)
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_corridor_snapshots_symbol_date
|
||||
ON corridor_snapshots(symbol, snapshot_date DESC);
|
||||
|
||||
-- Per-user corridor watchlist (tickers surfaced in the Corridor panel).
|
||||
CREATE TABLE IF NOT EXISTS corridor_watchlist (
|
||||
user_id TEXT NOT NULL REFERENCES users(id) ON DELETE CASCADE,
|
||||
symbol TEXT NOT NULL,
|
||||
added_at TEXT NOT NULL,
|
||||
PRIMARY KEY (user_id, symbol)
|
||||
);
|
||||
|
||||
-- Corridor Method backtest: grades @alojoh's weekly entry rankings against
|
||||
-- actual forward returns. An offline / periodic analysis, not live tracking.
|
||||
CREATE TABLE IF NOT EXISTS corridor_backtest (
|
||||
id TEXT PRIMARY KEY,
|
||||
article_date TEXT NOT NULL, -- YYYY-MM-DD of the ranked article
|
||||
article_id TEXT NOT NULL, -- X post id / source key
|
||||
ranking_type TEXT NOT NULL, -- entry_1y | entry_90d | upside
|
||||
top_symbols_json TEXT NOT NULL, -- JSON array of top-attractiveness names
|
||||
bottom_symbols_json TEXT NOT NULL, -- JSON array of least-attractive names
|
||||
top_avg_return REAL, -- mean N-day forward return (top)
|
||||
bottom_avg_return REAL, -- mean N-day forward return (bottom)
|
||||
spread REAL, -- top_avg_return - bottom_avg_return
|
||||
is_win INTEGER NOT NULL, -- 1 if spread > 0
|
||||
horizon_days INTEGER NOT NULL, -- forward window used (e.g. 7)
|
||||
graded_at TEXT NOT NULL,
|
||||
UNIQUE(article_date, ranking_type)
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_confluence_signal_symbol_slot ON confluence_signal_history(symbol, slot_id);
|
||||
|
||||
@@ -234,6 +234,39 @@ const housekeepTimer = setInterval(() => {
|
||||
}, DAILY_HOUSEKEEP_MS);
|
||||
housekeepTimer.unref();
|
||||
|
||||
// Confluence evaluation cycle: daily run of all wired slot family evaluators
|
||||
// over the confluence universe, persisting rack evaluations + signal history.
|
||||
// Idempotent per (symbol, asOf, rack) — cheap to run more often than daily.
|
||||
const CONFLUENCE_EVAL_MS = 60 * 60 * 1000; // hourly tick; re-eval only when stale
|
||||
const confluenceEvalTimer = setInterval(async () => {
|
||||
try {
|
||||
const { runConfluenceEvaluationCycle } = await import('./confluence/confluenceEngine.ts');
|
||||
const summary = await runConfluenceEvaluationCycle(database, cache);
|
||||
if (summary.symbolsEvaluated.length > 0) {
|
||||
console.log(
|
||||
`[confluence] evaluated=${summary.symbolsEvaluated.length} stored=${summary.evaluationsStored} ` +
|
||||
`reused=${summary.evaluationsReused} signals=${summary.signalsLogged} resolved=${summary.signalsResolved} ` +
|
||||
`corridorSnapshots=${summary.corridorSnapshots}`,
|
||||
);
|
||||
}
|
||||
} catch (e) {
|
||||
console.error('[confluence] evaluation cycle failed:', e);
|
||||
}
|
||||
}, CONFLUENCE_EVAL_MS);
|
||||
confluenceEvalTimer.unref();
|
||||
// First cycle shortly after boot once the EOD refresh has had a chance to land.
|
||||
setTimeout(async () => {
|
||||
try {
|
||||
const { runConfluenceEvaluationCycle } = await import('./confluence/confluenceEngine.ts');
|
||||
const summary = await runConfluenceEvaluationCycle(database, cache);
|
||||
if (summary.symbolsEvaluated.length > 0) {
|
||||
console.log(`[confluence] boot cycle: ${summary.symbolsEvaluated.length} symbols, ${summary.evaluationsStored} evaluations stored`);
|
||||
}
|
||||
} catch (e) {
|
||||
console.error('[confluence] boot cycle failed:', e);
|
||||
}
|
||||
}, 3 * 60 * 1000).unref();
|
||||
|
||||
|
||||
function readBody(req: IncomingMessage): Promise<string> {
|
||||
return new Promise((resolve, reject) => {
|
||||
|
||||
@@ -14,6 +14,7 @@ import type { CacheKey, SourceKind, CacheRepository, CacheScheduler } from '../c
|
||||
import {
|
||||
needsCandleRefresh,
|
||||
needsQuoteRefresh,
|
||||
needsTieredQuoteRefresh,
|
||||
needsSymbolMetaRefresh,
|
||||
parseCacheKey,
|
||||
} from '../cache/CacheRepository.ts';
|
||||
@@ -267,10 +268,10 @@ export class AdapterQueue implements CacheScheduler {
|
||||
}
|
||||
|
||||
/** One-time ops: remove junk test symbols and reset absurd refcounts. */
|
||||
cleanupDemandHygiene(): { removedJunk: number; cappedRefcounts: number; failedJunkCleared: number } {
|
||||
cleanupDemandHygiene(): { removedJunk: number; cappedRefcounts: number; failedJunkCleared: number; secPoisonCleared: number } {
|
||||
let removedJunk = 0;
|
||||
const demand = this._db.prepare('SELECT symbol, refcount, system_pin FROM symbol_demand').all() as Array<{
|
||||
symbol: string; refcount: number; system_pin: number | null;
|
||||
const demand = this._db.prepare('SELECT symbol, refcount, system_pin, ticker_kind FROM symbol_demand').all() as Array<{
|
||||
symbol: string; refcount: number; system_pin: number | null; ticker_kind: string | null;
|
||||
}>;
|
||||
for (const row of demand) {
|
||||
if (JUNK_SYMBOL_RE.test(row.symbol) && !row.system_pin) {
|
||||
@@ -287,10 +288,25 @@ export class AdapterQueue implements CacheScheduler {
|
||||
const failed = this._db.prepare(
|
||||
"DELETE FROM adapter_queue WHERE status='failed' AND (error LIKE 'quarantined:%' OR key LIKE '%:TEST%' OR key LIKE '%:FRESH%' OR key LIKE '%:ZZTEST%' OR key LIKE '%:FLOWTEST%' OR key LIKE '%:NEWTEST%')",
|
||||
).run();
|
||||
// Purge pending/backoff SEC jobs for ETF/index/crypto/fx symbols. These
|
||||
// never file 13F or SC 13G, so enqueuing them just manufactures a permanent
|
||||
// backlog (SPY/XLK/^VIX spin and 429 forever). secEligibleSymbols() blocks
|
||||
// new ones, but jobs queued before that guard land here and never die.
|
||||
let secPoisonCleared = 0;
|
||||
for (const row of demand) {
|
||||
const kind = (row.ticker_kind ?? 'equity').toLowerCase();
|
||||
if (kind === 'crypto' || kind === 'fx' || kind === 'index' || kind === 'etf') {
|
||||
const r = this._db.prepare(
|
||||
"DELETE FROM adapter_queue WHERE status IN ('pending','backoff') AND key LIKE '%' || ? || '%' AND (key LIKE 'sec-fetch:%' OR key LIKE 'sec-sc-fetch:%')",
|
||||
).run(row.symbol);
|
||||
secPoisonCleared += Number(r.changes);
|
||||
}
|
||||
}
|
||||
return {
|
||||
removedJunk,
|
||||
cappedRefcounts: Number(cap.changes),
|
||||
failedJunkCleared: Number(failed.changes),
|
||||
secPoisonCleared,
|
||||
};
|
||||
}
|
||||
|
||||
@@ -327,21 +343,64 @@ export class AdapterQueue implements CacheScheduler {
|
||||
).run(stuckCutoff);
|
||||
|
||||
// Prefer live marks (quote/candles) over secondary data so watchlist prices land first.
|
||||
// Sort by symbol tier (portfolio first) then cache kind then freshness.
|
||||
const tierMap = new Map<string, number>(
|
||||
(this._db.prepare('SELECT symbol, tier FROM symbol_demand').all() as Array<{ symbol: string; tier: number }>)
|
||||
.map((r) => [r.symbol, r.tier]),
|
||||
);
|
||||
/** Extract the symbol from a cache key (e.g. "yfinance:quote:AAPL" → "AAPL", "yfinance:candles:IREN:1d" → "IREN"). */
|
||||
const symFromKey = (key: string): string => {
|
||||
const idx = key.indexOf(':', key.indexOf(':') + 1);
|
||||
if (idx < 0) return '';
|
||||
const rest = key.slice(idx + 1);
|
||||
// Take everything up to the next colon (or end of string)
|
||||
const nextColon = rest.indexOf(':');
|
||||
return (nextColon > 0 ? rest.slice(0, nextColon) : rest).toUpperCase();
|
||||
};
|
||||
/** Kind rank: quote=0, candles=1, symbol=2, topHoldings=3, adjustments=4, else=5 */
|
||||
const kindRank = (key: string): number => {
|
||||
if (key.startsWith('yfinance:quote:')) return 0;
|
||||
if (key.startsWith('yfinance:candles:')) return 1;
|
||||
if (key.startsWith('yfinance:symbol:')) return 2;
|
||||
if (key.startsWith('yfinance:topHoldings:')) return 3;
|
||||
if (key.startsWith('yfinance:adjustments:')) return 4;
|
||||
return 5;
|
||||
};
|
||||
const jobs = this._db.prepare(
|
||||
`SELECT key, retry_count, backoff_until, scheduled_for FROM adapter_queue
|
||||
WHERE status IN ('pending','backoff')
|
||||
ORDER BY
|
||||
CASE
|
||||
WHEN key LIKE 'yfinance:quote:%' THEN 0
|
||||
WHEN key LIKE 'yfinance:candles:%' THEN 1
|
||||
WHEN key LIKE 'yfinance:symbol:%' THEN 2
|
||||
WHEN key LIKE 'yfinance:topHoldings:%' THEN 3
|
||||
WHEN key LIKE 'yfinance:adjustments:%' THEN 4
|
||||
ELSE 5
|
||||
END,
|
||||
(last_attempt IS NULL) DESC,
|
||||
last_attempt ASC`,
|
||||
).all() as Array<{ key: string; retry_count: number; backoff_until?: string | null; scheduled_for?: string | null }>;
|
||||
`SELECT key, retry_count, backoff_until, scheduled_for, last_attempt FROM adapter_queue
|
||||
WHERE status IN ('pending','backoff')`,
|
||||
).all() as Array<{ key: string; retry_count: number; backoff_until?: string | null; scheduled_for?: string | null; last_attempt?: string | null }>;
|
||||
|
||||
// Tier-aware sort: portfolio symbols (T0) first, then T1, T2, T3.
|
||||
// Within each tier: quote > candles > symbol > topHoldings > adjustments > other.
|
||||
// Within each kind: never-attempted first, then oldest first.
|
||||
jobs.sort((a, b) => {
|
||||
const ta = tierMap.get(symFromKey(a.key)) ?? 99;
|
||||
const tb = tierMap.get(symFromKey(b.key)) ?? 99;
|
||||
if (ta !== tb) return ta - tb;
|
||||
const ka = kindRank(a.key);
|
||||
const kb = kindRank(b.key);
|
||||
if (ka !== kb) return ka - kb;
|
||||
const na = a.last_attempt == null ? 0 : 1;
|
||||
const nb = b.last_attempt == null ? 0 : 1;
|
||||
if (na !== nb) return na - nb;
|
||||
if (a.last_attempt && b.last_attempt) return a.last_attempt.localeCompare(b.last_attempt);
|
||||
return 0;
|
||||
});
|
||||
|
||||
// Backlog pressure mode: when a critical quote backlog exists, defer
|
||||
// non-critical yfinance kinds so watchlist prices land first. Without
|
||||
// this, chain/shortinterest/topHoldings jobs steal the (serial) yfinance
|
||||
// source chain and quotes can stall for 26+ hours behind a 160-job pile.
|
||||
// Threshold: > 25 pending quote jobs means > 1 drain cycle of quotes alone.
|
||||
const pendingQuotes = this._db.prepare(
|
||||
"SELECT COUNT(*) AS c FROM adapter_queue WHERE status='pending' AND key LIKE 'yfinance:quote:%'",
|
||||
).get() as { c: number };
|
||||
const backlogPressure = pendingQuotes.c > 25;
|
||||
const NON_CRITICAL_YF_KINDS = new Set([
|
||||
'chain', 'shortinterest', 'dividendFundamentals',
|
||||
'topHoldings', 'expiry_dates', 'symbol', 'adjustments',
|
||||
]);
|
||||
|
||||
let processed = 0;
|
||||
const kindUsed: Record<string, number> = {};
|
||||
@@ -392,6 +451,11 @@ export class AdapterQueue implements CacheScheduler {
|
||||
const used = kindUsed[kind] ?? 0;
|
||||
if (used >= kindBudget) continue;
|
||||
|
||||
// Backlog pressure: defer non-critical yfinance kinds so quotes drain first.
|
||||
if (backlogPressure && source === 'yfinance' && NON_CRITICAL_YF_KINDS.has(kind)) continue;
|
||||
// Backlog pressure: skip T3 (background) symbols entirely so portfolio/alert symbols drain first.
|
||||
if (backlogPressure && sym && (tierMap.get(sym) ?? 99) >= 3) continue;
|
||||
|
||||
// Source-wide cool-down: skip all jobs for this vendor until the window ends.
|
||||
if (this.isSourceCoolingDown(source, Date.now())) continue;
|
||||
|
||||
@@ -422,6 +486,10 @@ export class AdapterQueue implements CacheScheduler {
|
||||
if (family && this.isVendorFamilyCoolingDown(family, Date.now())) continue;
|
||||
const kindBudget = DRAIN_KIND_BUDGET[kind] ?? DRAIN_KIND_BUDGET._default;
|
||||
if ((kindUsed[kind] ?? 0) >= kindBudget) continue;
|
||||
// Backlog pressure: defer non-critical yfinance kinds so quotes drain first.
|
||||
if (backlogPressure && source === 'yfinance' && NON_CRITICAL_YF_KINDS.has(kind)) continue;
|
||||
// Backlog pressure: skip T3 (background) symbols entirely so portfolio/alert symbols drain first.
|
||||
if (backlogPressure && sym && (tierMap.get(sym) ?? 99) >= 3) continue;
|
||||
specs.push({ key: job.key, source, family, sym, attempt: job.retry_count + 1 });
|
||||
processed += 1;
|
||||
kindUsed[kind] = (kindUsed[kind] ?? 0) + 1;
|
||||
@@ -454,36 +522,60 @@ export class AdapterQueue implements CacheScheduler {
|
||||
|
||||
// Re-check cool-downs now that it is this job's turn (another path may
|
||||
// have cooled the source/family while earlier jobs were running).
|
||||
if (this.isSourceCoolingDown(source, Date.now())) return;
|
||||
if (family && this.isVendorFamilyCoolingDown(family, Date.now())) return;
|
||||
if (this.isSourceCoolingDown(source, Date.now())) {
|
||||
this._db.prepare(
|
||||
"UPDATE adapter_queue SET status='backoff', last_attempt=?, retry_count=?, backoff_until=?, error=? WHERE key=?",
|
||||
).run(new Date().toISOString(), attempt, new Date(Date.now() + 30_000).toISOString(), 'source cooling', key);
|
||||
return;
|
||||
}
|
||||
if (family && this.isVendorFamilyCoolingDown(family, Date.now())) {
|
||||
this._db.prepare(
|
||||
"UPDATE adapter_queue SET status='backoff', last_attempt=?, retry_count=?, backoff_until=?, error=? WHERE key=?",
|
||||
).run(new Date().toISOString(), attempt, new Date(Date.now() + 30_000).toISOString(), 'vendor family cooling', key);
|
||||
return;
|
||||
}
|
||||
|
||||
const last = this._lastFetchAt[source] ?? 0;
|
||||
const wait = (this._rate[source] ?? 0) - (Date.now() - last);
|
||||
if (wait > 0) await sleep(wait);
|
||||
|
||||
// Re-check cool-down after sleep (another path may have set it).
|
||||
if (this.isSourceCoolingDown(source, Date.now())) return;
|
||||
if (family && this.isVendorFamilyCoolingDown(family, Date.now())) return;
|
||||
if (this.isSourceCoolingDown(source, Date.now())) {
|
||||
this._db.prepare(
|
||||
"UPDATE adapter_queue SET status='backoff', last_attempt=?, retry_count=?, backoff_until=?, error=? WHERE key=?",
|
||||
).run(new Date().toISOString(), attempt, new Date(Date.now() + 30_000).toISOString(), 'source cooling', key);
|
||||
return;
|
||||
}
|
||||
if (family && this.isVendorFamilyCoolingDown(family, Date.now())) {
|
||||
this._db.prepare(
|
||||
"UPDATE adapter_queue SET status='backoff', last_attempt=?, retry_count=?, backoff_until=?, error=? WHERE key=?",
|
||||
).run(new Date().toISOString(), attempt, new Date(Date.now() + 30_000).toISOString(), 'vendor family cooling', key);
|
||||
return;
|
||||
}
|
||||
|
||||
this._lastFetchAt[source] = Date.now();
|
||||
this._setStatus(key, 'in_flight');
|
||||
try {
|
||||
const res = await adapter.fetchOne(key);
|
||||
const FETCH_TIMEOUT_MS = 30_000;
|
||||
const fetchPromise = adapter.fetchOne(key);
|
||||
const timeoutPromise = new Promise<never>((_, reject) =>
|
||||
setTimeout(() => reject(new Error(`fetchOne timeout after ${FETCH_TIMEOUT_MS}ms for ${key}`)), FETCH_TIMEOUT_MS),
|
||||
);
|
||||
const res = await Promise.race([fetchPromise, timeoutPromise]);
|
||||
try { await this._cache?.set(key, res.value, res.ttlClass, res.provenance); } catch { /* adapter may persist directly */ }
|
||||
this.clearSourceCooldown(source);
|
||||
this._db.prepare("UPDATE adapter_queue SET status='done', last_attempt=?, error=NULL WHERE key=?").run(new Date().toISOString(), key);
|
||||
// After timeline posts land in x_cookie_posts, re-materialize fund captures.
|
||||
// (Schedule-time ingest runs *before* jobs finish and misses new posts.)
|
||||
if (key.startsWith('x:timeline:')) {
|
||||
const handle = key.slice('x:timeline:'.length);
|
||||
try {
|
||||
const { ingestAllFundCaptures } = await import('../services/captureIngest.ts');
|
||||
const stats = ingestAllFundCaptures(this._db);
|
||||
for (const s2 of stats) {
|
||||
if (s2.inserted > 0 || s2.captures > 0) {
|
||||
console.log(
|
||||
`[x-capture] post-timeline ${s2.fundId}: ${s2.captures} captures (${s2.inserted} new, ${s2.refreshed} refreshed)`,
|
||||
);
|
||||
}
|
||||
const { ingestByHandle } = await import('../services/captureIngest.ts');
|
||||
const s2 = ingestByHandle(this._db, handle);
|
||||
if (s2 && (s2.inserted > 0 || s2.captures > 0)) {
|
||||
console.log(
|
||||
`[x-capture] post-timeline ${s2.fundId}: ${s2.captures} captures (${s2.inserted} new, ${s2.refreshed} refreshed)`,
|
||||
);
|
||||
}
|
||||
} catch (capErr) {
|
||||
console.error('[x-capture] post-timeline ingest failed:', capErr instanceof Error ? capErr.message : capErr);
|
||||
@@ -610,6 +702,14 @@ this._db.prepare("UPDATE adapter_queue SET status='backoff', last_attempt=?, ret
|
||||
return [...new Set([...pins, ...rest])];
|
||||
}
|
||||
|
||||
/** Symbols at a specific tier — used by per-tier schedule branches. */
|
||||
private demandSymbolsByTier(tier: number): string[] {
|
||||
const rows = this._db.prepare(
|
||||
'SELECT symbol FROM symbol_demand WHERE tier = ? AND (in_demand = 1 OR COALESCE(system_pin, 0) = 1) ORDER BY symbol',
|
||||
).all(tier) as Array<{ symbol: string }>;
|
||||
return rows.filter((r) => !this.isSymbolQuarantined(r.symbol)).map((r) => r.symbol);
|
||||
}
|
||||
|
||||
/**
|
||||
* SEC equity-filings only touch equity-like tickers. ETFs/indexes/crypto/fx
|
||||
* never file 13F or SC 13G, so enqueueing them just manufactures a permanent
|
||||
@@ -629,6 +729,10 @@ this._db.prepare("UPDATE adapter_queue SET status='backoff', last_attempt=?, ret
|
||||
}
|
||||
|
||||
async enqueueDueSchedules(): Promise<void> {
|
||||
// Recompute symbol tiers from source tables (portfolio/alerts/watchlists/page views).
|
||||
// Self-healing: closes a holding → symbol decays to background on next tick.
|
||||
this.recomputeDemandTiers();
|
||||
|
||||
const now = new Date().toISOString();
|
||||
const due = this._db.prepare("SELECT * FROM queue_schedules WHERE next_enqueue IS NOT NULL AND next_enqueue <= ?").all(now) as Array<{ source_kind: string; interval_ms: number; last_enqueued: string | null; next_enqueue: string | null }>;
|
||||
for (const s of due) {
|
||||
@@ -663,12 +767,33 @@ this._db.prepare("UPDATE adapter_queue SET status='backoff', last_attempt=?, ret
|
||||
} else if (s.source_kind === 'sec-tickers') {
|
||||
await this.queue('sec-tickers:companyTickers:latest');
|
||||
} else if (s.source_kind === 'yfinance' || s.source_kind === 'yfinance-quote') {
|
||||
// Legacy 'yfinance' treated as quote tier
|
||||
// Legacy 'yfinance' / 'yfinance-quote' treated as quote tier (all demand symbols)
|
||||
for (const sym of symbols) {
|
||||
if (needsQuoteRefresh(d, sym)) {
|
||||
await this.queue(`yfinance:quote:${sym}`);
|
||||
}
|
||||
}
|
||||
} else if (s.source_kind === 'yfinance-quote-portfolio') {
|
||||
// Tier 0: portfolio holdings — tightest TTL, highest drain priority
|
||||
for (const sym of this.demandSymbolsByTier(0)) {
|
||||
if (needsTieredQuoteRefresh(d, sym, 0)) {
|
||||
await this.queue(`yfinance:quote:${sym}`);
|
||||
}
|
||||
}
|
||||
} else if (s.source_kind === 'yfinance-quote-priority') {
|
||||
// Tier 1: alert-critical symbols
|
||||
for (const sym of this.demandSymbolsByTier(1)) {
|
||||
if (needsTieredQuoteRefresh(d, sym, 1)) {
|
||||
await this.queue(`yfinance:quote:${sym}`);
|
||||
}
|
||||
}
|
||||
} else if (s.source_kind === 'yfinance-quote-watched') {
|
||||
// Tier 2: user watchlist symbols + recently viewed pages
|
||||
for (const sym of this.demandSymbolsByTier(2)) {
|
||||
if (needsTieredQuoteRefresh(d, sym, 2)) {
|
||||
await this.queue(`yfinance:quote:${sym}`);
|
||||
}
|
||||
}
|
||||
} else if (s.source_kind === 'yfinance-eod') {
|
||||
for (const sym of symbols) {
|
||||
if (needsCandleRefresh(d, sym)) {
|
||||
@@ -747,7 +872,12 @@ this._db.prepare("UPDATE adapter_queue SET status='backoff', last_attempt=?, ret
|
||||
}
|
||||
|
||||
// Alert-critical: re-queue SEC for demand symbols with broken/stale institutional pipeline.
|
||||
await this.requeueUnhealthySecSymbols();
|
||||
// Skip if SEC schedule is disabled (next_enqueue set far in the future).
|
||||
const secSchedule = this._db.prepare("SELECT next_enqueue FROM queue_schedules WHERE source_kind='sec-fetch'").get() as { next_enqueue: string | null } | undefined;
|
||||
const secDisabled = secSchedule?.next_enqueue && Date.parse(secSchedule.next_enqueue) > Date.now() + 365 * 86_400_000;
|
||||
if (!secDisabled) {
|
||||
await this.requeueUnhealthySecSymbols();
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -868,6 +998,56 @@ this._db.prepare("UPDATE adapter_queue SET status='backoff', last_attempt=?, ret
|
||||
return enqueued;
|
||||
}
|
||||
|
||||
/**
|
||||
* Recompute tier from source tables. Self-healing: if a holding is closed
|
||||
* or an alert subscription removed, the symbol naturally decays to background.
|
||||
* Tier 0 = portfolio (any open holding), 1 = alert-critical, 2 = watched
|
||||
* (user watchlist OR recently viewed page), 3 = background.
|
||||
*/
|
||||
recomputeDemandTiers(): void {
|
||||
const now = new Date();
|
||||
const VIEWED_DECAY_MS = 10 * 60_000; // page-view bump decays after 10 min
|
||||
const cutoff = new Date(now.getTime() - VIEWED_DECAY_MS).toISOString();
|
||||
|
||||
// Set in_demand=1 for all symbols that are still relevant (not orphaned).
|
||||
// Orphaned symbols: refcount=0, system_pin=0, no portfolio, no alert, no watchlist.
|
||||
// For those, set tier=3 and in_demand=0.
|
||||
this._db.prepare(`
|
||||
UPDATE symbol_demand SET tier = (
|
||||
SELECT MIN(t) FROM (
|
||||
SELECT 0 AS t WHERE EXISTS (
|
||||
SELECT 1 FROM portfolio_holdings WHERE symbol = symbol_demand.symbol AND status = 'open'
|
||||
)
|
||||
UNION ALL
|
||||
SELECT 1 WHERE EXISTS (
|
||||
SELECT 1 FROM alerts WHERE symbol = symbol_demand.symbol AND enabled = 1
|
||||
)
|
||||
UNION ALL
|
||||
SELECT 2 WHERE EXISTS (
|
||||
SELECT 1 FROM watchlists WHERE kind = 'user' AND symbols LIKE '%' || symbol_demand.symbol || '%'
|
||||
)
|
||||
UNION ALL
|
||||
SELECT 2 WHERE symbol_demand.last_viewed_at IS NOT NULL AND symbol_demand.last_viewed_at > ?
|
||||
UNION ALL
|
||||
SELECT 3
|
||||
)
|
||||
)
|
||||
`).run(cutoff);
|
||||
|
||||
// System pins that aren't boosted above T3 by portfolio/alerts/watchlists
|
||||
// stay at T3 (background, no schedule). They still exist in demand but
|
||||
// don't get scheduled - only fetched on-demand via page view bump.
|
||||
// Symbols with refcount=0, system_pin=0, and no source match become T3
|
||||
// and get in_demand=0 so they don't clutter the queue.
|
||||
this._db.prepare(`
|
||||
UPDATE symbol_demand SET in_demand = 0
|
||||
WHERE tier = 3
|
||||
AND COALESCE(system_pin, 0) = 0
|
||||
AND COALESCE(refcount, 0) = 0
|
||||
AND last_viewed_at IS NULL
|
||||
`).run();
|
||||
}
|
||||
|
||||
/**
|
||||
* Seed / migrate schedules to tiered yfinance plan.
|
||||
* Safe to call every boot: upserts missing tiers; migrates legacy single `yfinance` row.
|
||||
@@ -880,10 +1060,10 @@ this._db.prepare("UPDATE adapter_queue SET status='backoff', last_attempt=?, ret
|
||||
for (const [kind, ms] of defaults) {
|
||||
insert.run(kind, ms, null, new Date(Date.now() + Math.min(ms, 60_000)).toISOString());
|
||||
}
|
||||
// Drop legacy monolithic yfinance schedule if tiered ones exist
|
||||
const hasQuote = this._db.prepare("SELECT 1 FROM queue_schedules WHERE source_kind='yfinance-quote'").get();
|
||||
if (hasQuote) {
|
||||
this._db.prepare("DELETE FROM queue_schedules WHERE source_kind='yfinance'").run();
|
||||
// Drop legacy monolithic yfinance-quote schedule if tiered ones exist
|
||||
const hasTierSchedule = this._db.prepare("SELECT 1 FROM queue_schedules WHERE source_kind='yfinance-quote-portfolio'").get();
|
||||
if (hasTierSchedule) {
|
||||
this._db.prepare("DELETE FROM queue_schedules WHERE source_kind IN ('yfinance', 'yfinance-quote')").run();
|
||||
}
|
||||
// Ensure finra-bulk is not auto-seeded (403); delete if present from old seeds
|
||||
this._db.prepare("DELETE FROM queue_schedules WHERE source_kind='finra-bulk'").run();
|
||||
|
||||
@@ -213,10 +213,13 @@ test('seedDefaultSchedules creates tiered yfinance schedules and drops legacy',
|
||||
queue.seedDefaultSchedules();
|
||||
const kinds = (db.prepare('SELECT source_kind FROM queue_schedules ORDER BY source_kind').all() as Array<{ source_kind: string }>)
|
||||
.map((r) => r.source_kind);
|
||||
assert.ok(kinds.includes('yfinance-quote'));
|
||||
assert.ok(kinds.includes('yfinance-quote-portfolio'));
|
||||
assert.ok(kinds.includes('yfinance-quote-priority'));
|
||||
assert.ok(kinds.includes('yfinance-quote-watched'));
|
||||
assert.ok(kinds.includes('yfinance-eod'));
|
||||
assert.ok(kinds.includes('yfinance-meta'));
|
||||
assert.ok(kinds.includes('yfinance-holdings'));
|
||||
assert.ok(!kinds.includes('yfinance-quote'), 'legacy monolithic yfinance-quote schedule removed');
|
||||
assert.ok(!kinds.includes('yfinance'), 'legacy monolithic yfinance schedule removed');
|
||||
assert.ok(!kinds.includes('finra-bulk'), 'finra-bulk must not auto-seed');
|
||||
});
|
||||
@@ -228,12 +231,13 @@ test('TTL-aware quote schedule only enqueues stale quotes', async () => {
|
||||
await cache.set('yfinance:quote:NVDA', { symbol: 'NVDA', price: 100 } as import('../../cache/CacheRepository.ts').Quote, 'live_quote', {
|
||||
fetchedAt: new Date().toISOString(), sourceKind: 'yfinance',
|
||||
});
|
||||
await cache.subscribe('NVDA', 'equity');
|
||||
await cache.subscribe('AAPL', 'equity'); // no quote → stale
|
||||
// Bump both to watched tier (2) — the schedule that serves on-demand symbols
|
||||
await cache.bumpToWatched('NVDA', 'equity');
|
||||
await cache.bumpToWatched('AAPL', 'equity'); // no quote → stale
|
||||
// Clear any seed jobs
|
||||
db.prepare("DELETE FROM adapter_queue").run();
|
||||
// Force quote schedule due
|
||||
db.prepare("UPDATE queue_schedules SET next_enqueue=? WHERE source_kind='yfinance-quote'")
|
||||
db.prepare("UPDATE queue_schedules SET next_enqueue=? WHERE source_kind='yfinance-quote-watched'")
|
||||
.run(new Date(Date.now() - 1000).toISOString());
|
||||
await queue.enqueueDueSchedules();
|
||||
const pending = (db.prepare("SELECT key FROM adapter_queue WHERE status='pending'").all() as Array<{ key: string }>)
|
||||
|
||||
@@ -12,7 +12,7 @@ import type { SourceKind } from '../cache/CacheRepository.ts';
|
||||
/** Steady-state min gap between successful fetches for a source (ms). */
|
||||
export const DEFAULT_SOURCE_MIN_INTERVAL_MS: Record<SourceKind, number> = {
|
||||
// Job-level spacing (AdapterQueue). Per-request spacing is owned by vendorGate.
|
||||
yfinance: 2_500,
|
||||
yfinance: 2_000,
|
||||
sec: 400,
|
||||
'sec-fetch': 8_000,
|
||||
'sec-sc-fetch': 8_000,
|
||||
@@ -204,6 +204,22 @@ export function quoteTtlMs(now = new Date()): number {
|
||||
return 15 * 60_000;
|
||||
}
|
||||
|
||||
/** Per-tier quote TTL. Portfolio (T0) gets the tightest freshness target. */
|
||||
export function tieredQuoteTtlMs(tier: number, now = new Date()): number {
|
||||
if (tier === 0) {
|
||||
// Portfolio: 60s RTH, 120s extended, 600s closed
|
||||
if (isUsRegularHours(now)) return 60_000;
|
||||
if (isUsExtendedHours(now)) return 2 * 60_000;
|
||||
return 10 * 60_000;
|
||||
}
|
||||
if (tier === 1) {
|
||||
// Alert-critical: same as current quoteTtlMs
|
||||
return quoteTtlMs(now);
|
||||
}
|
||||
// Watched (T2): 10 min flat
|
||||
return 10 * 60_000;
|
||||
}
|
||||
|
||||
/** Candle bar considered current if last bar is within this age (ms). */
|
||||
export const CANDLE_FRESH_MS = 36 * 60 * 60_000; // 36h covers weekends lightly
|
||||
|
||||
@@ -212,7 +228,9 @@ export const SYMBOL_META_INCOMPLETE_TTL_MS = 6 * 60 * 60_000; // 6h
|
||||
|
||||
/** Default schedule intervals (ms). */
|
||||
export const SCHEDULE_INTERVALS = {
|
||||
'yfinance-quote': 5 * 60_000, // 5 min
|
||||
'yfinance-quote-portfolio': 60_000, // 1 min — portfolio holdings
|
||||
'yfinance-quote-priority': 5 * 60_000, // 5 min — alert-critical symbols
|
||||
'yfinance-quote-watched': 10 * 60_000, // 10 min — user watchlist symbols
|
||||
'yfinance-eod': 6 * 60 * 60_000, // 6h (incremental candles)
|
||||
'yfinance-meta': 24 * 60 * 60_000, // daily symbol meta
|
||||
'yfinance-holdings': 7 * 24 * 60 * 60_000, // weekly ETF holdings
|
||||
|
||||
@@ -144,6 +144,13 @@ function normalizePostedAtDate(postedAt: string): string {
|
||||
return postedAt.slice(0, 10);
|
||||
}
|
||||
|
||||
/** Materialize captures for a specific handle (called after its timeline job completes). */
|
||||
export function ingestByHandle(db: DatabaseSync, handle: string): CaptureIngestStats | null {
|
||||
const fund = db.prepare('SELECT id FROM tracked_funds WHERE lower(x_handle) = lower(?) AND enabled = 1').get(handle) as { id: string } | undefined;
|
||||
if (!fund) return null;
|
||||
return ingestFundCaptures(db, fund.id);
|
||||
}
|
||||
|
||||
/** Materialize captures for every enabled tracked fund with an x_handle. */
|
||||
export function ingestAllFundCaptures(db: DatabaseSync): CaptureIngestStats[] {
|
||||
return listTrackedFundHandles(db)
|
||||
|
||||
@@ -151,7 +151,7 @@ function seedBuiltIns(): void {
|
||||
registerVendorIntegration({
|
||||
family: 'yfinance',
|
||||
sourceKinds: ['yfinance', 'yfinance-quote', 'yfinance-eod', 'yfinance-meta', 'yfinance-holdings'],
|
||||
policy: { minIntervalMs: 400, maxInflight: 1, drainJobBudget: 3, hostPattern: 'yahoo|finance\\.yahoo' },
|
||||
policy: { minIntervalMs: 400, maxInflight: 1, drainJobBudget: 5, hostPattern: 'yahoo|finance\\.yahoo' },
|
||||
});
|
||||
registerVendorIntegration({
|
||||
family: 'sec',
|
||||
|
||||
+121
-14
@@ -38,6 +38,9 @@ import { CONFLUENCE_SLOT_IDS } from '../confluence/confluenceSlots.ts';
|
||||
import { ConfluenceRepository, rackFromSlots } from '../db/confluenceRepository.ts';
|
||||
import { runSlotBacktest, signalHistoryToStats, type ConfluenceFireEvent } from '../confluence/confluenceBacktest.ts';
|
||||
import { detectPictureChange } from '../confluence/confluenceRack.ts';
|
||||
import { runCorridorBacktest, aggregateBacktest } from '../analysis/corridorBacktest.ts';
|
||||
import { runConfluenceEvaluationCycle } from '../confluence/confluenceEngine.ts';
|
||||
import { CorridorRepository } from '../db/corridorRepository.ts';
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// X cookie credential helpers. Loads AES-256-GCM encrypted ct0/auth_token from
|
||||
@@ -452,8 +455,8 @@ const marketRouter = router({
|
||||
const { classifyRegime } = await import('../macro/MacroRegime.ts');
|
||||
|
||||
const symbol = input.symbol.toUpperCase();
|
||||
try { await ctx.cache.ensureInDemand(symbol, 'equity'); } catch { /* ignore */ }
|
||||
try { await ctx.cache.ensureInDemand(BENCHMARK_SYMBOL, 'etf'); } catch { /* ignore */ }
|
||||
try { await ctx.cache.bumpToWatched(symbol, 'equity'); } catch { /* ignore */ }
|
||||
try { await ctx.cache.bumpToWatched(BENCHMARK_SYMBOL, 'etf'); } catch { /* ignore */ }
|
||||
|
||||
const metaEntry = await ctx.cache.get<SymbolMeta>(`yfinance:symbol:${symbol}`);
|
||||
const meta = metaEntry.value;
|
||||
@@ -516,7 +519,7 @@ const marketRouter = router({
|
||||
|
||||
for (const s of symbolsToLoad) {
|
||||
const kind = s === BENCHMARK_SYMBOL || MARKET_ROTATION_UNIVERSE.some((u) => u.symbol === s) ? 'etf' : 'equity';
|
||||
try { await ctx.cache.ensureInDemand(s, kind); } catch { /* ignore */ }
|
||||
try { await ctx.cache.bumpToWatched(s, kind); } catch { /* ignore */ }
|
||||
}
|
||||
|
||||
const candleKeys = symbolsToLoad.map((s) => `yfinance:candles:${s}:1d`);
|
||||
@@ -776,8 +779,8 @@ const marketRouter = router({
|
||||
// SPY trend + VIX from cache/queue only (ADR-0009: no live Yahoo on request path).
|
||||
let spyCandles: PriceCandle[] = [];
|
||||
try {
|
||||
try { await ctx.cache.ensureInDemand('SPY', 'etf'); } catch { /* ignore */ }
|
||||
try { await ctx.cache.ensureInDemand('^VIX', 'index'); } catch { /* ignore */ }
|
||||
try { await ctx.cache.bumpToWatched('SPY', 'etf'); } catch { /* ignore */ }
|
||||
try { await ctx.cache.bumpToWatched('^VIX', 'index'); } catch { /* ignore */ }
|
||||
const spyEntry = await ctx.cache.get<PriceCandle[]>('yfinance:candles:SPY:1d');
|
||||
spyCandles = (spyEntry?.value ?? []) as PriceCandle[];
|
||||
factors.spy1M = totalReturnPct(spyCandles, 30 * 86_400_000);
|
||||
@@ -981,7 +984,7 @@ const marketRouter = router({
|
||||
|
||||
const symbols = [BENCHMARK_SYMBOL, ...MARKET_ROTATION_UNIVERSE.map((s) => s.symbol)];
|
||||
for (const sym of symbols) {
|
||||
try { await ctx.cache.ensureInDemand(sym, 'etf'); } catch { /* ignore */ }
|
||||
try { await ctx.cache.bumpToWatched(sym, 'etf'); } catch { /* ignore */ }
|
||||
}
|
||||
|
||||
const keys = symbols.map((s) => `yfinance:candles:${s}:1d`);
|
||||
@@ -1137,7 +1140,7 @@ const marketRouter = router({
|
||||
|
||||
const symbols = [BENCHMARK_SYMBOL, ...dedupedDefs.map((s) => s.symbol)];
|
||||
for (const sym of symbols) {
|
||||
try { await ctx.cache.ensureInDemand(sym, 'etf'); } catch { /* ignore */ }
|
||||
try { await ctx.cache.bumpToWatched(sym, 'etf'); } catch { /* ignore */ }
|
||||
}
|
||||
|
||||
const keys = symbols.map((s) => `yfinance:candles:${s}:1d`);
|
||||
@@ -1281,7 +1284,7 @@ const marketRouter = router({
|
||||
.query(async ({ ctx, input }) => {
|
||||
const { buildSeasonalitySnapshot, upcomingSimpleEvents } = await import('../analysis/seasonality.ts');
|
||||
const symbol = (input?.symbol ?? 'SPY').toUpperCase();
|
||||
try { await ctx.cache.ensureInDemand(symbol, symbol === 'SPY' ? 'etf' : 'equity'); } catch { /* ignore */ }
|
||||
try { await ctx.cache.bumpToWatched(symbol, symbol === 'SPY' ? 'etf' : 'equity'); } catch { /* ignore */ }
|
||||
const entry = await ctx.cache.get<PriceCandle[]>(`yfinance:candles:${symbol}:1d`);
|
||||
const candles = (entry?.value ?? []) as PriceCandle[];
|
||||
const snapshot = buildSeasonalitySnapshot(
|
||||
@@ -1429,7 +1432,7 @@ const marketRouter = router({
|
||||
for (const sym of symbols) {
|
||||
if (quoteMap.has(sym) && quoteMap.get(sym)!.price != null) continue;
|
||||
if (sym.includes('.') || sym.length > 5) continue;
|
||||
try { await ctx.cache.ensureInDemand(sym, 'equity'); } catch { /* ignore */ }
|
||||
try { await ctx.cache.bumpToWatched(sym, 'equity'); } catch { /* ignore */ }
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1467,7 +1470,7 @@ const marketRouter = router({
|
||||
ctx.db.prepare(
|
||||
'INSERT INTO rotation_custom_symbols (id, owner_id, symbol, name, grp, created_at) VALUES (?,?,?,?,?,?)'
|
||||
).run(crypto.randomUUID(), userId, input.symbol, input.name ?? input.symbol, input.group ?? 'Custom', new Date().toISOString());
|
||||
try { await ctx.cache.ensureInDemand(input.symbol, 'etf'); } catch { /* ignore */ }
|
||||
try { await ctx.cache.bumpToWatched(input.symbol, 'etf'); } catch { /* ignore */ }
|
||||
return { added: true };
|
||||
}),
|
||||
|
||||
@@ -3399,7 +3402,7 @@ const dealerStudyRouter = router({
|
||||
);
|
||||
// Ensure daily candles are in demand for later grading (ADR-0009).
|
||||
try {
|
||||
await ctx.cache.ensureInDemand(input.symbol.toUpperCase(), 'equity');
|
||||
await ctx.cache.bumpToWatched(input.symbol.toUpperCase(), 'equity');
|
||||
} catch { /* optional */ }
|
||||
return { ok: true as const, id, loggedAt };
|
||||
}),
|
||||
@@ -3835,7 +3838,7 @@ const mentorLedgerRouter = router({
|
||||
: row.logged_at;
|
||||
|
||||
try {
|
||||
await ctx.cache.ensureInDemand(row.symbol, 'equity');
|
||||
await ctx.cache.bumpToWatched(row.symbol, 'equity');
|
||||
} catch { /* optional */ }
|
||||
|
||||
const candles = ctx.db.prepare(
|
||||
@@ -5382,7 +5385,7 @@ const mirrorRouter = router({
|
||||
// ─── Confluence Signal Engine (M22) ─────────────────────────────────────────
|
||||
|
||||
const confluenceRouter = router({
|
||||
/** The read-only 34-slot catalog, grouped by family, with slot metadata. */
|
||||
/** The read-only slot catalog, grouped by family, with slot metadata. */
|
||||
slots: publicProcedure.query(async () => {
|
||||
return {
|
||||
slots: CONFLUENCE_SLOTS.map((s) => ({
|
||||
@@ -5476,7 +5479,7 @@ const confluenceRouter = router({
|
||||
id: z.string().min(1).max(80).optional(),
|
||||
name: z.string().min(1).max(80),
|
||||
description: z.string().max(300).optional(),
|
||||
slotIds: z.array(z.string()).min(1).max(34),
|
||||
slotIds: z.array(z.string()).min(1).max(CONFLUENCE_SLOT_IDS.length),
|
||||
}))
|
||||
.mutation(({ ctx, input }) => {
|
||||
const repo = new ConfluenceRepository(ctx.db);
|
||||
@@ -5498,6 +5501,110 @@ const confluenceRouter = router({
|
||||
}
|
||||
return rack;
|
||||
}),
|
||||
|
||||
// -------------------------------------------------------------------------
|
||||
// Price Corridor (M24): valuation-corridor snapshots + backtest ledger + link
|
||||
// -------------------------------------------------------------------------
|
||||
|
||||
/** Latest valuation-corridor snapshot for a symbol, with SPY market context. */
|
||||
corridorSnapshot: publicProcedure
|
||||
.input(z.object({ symbol: z.string().min(1).max(12) }))
|
||||
.query(({ ctx, input }) => {
|
||||
const repo = new CorridorRepository(ctx.db);
|
||||
const symbol = input.symbol.toUpperCase();
|
||||
const snapshot = repo.latestSnapshot(symbol);
|
||||
const spy = repo.latestSnapshot('SPY');
|
||||
const series = repo.snapshotsForSymbol(symbol, 90);
|
||||
const corridor = (snapshot?.corridor1yLow ?? null) !== null
|
||||
? {
|
||||
low: snapshot!.corridor1yLow,
|
||||
high: snapshot!.corridor1yHigh,
|
||||
median: snapshot!.corridor1yMedian,
|
||||
currentPE: snapshot!.trailingPE ?? snapshot!.forwardPE,
|
||||
fairValue1y: snapshot!.fairValue1y,
|
||||
impliedUpside1y: snapshot!.impliedUpside1y,
|
||||
pePercentile1y: snapshot!.pePercentile1y,
|
||||
}
|
||||
: null;
|
||||
return {
|
||||
symbol,
|
||||
snapshot,
|
||||
spy,
|
||||
series,
|
||||
corridor,
|
||||
};
|
||||
}),
|
||||
|
||||
/** Estuary user's corridor watchlist (tickers surfaced in the Corridor panel). */
|
||||
corridorWatchlist: protectedProcedure
|
||||
.query(({ ctx }) => {
|
||||
const repo = new CorridorRepository(ctx.db);
|
||||
return { symbols: repo.listWatchlist(ctx.userId!) };
|
||||
}),
|
||||
|
||||
/** Add a ticker to the corridor watchlist. */
|
||||
corridorWatchlistAdd: protectedProcedure
|
||||
.input(z.object({ symbol: z.string().min(1).max(12).transform((s) => s.toUpperCase()) }))
|
||||
.mutation(({ ctx, input }) => {
|
||||
const repo = new CorridorRepository(ctx.db);
|
||||
repo.addToWatchlist(ctx.userId!, input.symbol);
|
||||
return { symbols: repo.listWatchlist(ctx.userId!) };
|
||||
}),
|
||||
|
||||
/** Remove a ticker from the corridor watchlist. */
|
||||
corridorWatchlistRemove: protectedProcedure
|
||||
.input(z.object({ symbol: z.string().min(1).max(12).transform((s) => s.toUpperCase()) }))
|
||||
.mutation(({ ctx, input }) => {
|
||||
const repo = new CorridorRepository(ctx.db);
|
||||
repo.removeFromWatchlist(ctx.userId!, input.symbol);
|
||||
return { symbols: repo.listWatchlist(ctx.userId!) };
|
||||
}),
|
||||
|
||||
/** Corridor-method backtest scorecard: @alojoh's rated names vs actual forward returns. */
|
||||
corridorBacktest: publicProcedure
|
||||
.query(({ ctx }) => {
|
||||
const repo = new CorridorRepository(ctx.db);
|
||||
const rows = repo.listBacktests();
|
||||
return {
|
||||
rows,
|
||||
aggregate: aggregateBacktest(
|
||||
rows.map((r) => ({
|
||||
articleDate: r.articleDate,
|
||||
articleId: r.articleId,
|
||||
rankingType: r.rankingType as 'entry_1y' | 'entry_90d',
|
||||
topSymbols: r.topSymbols,
|
||||
bottomSymbols: r.bottomSymbols,
|
||||
topAvgReturn: r.topAvgReturn,
|
||||
bottomAvgReturn: r.bottomAvgReturn,
|
||||
spread: r.spread,
|
||||
isWin: r.isWin,
|
||||
horizonDays: r.horizonDays,
|
||||
resolvableCount: 0,
|
||||
})),
|
||||
),
|
||||
};
|
||||
}),
|
||||
|
||||
/** Replay the corridor-method backtest corpus against cached candles (offline; idempotent). */
|
||||
corridorBacktestRun: publicProcedure
|
||||
.mutation(async ({ ctx }) => {
|
||||
const { runCorridorBacktest } = await import('../analysis/corridorBacktest.ts');
|
||||
const result = await runCorridorBacktest(ctx.db, async (symbol) => {
|
||||
const entry = await ctx.cache.get<PriceCandle[]>(`yfinance:candles:${symbol.toUpperCase()}:1d`);
|
||||
return (entry?.value ?? []) as PriceCandle[];
|
||||
});
|
||||
return { grades: result.grades.length, resolved: result.aggregate.resolved, aggregate: result.aggregate };
|
||||
}),
|
||||
|
||||
/** Force a confluence evaluation cycle for the current symbol now (debug/admin). */
|
||||
runEvaluationNow: protectedProcedure
|
||||
.input(z.object({ symbol: z.string().min(1).max(12).optional() }))
|
||||
.mutation(async ({ ctx, input }) => {
|
||||
const summary = await runConfluenceEvaluationCycle(ctx.db, ctx.cache, {
|
||||
symbols: input?.symbol ? [input.symbol] : undefined,
|
||||
});
|
||||
return summary;
|
||||
}),
|
||||
});
|
||||
|
||||
export const appRouter = router({
|
||||
|
||||
Reference in New Issue
Block a user