148 lines
5.5 KiB
Python
148 lines
5.5 KiB
Python
"""Market-breadth state and early-warning indicators.
|
|
|
|
Breadth is a genuinely *leading* construct: a few mega-caps can keep an index
|
|
rising while participation narrows underneath — the classic pre-top divergence.
|
|
V2 measures an explicit, frozen basket rather than every ticker currently stored
|
|
in the database. That keeps the live series reproducible when the wider product
|
|
universe changes.
|
|
|
|
Two layers:
|
|
- breadth = % of the universe trading above its own 200-DMA (0-100).
|
|
- divergence = an early-warning score (0-100, high = fragile): the benchmark
|
|
price holding/rising *while* breadth falls. Absolute low breadth stays in the
|
|
State index so it is not counted twice.
|
|
|
|
The live monitor uses the breadth level in State and the pure divergence in
|
|
Warning. The event study evaluates the latter chronologically.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import logging
|
|
from datetime import date
|
|
|
|
from sqlalchemy import select
|
|
from sqlalchemy.ext.asyncio import AsyncSession
|
|
|
|
from app.models.ticker import Ticker
|
|
from app.services.price_service import query_ohlcv
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
Series = list[tuple[date, float]]
|
|
|
|
|
|
def _breadth_with_counts(
|
|
closes_by_symbol: dict[str, Series], window: int = 200, min_tickers: int = 20
|
|
) -> tuple[dict[date, float], dict[date, int]]:
|
|
"""Pure core: % of symbols above their own rolling SMA(window), per date.
|
|
|
|
Each symbol's SMA is computed once with a sliding sum (O(bars)); dates with
|
|
fewer than ``min_tickers`` qualifying names are dropped (too thin to trust).
|
|
"""
|
|
counts: dict[date, list[int]] = {} # date -> [above, total]
|
|
for series in closes_by_symbol.values():
|
|
ordered = sorted(series, key=lambda x: x[0])
|
|
dates = [d for d, _ in ordered]
|
|
closes = [c for _, c in ordered]
|
|
if len(closes) < window:
|
|
continue
|
|
running = sum(closes[:window])
|
|
for i in range(window - 1, len(closes)):
|
|
if i >= window:
|
|
running += closes[i] - closes[i - window]
|
|
sma = running / window
|
|
entry = counts.setdefault(dates[i], [0, 0])
|
|
entry[1] += 1
|
|
if closes[i] > sma:
|
|
entry[0] += 1
|
|
values = {
|
|
d: round(above / total * 100.0, 2)
|
|
for d, (above, total) in counts.items()
|
|
if total >= min_tickers
|
|
}
|
|
eligible = {d: total for d, (_, total) in counts.items() if total >= min_tickers}
|
|
return values, eligible
|
|
|
|
|
|
def _breadth_from_closes(
|
|
closes_by_symbol: dict[str, Series], window: int = 200, min_tickers: int = 20
|
|
) -> dict[date, float]:
|
|
"""Compatibility wrapper returning only the breadth percentage series."""
|
|
return _breadth_with_counts(closes_by_symbol, window, min_tickers)[0]
|
|
|
|
|
|
def compute_divergence_series(
|
|
breadth: dict[date, float], benchmark_closes: Series, lookback: int = 20
|
|
) -> dict[date, float]:
|
|
"""Early-warning score (0-100, high = fragile) per date.
|
|
|
|
This is deliberately a pure divergence: it is positive only when benchmark
|
|
price holds/rises while breadth falls. Absolute low breadth belongs in the
|
|
State score, so it is not counted again here. A 20 percentage-point breadth
|
|
deterioration maps to 100.
|
|
"""
|
|
bench = {d: c for d, c in benchmark_closes}
|
|
common = sorted(d for d in bench if d in breadth)
|
|
out: dict[date, float] = {}
|
|
for i in range(lookback, len(common)):
|
|
d, d0 = common[i], common[i - lookback]
|
|
price_past = bench[d0]
|
|
if price_past <= 0:
|
|
continue
|
|
price_ret = (bench[d] / price_past - 1.0) * 100.0 # %
|
|
breadth_chg = breadth[d] - breadth[d0] # percentage points
|
|
deterioration = max(0.0, -breadth_chg)
|
|
score = deterioration * 5.0 if price_ret >= 0 else 0.0
|
|
out[d] = max(0.0, min(100.0, round(score, 2)))
|
|
return out
|
|
|
|
|
|
async def _load_universe_closes(
|
|
db: AsyncSession, symbols: list[str] | None = None
|
|
) -> dict[str, Series]:
|
|
stmt = select(Ticker).order_by(Ticker.symbol)
|
|
if symbols is not None:
|
|
stmt = stmt.where(Ticker.symbol.in_(symbols))
|
|
result = await db.execute(stmt)
|
|
closes_by_symbol: dict[str, Series] = {}
|
|
for ticker in result.scalars().all():
|
|
try:
|
|
records = await query_ohlcv(db, ticker.symbol)
|
|
except Exception:
|
|
logger.exception("Breadth: OHLCV load failed for %s", ticker.symbol)
|
|
continue
|
|
if records:
|
|
closes_by_symbol[ticker.symbol] = [(r.date, float(r.close)) for r in records]
|
|
return closes_by_symbol
|
|
|
|
|
|
async def compute_breadth_series(
|
|
db: AsyncSession,
|
|
window: int = 200,
|
|
min_tickers: int = 20,
|
|
symbols: list[str] | None = None,
|
|
) -> dict[date, float]:
|
|
"""Historical breadth series across an explicit basket (or all stored names)."""
|
|
closes_by_symbol = await _load_universe_closes(db, symbols)
|
|
return _breadth_from_closes(closes_by_symbol, window, min_tickers)
|
|
|
|
|
|
async def compute_breadth_details(
|
|
db: AsyncSession,
|
|
symbols: list[str],
|
|
window: int = 200,
|
|
min_tickers: int = 20,
|
|
) -> tuple[dict[date, float], dict[date, int]]:
|
|
"""Breadth values plus the qualifying-member count for snapshot metadata."""
|
|
closes_by_symbol = await _load_universe_closes(db, symbols)
|
|
return _breadth_with_counts(closes_by_symbol, window, min_tickers)
|
|
|
|
|
|
async def compute_breadth_today(db: AsyncSession) -> float | None:
|
|
"""Latest breadth reading (thin wrapper, for future live use)."""
|
|
series = await compute_breadth_series(db)
|
|
if not series:
|
|
return None
|
|
return series[max(series)]
|