diff --git a/README.md b/README.md index a46a5f6..33ad501 100644 --- a/README.md +++ b/README.md @@ -124,7 +124,7 @@ Once a day (default 07:00). Steps run **in dependency order**, each consuming th 3. **R:R Scan** — persist clean Structural S/R for charts/alerts, recompute the 5-dimension scores, and build long/short setups from a transient Gate Target Ladder (ATR stops and nominal gate targets) for every ticker. Attach each ticker's residual 12‑1 momentum activation percentile plus the promoted 80/20 production rank. 4. **Outcome Eval** — resolve setups that hit target/stop or expired (default 30 trading days) and auto-close paper trades per the exit policy (default: 3x ATR trail with a 30-trading-day max hold). 5. **Market Regime** — recompute the regime index (breadth/trend). -6. **Regime Monitor** — observational early-warning snapshot (VIX, credit spreads via FRED); feeds nothing else. +6. **Regime Monitor** — separate v2 State/Warning risk thermometer with fixed-basket breadth, VIX, credit, and point-in-time fundamentals; feeds no trades. A failing step is logged; the pipeline continues with the next. @@ -279,7 +279,7 @@ Corollaries: never let an unvalidated score gate setups; the outcome evaluator m - Activation gate — qualifies setups on a residual-momentum percentile floor (the actual selection), a headline gate-target R:R floor (prod: 2.0) and a 20% primary-target reach-probability floor (validated long-only edge) - Recommendation layer — directional confidence, conflict detection, per-target reach-probability - Paper trading — take a setup, mark-to-market vs. latest close, auto-close per the exit policy (default: 3x ATR trail with a 30-trading-day max hold; time / percent-trailing / target-stop selectable), realized track record + outcome evaluation -- Market-regime index + FRED early-warning monitor (VIX, credit spreads); weekly backtest + manual event study +- Market-regime guard + observational State/Warning monitor (fixed-basket breadth, VIX, credit, PIT fundamentals) with a manual chronological correction study - Telegram alerts (e.g. regime-quadrant changes) - User-curated watchlist (cap: 20), enriched with composite score, R:R and S/R summary - JWT auth with admin role, configurable registration, user access control diff --git a/app/models/regime_snapshot.py b/app/models/regime_snapshot.py index 884a51b..2afbf0b 100644 --- a/app/models/regime_snapshot.py +++ b/app/models/regime_snapshot.py @@ -8,12 +8,12 @@ from app.database import Base class RegimeSnapshot(Base): - """Daily snapshot of the AI/Tech regime-change index. + """Daily point-in-time snapshot of the AI/Tech Regime Monitor. One row per calendar date (unique). ``breakdown_json`` holds the full - per-signal breakdown plus the raw inputs, so reads need no recomputation and - the 7/30-day trend is just a query over ``total_score``. Decoupled from the - rest of the platform: nothing reads this to gate or score trades. + ``breakdown_json`` is authoritative for v2 State, Warning, source dates, + coverage, and fixed-basket metadata. ``total_score``/``band`` retain the v2 + State reading for schema compatibility. Nothing reads this to gate trades. """ __tablename__ = "regime_snapshots" diff --git a/app/routers/market.py b/app/routers/market.py index a66d9af..e4fc038 100644 --- a/app/routers/market.py +++ b/app/routers/market.py @@ -1,7 +1,7 @@ """Market-level endpoints (benchmark regime + AI/Tech regime-change monitor).""" from fastapi import APIRouter, Depends, Query -from pydantic import BaseModel +from pydantic import BaseModel, Field, field_validator from sqlalchemy.ext.asyncio import AsyncSession from app.dependencies import get_db, require_access, require_admin @@ -40,12 +40,18 @@ async def backtest_report( class RegimeConfigUpdate(BaseModel): - weights: dict[str, float] | None = None - alert_threshold: float | None = None - tickers: dict | None = None - leader_weight: float | None = None - rs_lookback: int | None = None - fundamental_staleness_days: int | None = None + breadth_basket: list[str] | None = Field(default=None, min_length=20, max_length=100) + fundamental_staleness_days: int | None = Field(default=None, ge=30, le=180) + + @field_validator("breadth_basket") + @classmethod + def normalise_basket(cls, value: list[str] | None) -> list[str] | None: + if value is None: + return None + cleaned = [symbol.strip().upper().replace(".", "-") for symbol in value if symbol.strip()] + if len(cleaned) != len(set(cleaned)): + raise ValueError("breadth basket symbols must be unique") + return cleaned class RegimeFundamentalsUpdate(BaseModel): @@ -59,7 +65,7 @@ async def regime_monitor( _user: User = Depends(require_access), db: AsyncSession = Depends(get_db), ) -> APIEnvelope: - """Latest AI/Tech regime-change index (0-100) + per-signal breakdown + trend.""" + """Latest v2 State and Warning risk-thermometer readings.""" data = await regime_monitor_service.get_regime_monitor(db) return APIEnvelope(status="success", data=data) @@ -69,7 +75,7 @@ async def regime_config( _admin: User = Depends(require_admin), db: AsyncSession = Depends(get_db), ) -> APIEnvelope: - """Editable weights / thresholds / ticker lists for the regime monitor.""" + """Editable fixed breadth basket and fundamental freshness window.""" data = await regime_monitor_service.get_regime_config(db) return APIEnvelope(status="success", data=data) @@ -80,7 +86,7 @@ async def update_regime_config( _admin: User = Depends(require_admin), db: AsyncSession = Depends(get_db), ) -> APIEnvelope: - """Merge the supplied fields into the stored regime-monitor config.""" + """Update the deliberately small v2 operator configuration.""" updates = body.model_dump(exclude_none=True) data = await regime_monitor_service.update_regime_config(db, updates) return APIEnvelope(status="success", data=data) @@ -133,10 +139,10 @@ async def regime_event_study( @router.get("/regime/history", response_model=APIEnvelope) async def regime_history( - days: int = Query(default=400, ge=7, le=2000), + days: int = Query(default=800, ge=7, le=2000), _user: User = Depends(require_access), db: AsyncSession = Depends(get_db), ) -> APIEnvelope: - """Daily history of the index / early-warning / combined scores (for the chart).""" + """Point-in-time v2 State/Warning history. Legacy rows are excluded.""" data = await regime_monitor_service.get_regime_history(db, days=days) return APIEnvelope(status="success", data=data) diff --git a/app/scheduler.py b/app/scheduler.py index 5259af1..af61665 100644 --- a/app/scheduler.py +++ b/app/scheduler.py @@ -1001,12 +1001,20 @@ async def compute_regime_monitor() -> None: result = await update_regime_monitor(db) + state = result.get("state") or {} + warning = result.get("warning") or {} _runtime_progress(job_name, processed=1, total=1) _runtime_finish( job_name, "completed", processed=1, total=1, - message=f"Index: {result.get('total_score')} ({result.get('band')})", + message=f"State: {state.get('score')} · Warning: {warning.get('score')}", + ) + _log_event( + logging.INFO, + "job_complete", + job=job_name, + state=state.get("score"), + warning=warning.get("score"), ) - _log_event(logging.INFO, "job_complete", job=job_name, score=result.get("total_score")) except Exception as exc: _runtime_finish(job_name, "error", processed=0, total=1, message=str(exc)) _log_event(logging.ERROR, "job_error", job=job_name, error_type=type(exc).__name__, message=str(exc)) @@ -1086,7 +1094,11 @@ async def run_event_study_job() -> None: _runtime_progress(job_name, processed=1, total=1) if report.get("available"): - msg = f"{len(report.get('events', []))} events, lead Δ {report.get('lead_delta_days')}d" + metrics = report.get("metrics") or {} + msg = ( + f"{metrics.get('events_warned', 0)}/{metrics.get('events', 0)} warned, " + f"{metrics.get('false_alarms_per_year', 0)} false alarms/year" + ) else: msg = report.get("reason", "no data") _runtime_finish(job_name, "completed", processed=1, total=1, message=msg) diff --git a/app/services/alert_service.py b/app/services/alert_service.py index 9e4f4ea..828fad0 100644 --- a/app/services/alert_service.py +++ b/app/services/alert_service.py @@ -58,7 +58,9 @@ _BOOL_DEFAULTS = { KEY_SR: True, KEY_SCORE_DROP: True, KEY_DIGEST: True, - KEY_REGIME_QUADRANT: True, + # Experimental human-facing thermometer: opt in explicitly. Existing stored + # true values remain true; only missing/reset configurations default off. + KEY_REGIME_QUADRANT: False, KEY_TRADE_CLOSED: True, } @@ -90,19 +92,19 @@ SIGNAL_BUNDLE_SECTIONS = ( ) SIGNAL_BUNDLE_MAX_CHARS = 3900 # Telegram limit is 4096; keep room for HTML parsing -# Regime quadrant-change alert: (regime index x early-warning) quadrant. +# Regime quadrant-change alert: (State x Warning) quadrant. # Hysteresis (a deadband around each divider) stops a point sitting on a boundary # from flip-flopping; the cooldown caps how often a genuine change can re-alert. QUAD_TYPE = "regime_quadrant" -QUAD_X_DIV = 40.0 # regime index divider (matches the frontend quadrant) -QUAD_Y_DIV = 60.0 # early-warning divider +QUAD_X_DIV = 60.0 # v2 State divider (backend response is authoritative) +QUAD_Y_DIV = 60.0 # v2 Warning divider QUAD_MARGIN = 5.0 # half-width of the hysteresis deadband around each divider QUAD_COOLDOWN_DAYS = 3 # min days between quadrant-change alerts QUAD_LABELS = { - "1": "① Hot & brittle", - "2": "② Transition", - "3": "③ Healthy & broad", - "4": "④ Real downturn", + "1": "Early warning", + "2": "Active stress", + "3": "Healthy", + "4": "Stressed / stabilizing", } AlertItem = tuple[str, str, str] # alert_type, dedup_key, text @@ -693,49 +695,65 @@ def _closed_trade_bundle( def _bools_to_quadrant(x_high: bool, y_high: bool) -> str: if y_high: - return "2" if x_high else "1" # ② Transition / ① Hot & brittle - return "4" if x_high else "3" # ④ Real downturn / ③ Healthy & broad + return "2" if x_high else "1" # Active stress / Early warning + return "4" if x_high else "3" # Stressed/stabilizing / Healthy def _quadrant_to_bools(q: str) -> tuple[bool, bool]: return {"1": (False, True), "2": (True, True), "3": (False, False), "4": (True, False)}[q] -def _classify_quadrant(x: float, y: float, prev: str | None, margin: float = QUAD_MARGIN) -> str: - """Quadrant of (regime index x, early warning y), with per-axis hysteresis. +def _classify_quadrant( + x: float, + y: float, + prev: str | None, + margin: float = QUAD_MARGIN, + x_div: float = QUAD_X_DIV, + y_div: float = QUAD_Y_DIV, +) -> str: + """Quadrant of (State x, Warning y), with per-axis hysteresis. Each axis only flips once the value crosses its divider by ``margin`` in the new direction, so a point parked on a divider keeps its current quadrant instead of flip-flopping. ``prev`` None means a fresh (no-hysteresis) classify. """ if prev is None: - return _bools_to_quadrant(x >= QUAD_X_DIV, y >= QUAD_Y_DIV) + return _bools_to_quadrant(x >= x_div, y >= y_div) px, py = _quadrant_to_bools(prev) - x_high = (x >= QUAD_X_DIV - margin) if px else (x >= QUAD_X_DIV + margin) - y_high = (y >= QUAD_Y_DIV - margin) if py else (y >= QUAD_Y_DIV + margin) + x_high = (x >= x_div - margin) if px else (x >= x_div + margin) + y_high = (y >= y_div - margin) if py else (y >= y_div + margin) return _bools_to_quadrant(x_high, y_high) -def _quadrant_log_key(q: str, x: float, y: float) -> str: - return f"{q}:{x:.1f}:{y:.1f}" +def _quadrant_log_key(q: str, x: float, y: float, basket_hash: str | None = None) -> str: + return f"{basket_hash or 'legacy'}:{q}:{x:.1f}:{y:.1f}" -def _parse_quadrant_log_key(key: str | None) -> tuple[str | None, float | None, float | None]: +def _parse_quadrant_log_key( + key: str | None, +) -> tuple[str | None, str | None, float | None, float | None]: if not key: - return None, None, None + return None, None, None, None parts = key.split(":") - q = parts[0] + if parts[0] in QUAD_LABELS: + basket_hash, q, values = None, parts[0], parts[1:] + elif len(parts) >= 2: + basket_hash, q, values = parts[0], parts[1], parts[2:] + else: + return None, None, None, None if q not in QUAD_LABELS: - return None, None, None - if len(parts) >= 3: + return None, None, None, None + if len(values) >= 2: try: - return q, float(parts[1]), float(parts[2]) + return basket_hash, q, float(values[0]), float(values[1]) except ValueError: pass - return q, None, None + return basket_hash, q, None, None -async def _last_quadrant(db: AsyncSession) -> tuple[str | None, float | None, float | None, datetime | None]: +async def _last_quadrant( + db: AsyncSession, +) -> tuple[str | None, str | None, float | None, float | None, datetime | None]: """Most recently logged quadrant (and when), our baseline for change + cooldown.""" result = await db.execute( select(AlertLog.dedup_key, AlertLog.created_at) @@ -745,9 +763,9 @@ async def _last_quadrant(db: AsyncSession) -> tuple[str | None, float | None, fl ) row = result.first() if not row: - return None, None, None, None - prev_q, prev_x, prev_y = _parse_quadrant_log_key(row[0]) - return prev_q, prev_x, prev_y, row[1] + return None, None, None, None, None + basket_hash, prev_q, prev_x, prev_y = _parse_quadrant_log_key(row[0]) + return basket_hash, prev_q, prev_x, prev_y, row[1] async def _collect_regime_quadrant(db: AsyncSession) -> list[tuple[str, str]]: @@ -758,25 +776,64 @@ async def _collect_regime_quadrant(db: AsyncSession) -> list[tuple[str, str]]: cooldown has elapsed. The dispatch loop logs the new quadrant on send, which becomes the next baseline and resets the cooldown clock. """ - from app.services.regime_monitor_service import get_regime_monitor + from app.services.regime_monitor_service import get_regime_history, get_regime_monitor data = await get_regime_monitor(db) if not data.get("available"): return [] - x = data.get("total_score") - y = (data.get("early_warning") or {}).get("score") + state = data.get("state") or {} + warning = data.get("warning") or {} + x = state.get("score") + y = warning.get("score") if x is None or y is None: return [] - prev, prev_x, prev_y, prev_time = await _last_quadrant(db) - if prev is None: - _log_alert(db, QUAD_TYPE, _quadrant_log_key(_classify_quadrant(x, y, None), x, y)) # seed, no alert + quality = data.get("data_quality") or {} + if ( + float(state.get("coverage") or 0) < 75 + or float(warning.get("coverage") or 0) < 75 + or not quality.get("is_fresh") + ): return [] - new_q = _classify_quadrant(x, y, prev) + quadrant_cfg = data.get("quadrant_config") or {} + x_div = float(quadrant_cfg.get("state_divider", QUAD_X_DIV)) + y_div = float(quadrant_cfg.get("warning_divider", QUAD_Y_DIV)) + margin = float(quadrant_cfg.get("margin", QUAD_MARGIN)) + basket_hash = str((data.get("basket") or {}).get("hash") or "unknown") + + prev_hash, prev, prev_x, prev_y, prev_time = await _last_quadrant(db) + if prev is None or prev_hash != basket_hash: + seed = _classify_quadrant(x, y, None, margin, x_div, y_div) + _log_alert(db, QUAD_TYPE, _quadrant_log_key(seed, x, y, basket_hash)) + return [] + + new_q = _classify_quadrant(x, y, prev, margin, x_div, y_div) if new_q == prev: return [] + history = await get_regime_history(db, days=14) + valid = [ + point for point in history + if point.get("state") is not None + and point.get("warning") is not None + and float(point.get("state_coverage") or 0) >= 75 + and float(point.get("warning_coverage") or 0) >= 75 + ] + if len(valid) < 2: + return [] + prior = valid[-2] + prior_q = _classify_quadrant( + float(prior["state"]), + float(prior["warning"]), + prev, + margin, + x_div, + y_div, + ) + if prior_q != new_q: + return [] + if prev_time is not None: if prev_time.tzinfo is None: prev_time = prev_time.replace(tzinfo=timezone.utc) @@ -785,17 +842,19 @@ async def _collect_regime_quadrant(db: AsyncSession) -> list[tuple[str, str]]: if prev_x is not None and prev_y is not None: metrics = ( - f"regime {prev_x:.0f} → {x:.0f} ({x - prev_x:+.0f}) · " - f"early-warning {prev_y:.0f} → {y:.0f} ({y - prev_y:+.0f})" + f"State {prev_x:.0f} → {x:.0f} ({x - prev_x:+.0f}) · " + f"Warning {prev_y:.0f} → {y:.0f} ({y - prev_y:+.0f})" ) else: - metrics = f"regime {x:.0f} · early-warning {y:.0f}" + metrics = f"State {x:.0f} · Warning {y:.0f}" text = ( f"🧭 Regime quadrant change\n" f"{QUAD_LABELS.get(prev, prev)} → {QUAD_LABELS.get(new_q, new_q)}\n" - f"{metrics}" + f"{metrics}\n" + f"coverage: state {state.get('coverage'):.0f}% / warning {warning.get('coverage'):.0f}%\n" + f"Risk thermometer - not a trade signal." ) - return [(_quadrant_log_key(new_q, x, y), text)] + return [(_quadrant_log_key(new_q, x, y, basket_hash), text)] # --------------------------------------------------------------------------- diff --git a/app/services/breadth_service.py b/app/services/breadth_service.py index 1d3971d..dc1047b 100644 --- a/app/services/breadth_service.py +++ b/app/services/breadth_service.py @@ -1,18 +1,19 @@ -"""Market-breadth early-warning indicator (from the stored universe OHLCV). +"""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. -We measure it from the OHLCV we already store for the whole universe, so it costs -no new data source. +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 rising *while* breadth falls, plus a nudge for already-low breadth. + price holding/rising *while* breadth falls. Absolute low breadth stays in the + State index so it is not counted twice. -This module only *computes* the indicator. It is deliberately NOT wired into the -live regime index yet — the event study measures whether it actually leads before -it earns any weight. +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 @@ -31,9 +32,9 @@ logger = logging.getLogger(__name__) Series = list[tuple[date, float]] -def _breadth_from_closes( +def _breadth_with_counts( closes_by_symbol: dict[str, Series], window: int = 200, min_tickers: int = 20 -) -> dict[date, float]: +) -> 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 @@ -55,11 +56,20 @@ def _breadth_from_closes( entry[1] += 1 if closes[i] > sma: entry[0] += 1 - return { + 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( @@ -67,10 +77,10 @@ def compute_divergence_series( ) -> dict[date, float]: """Early-warning score (0-100, high = fragile) per date. - Fragility rises when the benchmark price climbs over ``lookback`` days while - breadth deteriorates over the same window, and is nudged up when the absolute - breadth level is already low. It is the *divergence* (not the level) that - makes this leading. + 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) @@ -82,14 +92,19 @@ def compute_divergence_series( continue price_ret = (bench[d] / price_past - 1.0) * 100.0 # % breadth_chg = breadth[d] - breadth[d0] # percentage points - raw = price_ret - breadth_chg # price up & breadth down -> large - score = 50.0 + raw * 2.0 + (50.0 - breadth[d]) * 0.4 + 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) -> dict[str, Series]: - result = await db.execute(select(Ticker).order_by(Ticker.symbol)) +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: @@ -103,13 +118,27 @@ async def _load_universe_closes(db: AsyncSession) -> dict[str, Series]: async def compute_breadth_series( - db: AsyncSession, window: int = 200, min_tickers: int = 20 + db: AsyncSession, + window: int = 200, + min_tickers: int = 20, + symbols: list[str] | None = None, ) -> dict[date, float]: - """Historical breadth series across the stored universe (for the event study).""" - closes_by_symbol = await _load_universe_closes(db) + """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) diff --git a/app/services/event_study_service.py b/app/services/event_study_service.py index 064d688..a822c2c 100644 --- a/app/services/event_study_service.py +++ b/app/services/event_study_service.py @@ -1,21 +1,9 @@ -"""Event study: does a candidate indicator actually *lead* regime breaks? +"""Compact chronological validation for the Regime Monitor warning score. -This is a backtest-style measurement, but the unit of analysis is **events** -(historical drawdowns), not trades. For each candidate indicator it answers: - - how many days of warning did it give before the break (event-centered)? - - at what false-alarm cost (signal-centered precision/recall vs. the base rate)? - -It compares the breadth-divergence early-warning candidate against a deterministic -**coincident** price composite (the existing regime price sub-scores), so you can -see whether the candidate crosses *earlier*. Everything is price/breadth only — -no LLM/FRED — so the result is reproducible. - -Honest caveat: with only a handful of real drawdowns in ~5y, the sample is tiny -and the numbers are noisy. Read the median lead time as an order of magnitude, and -do NOT overfit thresholds to this history. - -Report is cached in a SystemSetting (mirrors ``backtest_service``); a manual job -(Admin → Jobs) drives it. +The study calls its outcome a 10% correction, uses the first 70% of sessions to +freeze an 80th-percentile warning threshold, and reports alarm episodes only on +the final 30%. It is still labelled exploratory while the fixed breadth basket +is reconstructed before its freeze date. """ from __future__ import annotations @@ -34,304 +22,268 @@ logger = logging.getLogger(__name__) KEY_REPORT = "regime_event_study" -# Defaults. The 15% threshold gave only 2 events in 5y (statistically useless), -# so the default is lower with a cooldown-based dedup to surface more, cleaner -# events. Each indicator "warns" at its OWN 80th percentile rather than a shared -# absolute level, so the leading vs. coincident comparison is fair across scales. -EVENT_THRESHOLD_PCT = 10.0 # drawdown from the 52w high that counts as a "break" -COOLDOWN_DAYS = 40 # min trading days between event onsets (dedup) -DRAWDOWN_LOOKBACK = 252 # 52-week trailing high -HORIZON_DAYS = 20 # signal-centered prediction horizon -WARN_PERCENTILE = 80.0 # each indicator warns at its own Nth percentile -PRE, POST = 60, 20 # event-centered window (trading days) +EVENT_THRESHOLD_PCT = 10.0 +EVENT_COOLDOWN_DAYS = 40 +DRAWDOWN_LOOKBACK = 252 +HORIZON_DAYS = 20 +WARN_PERCENTILE = 80.0 +TRAIN_FRACTION = 0.70 def _median(values: list[float]) -> float | None: if not values: return None - s = sorted(values) - n = len(s) - mid = n // 2 - return float(s[mid]) if n % 2 else (s[mid - 1] + s[mid]) / 2.0 + ordered = sorted(values) + middle = len(ordered) // 2 + return ( + float(ordered[middle]) + if len(ordered) % 2 + else (ordered[middle - 1] + ordered[middle]) / 2.0 + ) def _percentile(values: list[float], pct: float) -> float | None: - """Linear-interpolated percentile of the non-None values.""" - vals = sorted(v for v in values if v is not None) - if not vals: + ordered = sorted(v for v in values if v is not None) + if not ordered: return None - k = (len(vals) - 1) * (pct / 100.0) - lo = int(k) - hi = min(lo + 1, len(vals) - 1) - return vals[lo] + (vals[hi] - vals[lo]) * (k - lo) + position = (len(ordered) - 1) * pct / 100.0 + lower = int(position) + upper = min(lower + 1, len(ordered) - 1) + return ordered[lower] + (ordered[upper] - ordered[lower]) * (position - lower) -# --------------------------------------------------------------------------- -# Event detection -# --------------------------------------------------------------------------- - def detect_events( closes: list[float], dates: list[date], threshold_pct: float = EVENT_THRESHOLD_PCT, lookback: int = DRAWDOWN_LOOKBACK, - cooldown: int = COOLDOWN_DAYS, + cooldown: int = EVENT_COOLDOWN_DAYS, ) -> list[dict]: - """Drawdown events: ``t0`` = a day the drawdown from the trailing 52w high - crosses up through ``threshold_pct`` (rising edge). De-duplicated by a - ``cooldown`` of trading days, so a continuous decline counts once but distinct - drawdowns separated by a recovery each register.""" + """Rising-edge corrections from the trailing 52-week high.""" events: list[dict] = [] - prev_dd = 0.0 + previous_drawdown = 0.0 last_event = -10**9 - for i in range(len(closes)): - window = closes[max(0, i - lookback + 1): i + 1] - hi = max(window) - dd = (hi - closes[i]) / hi * 100.0 if hi > 0 else 0.0 - if dd >= threshold_pct and prev_dd < threshold_pct and (i - last_event) >= cooldown: - events.append({"date": dates[i].isoformat(), "index": i, "depth_pct": round(dd, 1)}) - last_event = i - prev_dd = dd + for index, close in enumerate(closes): + high = max(closes[max(0, index - lookback + 1): index + 1]) + drawdown = (high - close) / high * 100.0 if high > 0 else 0.0 + if ( + drawdown >= threshold_pct + and previous_drawdown < threshold_pct + and index - last_event >= cooldown + ): + events.append({ + "date": dates[index].isoformat(), + "index": index, + "depth_pct": round(drawdown, 1), + }) + last_event = index + previous_drawdown = drawdown return events -# --------------------------------------------------------------------------- -# Event-centered: lead time + mean path -# --------------------------------------------------------------------------- - -def _lead(indicator: dict[date, float], t0: int, dates: list[date], pre: int, threshold: float) -> int | None: - """Earliest day within ``[t0-pre, t0]`` at which the indicator crosses - ``threshold`` — i.e. how many days of warning before the event, or None.""" - lead: int | None = None - for k in range(0, pre + 1): - idx = t0 - k - if idx < 0: - break - v = indicator.get(dates[idx]) - if v is not None and v >= threshold: - lead = k # keep going: the largest k = earliest warning in the window - return lead - - -def event_centered( +def alarm_episodes( indicator: dict[date, float], - events_idx: list[int], dates: list[date], - pre: int = PRE, - post: int = POST, - threshold: float = 60.0, -) -> dict: - """Align the indicator at each event's ``t0`` and measure how early it warned. + threshold: float, + start_index: int = 1, +) -> list[int]: + """Indices where the warning crosses upward; it must reset below first.""" + alarms: list[int] = [] + was_high = False + if start_index > 0: + previous = indicator.get(dates[start_index - 1]) + was_high = previous is not None and previous >= threshold + for index in range(start_index, len(dates)): + value = indicator.get(dates[index]) + if value is None: + continue + high = value >= threshold + if high and not was_high: + alarms.append(index) + was_high = high + return alarms - Lead time is measured against ``threshold`` (each indicator gets its own, - derived from its distribution). Also returns the cross-event mean path. - """ + +def evaluate_alarms( + alarm_indices: list[int], + event_indices: list[int], + dates: list[date], + horizon: int = HORIZON_DAYS, +) -> dict: + """Event recall, episode false alarms, and lead time for one holdout.""" leads: list[float] = [] - sums: dict[int, float] = {} - counts: dict[int, int] = {} - for t0 in events_idx: - lead = _lead(indicator, t0, dates, pre, threshold) + per_event: list[dict] = [] + warned = 0 + for event_index in event_indices: + matching = [ + alarm for alarm in alarm_indices if 0 < event_index - alarm <= horizon + ] + lead = max((event_index - alarm for alarm in matching), default=None) if lead is not None: - leads.append(lead) - for rel in range(-pre, post + 1): - idx = t0 + rel - if 0 <= idx < len(dates): - v = indicator.get(dates[idx]) - if v is not None: - sums[rel] = sums.get(rel, 0.0) + v - counts[rel] = counts.get(rel, 0) + 1 - mean_path = [ - {"rel_day": rel, "value": round(sums[rel] / counts[rel], 1)} for rel in sorted(sums) - ] + warned += 1 + leads.append(float(lead)) + per_event.append({ + "date": dates[event_index].isoformat(), + "warned": lead is not None, + "lead_days": lead, + }) + + false_alarms = sum( + 1 + for alarm in alarm_indices + if not any(0 < event - alarm <= horizon for event in event_indices) + ) return { + "events": len(event_indices), + "events_warned": warned, + "events_missed": len(event_indices) - warned, + "alarm_episodes": len(alarm_indices), + "false_alarms": false_alarms, "median_lead_days": _median(leads), - "events_with_signal": len(leads), - "events_total": len(events_idx), - "warn_threshold": round(threshold, 1), - "mean_path": mean_path, + "per_event": per_event, } -# --------------------------------------------------------------------------- -# Signal-centered: precision / recall vs. base rate -# --------------------------------------------------------------------------- - -def signal_centered( - indicator: dict[date, float], - events_idx: list[int], +def _warning_series( + prices: dict[str, rms.Series], + breadth_divergence: dict[date, float], dates: list[date], - horizon: int = HORIZON_DAYS, - thresholds: list[float] | None = None, -) -> dict: - """Treat ``indicator >= threshold`` as predicting a break within ``horizon`` - days. Sweep thresholds → precision/recall/alarm count, plus the base rate.""" - thresholds = thresholds or [50, 55, 60, 65, 70, 75, 80] - n = len(dates) - labels = [1 if any(i < e <= i + horizon for e in events_idx) else 0 for i in range(n)] - positives = sum(labels) - base_rate = positives / n if n else 0.0 - - rows: list[dict] = [] - for th in thresholds: - tp = fp = fn = 0 - for i in range(n): - v = indicator.get(dates[i]) - if v is None: - continue - pred = v >= th - if pred and labels[i]: - tp += 1 - elif pred and not labels[i]: - fp += 1 - elif not pred and labels[i]: - fn += 1 - precision = tp / (tp + fp) if (tp + fp) else None - recall = tp / (tp + fn) if (tp + fn) else None - rows.append({ - "threshold": th, - "precision": round(precision, 3) if precision is not None else None, - "recall": round(recall, 3) if recall is not None else None, - "alarms": tp + fp, - }) - return {"base_rate": round(base_rate, 3), "horizon_days": horizon, "rows": rows} - - -# --------------------------------------------------------------------------- -# Coincident baseline (deterministic price composite, reusing the regime sub-scores) -# --------------------------------------------------------------------------- - -def _coincident_series(prices: dict[str, list], dates: list[date], config: dict) -> dict[date, float]: - """Mean of the available price sub-scores (P1-P4) as-of each date — the - coincident baseline the leading candidate must beat on lead time.""" - lw = float(config.get("leader_weight", 2.0)) - lb = int(config.get("rs_lookback", 60)) - t = config["tickers"] - smh_full = prices.get(t["leaders"][0], []) if t["leaders"] else [] - qqq_full = prices.get(t["confirm"][0], []) if t["confirm"] else [] - spy_full = prices.get(t["market"], []) + config: dict, +) -> dict[date, float]: + """Technical Warning score used historically (fundamentals have no PIT history).""" + tickers = config["tickers"] + smh_full = prices.get(tickers["leaders"][0], []) + spy_full = prices.get(tickers["market"], []) out: dict[date, float] = {} - for d in dates: - smh = rms._closes_asof(smh_full, d) - qqq = rms._closes_asof(qqq_full, d) - spy = rms._closes_asof(spy_full, d) - subs = [ - rms.p1_trend_break(smh, qqq, lw), - rms.p2_death_cross(smh, qqq, lw), - rms.p3_drawdown(smh, qqq), - rms.p4_relative_strength(smh, spy, lb), - ] - vals = [v for v in subs if v is not None] - if vals: - out[d] = round(sum(vals) / len(vals), 2) + for session in dates: + divergence = breadth_divergence.get(session) + relative = rms.p4_relative_strength( + rms._closes_asof(smh_full, session), + rms._closes_asof(spy_full, session), + ) + values: list[tuple[float, float]] = [] + if divergence is not None: + values.append((divergence, rms.WARNING_WEIGHTS["breadth_divergence"])) + if relative is not None: + values.append((relative, rms.WARNING_WEIGHTS["relative_strength"])) + if values: + out[session] = round( + sum(value * weight for value, weight in values) + / sum(weight for _, weight in values), + 2, + ) return out -# --------------------------------------------------------------------------- -# Orchestration -# --------------------------------------------------------------------------- - async def run_event_study( db: AsyncSession, threshold_pct: float = EVENT_THRESHOLD_PCT, horizon: int = HORIZON_DAYS, - cooldown: int = COOLDOWN_DAYS, - warn_percentile: float = WARN_PERCENTILE, ) -> dict: - """Run the study: detect events on the benchmark, then measure breadth-divergence - vs. the coincident price composite. Best-effort; returns available=False on no data.""" config = await rms.get_regime_config(db) end = date.today() start = end - timedelta(days=5 * 365 + 30) - prices = await rms._fetch_prices(config, start, end) - leader = config["tickers"]["leaders"][0] if config["tickers"]["leaders"] else "SMH" - bench = sorted(prices.get(leader, []), key=lambda x: x[0]) - if len(bench) < 260: + leader = config["tickers"]["leaders"][0] + benchmark = sorted(prices.get(leader, []), key=lambda item: item[0]) + if len(benchmark) < 500: return {"available": False, "reason": "insufficient benchmark history"} - dates = [d for d, _ in bench] - closes = [c for _, c in bench] - events = detect_events(closes, dates, threshold_pct, cooldown=cooldown) - events_idx = [e["index"] for e in events] + dates = [d for d, _ in benchmark] + closes = [value for _, value in benchmark] + breadth, _ = await breadth_service.compute_breadth_details( + db, config["breadth_basket"], window=200, min_tickers=20 + ) + divergence = breadth_service.compute_divergence_series(breadth, benchmark) + warning = _warning_series(prices, divergence, dates, config) - breadth = await breadth_service.compute_breadth_series(db) - divergence = breadth_service.compute_divergence_series(breadth, bench) - coincident = _coincident_series(prices, dates, config) + split = max(1, min(len(dates) - 1, int(len(dates) * TRAIN_FRACTION))) + train_values = [warning[d] for d in dates[:split] if d in warning] + warn_threshold = _percentile(train_values, WARN_PERCENTILE) + if warn_threshold is None: + return {"available": False, "reason": "insufficient warning history"} - # Each indicator warns at its OWN distribution's percentile, so a leading - # indicator isn't penalised for living on a different scale than the baseline. - warn = { - "breadth_divergence": _percentile(list(divergence.values()), warn_percentile) or 60.0, - "coincident_price": _percentile(list(coincident.values()), warn_percentile) or 60.0, - } - series_by_key = {"breadth_divergence": divergence, "coincident_price": coincident} + all_events = detect_events(closes, dates, threshold_pct) + holdout_events = [event["index"] for event in all_events if event["index"] >= split] + alarms = alarm_episodes(warning, dates, warn_threshold, start_index=split) + metrics = evaluate_alarms(alarms, holdout_events, dates, horizon) + holdout_sessions = max(1, len(dates) - split) + metrics["false_alarms_per_year"] = round( + metrics["false_alarms"] / (holdout_sessions / 252.0), 2 + ) - def _evaluate(series: dict[date, float], threshold: float) -> dict: - return { - **event_centered(series, events_idx, dates, threshold=threshold), - "signal": signal_centered(series, events_idx, dates, horizon), - } - - indicators = {key: _evaluate(series_by_key[key], warn[key]) for key in series_by_key} - - # Per-event comparison: which event, and each indicator's lead on THAT event — - # so a median over a tiny sample can't hide an apples-to-oranges comparison. - per_event = [ - { - "date": e["date"], - "depth_pct": e["depth_pct"], - "breadth_lead": _lead(divergence, e["index"], dates, PRE, warn["breadth_divergence"]), - "coincident_lead": _lead(coincident, e["index"], dates, PRE, warn["coincident_price"]), - } - for e in events - ] - - bd = indicators["breadth_divergence"]["median_lead_days"] - cd = indicators["coincident_price"]["median_lead_days"] - lead_delta = (bd - cd) if (bd is not None and cd is not None) else None - - recent_breadth = [ - {"date": d.isoformat(), "breadth": breadth[d], "divergence": divergence.get(d)} - for d in dates[-90:] - if d in breadth - ] + basket_asof = date.fromisoformat(config["basket_asof"]) + retrospective = dates[split] < basket_asof + evaluation = "exploratory" if retrospective else "holdout" + lead_text = ( + f"median lead {metrics['median_lead_days']:.0f} sessions" + if metrics["median_lead_days"] is not None + else "no successful warning lead" + ) + summary = ( + f"{evaluation.capitalize()} chronological test: warning episodes preceded " + f"{metrics['events_warned']}/{metrics['events']} 10% corrections; " + f"{metrics['events_missed']} missed, {metrics['false_alarms_per_year']:.1f} " + f"false alarms/year, {lead_text}." + ) + per_event = metrics.pop("per_event") report = { "available": True, + "methodology": rms.METHODOLOGY, "generated_at": datetime.now(timezone.utc).isoformat(), + "evaluation": evaluation, + "summary": summary, "params": { "benchmark": leader, + "outcome": "10% correction from trailing 52-week high", "event_threshold_pct": threshold_pct, - "cooldown_days": cooldown, + "event_cooldown_days": EVENT_COOLDOWN_DAYS, "horizon_days": horizon, - "warn_percentile": warn_percentile, + "train_fraction": TRAIN_FRACTION, + "warn_percentile": WARN_PERCENTILE, + "warn_threshold": round(warn_threshold, 1), + "basket_hash": rms._basket_hash(config["breadth_basket"]), + "basket_asof": config["basket_asof"], }, - "events": events, - "indicators": indicators, - "per_event": per_event, - "lead_delta_days": lead_delta, - "recent_breadth": recent_breadth, + "sample": { + "start": dates[0].isoformat(), + "end": dates[-1].isoformat(), + "train_end": dates[split - 1].isoformat(), + "test_start": dates[split].isoformat(), + "sessions": len(dates), + "holdout_sessions": holdout_sessions, + }, + "metrics": metrics, + "events": per_event, + "recent_breadth": [ + {"date": d.isoformat(), "breadth": breadth[d], "warning": warning.get(d)} + for d in dates[-90:] + if d in breadth + ], } logger.info(json.dumps({ - "event": "event_study_complete", "events": len(events), - "breadth_lead": bd, "coincident_lead": cd, + "event": "regime_event_study_complete", + "evaluation": evaluation, + "events": metrics["events"], + "warned": metrics["events_warned"], + "false_alarms_per_year": metrics["false_alarms_per_year"], })) return report async def run_and_store(db: AsyncSession) -> dict: - """Run the event study and cache the report in a SystemSetting. Job entrypoint.""" report = await run_event_study(db) await update_setting(db, KEY_REPORT, json.dumps(report)) return report async def get_event_study_report(db: AsyncSession) -> dict | None: - """Return the last cached event-study report, or None if never run.""" setting = await settings_store.get_setting(db, KEY_REPORT) if setting is None: return None try: - return json.loads(setting.value) + report = json.loads(setting.value) except (TypeError, ValueError): return None + return report if report.get("methodology") == rms.METHODOLOGY else None diff --git a/app/services/regime_monitor_service.py b/app/services/regime_monitor_service.py index 7975420..f1df46f 100644 --- a/app/services/regime_monitor_service.py +++ b/app/services/regime_monitor_service.py @@ -1,29 +1,21 @@ -"""AI/Tech Regime-Change Monitor. +"""AI/Tech Regime Monitor v2. -A standalone, observational tool: it scores how far the AI/Tech bull regime has -deteriorated toward a re-rating, as a single 0-100 **index** (not a calibrated -probability), broken down per signal. It is intentionally decoupled — nothing -here feeds gates, scoring, alerts, or trade logic. It only computes a number for -its own tab. +The monitor is a risk thermometer, not a probability or trading rule. It keeps +two deliberately separate outputs: -Design mirrors ``market_regime_service``: benchmark/sector bars are pulled -directly via Alpaca (no Universe membership needed), macro inputs (VIX, HY credit -spreads) come from FRED, and the daily result is persisted as one -``RegimeSnapshot`` row per date so the UI can show a 7/30-day trend. On the first -run the history is backfilled by replaying the (already-fetched) price/FRED series -as-of each past day, so the trend is populated immediately. +* State: current structural stress (price, breadth, credit, volatility). +* Warning: deterioration/divergence that may precede State (breadth, relative + strength, and sourced fundamental observations). -Signals (sub-score 0 = healthy … 100 = regime breaking): - P1 trend break (% under 200-DMA, SMH-led) P2 death cross + 200-slope - P3 drawdown from 52w high P4 relative strength SMH/SPY - P5 volatility (VIX) P6 NVDA canary divergence (opt.) - F1 hyperscaler capex guidance (LLM/manual) F2 HY credit-spread percentile - F3 "good news, stock down" (LLM/manual) F4 market breadth RSP/SPY +Daily snapshots are the point-in-time record. The first v2 run rewrites the +latest ``REBUILD_SESSIONS`` trading sessions once; ordinary runs thereafter only +upsert the latest trading date. Fundamental observations are never replayed +before their effective date. """ from __future__ import annotations -import copy +import hashlib import json import logging import os @@ -31,11 +23,11 @@ from datetime import date, datetime, timedelta, timezone from pathlib import Path import httpx -from sqlalchemy import func, select +from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession from app.config import settings -from app.exceptions import ProviderError +from app.exceptions import ProviderError, ValidationError from app.models.regime_snapshot import RegimeSnapshot from app.providers.alpaca import AlpacaOHLCVProvider from app.services import breadth_service, settings_store @@ -49,50 +41,60 @@ _CA_BUNDLE = os.environ.get("SSL_CERT_FILE", "") KEY_CONFIG = "regime_monitor_config" KEY_FUNDAMENTALS = "regime_fundamental_overrides" -# All weights/thresholds are admin-editable via the KEY_CONFIG SystemSetting. -# Default weights sum to 100 (P6 off). SMH is the leading sensor, QQQ confirms. +METHODOLOGY = "v2" +REBUILD_SESSIONS = 400 +MIN_COVERAGE = 75.0 +SOURCE_MAX_LAG_DAYS = 7 + +QUADRANT_STATE_DIVIDER = 60.0 +QUADRANT_WARNING_DIVIDER = 60.0 +QUADRANT_MARGIN = 5.0 + +HY_OAS_MILD = 3.5 +HY_OAS_ELEVATED = 5.0 +HY_OAS_STRESSED = 7.0 +HY_OAS_REFERENCE_YEARS = 10.0 + +STATE_WEIGHTS = { + "price": 40.0, + "breadth": 25.0, + "credit": 20.0, + "volatility": 15.0, +} +WARNING_WEIGHTS = { + "breadth_divergence": 50.0, + "relative_strength": 30.0, + "capex": 12.0, + "earnings_reaction": 8.0, +} + +# Fixed at the v2 launch. These are liquid S&P 500/Nasdaq AI, semiconductor, +# infrastructure, cloud, and enterprise-software names that the platform's +# normal universe sync already stores. +DEFAULT_BREADTH_BASKET = [ + "AAPL", "MSFT", "NVDA", "AMZN", "META", "GOOGL", "AVGO", "AMD", + "ORCL", "CRM", "NOW", "PLTR", "ANET", "DELL", "SMCI", "MU", + "QCOM", "INTC", "AMAT", "LRCX", "KLAC", "SNPS", "CDNS", "ADI", + "TXN", "IBM", "CSCO", "PANW", "CRWD", "VRT", +] + DEFAULT_CONFIG: dict = { "tickers": { - "leaders": ["SMH"], # semis — fast early signal - "confirm": ["QQQ"], # broad tech — confirmation + "leaders": ["SMH"], + "confirm": ["QQQ"], "market": "SPY", - "breadth": "RSP", # equal-weight S&P for breadth - "canary": "NVDA", # sector lead-stock (optional early warning) "hyperscalers": ["GOOGL", "AMZN", "META", "MSFT"], }, - "weights": { - "P1": 12, "P2": 8, "P3": 10, "P4": 8, "P5": 7, "P6": 0, - "F1": 25, "F2": 15, "F3": 8, "F4": 7, - }, - "alert_threshold": 65, - # Observational early-warning blend: a small Combined score = weighted mean of - # the coincident index and the breadth-divergence early-warning score. Kept - # separate from the index weights above so the early-warning side stays - # decoupled until proven. Tunable; need not sum to 1 (normalised). - "combined_weights": {"coincident": 0.6, "early_warning": 0.4}, - "leader_weight": 2.0, # SMH counts 2x vs QQQ where both feed a signal - "rs_lookback": 60, # trading days for relative-strength / breadth trend + "breadth_basket": DEFAULT_BREADTH_BASKET, + "basket_asof": "2026-07-15", "fundamental_staleness_days": 80, } -SIGNAL_LABELS: dict[str, str] = { - "P1": "Trend break (200-DMA)", - "P2": "Death cross + slope", - "P3": "Drawdown from 52w high", - "P4": "Relative strength SMH/SPY", - "P5": "Volatility (VIX)", - "P6": "NVDA canary divergence", - "F1": "Hyperscaler capex guidance", - "F2": "Credit spreads (HY OAS)", - "F3": "Good news, stock down", - "F4": "Market breadth RSP/SPY", -} - -_PRICE_SIGNALS = {"P1", "P2", "P3", "P4", "P6"} +Series = list[tuple[date, float]] # --------------------------------------------------------------------------- -# Small numeric helpers +# Pure numeric helpers and sensors # --------------------------------------------------------------------------- def _clamp(x: float, lo: float = 0.0, hi: float = 100.0) -> float: @@ -109,8 +111,7 @@ def _mean(values: list[float]) -> float | None: return sum(values) / len(values) if values else None -def _blend(leader: float | None, confirm: float | None, leader_weight: float) -> float | None: - """Weighted blend of a leading vs a confirming sub-score (SMH vs QQQ).""" +def _blend(leader: float | None, confirm: float | None, leader_weight: float = 2.0) -> float | None: parts: list[tuple[float, float]] = [] if leader is not None: parts.append((leader, leader_weight)) @@ -118,13 +119,10 @@ def _blend(leader: float | None, confirm: float | None, leader_weight: float) -> parts.append((confirm, 1.0)) if not parts: return None - num = sum(v * w for v, w in parts) - den = sum(w for _, w in parts) - return num / den + return sum(v * w for v, w in parts) / sum(w for _, w in parts) def band_for(score: float) -> str: - """Map the 0-100 index onto its label band.""" if score < 30: return "stable" if score < 60: @@ -134,10 +132,6 @@ def band_for(score: float) -> str: return "breaking" -# --------------------------------------------------------------------------- -# Pure sub-score functions (0 = healthy, 100 = regime breaking). None = no data. -# --------------------------------------------------------------------------- - def _under_200(closes: list[float]) -> float | None: sma200 = _sma(closes, 200) if sma200 is None: @@ -145,8 +139,7 @@ def _under_200(closes: list[float]) -> float | None: return 100.0 if closes[-1] < sma200 else 0.0 -def p1_trend_break(smh: list[float], qqq: list[float], leader_weight: float) -> float | None: - """Weighted share trading below the 200-DMA (SMH leads, QQQ confirms).""" +def p1_trend_break(smh: list[float], qqq: list[float], leader_weight: float = 2.0) -> float | None: return _blend(_under_200(smh), _under_200(qqq), leader_weight) @@ -156,28 +149,27 @@ def _death_cross(closes: list[float]) -> float | None: if sma50 is None or sma200 is None or len(closes) < 221 or sma200 == 0: return None gap_pct = (sma50 / sma200 - 1.0) * 100.0 - severity = 0.0 if gap_pct >= 0 else _clamp(-gap_pct * 20.0) # -5% gap -> 100 + severity = 0.0 if gap_pct >= 0 else _clamp(-gap_pct * 20.0) sma200_past = _sma(closes[:-20], 200) - slope_factor = 1.0 if sma200_past: slope_pct = (sma200 / sma200_past - 1.0) * 100.0 - slope_factor = 1.0 if slope_pct < 0 else 0.5 # damp if 200 still rising - return severity * slope_factor + if slope_pct >= 0: + severity *= 0.5 + return severity -def p2_death_cross(smh: list[float], qqq: list[float], leader_weight: float) -> float | None: +def p2_death_cross(smh: list[float], qqq: list[float], leader_weight: float = 2.0) -> float | None: return _blend(_death_cross(smh), _death_cross(qqq), leader_weight) def _drawdown(closes: list[float]) -> float | None: if len(closes) < 30: return None - window = closes[-252:] - peak = max(window) + peak = max(closes[-252:]) if peak <= 0: return None dd_pct = (peak - closes[-1]) / peak * 100.0 - return _clamp(dd_pct * 5.0) # 20% below 52w high -> 100 + return _clamp(dd_pct * 5.0) def p3_drawdown(smh: list[float], qqq: list[float]) -> float | None: @@ -185,119 +177,188 @@ def p3_drawdown(smh: list[float], qqq: list[float]) -> float | None: return max(vals) if vals else None -def _ratio_trend(a: list[float], b: list[float], lookback: int) -> float | None: - """Falling a/b (a underperforming b) -> higher score. Flat -> 50.""" - if len(a) < lookback + 1 or len(b) < lookback + 1: +def p4_relative_strength(smh: list[float], spy: list[float], lookback: int = 60) -> float | None: + """Stress-only SMH/SPY rollover: flat/outperformance=0, -10%=100.""" + if len(smh) < lookback + 1 or len(spy) < lookback + 1: return None - if b[-1] == 0 or b[-lookback - 1] == 0: + if spy[-1] == 0 or spy[-lookback - 1] == 0: return None - now = a[-1] / b[-1] - past = a[-lookback - 1] / b[-lookback - 1] + now = smh[-1] / spy[-1] + past = smh[-lookback - 1] / spy[-lookback - 1] if past == 0: return None chg_pct = (now / past - 1.0) * 100.0 - return _clamp(50.0 - chg_pct * 5.0) # -10% -> 100, +10% -> 0 - - -def p4_relative_strength(smh: list[float], spy: list[float], lookback: int) -> float | None: - return _ratio_trend(smh, spy, lookback) - - -def f4_breadth(rsp: list[float], spy: list[float], lookback: int) -> float | None: - """Narrowing breadth (equal-weight lagging cap-weight) -> RSP/SPY falls -> higher.""" - return _ratio_trend(rsp, spy, lookback) + return _clamp(-chg_pct * 10.0) def p5_volatility(vix: float | None) -> float | None: if vix is None: return None - return _clamp((vix - 15.0) / 15.0 * 100.0) # <=15 -> 0, >=30 -> 100 + return _clamp((vix - 15.0) / 15.0 * 100.0) + + +def breadth_level_score(pct_above_200: float | None) -> float | None: + """Broad >=60%=healthy; <=20%=full breadth stress; linear between.""" + if pct_above_200 is None: + return None + return _clamp((60.0 - pct_above_200) / 40.0 * 100.0) + + +def _oas_absolute_score(value: float) -> float: + if value <= HY_OAS_MILD: + return 0.0 + if value <= HY_OAS_ELEVATED: + return (value - HY_OAS_MILD) / (HY_OAS_ELEVATED - HY_OAS_MILD) * 50.0 + if value < HY_OAS_STRESSED: + return 50.0 + (value - HY_OAS_ELEVATED) / (HY_OAS_STRESSED - HY_OAS_ELEVATED) * 50.0 + return 100.0 def f2_credit_spreads(oas_values: list[float]) -> float | None: - """Percentile rank of the latest HY OAS within the window. Wider = higher.""" - if len(oas_values) < 30: + """HY OAS stress: 70% named absolute anchors + 30% upper-tail percentile.""" + if not oas_values: return None latest = oas_values[-1] - below = sum(1 for v in oas_values if v <= latest) - return _clamp(below / len(oas_values) * 100.0) + absolute = _oas_absolute_score(latest) + if len(oas_values) < 30: + return round(absolute, 2) + less = sum(1 for v in oas_values if v < latest) + equal = sum(1 for v in oas_values if v == latest) + percentile = (less + 0.5 * equal) / len(oas_values) * 100.0 + relative = _clamp((percentile - 50.0) / 45.0 * 100.0) + return round(absolute * 0.7 + relative * 0.3, 2) -def p6_canary(nvda: list[float], smh: list[float]) -> float | None: - """NVDA below its 50-DMA while SMH's trend is still intact = lead divergence.""" - sma50 = _sma(nvda, 50) - if sma50 is None: - return None - nvda_weak = nvda[-1] < sma50 - sma200 = _sma(smh, 200) - smh_intact = sma200 is not None and smh[-1] > sma200 - if nvda_weak and smh_intact: - return 100.0 - return 50.0 if nvda_weak else 0.0 +def _sensor(sensor_id: str, label: str, score: float | None, **details: object) -> dict: + return { + "id": sensor_id, + "label": label, + "score": round(score, 1) if score is not None else None, + "available": score is not None, + "details": details, + } -# --------------------------------------------------------------------------- -# Aggregation -# --------------------------------------------------------------------------- - -def compute_regime_score(sub_scores: dict[str, float | None], weights: dict[str, float]) -> dict: - """Weighted mean over the *available* signals (weight>0 and score present). - - Missing-data signals drop out of both numerator and denominator and are - reported with ``available=False``. Contributions sum to the total. - """ - denom = sum( - weights.get(sid, 0) - for sid, score in sub_scores.items() - if score is not None and weights.get(sid, 0) > 0 +def _score_pillars(pillars: list[dict], weights: dict[str, float]) -> dict: + expected = sum(max(0.0, float(w)) for w in weights.values()) + available_weight = sum( + max(0.0, float(weights.get(p["id"], 0.0))) + for p in pillars + if p.get("score") is not None ) - total = 0.0 - breakdown: list[dict] = [] - for sid in SIGNAL_LABELS: - weight = weights.get(sid, 0) - if weight <= 0: - continue - score = sub_scores.get(sid) - available = score is not None - contribution = (score * weight / denom) if (available and denom > 0) else 0.0 - if available: - total += contribution - breakdown.append({ - "id": sid, - "label": SIGNAL_LABELS[sid], - "sub_score": round(score, 1) if available else None, - "weight": weight, - "available": available, - "contribution": round(contribution, 2), - }) - return {"total_score": round(total, 1), "band": band_for(total), "breakdown": breakdown} + coverage = available_weight / expected * 100.0 if expected else 0.0 + score = None + if available_weight: + score = sum( + float(p["score"]) * float(weights.get(p["id"], 0.0)) + for p in pillars + if p.get("score") is not None + ) / available_weight + + rows: list[dict] = [] + for pillar in pillars: + row = dict(pillar) + weight = float(weights.get(row["id"], 0.0)) + row["weight"] = weight + row["available"] = row.get("score") is not None + row["contribution"] = ( + round(float(row["score"]) * weight / available_weight, 2) + if row["available"] and available_weight + else 0.0 + ) + rows.append(row) + + rounded = round(score, 1) if score is not None else None + return { + "score": rounded, + "band": band_for(rounded) if rounded is not None and coverage >= MIN_COVERAGE else None, + "coverage": round(coverage, 1), + "minimum_coverage": MIN_COVERAGE, + "available_pillars": [p["id"] for p in rows if p["available"]], + "pillars": rows, + } # --------------------------------------------------------------------------- -# As-of series helpers (for backfill replay) +# Point-in-time helpers # --------------------------------------------------------------------------- -Series = list[tuple[date, float]] - - def _closes_asof(series: Series, as_of: date) -> list[float]: return [v for d, v in series if d <= as_of] -def _value_asof(series: Series | None, as_of: date) -> float | None: +def _item_asof(series: Series | None, as_of: date) -> tuple[date, float] | None: if not series: return None - vals = [v for d, v in series if d <= as_of] - return vals[-1] if vals else None + chosen: tuple[date, float] | None = None + for item in series: + if item[0] <= as_of: + chosen = item + else: + break + return chosen + + +def _value_asof(series: Series | None, as_of: date) -> float | None: + item = _item_asof(series, as_of) + return item[1] if item else None def _window_asof(series: Series | None, as_of: date, years: float) -> list[float]: if not series: return [] - start = as_of - timedelta(days=int(365 * years)) + start = as_of - timedelta(days=int(365.25 * years)) return [v for d, v in series if start <= d <= as_of] +def _next_weekday(d: date) -> date: + candidate = d + timedelta(days=1) + while candidate.weekday() >= 5: + candidate += timedelta(days=1) + return candidate + + +def _parse_date(value: object) -> date | None: + if not value: + return None + try: + return date.fromisoformat(str(value)[:10]) + except ValueError: + return None + + +def _fundamental_effective_date(overrides: dict) -> date | None: + explicit = _parse_date(overrides.get("effective_date")) + if explicit: + return explicit + fetched = _parse_date(overrides.get("fetched_at")) + return _next_weekday(fetched) if fetched else None + + +def _fundamental_scores_asof(overrides: dict, config: dict, as_of: date) -> tuple[float | None, float | None, dict]: + effective = _fundamental_effective_date(overrides) + if effective is None or as_of < effective: + return None, None, {"effective_date": effective.isoformat() if effective else None, "age_days": None} + age = (as_of - effective).days + stale = age > int(config.get("fundamental_staleness_days", 80)) + f1 = overrides.get("f1_score") + f3 = overrides.get("f3_score") + return ( + None if stale or f1 is None else _clamp(float(f1)), + None if stale or f3 is None else _clamp(float(f3)), + {"effective_date": effective.isoformat(), "age_days": age, "stale": stale}, + ) + + +def _basket_hash(symbols: list[str]) -> str: + canonical = ",".join(sorted({s.strip().upper() for s in symbols if s.strip()})) + return hashlib.sha256(canonical.encode("utf-8")).hexdigest()[:12] + + +def _mapping_series(values: dict[date, float]) -> Series: + return sorted(values.items(), key=lambda item: item[0]) + + def _compute_index( prices: dict[str, Series], vix_series: Series | None, @@ -305,82 +366,206 @@ def _compute_index( overrides: dict, config: dict, as_of: date, + breadth_series: Series | None = None, + divergence_series: Series | None = None, + breadth_counts: dict[date, int] | None = None, ) -> dict: - """Compute the full index result as-of *as_of* from raw series.""" - t = config["tickers"] - lw = float(config.get("leader_weight", 2.0)) - lb = int(config.get("rs_lookback", 60)) + """Compute the complete v2 State/Warning snapshot as of one trading date.""" + tickers = config["tickers"] + smh = _closes_asof(prices.get(tickers["leaders"][0], []), as_of) + qqq = _closes_asof(prices.get(tickers["confirm"][0], []), as_of) + spy = _closes_asof(prices.get(tickers["market"], []), as_of) - smh = _closes_asof(prices.get(t["leaders"][0], []), as_of) if t["leaders"] else [] - qqq = _closes_asof(prices.get(t["confirm"][0], []), as_of) if t["confirm"] else [] - spy = _closes_asof(prices.get(t["market"], []), as_of) - rsp = _closes_asof(prices.get(t["breadth"], []), as_of) - nvda = _closes_asof(prices.get(t["canary"], []), as_of) - vix = _value_asof(vix_series, as_of) - oas = _window_asof(oas_series, as_of, 3) + p1 = p1_trend_break(smh, qqq) + p2 = p2_death_cross(smh, qqq) + p3 = p3_drawdown(smh, qqq) + price_values = [v for v in (p1, p2, p3) if v is not None] + price_score = max(price_values) if price_values else None - sub_scores: dict[str, float | None] = { - "P1": p1_trend_break(smh, qqq, lw), - "P2": p2_death_cross(smh, qqq, lw), - "P3": p3_drawdown(smh, qqq), - "P4": p4_relative_strength(smh, spy, lb), - "P5": p5_volatility(vix), - "P6": p6_canary(nvda, smh), - "F1": overrides.get("f1_score"), - "F2": f2_credit_spreads(oas), - "F3": overrides.get("f3_score"), - "F4": f4_breadth(rsp, spy, lb), + breadth_item = _item_asof(breadth_series, as_of) + breadth_pct = breadth_item[1] if breadth_item else None + breadth_score = breadth_level_score(breadth_pct) + vix_item = _item_asof(vix_series, as_of) + vix_score = p5_volatility(vix_item[1] if vix_item else None) + oas_item = _item_asof(oas_series, as_of) + oas_window = _window_asof(oas_series, as_of, HY_OAS_REFERENCE_YEARS) + credit_score = f2_credit_spreads(oas_window) + + divergence = _value_asof(divergence_series, as_of) + relative_strength = p4_relative_strength(smh, spy) + f1, f3, fundamental_meta = _fundamental_scores_asof(overrides, config, as_of) + + state_pillars = [ + { + "id": "price", + "label": "Price structure", + "score": round(price_score, 1) if price_score is not None else None, + "sensors": [ + _sensor("P1", "Trend break (200-DMA)", p1), + _sensor("P2", "Death cross + slope", p2), + _sensor("P3", "Drawdown from 52w high", p3), + ], + }, + { + "id": "breadth", + "label": "Breadth level", + "score": round(breadth_score, 1) if breadth_score is not None else None, + "sensors": [_sensor("B1", "% basket above 200-DMA", breadth_score, pct_above_200=breadth_pct)], + }, + { + "id": "credit", + "label": "Credit level", + "score": round(credit_score, 1) if credit_score is not None else None, + "sensors": [_sensor("C1", "HY option-adjusted spread", credit_score, oas=oas_item[1] if oas_item else None)], + }, + { + "id": "volatility", + "label": "Volatility level", + "score": round(vix_score, 1) if vix_score is not None else None, + "sensors": [_sensor("V1", "VIX level", vix_score, vix=vix_item[1] if vix_item else None)], + }, + ] + warning_pillars = [ + { + "id": "breadth_divergence", + "label": "Breadth divergence", + "score": round(divergence, 1) if divergence is not None else None, + "sensors": [_sensor("W1", "Price holding while breadth narrows", divergence)], + }, + { + "id": "relative_strength", + "label": "SMH/SPY rollover", + "score": round(relative_strength, 1) if relative_strength is not None else None, + "sensors": [_sensor("W2", "60-session relative-strength deterioration", relative_strength)], + }, + { + "id": "capex", + "label": "Hyperscaler capex revisions", + "score": round(f1, 1) if f1 is not None else None, + "sensors": [_sensor("F1", "Capex guidance cuts", f1)], + }, + { + "id": "earnings_reaction", + "label": "Good news, stock down", + "score": round(f3, 1) if f3 is not None else None, + "sensors": [_sensor("F3", "Abnormal earnings reaction", f3)], + }, + ] + + state = _score_pillars(state_pillars, STATE_WEIGHTS) + warning = _score_pillars(warning_pillars, WARNING_WEIGHTS) + + price_item = _item_asof(prices.get(tickers["leaders"][0]), as_of) + dated_sources = { + "price": price_item[0] if price_item else None, + "breadth": breadth_item[0] if breadth_item else None, + "vix": vix_item[0] if vix_item else None, + "credit": oas_item[0] if oas_item else None, } - - result = compute_regime_score(sub_scores, config["weights"]) - result["date"] = as_of.isoformat() - result["alert_threshold"] = config.get("alert_threshold", 65) - result["inputs"] = { - "vix": round(vix, 2) if vix is not None else None, - "hy_oas": round(oas[-1], 2) if oas else None, - "fundamentals_fetched_at": overrides.get("fetched_at"), + source_ages = { + key: (as_of - d).days for key, d in dated_sources.items() if d is not None + } + stale_inputs = [key for key, age in source_ages.items() if age > SOURCE_MAX_LAG_DAYS] + basket = list(config["breadth_basket"]) + basket_count = None + if breadth_counts and breadth_item: + basket_count = breadth_counts.get(breadth_item[0]) + + return { + "methodology": METHODOLOGY, + "date": as_of.isoformat(), + "state": state, + "warning": warning, + "quadrant_config": { + "state_divider": QUADRANT_STATE_DIVIDER, + "warning_divider": QUADRANT_WARNING_DIVIDER, + "margin": QUADRANT_MARGIN, + }, + "basket": { + "symbols": basket, + "hash": _basket_hash(basket), + "basket_asof": config["basket_asof"], + "members_available": basket_count, + "members_expected": len(basket), + "history_kind": "forward" if as_of >= date.fromisoformat(config["basket_asof"]) else "retrospective", + }, + "inputs": { + "vix": round(vix_item[1], 2) if vix_item else None, + "vix_date": vix_item[0].isoformat() if vix_item else None, + "hy_oas": round(oas_item[1], 2) if oas_item else None, + "hy_oas_date": oas_item[0].isoformat() if oas_item else None, + "breadth_pct_above_200": round(breadth_pct, 1) if breadth_pct is not None else None, + "breadth_date": breadth_item[0].isoformat() if breadth_item else None, + "fundamentals_fetched_at": overrides.get("fetched_at"), + "fundamentals_effective_date": fundamental_meta.get("effective_date"), + "fundamentals_age_days": fundamental_meta.get("age_days"), + }, + "data_quality": { + "minimum_coverage": MIN_COVERAGE, + "oldest_market_input_age_days": max(source_ages.values()) if source_ages else None, + "stale_inputs": stale_inputs, + "inputs_fresh": not stale_inputs, + }, } - return result # --------------------------------------------------------------------------- -# Config + fundamental-override storage +# Configuration and fundamental storage # --------------------------------------------------------------------------- +def _normalise_basket(symbols: list[str]) -> list[str]: + cleaned = [str(s).strip().upper().replace(".", "-") for s in symbols if str(s).strip()] + if len(cleaned) != len(set(cleaned)): + raise ValidationError("Breadth basket symbols must be unique") + if not 20 <= len(cleaned) <= 100: + raise ValidationError("Breadth basket must contain between 20 and 100 symbols") + return cleaned + + async def get_regime_config(db: AsyncSession) -> dict: - """DEFAULT_CONFIG deep-merged with the stored override (nested for dicts).""" - cfg = copy.deepcopy(DEFAULT_CONFIG) + cfg = json.loads(json.dumps(DEFAULT_CONFIG)) raw = await settings_store.get_value(db, KEY_CONFIG) if raw: try: stored = json.loads(raw) - for k, v in stored.items(): - if isinstance(v, dict) and isinstance(cfg.get(k), dict): - cfg[k].update(v) - else: - cfg[k] = v - except (TypeError, ValueError): - logger.warning("Corrupt %s; using defaults", KEY_CONFIG) + if isinstance(stored.get("breadth_basket"), list): + cfg["breadth_basket"] = _normalise_basket(stored["breadth_basket"]) + if stored.get("basket_asof"): + cfg["basket_asof"] = str(stored["basket_asof"]) + if stored.get("fundamental_staleness_days") is not None: + cfg["fundamental_staleness_days"] = int(stored["fundamental_staleness_days"]) + except (TypeError, ValueError, ValidationError): + logger.warning("Corrupt %s; using v2 defaults", KEY_CONFIG) return cfg async def update_regime_config(db: AsyncSession, updates: dict) -> dict: - """Merge *updates* into the stored config and persist. Returns the new config.""" cfg = await get_regime_config(db) - for k, v in (updates or {}).items(): - if isinstance(v, dict) and isinstance(cfg.get(k), dict): - cfg[k].update(v) - else: - cfg[k] = v + if "breadth_basket" in updates: + basket = _normalise_basket(updates["breadth_basket"]) + if basket != cfg["breadth_basket"]: + cfg["breadth_basket"] = basket + cfg["basket_asof"] = date.today().isoformat() + if "fundamental_staleness_days" in updates: + days = int(updates["fundamental_staleness_days"]) + if not 30 <= days <= 180: + raise ValidationError("Fundamental staleness must be between 30 and 180 days") + cfg["fundamental_staleness_days"] = days await update_setting(db, KEY_CONFIG, json.dumps(cfg)) return cfg async def get_fundamental_overrides(db: AsyncSession) -> dict: - """Current F1/F3 override (LLM-proposed or manual). Defaults to neutral 50.""" + default = { + "f1_score": None, + "f3_score": None, + "locked": False, + "reasoning": None, + "fetched_at": None, + "effective_date": None, + "source": "default", + } raw = await settings_store.get_value(db, KEY_FUNDAMENTALS) - default = {"f1_score": 50.0, "f3_score": 50.0, "locked": False, - "reasoning": None, "fetched_at": None, "source": "default"} if not raw: return default try: @@ -396,35 +581,35 @@ async def set_fundamental_overrides( f3_score: float | None = None, locked: bool | None = None, ) -> dict: - """Manual override of F1/F3. Setting any value locks out the LLM refresh - unless ``locked`` is explicitly cleared.""" current = await get_fundamental_overrides(db) + observation_changed = f1_score is not None or f3_score is not None if f1_score is not None: current["f1_score"] = _clamp(float(f1_score)) if f3_score is not None: current["f3_score"] = _clamp(float(f3_score)) if locked is not None: current["locked"] = bool(locked) - elif f1_score is not None or f3_score is not None: + elif observation_changed: current["locked"] = True - current["source"] = "manual" - current["fetched_at"] = datetime.now(timezone.utc).isoformat() + if observation_changed: + now = datetime.now(timezone.utc) + current.update({ + "source": "manual", + "fetched_at": now.isoformat(), + "effective_date": _next_weekday(now.date()).isoformat(), + }) await update_setting(db, KEY_FUNDAMENTALS, json.dumps(current)) return current # --------------------------------------------------------------------------- -# Data fetching: Alpaca prices + FRED macro +# External data fetching # --------------------------------------------------------------------------- def _price_symbols(config: dict) -> list[str]: - t = config["tickers"] - syms = list(t["leaders"]) + list(t["confirm"]) + [t["market"], t["breadth"], t["canary"]] - seen: list[str] = [] - for s in syms: - if s and s not in seen: - seen.append(s) - return seen + tickers = config["tickers"] + symbols = list(tickers["leaders"]) + list(tickers["confirm"]) + [tickers["market"]] + return list(dict.fromkeys(s for s in symbols if s)) async def _fetch_prices(config: dict, start: date, end: date) -> dict[str, Series]: @@ -432,17 +617,16 @@ async def _fetch_prices(config: dict, start: date, end: date) -> dict[str, Serie return {} provider = AlpacaOHLCVProvider(settings.alpaca_api_key, settings.alpaca_api_secret) out: dict[str, Series] = {} - for sym in _price_symbols(config): + for symbol in _price_symbols(config): try: - bars = await provider.fetch_ohlcv(sym, start, end) - out[sym] = sorted(((b.date, float(b.close)) for b in bars), key=lambda x: x[0]) + bars = await provider.fetch_ohlcv(symbol, start, end) + out[symbol] = sorted(((b.date, float(b.close)) for b in bars), key=lambda item: item[0]) except Exception as exc: - logger.warning("Regime monitor: price fetch failed for %s: %s", sym, exc) + logger.warning("Regime monitor: price fetch failed for %s: %s", symbol, exc) return out async def _fetch_fred_series(series_id: str, start: date, end: date) -> Series | None: - """Fetch a FRED series as [(date, value)]. None if no API key configured.""" if not settings.fred_api_key: return None verify = _CA_BUNDLE if (_CA_BUNDLE and Path(_CA_BUNDLE).exists()) else True @@ -455,68 +639,82 @@ async def _fetch_fred_series(series_id: str, start: date, end: date) -> Series | } try: async with httpx.AsyncClient(timeout=30, verify=verify) as client: - resp = await client.get( + response = await client.get( "https://api.stlouisfed.org/fred/series/observations", params=params ) - resp.raise_for_status() - payload = resp.json() + response.raise_for_status() + payload = response.json() except Exception as exc: logger.warning("Regime monitor: FRED fetch failed for %s: %s", series_id, exc) return None out: Series = [] - for obs in payload.get("observations", []): - value = obs.get("value") + for observation in payload.get("observations", []): + value = observation.get("value") if value in (None, ".", ""): continue try: - out.append((date.fromisoformat(obs["date"]), float(value))) + out.append((date.fromisoformat(observation["date"]), float(value))) except (TypeError, ValueError): continue - return sorted(out, key=lambda x: x[0]) + return sorted(out, key=lambda item: item[0]) # --------------------------------------------------------------------------- -# Snapshot persistence +# Snapshot persistence and reads # --------------------------------------------------------------------------- -async def _upsert_snapshot(db: AsyncSession, result: dict) -> None: - d = date.fromisoformat(result["date"]) - existing = await db.execute(select(RegimeSnapshot).where(RegimeSnapshot.date == d)) +async def _upsert_snapshot( + db: AsyncSession, + result: dict, + *, + rewrite_existing_v2: bool, +) -> tuple[bool, dict]: + snapshot_date = date.fromisoformat(result["date"]) + existing = await db.execute(select(RegimeSnapshot).where(RegimeSnapshot.date == snapshot_date)) row = existing.scalar_one_or_none() + state_score = (result.get("state") or {}).get("score") + state_band = (result.get("state") or {}).get("band") payload = json.dumps(result) if row is None: db.add(RegimeSnapshot( - date=d, - total_score=result["total_score"], - band=result["band"], + date=snapshot_date, + total_score=float(state_score or 0.0), + band=state_band or "unavailable", breakdown_json=payload, created_at=datetime.now(timezone.utc), )) else: - row.total_score = result["total_score"] - row.band = result["band"] + existing_v2 = _parse_v2(row.breakdown_json) + if existing_v2 is not None and not rewrite_existing_v2: + return False, existing_v2 + row.total_score = float(state_score or 0.0) + row.band = state_band or "unavailable" row.breakdown_json = payload + return True, result -async def _snapshot_count(db: AsyncSession) -> int: - res = await db.execute(select(func.count()).select_from(RegimeSnapshot)) - return int(res.scalar() or 0) +def _parse_v2(raw: str) -> dict | None: + try: + parsed = json.loads(raw) + except (TypeError, ValueError): + return None + return parsed if parsed.get("methodology") == METHODOLOGY else None -# --------------------------------------------------------------------------- -# Job entrypoint + reads -# --------------------------------------------------------------------------- +async def _latest_v2_row(db: AsyncSession) -> tuple[RegimeSnapshot, dict] | None: + result = await db.execute( + select(RegimeSnapshot).order_by(RegimeSnapshot.date.desc()).limit(1000) + ) + for row in result.scalars().all(): + parsed = _parse_v2(row.breakdown_json) + if parsed is not None: + return row, parsed + return None -async def update_regime_monitor(db: AsyncSession, backfill_days: int = 90) -> dict: - """Compute the latest index, persist it, and backfill history on first run. - Job entrypoint (daily-pipeline step). Best-effort throughout: missing keys or - a failed source degrade gracefully (signals drop to n/a) rather than abort. - """ +async def update_regime_monitor(db: AsyncSession, rebuild_sessions: int = REBUILD_SESSIONS) -> dict: config = await get_regime_config(db) - - # Refresh the LLM fundamentals if stale (and not manually locked). Best-effort. overrides = await get_fundamental_overrides(db) if _fundamentals_stale(overrides, config) and not overrides.get("locked"): try: @@ -525,204 +723,175 @@ async def update_regime_monitor(db: AsyncSession, backfill_days: int = 90) -> di logger.warning("Regime monitor: fundamentals refresh skipped: %s", exc) end = date.today() - start = end - timedelta(days=400) - prices = await _fetch_prices(config, start, end) - vix_series = await _fetch_fred_series("VIXCLS", start, end) - oas_series = await _fetch_fred_series("BAMLH0A0HYM2", end - timedelta(days=1200), end) + prices = await _fetch_prices(config, end - timedelta(days=1200), end) + leader = config["tickers"]["leaders"][0] + leader_series = prices.get(leader, []) + if not leader_series: + return {"available": False, "reason": "no benchmark price data"} + latest_date = leader_series[-1][0] - # Anchor "today" on the latest actual trading day we have prices for. - leader = config["tickers"]["leaders"][0] if config["tickers"]["leaders"] else None - leader_series = prices.get(leader or "", []) - latest_date = leader_series[-1][0] if leader_series else end + vix_series = await _fetch_fred_series("VIXCLS", end - timedelta(days=1200), end) + oas_series = await _fetch_fred_series( + "BAMLH0A0HYM2", end - timedelta(days=int(365.25 * 13)), end + ) - # Early-warning signal: breadth-divergence over the stored universe (leads but - # noisy). Computed once here so the daily job carries it live, as a SEPARATE - # score next to the coincident index — not folded into the index weights. - # Best-effort: a breadth failure must not stop the index update. + basket = config["breadth_basket"] try: - breadth = await breadth_service.compute_breadth_series(db) - divergence = breadth_service.compute_divergence_series(breadth, sorted(leader_series)) - except Exception as exc: - logger.warning("Regime monitor: breadth/divergence skipped: %s", exc) - divergence = {} - # As-of lookup: the stored universe (breadth) can lag the live benchmark date - # by a day or two, so an exact-date match would blank the newest snapshot. - div_items = sorted(divergence.items()) - cw = config.get("combined_weights") or {"coincident": 0.6, "early_warning": 0.4} - - dates = {latest_date} - if await _snapshot_count(db) < 5 and leader_series: - cutoff = end - timedelta(days=backfill_days) - dates |= {d for d, _ in leader_series if d >= cutoff} - - latest_result: dict | None = None - for d in sorted(dates): - result = _compute_index(prices, vix_series, oas_series, overrides, config, d) - _attach_early_warning(result, _divergence_asof(div_items, d), cw) - await _upsert_snapshot(db, result) - latest_result = result - - # Backfill early-warning + combined onto recent existing snapshots (e.g. the - # index history written before this signal existed) so their 7/30-day trends - # populate immediately rather than only filling in over the coming weeks. - if div_items: - recent = await db.execute( - select(RegimeSnapshot).where(RegimeSnapshot.date >= end - timedelta(days=120)) + breadth, breadth_counts = await breadth_service.compute_breadth_details( + db, basket, window=200, min_tickers=20 ) - for row in recent.scalars().all(): - try: - res = json.loads(row.breakdown_json) - except (TypeError, ValueError): - continue - if (res.get("early_warning") or {}).get("score") is not None: - continue - _attach_early_warning(res, _divergence_asof(div_items, row.date), cw) - row.breakdown_json = json.dumps(res) + divergence = breadth_service.compute_divergence_series(breadth, leader_series) + except Exception as exc: + logger.warning("Regime monitor: fixed-basket breadth skipped: %s", exc) + breadth, breadth_counts, divergence = {}, {}, {} + + latest_v2 = await _latest_v2_row(db) + rebuilding = latest_v2 is None and bool(leader_series) + if rebuilding: + dates = [d for d, _ in leader_series[-max(1, rebuild_sessions):]] + else: + # Routine PIT rule: only the latest trading date may be inserted/updated. + dates = [latest_date] + + breadth_series = _mapping_series(breadth) + divergence_series = _mapping_series(divergence) + latest_result: dict | None = None + snapshots_written = 0 + for snapshot_date in dates: + computed = _compute_index( + prices, + vix_series, + oas_series, + overrides, + config, + snapshot_date, + breadth_series, + divergence_series, + breadth_counts, + ) + written, latest_result = await _upsert_snapshot( + db, + computed, + rewrite_existing_v2=rebuilding or snapshot_date == date.today(), + ) + snapshots_written += int(written) await db.commit() logger.info(json.dumps({ "event": "regime_monitor_updated", - "date": latest_result["date"] if latest_result else None, - "score": latest_result["total_score"] if latest_result else None, - "snapshots_written": len(dates), + "methodology": METHODOLOGY, + "date": latest_result.get("date") if latest_result else None, + "state": ((latest_result or {}).get("state") or {}).get("score"), + "warning": ((latest_result or {}).get("warning") or {}).get("score"), + "snapshots_written": snapshots_written, })) return latest_result or {"available": False, "reason": "no data"} -def _divergence_asof(div_items: list[tuple[date, float]], as_of: date, max_lag_days: int = 7) -> float | None: - """Latest divergence value on/before ``as_of``, tolerating a small data lag - between the live benchmark and the stored universe. None if too stale/absent.""" - chosen: tuple[date, float] | None = None - for d, v in div_items: - if d <= as_of: - chosen = (d, v) - else: - break - if chosen is None or (as_of - chosen[0]).days > max_lag_days: - return None - return chosen[1] - - -def _attach_early_warning(result: dict, ew: float | None, weights: dict) -> None: - """Attach the separate early-warning score and a combined blend to a snapshot. - - ``ew`` is the breadth-divergence value as-of this date (or None). The combined - score is a normalised weighted mean of the coincident index and the early - warning — observational, kept apart from the index itself. - """ - result["early_warning"] = { - "score": round(ew, 1) if ew is not None else None, - "band": band_for(ew) if ew is not None else None, - } - if ew is None: - combined = result["total_score"] - else: - wc = float(weights.get("coincident", 0.6)) - we = float(weights.get("early_warning", 0.4)) - wsum = (wc + we) or 1.0 - combined = (result["total_score"] * wc + ew * we) / wsum - result["combined"] = {"score": round(combined, 1), "band": band_for(combined)} - - -async def _result_at_or_before(db: AsyncSession, target: date) -> dict | None: - """Parsed snapshot result for the latest date on/before ``target``.""" - res = await db.execute( +async def _result_at_or_before( + db: AsyncSession, + target: date, + basket_hash: str | None = None, +) -> dict | None: + result = await db.execute( select(RegimeSnapshot.breakdown_json) .where(RegimeSnapshot.date <= target) .order_by(RegimeSnapshot.date.desc()) - .limit(1) + .limit(1000) ) - raw = res.scalar_one_or_none() - if raw is None: - return None - try: - return json.loads(raw) - except (TypeError, ValueError): - return None + for raw in result.scalars().all(): + parsed = _parse_v2(raw) + parsed_hash = ((parsed or {}).get("basket") or {}).get("hash") + if parsed is not None and (basket_hash is None or parsed_hash == basket_hash): + return parsed + return None -def _delta(curr: float | None, prev: float | None) -> float | None: - return round(curr - prev, 1) if (curr is not None and prev is not None) else None +def _delta(current: dict, previous: dict | None) -> float | None: + if not previous: + return None + if current.get("available_pillars") != previous.get("available_pillars"): + return None + a, b = current.get("score"), previous.get("score") + return round(a - b, 1) if a is not None and b is not None else None async def get_regime_monitor(db: AsyncSession) -> dict: - """Latest snapshot + 7/30-day trend deltas for the index, early-warning, and - combined scores. Cheap (a few row reads).""" - res = await db.execute( - select(RegimeSnapshot).order_by(RegimeSnapshot.date.desc()).limit(1) - ) - latest = res.scalar_one_or_none() + latest = await _latest_v2_row(db) if latest is None: - return {"available": False, "reason": "not computed yet"} + return {"available": False, "reason": "v2 not computed yet"} + row, result = latest + basket_hash = (result.get("basket") or {}).get("hash") + previous_7 = await _result_at_or_before( + db, row.date - timedelta(days=7), basket_hash + ) + previous_30 = await _result_at_or_before( + db, row.date - timedelta(days=30), basket_hash + ) - try: - result = json.loads(latest.breakdown_json) - except (TypeError, ValueError): - result = {"date": latest.date.isoformat(), "total_score": latest.total_score, - "band": latest.band, "breakdown": []} - - r7 = await _result_at_or_before(db, latest.date - timedelta(days=7)) - r30 = await _result_at_or_before(db, latest.date - timedelta(days=30)) - - def _nested(r: dict | None, key: str) -> float | None: - return (r.get(key) or {}).get("score") if r else None - - result["available"] = True - cur_total = result.get("total_score") - result["trend"] = { - "delta_7": _delta(cur_total, (r7 or {}).get("total_score")), - "delta_30": _delta(cur_total, (r30 or {}).get("total_score")), - } - for key in ("early_warning", "combined"): - block = result.get(key) or {"score": None, "band": None} - block["delta_7"] = _delta(block.get("score"), _nested(r7, key)) - block["delta_30"] = _delta(block.get("score"), _nested(r30, key)) + for key in ("state", "warning"): + block = result.get(key) or {} + block["trend"] = { + "delta_7": _delta(block, (previous_7 or {}).get(key)), + "delta_30": _delta(block, (previous_30 or {}).get(key)), + } result[key] = block + + snapshot_age = (date.today() - row.date).days + quality = result.get("data_quality") or {} + quality["snapshot_age_days"] = snapshot_age + quality["is_fresh"] = bool(quality.get("inputs_fresh")) and snapshot_age <= 4 + result["data_quality"] = quality + result["available"] = True return result -async def get_regime_history(db: AsyncSession, days: int = 400) -> list[dict]: - """Daily history of the index, early-warning, and combined scores for the - score-over-time chart. One point per snapshot date, ascending.""" +async def get_regime_history(db: AsyncSession, days: int = 800) -> list[dict]: cutoff = date.today() - timedelta(days=days) - res = await db.execute( + result = await db.execute( select(RegimeSnapshot) .where(RegimeSnapshot.date >= cutoff) .order_by(RegimeSnapshot.date.asc()) ) out: list[dict] = [] - for row in res.scalars().all(): - try: - data = json.loads(row.breakdown_json) - except (TypeError, ValueError): - data = {} + for row in result.scalars().all(): + data = _parse_v2(row.breakdown_json) + if data is None: + continue + state, warning = data.get("state") or {}, data.get("warning") or {} out.append({ "date": row.date.isoformat(), - "index": row.total_score, - "early_warning": (data.get("early_warning") or {}).get("score"), - "combined": (data.get("combined") or {}).get("score"), + "state": state.get("score") if state.get("band") is not None else None, + "warning": warning.get("score") if warning.get("band") is not None else None, + "state_coverage": state.get("coverage"), + "warning_coverage": warning.get("coverage"), + "basket_hash": (data.get("basket") or {}).get("hash"), }) - return out + if not out: + return out + latest_hash = out[-1]["basket_hash"] + if latest_hash is None: + return out + return [point for point in out if point["basket_hash"] == latest_hash] # --------------------------------------------------------------------------- -# F1/F3 via grounded LLM (reuses the configured sentiment provider) +# Grounded fundamental extraction # --------------------------------------------------------------------------- _CAPEX_PROMPT = """\ You are a markets analyst. Search the web for the MOST RECENT (last reported \ quarter) capital-expenditure (capex) guidance from these hyperscalers: {names}. -For each name, classify the direction of its forward capex/AI-infrastructure \ -guidance vs. the prior quarter as exactly one of: "raising", "holding", "cutting". +For each name, classify forward capex/AI-infrastructure guidance vs. the prior \ +quarter as exactly one of: "raising", "holding", "cutting", "unknown". -Also judge the recent "good news, stock down" dynamic: across these names and \ -the semiconductor sector, did stocks broadly FALL despite earnings/revenue beats \ -in the last reporting season? Answer "yes", "no", or "mixed". +Also judge the recent good-news-stock-down dynamic across these names and the \ +semiconductor sector after earnings/revenue beats. Answer "yes", "no", or "mixed". -Respond ONLY with a JSON object (no markdown): +Respond ONLY with JSON (no markdown): {{"capex": {{ {example} }}, "good_news_stock_down": "yes|no|mixed", \ -"reasoning": "<2-3 sentences citing the specific guidance you found>"}} +"reasoning": "<2-3 sourced sentences>"}} """ @@ -731,13 +900,14 @@ def _fundamentals_stale(overrides: dict, config: dict) -> bool: if not fetched: return True try: - ts = datetime.fromisoformat(fetched) + timestamp = datetime.fromisoformat(fetched) except (TypeError, ValueError): return True - if ts.tzinfo is None: - ts = ts.replace(tzinfo=timezone.utc) - max_age = timedelta(days=int(config.get("fundamental_staleness_days", 80))) - return datetime.now(timezone.utc) - ts > max_age + if timestamp.tzinfo is None: + timestamp = timestamp.replace(tzinfo=timezone.utc) + return datetime.now(timezone.utc) - timestamp > timedelta( + days=int(config.get("fundamental_staleness_days", 80)) + ) def _strip_fences(text: str) -> str: @@ -759,15 +929,14 @@ def _extract_responses_text(response: object) -> str: async def _call_llm_json(cfg: dict, prompt: str) -> dict: - """Send one grounded prompt via the configured LLM and parse its JSON reply.""" provider, model, api_key = cfg["provider"], cfg["model"], cfg["api_key"] base_url = cfg.get("base_url") - if provider == "gemini": from google import genai from google.genai import types + client = genai.Client(api_key=api_key) - resp = await client.aio.models.generate_content( + response = await client.aio.models.generate_content( model=model, contents=prompt, config=types.GenerateContentConfig( @@ -775,9 +944,10 @@ async def _call_llm_json(cfg: dict, prompt: str) -> dict: response_mime_type="application/json", ), ) - return json.loads(_strip_fences(resp.text)) + return json.loads(_strip_fences(response.text)) from openai import AsyncOpenAI + verify = _CA_BUNDLE if (_CA_BUNDLE and Path(_CA_BUNDLE).exists()) else True client = AsyncOpenAI( api_key=api_key, @@ -786,68 +956,69 @@ async def _call_llm_json(cfg: dict, prompt: str) -> dict: ) if provider in ("openai", "xai"): tool = "web_search_preview" if provider == "openai" else "web_search" - resp = await client.responses.create( + response = await client.responses.create( model=model, tools=[{"type": tool}], instructions="Respond with valid JSON only, no markdown fences.", input=prompt, ) - return json.loads(_strip_fences(_extract_responses_text(resp))) + return json.loads(_strip_fences(_extract_responses_text(response))) - # deepseek / generic OpenAI-compatible: no web search, knowledge-based. - resp = await client.chat.completions.create( + response = await client.chat.completions.create( model=model, messages=[{"role": "user", "content": prompt}], response_format={"type": "json_object"}, ) - return json.loads(_strip_fences(resp.choices[0].message.content)) + return json.loads(_strip_fences(response.choices[0].message.content)) -_CAPEX_STATE_SCORES = {"raising": 0.0, "holding": 50.0, "cutting": 100.0} -_GNSD_SCORES = {"yes": 100.0, "mixed": 50.0, "no": 0.0} +_CAPEX_STATE_SCORES = {"raising": 0.0, "holding": 0.0, "cutting": 100.0} +_GNSD_SCORES = {"yes": 100.0, "no": 0.0} async def refresh_fundamental_overrides( db: AsyncSession, config: dict | None = None, force: bool = False ) -> dict: - """Ask the configured LLM to propose F1 (capex) and F3 (earnings reaction). - - Skips (returns current) if a manual override is locked, unless ``force``. - """ current = await get_fundamental_overrides(db) if current.get("locked") and not force: return current config = config or await get_regime_config(db) - cfg = await resolve_llm_config(db) - if not cfg.get("api_key"): - raise ProviderError(f"No API key configured for LLM provider '{cfg.get('provider')}'") + llm = await resolve_llm_config(db) + if not llm.get("api_key"): + raise ProviderError(f"No API key configured for LLM provider '{llm.get('provider')}'") names = config["tickers"]["hyperscalers"] - example = ", ".join(f'"{n}": "holding"' for n in names) - prompt = _CAPEX_PROMPT.format(names=", ".join(names), example=example) - parsed = await _call_llm_json(cfg, prompt) - + example = ", ".join(f'"{name}": "holding"' for name in names) + parsed = await _call_llm_json( + llm, _CAPEX_PROMPT.format(names=", ".join(names), example=example) + ) capex = parsed.get("capex", {}) if isinstance(parsed, dict) else {} scores = [ - _CAPEX_STATE_SCORES[str(capex.get(n, "")).strip().lower()] - for n in names - if str(capex.get(n, "")).strip().lower() in _CAPEX_STATE_SCORES + _CAPEX_STATE_SCORES[value] + for name in names + if (value := str(capex.get(name, "")).strip().lower()) in _CAPEX_STATE_SCORES ] - f1 = _mean(scores) if scores else 50.0 - gnsd = str(parsed.get("good_news_stock_down", "")).strip().lower() - f3 = _GNSD_SCORES.get(gnsd, 50.0) - + f1 = _mean(scores) if len(scores) >= 3 else None + reaction = str(parsed.get("good_news_stock_down", "")).strip().lower() + f3 = _GNSD_SCORES.get(reaction) + now = datetime.now(timezone.utc) result = { - "f1_score": round(f1, 1), + "f1_score": round(f1, 1) if f1 is not None else None, "f3_score": f3, "capex": capex, - "good_news_stock_down": gnsd or None, + "good_news_stock_down": reaction or None, "reasoning": parsed.get("reasoning") if isinstance(parsed, dict) else None, - "fetched_at": datetime.now(timezone.utc).isoformat(), + "fetched_at": now.isoformat(), + "effective_date": _next_weekday(now.date()).isoformat(), "locked": False, - "source": cfg.get("provider"), + "source": llm.get("provider"), } await update_setting(db, KEY_FUNDAMENTALS, json.dumps(result)) - logger.info(json.dumps({"event": "regime_fundamentals_refreshed", "f1": result["f1_score"], "f3": result["f3_score"]})) + logger.info(json.dumps({ + "event": "regime_fundamentals_refreshed", + "f1": result["f1_score"], + "f3": result["f3_score"], + "effective_date": result["effective_date"], + })) return result diff --git a/docs/research/regime-monitor-v2.md b/docs/research/regime-monitor-v2.md new file mode 100644 index 0000000..85832c1 --- /dev/null +++ b/docs/research/regime-monitor-v2.md @@ -0,0 +1,66 @@ +# Regime Monitor v2 methodology + +The Regime Monitor is an observational AI/Tech risk thermometer. It does not +gate entries, exits, position size, ranking, or alerts about individual setups. + +## Outputs + +**State** measures current structural stress: + +- Price structure, 40%: `max(P1, P2, P3)`, so the correlated 200-DMA, death-cross, + and drawdown readings receive one capped vote. +- Fixed-basket breadth level, 25%. +- HY option-adjusted credit spread, 20%. +- VIX level, 15%. + +**Warning** measures deterioration and divergence: + +- Fixed-basket breadth divergence while SMH holds/rises, 50%. +- 60-session SMH/SPY relative-strength deterioration, 30%. +- Hyperscaler capex cuts, 12%. +- Good-news-stock-down earnings reactions, 8%. + +Combined, RSP/SPY (former F4), and the NVDA canary (former P6) do not enter v2. + +## Scale and missing data + +Zero means ordinary/healthy, and only stress contributes positively. Automated +capex `raising`/`holding` and no good-news-stock-down pattern map to zero; +`mixed`, unknown, and stale observations are unavailable rather than neutral 50. + +Scores renormalize over available fixed weights, but a band is published only at +75% or greater coverage. Trend deltas are suppressed when the participating +pillar set changes. Bands are stable `<30`, watch `<60`, elevated `<80`, and +breaking `>=80`. + +Credit uses named HY OAS anchors (3.5 mild, 5.0 elevated, 7.0 stressed) for 70% +of its score and a ten-year upper-tail percentile for 30%. + +## Point-in-time record + +The first v2 run rebuilds the latest 400 trading sessions with sufficient sensor +warm-up. Routine runs thereafter insert/update only the latest trading date. +Fundamental observations have an effective date (normally the next session after +collection) and are never replayed backward. The history API and main chart show +only snapshots marked `methodology: v2`. + +Each snapshot stores the fixed basket symbols, hash, and freeze date. Reconstructed +history before that freeze date is retrospective/exploratory; readings after it +form the forward record. + +## Warning study + +The study calls the outcome a **10% correction**, not a regime break. The first +70% of sessions freezes the 80th-percentile warning threshold; alarm episodes are +measured on the final 30%. An alarm requires an upward crossing and another alarm +requires a reset below the threshold. The report exposes warned/missed events, +false alarms per year, median lead, sample dates, event count, report date, and +whether the result is exploratory or a true forward holdout. UI claims are +generated from that report; no performance sentence is hard-coded. + +## Operator rule + +Quadrant alerts default off for new/reset configurations. When enabled they +require fresh inputs, at least 75% coverage on both axes, two consecutive daily +confirmations, hysteresis, and cooldown. Every alert states: **Risk thermometer — +not a trade signal.** diff --git a/frontend/src/api/regime.ts b/frontend/src/api/regime.ts index 6a5fba8..d90ff09 100644 --- a/frontend/src/api/regime.ts +++ b/frontend/src/api/regime.ts @@ -11,7 +11,7 @@ export function getRegimeMonitor() { return apiClient.get('regime/monitor').then((r) => r.data); } -export function getRegimeHistory(days = 400) { +export function getRegimeHistory(days = 800) { return apiClient .get('regime/history', { params: { days } }) .then((r) => r.data); diff --git a/frontend/src/components/regime/RegimeQuadrant.tsx b/frontend/src/components/regime/RegimeQuadrant.tsx index ae65e71..3fa9724 100644 --- a/frontend/src/components/regime/RegimeQuadrant.tsx +++ b/frontend/src/components/regime/RegimeQuadrant.tsx @@ -13,15 +13,13 @@ import { ReferenceLine, ReferenceArea, } from 'recharts'; -import { getRegimeHistory } from '../../api/regime'; +import { getRegimeHistory, getRegimeMonitor } from '../../api/regime'; import { Callout } from '../ui/Callout'; import { SkeletonCard } from '../ui/Skeleton'; // Lazy-loaded (see RegimePage) so recharts stays in the regime-tab chunk. -// Quadrant dividers. Regime < 40 ≈ intact; early-warning > 60 ≈ elevated. -const X_DIV = 40; // regime index -const Y_DIV = 60; // early warning +// Quadrant boundaries come from the backend v2 methodology response. const TRAIL = 60; // sessions shown interface QPoint { @@ -64,7 +62,7 @@ function QuadrantTip({ active, payload }: { active?: boolean; payload?: { payloa
{p.date}
- Regime {Math.round(p.x)} · Early warning{' '} + State {Math.round(p.x)} · Warning{' '} {Math.round(p.y)}
@@ -72,14 +70,17 @@ function QuadrantTip({ active, payload }: { active?: boolean; payload?: { payloa } export default function RegimeQuadrant() { - const history = useQuery({ queryKey: ['regime', 'history'], queryFn: () => getRegimeHistory(400) }); + const history = useQuery({ queryKey: ['regime', 'history'], queryFn: () => getRegimeHistory(800) }); + const monitor = useQuery({ queryKey: ['regime', 'monitor'], queryFn: getRegimeMonitor }); + const xDiv = monitor.data?.quadrant_config?.state_divider ?? 60; + const yDiv = monitor.data?.quadrant_config?.warning_divider ?? 60; const points = useMemo(() => { const data = history.data ?? []; return data - .filter((p) => p.early_warning != null) + .filter((p) => p.state != null && p.warning != null) .slice(-TRAIL) - .map((p) => ({ x: p.index, y: p.early_warning as number, date: p.date })); + .map((p) => ({ x: p.state as number, y: p.warning as number, date: p.date })); }, [history.data]); const trail = useMemo(() => smoothTrail(points), [points]); @@ -89,11 +90,11 @@ export default function RegimeQuadrant() {
- Regime quadrant — last {TRAIL} sessions + State × Warning quadrant — last {TRAIL} sessions
{latest && (
- now: regime {Math.round(latest.x)} · warning{' '} + now: State {Math.round(latest.x)} · Warning{' '} {Math.round(latest.y)}
)} @@ -103,7 +104,7 @@ export default function RegimeQuadrant() { ) : !points.length ? ( - Not enough history yet — the early-warning fills in as the daily job runs. + Not enough coverage-qualified v2 history yet. ) : ( <> @@ -111,13 +112,13 @@ export default function RegimeQuadrant() { {/* Quadrant shading (drawn first, behind everything) */} - - - - + + + + - - + + } /> @@ -166,15 +167,15 @@ export default function RegimeQuadrant() {
- ① Hot & brittle — narrow melt-up, shakeout risk - ② Transition — break may be starting - ③ Healthy & broad — calm uptrend - ④ Real downturn — regime breaking, broad + Early warning — state calm, fragility rising + Active stress — damaged and deteriorating + Healthy — calm and broadly supported + Stressed / stabilizing — damage remains, warning lower

White dot = today; the trail fades from muted (older) to bright blue (newer) over the last {TRAIL}{' '} - sessions, smoothed. The tell isn't a single spot but the move ①→④ (early warning rolling over while - the regime index climbs = divergence resolving downward). Observational — not wired into trades. + sessions, smoothed. The path matters more than a single point. Risk thermometer — not an entry, exit, + or sizing signal.

)} diff --git a/frontend/src/components/regime/ScoreHistoryChart.tsx b/frontend/src/components/regime/ScoreHistoryChart.tsx index 13a7530..8c3df78 100644 --- a/frontend/src/components/regime/ScoreHistoryChart.tsx +++ b/frontend/src/components/regime/ScoreHistoryChart.tsx @@ -26,14 +26,13 @@ const HISTORY_RANGES = [ type HistoryRange = (typeof HISTORY_RANGES)[number]['key']; const HISTORY_SERIES = [ - { key: 'index', label: 'Index', color: '#60a5fa' }, - { key: 'early_warning', label: 'Early warning', color: '#fb923c' }, - { key: 'combined', label: 'Combined', color: '#a78bfa' }, + { key: 'state', label: 'State', color: '#60a5fa' }, + { key: 'warning', label: 'Warning', color: '#fb923c' }, ] as const; export default function ScoreHistoryChart() { const [range, setRange] = useState('3M'); - const history = useQuery({ queryKey: ['regime', 'history'], queryFn: () => getRegimeHistory(400) }); + const history = useQuery({ queryKey: ['regime', 'history'], queryFn: () => getRegimeHistory(800) }); const filtered = useMemo(() => { const data = history.data ?? []; @@ -113,7 +112,6 @@ export default function ScoreHistoryChart() { stroke={s.color} dot={false} strokeWidth={1.5} - connectNulls isAnimationActive={false} /> ))} diff --git a/frontend/src/lib/types.ts b/frontend/src/lib/types.ts index e8d017f..1886164 100644 --- a/frontend/src/lib/types.ts +++ b/frontend/src/lib/types.ts @@ -439,106 +439,134 @@ export type RegimeBand = 'stable' | 'watch' | 'elevated' | 'breaking'; export interface RegimeSignal { id: string; label: string; - sub_score: number | null; - weight: number; + score: number | null; available: boolean; - contribution: number; + details?: Record; } -export interface RegimeSubScore { +export interface RegimePillar { + id: string; + label: string; + score: number | null; + weight: number; + contribution: number; + available: boolean; + sensors: RegimeSignal[]; +} + +export interface RegimeReading { score: number | null; band: RegimeBand | null; - delta_7?: number | null; - delta_30?: number | null; + coverage: number; + minimum_coverage: number; + available_pillars: string[]; + pillars: RegimePillar[]; + trend?: { delta_7: number | null; delta_30: number | null }; } export interface RegimeHistoryPoint { date: string; - index: number; - early_warning: number | null; - combined: number | null; + state: number | null; + warning: number | null; + state_coverage: number | null; + warning_coverage: number | null; + basket_hash: string | null; } export interface RegimeMonitor { available: boolean; reason?: string; + methodology?: string; date?: string; - total_score?: number; - band?: RegimeBand; - alert_threshold?: number; - breakdown?: RegimeSignal[]; + state?: RegimeReading; + warning?: RegimeReading; inputs?: { vix: number | null; + vix_date: string | null; hy_oas: number | null; + hy_oas_date: string | null; + breadth_pct_above_200: number | null; + breadth_date: string | null; fundamentals_fetched_at: string | null; + fundamentals_effective_date: string | null; + fundamentals_age_days: number | null; }; - trend?: { delta_7: number | null; delta_30: number | null }; - // Separate, observational early-warning score (breadth divergence) + a small - // combined blend. Decoupled from the index above. - early_warning?: RegimeSubScore; - combined?: RegimeSubScore; + basket?: { + symbols: string[]; + hash: string; + basket_asof: string; + members_available: number | null; + members_expected: number; + history_kind: 'forward' | 'retrospective'; + }; + data_quality?: { + minimum_coverage: number; + oldest_market_input_age_days: number | null; + stale_inputs: string[]; + inputs_fresh: boolean; + snapshot_age_days?: number; + is_fresh?: boolean; + }; + quadrant_config?: { state_divider: number; warning_divider: number; margin: number }; } export interface RegimeFundamentals { - f1_score: number; - f3_score: number; + f1_score: number | null; + f3_score: number | null; locked: boolean; reasoning: string | null; fetched_at: string | null; + effective_date: string | null; source: string; capex?: Record; good_news_stock_down?: string | null; } export interface RegimeConfig { - weights: Record; - alert_threshold: number; - tickers: Record; - leader_weight: number; - rs_lookback: number; + breadth_basket: string[]; + basket_asof: string; fundamental_staleness_days: number; } // Event study — measured lead time of early-warning indicators vs. drawdowns -export interface EventStudyLeadStats { - median_lead_days: number | null; - events_with_signal: number; - events_total: number; - warn_threshold: number; - mean_path: { rel_day: number; value: number }[]; - signal: { - base_rate: number; - horizon_days: number; - rows: { threshold: number; precision: number | null; recall: number | null; alarms: number }[]; - }; -} - -export interface EventStudyPerEvent { - date: string; - depth_pct: number; - breadth_lead: number | null; - coincident_lead: number | null; -} - export interface EventStudyReport { available: boolean; reason?: string; + methodology?: string; generated_at?: string; + evaluation?: 'exploratory' | 'holdout'; + summary?: string; params?: { benchmark: string; + outcome: string; event_threshold_pct: number; - cooldown_days: number; + event_cooldown_days: number; horizon_days: number; + train_fraction: number; warn_percentile: number; + warn_threshold: number; + basket_hash: string; + basket_asof: string; }; - events?: { date: string; index: number; depth_pct: number }[]; - indicators?: { - breadth_divergence: EventStudyLeadStats; - coincident_price: EventStudyLeadStats; + sample?: { + start: string; + end: string; + train_end: string; + test_start: string; + sessions: number; + holdout_sessions: number; }; - per_event?: EventStudyPerEvent[]; - lead_delta_days?: number | null; - recent_breadth?: { date: string; breadth: number; divergence: number | null }[]; + metrics?: { + events: number; + events_warned: number; + events_missed: number; + alarm_episodes: number; + false_alarms: number; + false_alarms_per_year: number; + median_lead_days: number | null; + }; + events?: { date: string; warned: boolean; lead_days: number | null }[]; + recent_breadth?: { date: string; breadth: number; warning: number | null }[]; } export interface AlertConfig { diff --git a/frontend/src/pages/RegimePage.tsx b/frontend/src/pages/RegimePage.tsx index 1618515..053d764 100644 --- a/frontend/src/pages/RegimePage.tsx +++ b/frontend/src/pages/RegimePage.tsx @@ -1,5 +1,5 @@ -import { useState, lazy, Suspense, type ReactNode } from 'react'; -import { useQuery, useMutation, useQueryClient } from '@tanstack/react-query'; +import { lazy, Suspense, useState, type ReactNode } from 'react'; +import { useMutation, useQuery, useQueryClient } from '@tanstack/react-query'; import { PageHeader } from '../components/ui/PageHeader'; import { Callout } from '../components/ui/Callout'; import { Disclosure } from '../components/ui/Disclosure'; @@ -7,44 +7,38 @@ import { Badge } from '../components/ui/Badge'; import { SkeletonCard, SkeletonTable } from '../components/ui/Skeleton'; import { useAuthStore } from '../stores/authStore'; import { - getRegimeMonitor, - getRegimeConfig, - updateRegimeConfig, - getRegimeFundamentals, - updateRegimeFundamentals, - refreshRegimeFundamentals, getEventStudy, + getRegimeConfig, + getRegimeFundamentals, + getRegimeMonitor, + refreshRegimeFundamentals, + updateRegimeConfig, + updateRegimeFundamentals, } from '../api/regime'; - -// Lazy so recharts (heavy) ships in its own chunk, loaded only on this tab. -const ScoreHistoryChart = lazy(() => import('../components/regime/ScoreHistoryChart')); -const RegimeQuadrant = lazy(() => import('../components/regime/RegimeQuadrant')); import type { + EventStudyReport, RegimeBand, - RegimeSignal, RegimeConfig, RegimeFundamentals, - EventStudyReport, - EventStudyLeadStats, - EventStudyPerEvent, + RegimeReading, } from '../lib/types'; +const ScoreHistoryChart = lazy(() => import('../components/regime/ScoreHistoryChart')); +const RegimeQuadrant = lazy(() => import('../components/regime/RegimeQuadrant')); + const BAND_STYLES: Record = { stable: { text: 'text-emerald-400', bar: 'bg-emerald-400', ring: 'border-emerald-400/30', label: 'Stable' }, watch: { text: 'text-amber-400', bar: 'bg-amber-400', ring: 'border-amber-400/30', label: 'Watch' }, elevated: { text: 'text-orange-400', bar: 'bg-orange-400', ring: 'border-orange-400/30', label: 'Elevated' }, - breaking: { text: 'text-red-400', bar: 'bg-red-400', ring: 'border-red-400/30', label: 'Breaking' }, + breaking: { text: 'text-red-400', bar: 'bg-red-400', ring: 'border-red-400/30', label: 'High stress' }, }; function TrendChip({ label, delta }: { label: string; delta: number | null | undefined }) { if (delta == null) { return {label}: n/a; } - const rising = delta > 0; - const flat = delta === 0; - // Higher index = worse, so a rising score is the warning direction. - const color = flat ? 'text-gray-400' : rising ? 'text-red-400' : 'text-emerald-400'; - const arrow = flat ? '→' : rising ? '↑' : '↓'; + const color = delta === 0 ? 'text-gray-400' : delta > 0 ? 'text-red-400' : 'text-emerald-400'; + const arrow = delta === 0 ? '→' : delta > 0 ? '↑' : '↓'; return ( {label}: {arrow} {delta > 0 ? '+' : ''}{delta} @@ -54,61 +48,51 @@ function TrendChip({ label, delta }: { label: string; delta: number | null | und function ScoreGauge({ label, - score, - band, - trend, - threshold, + reading, + divider, footnote, - size = 'lg', }: { label: string; - score: number | null | undefined; - band: RegimeBand | null | undefined; - trend?: { delta_7?: number | null; delta_30?: number | null }; - threshold?: number; - footnote?: ReactNode; - size?: 'lg' | 'md'; + reading: RegimeReading | undefined; + divider?: number; + footnote: ReactNode; }) { - const naa = score == null; - const style = BAND_STYLES[(band ?? 'stable') as RegimeBand]; - const s = score ?? 0; - const clamp = (v: number) => Math.min(100, Math.max(0, v)); - const numCls = size === 'lg' ? 'text-6xl' : 'text-4xl'; + const score = reading?.score; + const complete = reading?.band != null; + const style = complete ? BAND_STYLES[reading.band as RegimeBand] : null; + const position = Math.min(100, Math.max(0, score ?? 0)); return ( -
+
{label}
- - {naa ? '—' : Math.round(s)} + + {score == null ? '—' : Math.round(score)} - {!naa && / 100} + {score != null && / 100} +
+
+ + {style?.label ?? 'Incomplete'} + + coverage {Math.round(reading?.coverage ?? 0)}%
- {!naa &&

{style.label}

}
- {trend && ( -
- - -
- )} +
+ + +
- - {!naa && ( + {score != null && ( <> - {/* Band track with score (+ optional threshold) markers */} -
- {threshold != null && ( -
+
+ {divider != null && ( +
)}
@@ -116,68 +100,110 @@ function ScoreGauge({
)} - {footnote &&

{footnote}

} +

{footnote}

); } -function Breakdown({ breakdown }: { breakdown: RegimeSignal[] }) { +function PillarBreakdown({ title, reading }: { title: string; reading: RegimeReading }) { return ( -
- - - - - - - - - - - {breakdown.map((s) => ( - - - - - + +
+
SignalSub-scoreWeightContribution
- {s.id}{' '} - {s.label} - - {s.available && s.sub_score != null ? ( -
-
-
-
- {s.sub_score} -
- ) : ( - n/a - )} -
{s.weight} - {s.available ? s.contribution.toFixed(1) : '—'} -
+ + + + + + + + + {reading.pillars.map((pillar) => ( + + + + + + + ))} + +
Pillar / sensorScoreWeightContribution
+
{pillar.label}
+
+ {pillar.sensors.map((sensor) => ( +
+ {sensor.id} {sensor.label}:{' '} + {sensor.score == null ? 'n/a' : sensor.score} +
+ ))} +
+
{pillar.score ?? '—'}{pillar.weight}{pillar.available ? pillar.contribution.toFixed(1) : '—'}
+
+ + ); +} + +function EventStudyBody({ report }: { report: EventStudyReport }) { + const metrics = report.metrics; + return ( +
+
+ + {report.generated_at && generated {new Date(report.generated_at).toLocaleDateString()}} + {report.sample && test {report.sample.test_start} → {report.sample.end}} +
+

{report.summary}

+ {metrics && ( +
+ {[ + ['Warned', `${metrics.events_warned}/${metrics.events}`], + ['Missed', metrics.events_missed], + ['False alarms/year', metrics.false_alarms_per_year.toFixed(1)], + ['Median lead', metrics.median_lead_days == null ? '—' : `${metrics.median_lead_days}d`], + ].map(([label, value]) => ( +
+
{label}
+
{value}
+
))} - - +
+ )} + {report.events && report.events.length > 0 && ( +
+ + + + + + + {report.events.map((event) => ( + + + + + + ))} +
CorrectionWarnedLead
{event.date}{event.warned ? 'yes' : 'no'}{event.lead_days == null ? '—' : `${event.lead_days}d`}
+
+ )} +

+ The threshold is frozen on the training period and measured on the chronological test period. Reconstructed + pre-freeze basket history remains exploratory. +

); } -function SliderRow({ label, value, onChange }: { label: string; value: number; onChange: (v: number) => void }) { +function EventStudyPanel() { + const study = useQuery({ queryKey: ['regime', 'event-study'], queryFn: getEventStudy }); return ( - + + {study.isLoading && } + {study.data === null && Not run yet — trigger “Event Study” in Admin → Jobs.} + {study.data && !study.data.available && {study.data.reason ?? 'No data'}} + {study.data?.available && } + ); } @@ -194,379 +220,130 @@ function FundamentalsEditor({ saving: boolean; refreshing: boolean; }) { - const [f1, setF1] = useState(Math.round(data.f1_score)); - const [f3, setF3] = useState(Math.round(data.f3_score)); + const [f1, setF1] = useState(data.f1_score ?? 0); + const [f3, setF3] = useState(data.f3_score ?? 0); return (
Source: {data.source} - {data.fetched_at && · {new Date(data.fetched_at).toLocaleDateString()}} + {data.fetched_at && · fetched {new Date(data.fetched_at).toLocaleDateString()}} + {data.effective_date && · effective {data.effective_date}} {data.locked && }
{data.reasoning &&

{data.reasoning}

} - - -
- - - {data.locked && ( - - )} + {[ + ['F1 · Capex cuts', f1, setF1], + ['F3 · Good news, stock down', f3, setF3], + ].map(([label, value, setter]) => ( + + ))} +
+ + + {data.locked && }
); } -function WeightsEditor({ - data, - onSave, - saving, -}: { - data: RegimeConfig; - onSave: (updates: Partial) => void; - saving: boolean; -}) { - const [weights, setWeights] = useState>(() => ({ ...data.weights })); - const [threshold, setThreshold] = useState(data.alert_threshold); - - const setWeight = (key: string, value: string) => { - const num = parseFloat(value); - setWeights((prev) => ({ ...prev, [key]: isNaN(num) ? 0 : num })); - }; - +function ConfigEditor({ data, onSave, saving }: { data: RegimeConfig; onSave: (updates: Partial) => void; saving: boolean }) { + const [basket, setBasket] = useState(data.breadth_basket.join(', ')); + const [staleness, setStaleness] = useState(data.fundamental_staleness_days); + const symbols = basket.split(/[\s,]+/).map((symbol) => symbol.trim().toUpperCase()).filter(Boolean); return (
-
- {Object.keys(weights).map((key) => ( - - ))} -
-