Download monitoring/web_api.py from raghava4u/Trading-Bot-M20: direct link, hf CLI and curl.
- Browser
- Download file 17.9 kB
-
https://huggingface.co/raghava4u/Trading-Bot-M20/resolve/main/monitoring/web_api.py
- Command line
-
hf download hf://raghava4u/Trading-Bot-M20/monitoring/web_api.py
-
curl -L -o web_api.py https://huggingface.co/raghava4u/Trading-Bot-M20/resolve/main/monitoring/web_api.py
17.9 kB
| """ | |
| 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 | |
| def index(): | |
| return send_from_directory(str(STATIC_DIR), "dashboard.html") | |
| 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, | |
| }) | |
| 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(), | |
| }) | |
| 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) | |
| 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, | |
| }) | |
| 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)}) | |
| 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)}) | |
| 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) | |
| def api_trades(): | |
| """Return recent trade events.""" | |
| with _trade_log_lock: | |
| return jsonify(_trade_log[-50:]) | |
| 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" | |
| def backtest_viewer(): | |
| return send_from_directory(str(STATIC_DIR), "backtest_dashboard.html") | |
| 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())) | |
| 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) | |
| 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 | |
| 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, | |
| }) | |
| 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"}) | |