""" 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()