""" monitoring/web_api.py — Lightweight Flask API for the HTML dashboard. Serves JSON endpoints for portfolio, predictions, signals, risk state, and chart data. """ from __future__ import annotations import csv import datetime import json import logging import os import threading from pathlib import Path from zoneinfo import ZoneInfo _ET = ZoneInfo("America/New_York") from flask import Flask, jsonify, send_from_directory, request from flask_cors import CORS import config from contracts import PriceTarget, FinalScore from data import storage from execution import broker from execution import risk as risk_module from execution.portfolio import Portfolio from execution.pdt_tracker import count_day_trades from signals.sentiment import is_warming_up, get_last_refresh, get_sentiment logger = logging.getLogger("trading_system.web_api") app = Flask(__name__, static_folder=None) CORS(app) # ── Shared state (set by main.py) ─────────────────────────────────────────── _portfolio: Portfolio | None = None _predictions: dict[str, PriceTarget] = {} _signals: dict[str, FinalScore] = {} _trade_log: list[dict] = [] _trade_log_lock = threading.Lock() MAX_TRADE_LOG = 500 def init(portfolio: Portfolio): """Set the portfolio reference from main.py.""" global _portfolio _portfolio = portfolio def update_prediction(symbol: str, target: PriceTarget): _predictions[symbol] = target def update_signal(symbol: str, signal: FinalScore): _signals[symbol] = signal def log_trade_event(event: dict): """Append a trade event to the rolling log.""" with _trade_log_lock: _trade_log.append({ "time": datetime.datetime.now(datetime.timezone.utc).isoformat(), **event, }) if len(_trade_log) > MAX_TRADE_LOG: del _trade_log[:len(_trade_log) - MAX_TRADE_LOG] # ── Routes ─────────────────────────────────────────────────────────────────── STATIC_DIR = Path(__file__).resolve().parent.parent @app.route("/") def index(): return send_from_directory(str(STATIC_DIR), "dashboard.html") @app.route("/api/config") def api_config(): return jsonify({ "universe": config.UNIVERSE, "trading_mode": config.TRADING_MODE, "dry_run": config.DRY_RUN, "target_daily_profit": config.TARGET_DAILY_PROFIT_USD, "max_daily_loss": config.MAX_DAILY_LOSS_USD, "max_positions": config.MAX_OPEN_POSITIONS, "risk_per_trade_pct": config.RISK_PER_TRADE_PCT, "technical_weight": config.TECHNICAL_WEIGHT, "sentiment_weight": config.SENTIMENT_WEIGHT, "buy_threshold": config.SIGNAL_BUY_THRESHOLD, "sell_threshold": config.SIGNAL_SELL_THRESHOLD, "allow_short": config.ALLOW_SHORT, "allow_overnight": config.ALLOW_OVERNIGHT_POSITIONS, "gpu": config.GPU_AVAILABLE, "above_25k": config.ACCOUNT_BALANCE_ABOVE_25K, }) @app.route("/api/portfolio") def api_portfolio(): if not _portfolio: return jsonify({"positions": {}, "realized": 0, "unrealized": 0, "total": 0, "count": 0}) positions = {} for sym, pos in _portfolio.positions.items(): positions[sym] = { "symbol": sym, "side": pos.get("side", ""), "qty": pos.get("qty", 0), "entry_price": pos.get("entry_price", 0), "current_price": pos.get("current_price", 0), "unrealized_pl": pos.get("unrealized_pl", 0), "market_value": pos.get("market_value", 0), } return jsonify({ "positions": positions, "realized": _portfolio.daily_realized_pnl, "unrealized": _portfolio.total_unrealized_pnl, "total": _portfolio.total_pnl, "count": _portfolio.open_count, "closed_trades": _portfolio.closed_trades[-20:], "equity": broker.get_equity(), "buying_power": broker.get_buying_power(), "cash": broker.get_cash(), }) @app.route("/api/predictions") def api_predictions(): result = {} for sym in config.UNIVERSE: pred = _predictions.get(sym) sig = _signals.get(sym) if pred: result[sym] = { "symbol": sym, "current_price": pred.current_price, "direction": pred.predicted_direction, "entry_price": pred.entry_price, "stop_loss": pred.stop_loss_price, "take_profit": pred.take_profit_price, "expected_profit": pred.expected_profit_usd, "confidence": pred.confidence, "decision": sig.decision if sig else "—", "final_score": sig.final if sig else 0, "technical": sig.technical if sig else 0, "sentiment": sig.sentiment if sig else 0, } else: result[sym] = {"symbol": sym, "direction": "—", "decision": "—"} return jsonify(result) @app.route("/api/risk") def api_risk(): now_et = datetime.datetime.now(_ET) pdt_count = count_day_trades() if not config.ACCOUNT_BALANCE_ABOVE_25K else -1 sentiment_status = "warming_up" sentiment_age = None if not is_warming_up(): last = get_last_refresh() if last: sentiment_age = (datetime.datetime.now(datetime.timezone.utc) - last).total_seconds() / 60 sentiment_status = "stale" if sentiment_age > 35 else "fresh" else: sentiment_status = "no_data" circuits = [] if risk_module.profit_locked: circuits.append("profit_locked") if risk_module.half_size_mode: circuits.append("half_size") if risk_module.consecutive_losses >= 3: circuits.append(f"losing_streak_{risk_module.consecutive_losses}") if broker.safe_mode_active: circuits.append("safe_mode") overnight_risk = _portfolio.has_overnight_risk() if _portfolio else False return jsonify({ "daily_realized": risk_module.daily_realized_pnl, "daily_unrealized": risk_module.daily_unrealized_pnl, "daily_open_equity": risk_module.daily_open_equity, "consecutive_losses": risk_module.consecutive_losses, "half_size_mode": risk_module.half_size_mode, "profit_locked": risk_module.profit_locked, "safe_mode": broker.safe_mode_active, "pdt_count": pdt_count, "pdt_max": config.PDT_MAX_DAY_TRADES, "circuits_active": circuits, "sentiment_status": sentiment_status, "sentiment_age_min": sentiment_age, "overnight_risk": overnight_risk, "time_et": now_et.strftime("%H:%M:%S"), "market_open": 9 * 60 + 30 <= now_et.hour * 60 + now_et.minute <= 16 * 60 and now_et.weekday() < 5, }) @app.route("/api/chart/") def api_chart(symbol: str): """Return recent 5-min bars for mini charts.""" symbol = symbol.upper() if symbol not in config.UNIVERSE: return jsonify({"error": "symbol not in universe"}), 400 try: df = storage.get_all_bars(symbol, "5Min") if df.empty: return jsonify({"bars": []}) # Last 78 bars (~1 trading day) tail = df.tail(78) bars = [] for _, row in tail.iterrows(): bars.append({ "t": str(row.get("timestamp", "")), "o": round(float(row["open"]), 2), "h": round(float(row["high"]), 2), "l": round(float(row["low"]), 2), "c": round(float(row["close"]), 2), "v": int(row["volume"]), }) return jsonify({"bars": bars, "symbol": symbol}) except Exception as e: return jsonify({"bars": [], "error": str(e)}) @app.route("/api/chart_daily/") def api_chart_daily(symbol: str): """Return recent daily bars for trend chart.""" symbol = symbol.upper() if symbol not in config.UNIVERSE: return jsonify({"error": "symbol not in universe"}), 400 try: df = storage.get_all_bars(symbol, "1Day") if df.empty: return jsonify({"bars": []}) tail = df.tail(60) bars = [] for _, row in tail.iterrows(): bars.append({ "t": str(row.get("timestamp", ""))[:10], "o": round(float(row["open"]), 2), "h": round(float(row["high"]), 2), "l": round(float(row["low"]), 2), "c": round(float(row["close"]), 2), "v": int(row["volume"]), }) return jsonify({"bars": bars, "symbol": symbol}) except Exception as e: return jsonify({"bars": [], "error": str(e)}) @app.route("/api/sentiment") def api_sentiment(): """Return sentiment scores for all symbols.""" result = {} for sym in config.UNIVERSE: score = get_sentiment(sym) # Prefer the in-memory cache if it has data; otherwise fall back to DB cached value. if score and getattr(score, "source_count", 0) > 0: result[sym] = { "score": score.score, "cached_at": score.cached_at.isoformat(), "source_count": score.source_count, "stale": score.stale, } else: # Try DB-stored cached sentiment as a fallback so the UI can display values # even if the in-memory cache hasn't been populated yet. cached = storage.get_cached_sentiment(sym) if cached: # `cached` uses ISO timestamp strings from the DB. result[sym] = { "score": float(cached.get("score", 0.0)), "cached_at": cached.get("cached_at"), "source_count": int(cached.get("source_count", 0)), "stale": False, } else: result[sym] = {"score": 0.0, "stale": True, "source_count": 0} return jsonify(result) @app.route("/api/trades") def api_trades(): """Return recent trade events.""" with _trade_log_lock: return jsonify(_trade_log[-50:]) @app.route("/api/system_info") def api_system_info(): """Return model info, data download status, and system metrics.""" gpu_available = False gpu_name = None gpu_vram = None try: import torch gpu_available = torch.cuda.is_available() gpu_name = torch.cuda.get_device_name(0) if gpu_available else None if gpu_available: props = torch.cuda.get_device_properties(0) total = getattr(props, 'total_memory', None) or getattr(props, 'total_mem', 0) gpu_vram = round(total / (1024**3), 1) except Exception as e: logger.warning("Could not load torch for system info: %s", e) model_info = { "name": "ProsusAI/finbert", "type": "FinBERT (BERT fine-tuned for financial sentiment)", "labels": ["positive", "negative", "neutral"], "score_formula": "positive - negative → [-1, +1]", "max_tokens": 128, "batch_size": 16, "fallback": "VADER (no GPU)", "device": f"cuda ({gpu_name})" if gpu_available else "cpu (VADER fallback)", "gpu_available": gpu_available, "gpu_name": gpu_name, "gpu_vram_gb": gpu_vram, "news_sources": [], } if config.NEWS_API_KEY: model_info["news_sources"].append("NewsAPI") if config.GNEWS_API_KEY: model_info["news_sources"].append("GNews") # ── Technical indicators ── indicators = [ {"name": "RSI(14)", "weight": 0.20, "type": "Momentum"}, {"name": "MACD(12,26,9)", "weight": 0.20, "type": "Trend"}, {"name": "Bollinger(20,2σ)", "weight": 0.15, "type": "Volatility"}, {"name": "VWAP", "weight": 0.20, "type": "Volume-Price"}, {"name": "EMA-200", "weight": 0.15, "type": "Trend"}, {"name": "ATR Percentile", "weight": 0.10, "type": "Volatility Filter"}, ] # ── Data download status ── data_status = {} for sym in config.UNIVERSE: sym_data = {} for tf, label in [("1Day", "daily"), ("1Hour", "hourly"), ("5Min", "5min")]: count = storage.count_bars(sym, tf) last_ts = storage.get_last_timestamp(sym, tf) sym_data[label] = { "bars": count, "last_update": last_ts.isoformat() if last_ts else None, } data_status[sym] = sym_data total_bars = sum( d[tf]["bars"] for d in data_status.values() for tf in ["daily", "hourly", "5min"] ) # ── Backtest metrics (if available from last run) ── metrics = _last_backtest_metrics.copy() if _last_backtest_metrics else None return jsonify({ "model": model_info, "indicators": indicators, "data_status": data_status, "total_bars": total_bars, "universe": config.UNIVERSE, "metrics": metrics, }) # ── Backtest metrics store ── _last_backtest_metrics: dict = {} def update_backtest_metrics(metrics: dict): """Store the latest backtest metrics for display.""" global _last_backtest_metrics _last_backtest_metrics = metrics # ── Backtest results viewer ───────────────────────────────────────────────── BACKTEST_DIR = Path(__file__).resolve().parent.parent / "backtest_results" @app.route("/backtest") def backtest_viewer(): return send_from_directory(str(STATIC_DIR), "backtest_dashboard.html") @app.route("/api/backtest/list") def api_backtest_list(): """List all available backtest result sets.""" if not BACKTEST_DIR.exists(): return jsonify([]) sets: dict[str, dict] = {} for f in sorted(BACKTEST_DIR.glob("*_trades.csv")): key = f.stem.replace("_trades", "") parts = key.rsplit("_", 2) # symbols_start_end if len(parts) >= 3: symbols_str, start, end = parts[0], parts[1], parts[2] else: symbols_str, start, end = key, "", "" symbols = symbols_str.split("+") pnl_file = BACKTEST_DIR / f"{key}_daily_pnl.csv" sets[key] = { "key": key, "symbols": symbols, "symbol_count": len(symbols), "start": start, "end": end, "has_pnl": pnl_file.exists(), "trades_file": f.name, } return jsonify(list(sets.values())) @app.route("/api/backtest/trades/") def api_backtest_trades(key: str): """Return trades for a backtest run.""" trades_file = BACKTEST_DIR / f"{key}_trades.csv" if not trades_file.exists(): return jsonify({"error": "not found"}), 404 rows = [] with open(trades_file, newline="") as fh: reader = csv.DictReader(fh) for row in reader: for num_col in ("entry_price", "exit_price", "qty", "pnl", "duration_min", "confidence", "conf_multiplier"): if num_col in row and row[num_col]: try: row[num_col] = float(row[num_col]) except ValueError: pass rows.append(row) return jsonify(rows) @app.route("/api/backtest/daily_pnl/") def api_backtest_daily_pnl(key: str): """Return daily P&L for a backtest run.""" pnl_file = BACKTEST_DIR / f"{key}_daily_pnl.csv" if not pnl_file.exists(): return jsonify({"error": "not found"}), 404 rows = [] with open(pnl_file, newline="") as fh: reader = csv.DictReader(fh) for row in reader: if "pnl" in row: try: row["pnl"] = float(row["pnl"]) except ValueError: pass rows.append(row) return jsonify(rows) def start_server(portfolio: Portfolio, host: str = "127.0.0.1", port: int = 5000): """Start the Flask API server in a daemon thread.""" init(portfolio) thread = threading.Thread( target=lambda: app.run(host=host, port=port, debug=False, use_reloader=False), daemon=True, name="WebDashboardAPI", ) thread.start() logger.info("Web dashboard API started on http://localhost:%d", port) return thread @app.route("/health") def health(): """Health check endpoint for external monitoring.""" return jsonify({ "status": "ok", "timestamp": datetime.datetime.now(datetime.timezone.utc).isoformat(), "trading_mode": config.TRADING_MODE, "safe_mode": broker.safe_mode_active, }) @app.route("/kill", methods=["POST"]) def kill_switch(): """Emergency kill switch. Requires KILL_TOKEN header.""" expected_token = os.environ.get("KILL_TOKEN", "") if not expected_token: return jsonify({"error": "KILL_TOKEN not configured"}), 503 provided = request.headers.get("X-Kill-Token", "") if provided != expected_token: return jsonify({"error": "unauthorized"}), 403 broker.activate_safe_mode("Remote kill switch activated") return jsonify({"status": "safe_mode_activated"})