Trading-Bot-M20 / ml /ensemble_engine.py
raghava4u's picture
Upload folder using huggingface_hub
d53dc44 verified
Raw History Blame Contribute Delete
40.1 kB
"""
ml/ensemble_engine.py — Institution-grade ensemble backtest engine.
Integrates ALL new modules for a 3-way comparison:
A) Rule-based (original regime_engine)
B) ML-enhanced (existing ml_regime_engine)
C) ENSEMBLE (new): calibration + ensemble models + adaptive learning +
feature selection + regime-specific models + advanced risk + monitoring
Usage:
python -m ml.ensemble_engine --compare # 3-way A/B/C
python -m ml.ensemble_engine --ensemble # ensemble only
python -m ml.ensemble_engine --robustness # + robustness testing
python -m ml.ensemble_engine --timeframe 1Hour # hourly
"""
from __future__ import annotations
import argparse
import logging
import sys
import time
from pathlib import Path
import numpy as np
import pandas as pd
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
import config
from data.downloader import download_bars
from data.storage import get_all_bars
from backtester.strategies import generate_signals, STRATEGY_PARAMS
from backtester.regime_detector import classify_regimes, map_regime_to_intraday, regime_to_strategy
from backtester.portfolio import PortfolioManager, PortfolioRiskLimits, PortfolioPosition
from backtester.report import compute_metrics, generate_report
# Original ML modules
from ml.regime_classifier import MLRegimeClassifier
from ml.signal_model import MLSignalModel
from ml.strategy_blender import StrategyBlender
from ml.position_sizer import ProbabilisticSizer
from ml.portfolio_optimizer import DynamicPortfolioOptimizer
from ml.execution import ExecutionOptimizer
# New institution-grade modules
from ml.calibration import ProbabilityCalibrator, WalkForwardCalibrator
from ml.ensemble import EnsembleModel, EnsemblePrediction
from ml.adaptive import AdaptiveLearner, AdaptiveConfig
from ml.feature_selection import SHAPFeatureSelector, print_feature_report
from ml.regime_models import RegimeSpecificModels
from ml.risk_manager import AdvancedRiskManager, CorrelationRiskFilter
from ml.robustness import MonteCarloSimulator, StressTester, RobustnessReport, compute_robustness_score, print_robustness_report
from ml.monitoring import PerformanceMonitor, print_monitoring_dashboard
logger = logging.getLogger("trading_system.ml.ensemble_engine")
SPREAD_COST_PCT = 0.02
SLIPPAGE_PCT = 0.01
# ═══════════════════════════════════════════════════════════════════════════════
# ENSEMBLE BACKTEST
# ═══════════════════════════════════════════════════════════════════════════════
def ensemble_portfolio_backtest(
symbols: list[str],
data: dict[str, pd.DataFrame],
daily_data: dict[str, pd.DataFrame],
hourly_data: dict[str, pd.DataFrame] | None = None,
account_size: float = 100_000.0,
risk_per_trade_pct: float = 1.0,
limits: PortfolioRiskLimits | None = None,
regime_proxy_symbol: str = "SPY",
cooldown_bars: int = 6,
use_ensemble: bool = True,
use_lstm: bool = True,
) -> dict:
"""Institution-grade ensemble portfolio backtest.
Integrates: ensemble models, calibration, adaptive learning,
regime-specific models, advanced risk, execution optimization,
and performance monitoring.
"""
from ml.features import (
compute_bar_features, compute_regime_features,
merge_multi_tf_features, make_return_labels, make_regime_labels,
)
from signals.technical import compute_atr
limits = limits or PortfolioRiskLimits()
pm = PortfolioManager(account_size, limits)
regime_daily = daily_data.get(regime_proxy_symbol)
if regime_daily is None:
regime_daily = next(iter(daily_data.values()))
# ── Initialize ML Components ──
logger.info("Initializing ensemble system...")
regime_clf = MLRegimeClassifier(model_type="xgboost", min_train_days=60)
regime_result = regime_clf.walk_forward_predict(regime_daily)
blender = StrategyBlender()
exec_opt = ExecutionOptimizer(latency_bars=1, partial_fill_rate=0.95)
sizer = ProbabilisticSizer(base_risk_pct=risk_per_trade_pct)
risk_mgr = AdvancedRiskManager(target_annual_vol=0.15, kill_switch_dd=0.10)
corr_filter = CorrelationRiskFilter(max_correlation=0.70)
monitor = PerformanceMonitor(rolling_window=50)
adaptive = AdaptiveLearner(AdaptiveConfig(retrain_interval_bars=200, gap_bars=78))
calibrator = WalkForwardCalibrator(method="platt")
# Ensemble model per symbol
ensemble_models: dict[str, EnsembleModel] = {}
regime_specific: dict[str, RegimeSpecificModels] = {}
exec_masks: dict[str, pd.Series] = {}
ml_regime_probs: dict[str, pd.DataFrame] = {}
strategy_weights_map: dict[str, pd.DataFrame] = {}
# ── Pre-compute features and train models per symbol ──
symbol_features: dict[str, pd.DataFrame] = {}
symbol_labels: dict[str, np.ndarray] = {}
symbol_regime_labels: dict[str, np.ndarray] = {}
symbol_mom_signals: dict[str, pd.DataFrame] = {}
symbol_mr_signals: dict[str, pd.DataFrame] = {}
symbol_mom_atr: dict[str, pd.Series] = {}
symbol_mr_atr: dict[str, pd.Series] = {}
symbol_dfs: dict[str, pd.DataFrame] = {}
symbol_regimes: dict[str, pd.Series] = {}
for sym in symbols:
df = data.get(sym)
if df is None or len(df) < 200:
continue
symbol_dfs[sym] = df
daily_df = daily_data.get(sym)
hourly_df = hourly_data.get(sym) if hourly_data else None
# Compute features
try:
bar_feats = compute_bar_features(df, prefix="5m")
if daily_df is not None and len(daily_df) > 20:
regime_feats = compute_regime_features(daily_df)
features = merge_multi_tf_features(bar_feats, regime_feats, None)
else:
features = bar_feats
labels = make_return_labels(df, horizon=6, threshold=0.15)
features = features.iloc[:len(labels)]
labels = labels[:len(features)]
# Align and clean
valid = ~features.isna().any(axis=1) & ~np.isnan(labels)
features = features[valid]
labels = labels[valid.values]
symbol_features[sym] = features
symbol_labels[sym] = labels.astype(int) + 1 # {-1,0,1} → {0,1,2}
# Regime labels for regime-specific models
if daily_df is not None and len(daily_df) > 20:
r_labels = make_regime_labels(daily_df)
regime_daily_series = pd.Series(r_labels, index=daily_df.index[:len(r_labels)])
intra_regimes = regime_daily_series.reindex(features.index, method='ffill')
symbol_regime_labels[sym] = intra_regimes.fillna(0).astype(int).values
except Exception as e:
logger.warning("Feature computation failed for %s: %s", sym, e)
# Map regime probs to intraday
intra_probs = regime_clf.map_to_intraday(regime_result, df)
ml_regime_probs[sym] = intra_probs
strategy_weights_map[sym] = blender.compute_strategy_weights(intra_probs)
# Generate rule-based signals
mom_sig, mom_atr, _ = generate_signals("momentum", df, daily_df)
symbol_mom_signals[sym] = mom_sig
symbol_mom_atr[sym] = mom_atr
mr_sig, mr_atr, _ = generate_signals("mean_reversion", df, daily_df)
symbol_mr_signals[sym] = mr_sig
symbol_mr_atr[sym] = mr_atr
regime_series = map_regime_to_intraday(regime_daily, df)
symbol_regimes[sym] = regime_series
# Execution mask
atr = compute_atr(df, 14)
exec_masks[sym] = exec_opt.compute_execution_mask(df, atr)
if not symbol_dfs:
return {"error": "no_data", "trades": [], "daily_pnl": {}}
# ── Train Ensemble Models (Walk-Forward) ──
if use_ensemble:
for sym in symbol_features:
feats = symbol_features[sym]
labels = symbol_labels[sym]
if len(feats) < 500:
logger.warning("Insufficient data for ensemble on %s (%d bars)", sym, len(feats))
continue
logger.info("Training ensemble for %s (%d samples)...", sym, len(feats))
# Split walk-forward: 70% train, 30% test
train_end = int(len(feats) * 0.70)
X_train = feats.iloc[:train_end]
y_train = labels[:train_end]
# Train ensemble
ens = EnsembleModel(
combination="average",
use_lstm=use_lstm,
n_classes=3,
xgb_weight=0.35,
lgb_weight=0.35,
lstm_weight=0.30,
)
ens.fit(X_train, y_train)
ensemble_models[sym] = ens
# Train regime-specific models
if sym in symbol_regime_labels:
rsm = RegimeSpecificModels(min_regime_samples=100)
reg_labels = symbol_regime_labels[sym][:train_end]
rsm.fit(X_train, y_train, reg_labels)
regime_specific[sym] = rsm
# Calibrate on validation portion
val_probs = ens.predict_proba(X_train.iloc[-200:]).probs
val_labels = y_train[-200:]
calibrator.update(val_probs.max(axis=1), (val_labels > 0).astype(int))
logger.info(" %s ensemble: %d members trained", sym, len(ens._fitted_members))
# ── Dynamic portfolio optimizer ──
port_opt = DynamicPortfolioOptimizer()
daily_returns = {}
for sym in symbols:
if sym in daily_data:
daily_returns[sym] = daily_data[sym]["close"].pct_change().dropna()
if daily_returns:
port_opt.update_correlations(daily_returns)
# ── Simulate bar by bar ──
all_timestamps = sorted(set().union(*(df.index for df in symbol_dfs.values())))
cooldown_tracker: dict[str, int] = {}
gap_losses = 0.0
peak_equity = account_size
portfolio_returns: list[float] = []
prev_equity = account_size
regime_trade_count = {"BULL": 0, "BEAR": 0, "SIDEWAYS": 0}
strategy_trade_count = {"ensemble": 0, "momentum": 0, "mean_reversion": 0}
ensemble_stats = {"total_signals": 0, "ensemble_override": 0, "kill_switch_blocks": 0}
risk_mgr.update_equity(account_size)
for ti, ts in enumerate(all_timestamps):
date_str = str(ts.date()) if hasattr(ts, 'date') else str(ts)[:10]
pm.update_day(date_str)
# Track daily returns for risk manager
current_equity = pm.equity
if prev_equity > 0:
bar_return = (current_equity - prev_equity) / prev_equity
portfolio_returns.append(bar_return)
prev_equity = current_equity
# Risk state
risk_state = risk_mgr.get_risk_state(
current_equity,
np.array(portfolio_returns[-100:]) if portfolio_returns else np.array([0.0]),
)
# ── Exit checks (same as original) ──
for sym in list(pm.positions.keys()):
df = symbol_dfs[sym]
if ts not in df.index:
continue
idx = df.index.get_loc(ts)
pos = pm.positions[sym]
bar_low = float(df["low"].iloc[idx])
bar_high = float(df["high"].iloc[idx])
bar_open = float(df["open"].iloc[idx])
current_price = float(df["close"].iloc[idx])
# Dynamic stop tightening in high vol
stop_mult = risk_mgr.get_stop_multiplier(risk_state.realized_vol)
effective_stop = pos.stop
if stop_mult < 1.0:
if pos.side == "buy":
tighter = pos.entry_price - (pos.entry_price - pos.stop) * stop_mult
effective_stop = max(pos.stop, tighter)
else:
tighter = pos.entry_price + (pos.stop - pos.entry_price) * stop_mult
effective_stop = min(pos.stop, tighter)
stop_hit = False
exit_price = 0.0
if pos.side == "buy" and bar_low <= effective_stop:
exit_price = min(bar_open, effective_stop)
if bar_open < effective_stop:
gap_losses += abs(effective_stop - bar_open) * pos.qty
stop_hit = True
elif pos.side == "sell" and bar_high >= effective_stop:
exit_price = max(bar_open, effective_stop)
if bar_open > effective_stop:
gap_losses += abs(bar_open - effective_stop) * pos.qty
stop_hit = True
tp_hit = False
if not stop_hit:
if pos.side == "buy" and bar_high >= pos.tp:
exit_price = pos.tp
tp_hit = True
elif pos.side == "sell" and bar_low <= pos.tp:
exit_price = pos.tp
tp_hit = True
eod_close = False
if not stop_hit and not tp_hit:
if hasattr(ts, 'hour'):
et = ts.tz_convert("US/Eastern") if ts.tzinfo else ts
if et.hour >= 15 and et.minute >= 45:
exit_price = current_price
eod_close = True
if stop_hit or tp_hit or eod_close:
reason = "stop_loss" if stop_hit else "take_profit" if tp_hit else "end_of_day"
pnl = pm.close_position(sym, exit_price, ts, reason,
SPREAD_COST_PCT, SLIPPAGE_PCT)
pm.add_realized_pnl(pnl)
# Record to monitor
trade_ret = pnl / (pos.dollar_risk + 1e-9)
monitor.record_trade(trade_ret / 100, won=pnl > 0)
adaptive.record_trade_return(trade_ret)
# ── Kill switch check ──
if risk_state.kill_switch_active:
ensemble_stats["kill_switch_blocks"] += 1
continue
# ── Entry checks ──
for sym in symbols:
if sym in pm.positions:
continue
if sym not in symbol_dfs:
continue
df = symbol_dfs[sym]
if ts not in df.index:
continue
idx = df.index.get_loc(ts)
if idx >= len(df) - 1:
continue
last = cooldown_tracker.get(sym, -999)
if ti - last < cooldown_bars:
continue
# Execution filter
if sym in exec_masks and not exec_masks[sym].iloc[idx]:
continue
# ── ENSEMBLE SIGNAL ──
sig = 0
conf = 0.0
strategy_name = "momentum"
if use_ensemble and sym in ensemble_models:
ens = ensemble_models[sym]
feats = symbol_features.get(sym)
if feats is not None and idx < len(feats):
ensemble_stats["total_signals"] += 1
# Get ensemble prediction
try:
row = feats.iloc[[idx]]
ens_pred = ens.predict_proba(row)
probs = ens_pred.probs[0] # [p_down, p_flat, p_up]
# Calibrate
raw_conf = float(probs.max())
cal_conf = calibrator.calibrate(np.array([raw_conf]))[0]
# Regime-specific model overlay
if sym in regime_specific and sym in ml_regime_probs:
rp = ml_regime_probs[sym]
if idx < len(rp):
regime_p = np.array([[
float(rp["p_sideways"].iloc[idx]),
float(rp["p_bull"].iloc[idx]),
float(rp["p_bear"].iloc[idx]),
]])
rsm_pred = regime_specific[sym].predict(row, regime_p)
# Blend ensemble + regime-specific 60/40
probs = 0.6 * probs + 0.4 * rsm_pred.blended_probs[0]
predicted_class = int(np.argmax(probs))
max_prob = float(probs[predicted_class])
# Conservative: only trade if calibrated confidence > 0.40
if cal_conf < 0.40 or max_prob < 0.38:
continue
if predicted_class == 2: # UP
sig = 1
conf = cal_conf * 10 # scale to ~5-7 range
elif predicted_class == 0: # DOWN
sig = -1
conf = cal_conf * 10
# class 1 (FLAT) → no trade
ensemble_stats["ensemble_override"] += 1
except Exception as e:
logger.debug("Ensemble predict failed for %s: %s", sym, e)
sig = 0
# Fallback: rule-based blended signal
if sig == 0 and sym in symbol_mom_signals:
regime = symbol_regimes[sym].iloc[idx] if sym in symbol_regimes else "SIDEWAYS"
strategy_name = regime_to_strategy(regime)
strat_wts = strategy_weights_map.get(sym)
if strat_wts is not None:
blended = blender.blend_signals(
symbol_mom_signals[sym],
symbol_mr_signals[sym],
strat_wts,
)
sig = int(blended["entry"].iloc[idx])
conf = float(blended["confidence"].iloc[idx])
else:
if strategy_name == "mean_reversion":
sig_df = symbol_mr_signals[sym]
else:
sig_df = symbol_mom_signals[sym]
sig = int(sig_df["entry"].iloc[idx])
conf = float(sig_df["confidence"].iloc[idx])
if sig == 0:
continue
# ATR
if strategy_name == "mean_reversion":
atr_series = symbol_mr_atr.get(sym)
else:
atr_series = symbol_mom_atr.get(sym)
if atr_series is None:
continue
current_atr = float(atr_series.iloc[idx])
if np.isnan(current_atr) or current_atr <= 0:
continue
side = "buy" if sig == 1 else "sell"
entry_price = float(df["open"].iloc[idx + 1])
# Apply latency
exec_idx = exec_opt.apply_latency(idx + 1, len(df) - 1)
if exec_idx != idx + 1:
entry_price = float(df["open"].iloc[exec_idx])
entry_slippage = entry_price * SLIPPAGE_PCT / 100
entry_price += entry_slippage if side == "buy" else -entry_slippage
# Strategy parameters
params = STRATEGY_PARAMS.get(strategy_name, STRATEGY_PARAMS["momentum"])
stop_distance = params["stop_atr_mult"] * current_atr
# Vol-targeting scale
vol_scale = risk_state.vol_scale
# Confidence multiplier
conf_lo, conf_hi = params["conf_range"]
conf_multiplier = 0.5 + (min(conf, conf_hi) - conf_lo) / (conf_hi - conf_lo)
conf_multiplier = max(0.5, min(1.5, conf_multiplier))
risk_scale = params.get("risk_scale", 1.0)
peak_equity = max(peak_equity, pm.equity)
# Position sizing with vol targeting
atr_avg = float(atr_series.rolling(50).mean().iloc[idx])
if np.isnan(atr_avg) or atr_avg <= 0:
atr_avg = current_atr
ml_conf = conf / 10.0 # normalize back to 0-1 range
dollar_risk = sizer.compute_position_risk(
equity=pm.equity,
peak_equity=peak_equity,
ml_confidence=ml_conf,
conf_multiplier=conf_multiplier,
atr_current=current_atr,
atr_avg=atr_avg,
strategy_risk_scale=risk_scale,
)
# Apply vol-targeting scale
dollar_risk *= vol_scale
qty = dollar_risk / stop_distance
if qty < 0.01:
continue
# Portfolio check
if hasattr(port_opt, 'can_enter'):
regime_p = {}
if sym in ml_regime_probs:
rp = ml_regime_probs[sym]
if idx < len(rp):
regime_p = {
"p_bull": float(rp["p_bull"].iloc[idx]),
"p_bear": float(rp["p_bear"].iloc[idx]),
"p_sideways": float(rp["p_sideways"].iloc[idx]),
}
can, reason = port_opt.can_enter(
sym, dollar_risk, pm.equity, pm.positions,
max_per_cluster=2, regime_probs=regime_p,
)
else:
can, reason = pm.can_enter(sym, dollar_risk)
if not can:
continue
# Correlation filter
if daily_returns:
corr_matrix = corr_filter.compute_correlation_matrix(daily_returns)
pos_exposures = {
s: p.dollar_risk / pm.equity
for s, p in pm.positions.items()
}
if not corr_filter.check_position_allowed(sym, pos_exposures, corr_matrix):
continue
# Stops and targets
if side == "buy":
stop = entry_price - stop_distance
tp = entry_price + params["tp_rr_ratio"] * stop_distance
else:
stop = entry_price + stop_distance
tp = entry_price - params["tp_rr_ratio"] * stop_distance
pos = PortfolioPosition(
symbol=sym, side=side, entry_price=entry_price,
qty=qty, stop=stop, tp=tp, initial_risk=stop_distance,
entry_time=df.index[min(exec_idx, len(df) - 1)],
confidence=conf,
conf_multiplier=conf_multiplier, dollar_risk=dollar_risk,
)
pm.open_position(pos)
cooldown_tracker[sym] = ti
regime = symbol_regimes[sym].iloc[idx] if sym in symbol_regimes else "SIDEWAYS"
regime_trade_count[regime] = regime_trade_count.get(regime, 0) + 1
strategy_trade_count["ensemble"] = strategy_trade_count.get("ensemble", 0) + 1
# ── Close remaining positions ──
for sym in list(pm.positions.keys()):
df = symbol_dfs[sym]
last_price = float(df["close"].iloc[-1])
pnl = pm.close_position(sym, last_price, df.index[-1], "end_of_backtest",
SPREAD_COST_PCT, SLIPPAGE_PCT)
pm.add_realized_pnl(pnl)
return {
"account_size": account_size,
"final_equity": round(pm.equity, 2),
"total_trades": len(pm.trades),
"trades": pm.trades,
"daily_pnl": pm.daily_pnl,
"gap_losses_usd": round(gap_losses, 2),
"pdt_blocked_count": 0,
"symbol": "+".join(symbols),
"strategy": "ensemble_institution",
"start_date": str(all_timestamps[0].date()) if all_timestamps else "N/A",
"end_date": str(all_timestamps[-1].date()) if all_timestamps else "N/A",
"regime_trade_count": regime_trade_count,
"strategy_trade_count": strategy_trade_count,
"ensemble_stats": ensemble_stats,
"monitor": monitor,
}
# ═══════════════════════════════════════════════════════════════════════════════
# 3-WAY COMPARISON: RULE vs ML vs ENSEMBLE
# ═══════════════════════════════════════════════════════════════════════════════
def compare_three_way(
symbols: list[str],
data: dict[str, pd.DataFrame],
daily_data: dict[str, pd.DataFrame],
hourly_data: dict[str, pd.DataFrame] | None = None,
account_size: float = 100_000.0,
risk_per_trade_pct: float = 1.0,
limits: PortfolioRiskLimits | None = None,
regime_proxy_symbol: str = "SPY",
cooldown_bars: int = 6,
run_robustness: bool = False,
) -> dict:
"""Run 3-way comparison: Rule-based vs ML-enhanced vs Ensemble."""
from ml.ml_regime_engine import ml_regime_portfolio_backtest
print("\n" + "=" * 100)
print(" 3-WAY COMPARISON: Rule-Based vs ML-Enhanced vs ENSEMBLE")
print("=" * 100)
results = {}
# A) Rule-based
print("\n [1/3] Running RULE-BASED backtest...")
t0 = time.time()
rule_result = ml_regime_portfolio_backtest(
symbols=symbols, data=data, daily_data=daily_data,
hourly_data=hourly_data, account_size=account_size,
risk_per_trade_pct=risk_per_trade_pct, limits=limits,
regime_proxy_symbol=regime_proxy_symbol,
cooldown_bars=cooldown_bars, use_ml=False,
)
rule_time = time.time() - t0
rule_metrics = compute_metrics(rule_result)
results["rule_based"] = {"result": rule_result, "metrics": rule_metrics, "time": rule_time}
print(f" Done in {rule_time:.1f}s — Return: {rule_metrics.get('total_return_pct', 0):+.2f}%")
# B) ML-enhanced
print(f" [2/3] Running ML-ENHANCED backtest...")
t0 = time.time()
ml_result = ml_regime_portfolio_backtest(
symbols=symbols, data=data, daily_data=daily_data,
hourly_data=hourly_data, account_size=account_size,
risk_per_trade_pct=risk_per_trade_pct, limits=limits,
regime_proxy_symbol=regime_proxy_symbol,
cooldown_bars=cooldown_bars, use_ml=True,
)
ml_time = time.time() - t0
ml_metrics = compute_metrics(ml_result)
results["ml_enhanced"] = {"result": ml_result, "metrics": ml_metrics, "time": ml_time}
print(f" Done in {ml_time:.1f}s — Return: {ml_metrics.get('total_return_pct', 0):+.2f}%")
# C) Ensemble
print(f" [3/3] Running ENSEMBLE backtest...")
t0 = time.time()
ens_result = ensemble_portfolio_backtest(
symbols=symbols, data=data, daily_data=daily_data,
hourly_data=hourly_data, account_size=account_size,
risk_per_trade_pct=risk_per_trade_pct, limits=limits,
regime_proxy_symbol=regime_proxy_symbol,
cooldown_bars=cooldown_bars, use_ensemble=True,
)
ens_time = time.time() - t0
ens_metrics = compute_metrics(ens_result)
results["ensemble"] = {"result": ens_result, "metrics": ens_metrics, "time": ens_time}
print(f" Done in {ens_time:.1f}s — Return: {ens_metrics.get('total_return_pct', 0):+.2f}%")
# ── Print comparison table ──
print("\n" + "=" * 100)
print(" 3-WAY COMPARISON RESULTS")
print("=" * 100)
print(f" {'Metric':<25} {'Rule-Based':>15} {'ML-Enhanced':>15} {'ENSEMBLE':>15} {'Best':>10}")
print(" " + "-" * 80)
comparisons = [
("Total Return %", "total_return_pct", "+"),
("Sharpe Ratio", "sharpe_ratio", "+"),
("Sortino Ratio", "sortino_ratio", "+"),
("Max Drawdown %", "max_drawdown_pct", "-"),
("Win Rate %", "win_rate_pct", "+"),
("Expectancy $/trade", "expectancy_usd", "+"),
("Total Trades", "total_trades", None),
("CAGR %", "cagr_pct", "+"),
("Calmar Ratio", "calmar_ratio", "+"),
]
for label, key, direction in comparisons:
rv = rule_metrics.get(key, 0)
mv = ml_metrics.get(key, 0)
ev = ens_metrics.get(key, 0)
if direction == "+":
best = "ENS" if ev >= mv and ev >= rv else ("ML" if mv >= rv else "RULE")
elif direction == "-":
best = "ENS" if ev <= mv and ev <= rv else ("ML" if mv <= rv else "RULE")
else:
best = ""
print(f" {label:<25} {rv:>15.2f} {mv:>15.2f} {ev:>15.2f} {best:>10}")
# Profit factors
for label, result in [("Rule", rule_result), ("ML", ml_result), ("Ensemble", ens_result)]:
trades = result.get("trades", [])
gp = sum(t["pnl"] for t in trades if t["pnl"] > 0)
gl = abs(sum(t["pnl"] for t in trades if t["pnl"] <= 0))
pf = gp / gl if gl > 0 else 0
print(f" {'Profit Factor (' + label + ')':<25} {pf:>15.2f}")
# Monthly comparison
print("\n MONTHLY RETURN COMPARISON:")
print(f" {'Month':>10} {'Rule %':>10} {'ML %':>10} {'ENS %':>10}")
print(" " + "-" * 50)
rule_monthly = _compute_monthly(rule_result)
ml_monthly = _compute_monthly(ml_result)
ens_monthly = _compute_monthly(ens_result)
all_months = sorted(set(list(rule_monthly) + list(ml_monthly) + list(ens_monthly)))
for month in all_months:
rp = rule_monthly.get(month, 0) / account_size * 100
mp = ml_monthly.get(month, 0) / account_size * 100
ep = ens_monthly.get(month, 0) / account_size * 100
print(f" {month:>10} {rp:>+9.2f}% {mp:>+9.2f}% {ep:>+9.2f}%")
# Ensemble stats
ens_stats = ens_result.get("ensemble_stats", {})
if ens_stats:
print(f"\n ENSEMBLE STATS:")
for k, v in ens_stats.items():
print(f" {k:<25} {v:>6}")
print(f"\n Execution time: Rule={rule_time:.1f}s, ML={ml_time:.1f}s, Ensemble={ens_time:.1f}s")
# ── Robustness testing on ensemble results ──
if run_robustness and ens_result.get("trades"):
print("\n" + "=" * 100)
print(" ROBUSTNESS TESTING (Ensemble)")
print("=" * 100)
trade_returns = np.array([
t["pnl"] / max(abs(t.get("dollar_risk", t["pnl"])), 1)
for t in ens_result["trades"]
])
mc_sim = MonteCarloSimulator(n_simulations=1000)
mc_result = mc_sim.simulate(trade_returns)
stress = StressTester()
stress_results = stress.run_all(trade_returns)
score = compute_robustness_score(mc_result, stress_results)
report = RobustnessReport(
monte_carlo=mc_result,
stress_tests=stress_results,
overall_score=score,
)
print_robustness_report(report)
# ── Performance monitoring dashboard ──
if "monitor" in ens_result:
ens_monitor = ens_result["monitor"]
ens_monitor.set_backtest_baseline(
sharpe=ens_metrics.get("sharpe_ratio", 0),
winrate=ens_metrics.get("win_rate_pct", 0) / 100,
total_return=ens_metrics.get("total_return_pct", 0) / 100,
max_dd=ens_metrics.get("max_drawdown_pct", 0) / 100,
)
print_monitoring_dashboard(ens_monitor)
print("=" * 100)
return results
def _compute_monthly(result: dict) -> dict:
daily_pnl = result.get("daily_pnl", {})
monthly = {}
for date_str, pnl in daily_pnl.items():
month_key = date_str[:7]
monthly[month_key] = monthly.get(month_key, 0) + pnl
return monthly
# ═══════════════════════════════════════════════════════════════════════════════
# CLI
# ═══════════════════════════════════════════════════════════════════════════════
def main():
parser = argparse.ArgumentParser(description="Institution-Grade Ensemble Backtester")
parser.add_argument("--symbols", type=str, default=None)
parser.add_argument("--timeframe", type=str, default="5Min", choices=["5Min", "1Hour"])
parser.add_argument("--account-size", type=float, default=100_000)
parser.add_argument("--risk-pct", type=float, default=1.0)
parser.add_argument("--cooldown", type=int, default=6)
parser.add_argument("--download-days", type=int, default=30)
parser.add_argument("--compare", action="store_true", help="3-way A/B/C comparison")
parser.add_argument("--ensemble", action="store_true", help="Ensemble only")
parser.add_argument("--robustness", action="store_true", help="Include robustness testing")
parser.add_argument("--no-lstm", action="store_true", help="Disable LSTM in ensemble")
args = parser.parse_args()
from monitoring.logger import setup_logging
setup_logging("INFO")
symbols = ([s.strip().upper() for s in args.symbols.split(",")]
if args.symbols else config.UNIVERSE)
# Download data
logger.info("Downloading data for %d symbols...", len(symbols))
intraday_data = {}
daily_data = {}
hourly_data = {}
download_days = args.download_days
if args.timeframe == "1Hour":
download_days = min(download_days, 730)
else:
download_days = min(download_days, 30)
for sym in symbols:
logger.info("Downloading %s...", sym)
try:
download_bars(sym, args.timeframe, download_days)
download_bars(sym, "1Day", 500)
if args.timeframe == "5Min":
try:
download_bars(sym, "1Hour", 180)
except Exception:
pass
except Exception as e:
logger.warning("Download failed for %s: %s", sym, e)
continue
df = get_all_bars(sym, args.timeframe)
ddf = get_all_bars(sym, "1Day")
if not df.empty and len(df) >= 100:
intraday_data[sym] = df
logger.info(" %s %s: %d bars (%s -> %s)",
sym, args.timeframe, len(df),
df.index[0].date(), df.index[-1].date())
if not ddf.empty:
daily_data[sym] = ddf
hdf = get_all_bars(sym, "1Hour")
if not hdf.empty:
hourly_data[sym] = hdf
if not intraday_data:
logger.error("No data available")
sys.exit(1)
regime_proxy = "SPY" if "SPY" in daily_data else next(iter(daily_data.keys()))
limits = PortfolioRiskLimits(
max_total_exposure_pct=60.0,
max_concurrent_positions=4,
max_per_group=2,
max_daily_loss_pct=2.0,
max_per_symbol_risk_pct=2.0,
)
if args.compare:
compare_three_way(
symbols=list(intraday_data.keys()),
data=intraday_data,
daily_data=daily_data,
hourly_data=hourly_data or None,
account_size=args.account_size,
risk_per_trade_pct=args.risk_pct,
limits=limits,
regime_proxy_symbol=regime_proxy,
cooldown_bars=args.cooldown,
run_robustness=args.robustness,
)
else:
# Ensemble only
logger.info("Running ENSEMBLE backtest...")
result = ensemble_portfolio_backtest(
symbols=list(intraday_data.keys()),
data=intraday_data,
daily_data=daily_data,
hourly_data=hourly_data or None,
account_size=args.account_size,
risk_per_trade_pct=args.risk_pct,
limits=limits,
regime_proxy_symbol=regime_proxy,
cooldown_bars=args.cooldown,
use_ensemble=True,
use_lstm=not args.no_lstm,
)
if "error" not in result:
metrics = compute_metrics(result)
_print_ensemble_summary(result, metrics)
generate_report(result)
if args.robustness and result.get("trades"):
trade_returns = np.array([
t["pnl"] / max(abs(t.get("dollar_risk", t["pnl"])), 1)
for t in result["trades"]
])
mc_sim = MonteCarloSimulator(n_simulations=1000)
mc_result = mc_sim.simulate(trade_returns)
stress = StressTester()
stress_results = stress.run_all(trade_returns)
score = compute_robustness_score(mc_result, stress_results)
report = RobustnessReport(
monte_carlo=mc_result,
stress_tests=stress_results,
overall_score=score,
)
print_robustness_report(report)
if "monitor" in result:
print_monitoring_dashboard(result["monitor"])
else:
logger.error("Backtest error: %s", result.get("error"))
def _print_ensemble_summary(result: dict, metrics: dict):
"""Pretty-print ensemble backtest results."""
print()
print("=" * 90)
print(f" ENSEMBLE INSTITUTION-GRADE BACKTEST RESULTS")
print("=" * 90)
print(f" Period: {result['start_date']} -> {result['end_date']}")
print(f" Symbols: {result['symbol']}")
print(f" Account: ${result['account_size']:,.0f} -> ${result['final_equity']:,.2f}")
print(f" Total Return: {metrics['total_return_pct']:+.2f}%")
print(f" Total Trades: {metrics['total_trades']}")
print(f" Win Rate: {metrics['win_rate_pct']:.1f}%")
print(f" Sharpe Ratio: {metrics['sharpe_ratio']:.2f}")
print(f" Sortino Ratio: {metrics['sortino_ratio']:.2f}")
print(f" Max Drawdown: {metrics['max_drawdown_pct']:.2f}%")
print(f" Calmar Ratio: {metrics['calmar_ratio']:.2f}")
print(f" Expectancy: ${metrics['expectancy_usd']:.2f}/trade")
trades = result.get("trades", [])
gp = sum(t["pnl"] for t in trades if t["pnl"] > 0)
gl = abs(sum(t["pnl"] for t in trades if t["pnl"] <= 0))
pf = gp / gl if gl > 0 else 0
print(f" Profit Factor: {pf:.2f}")
ens_stats = result.get("ensemble_stats", {})
if ens_stats:
print(f"\n ENSEMBLE STATS:")
for k, v in ens_stats.items():
print(f" {k:<25} {v:>6}")
# Monthly
daily_pnl = result.get("daily_pnl", {})
if daily_pnl:
monthly = {}
for date_str, pnl in daily_pnl.items():
month_key = date_str[:7]
monthly[month_key] = monthly.get(month_key, 0) + pnl
print(f"\n MONTHLY P&L:")
for month in sorted(monthly.keys()):
pnl = monthly[month]
pct = pnl / result["account_size"] * 100
bar = "█" * max(1, int(abs(pct) * 5))
sign = "+" if pnl >= 0 else ""
print(f" {month} {sign}${pnl:>8,.2f} ({sign}{pct:.2f}%) {'▓' if pnl >= 0 else '░'}{bar}")
print("=" * 90)
if __name__ == "__main__":
main()