raghava4u's picture
Upload folder using huggingface_hub
d53dc44 verified
Raw History Blame Contribute Delete
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
@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/<symbol>")
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/<symbol>")
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/<path:key>")
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/<path:key>")
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"})