Download ml/ensemble_engine.py from raghava4u/Trading-Bot-M20: direct link, hf CLI and curl.
- Browser
- Download file 40.1 kB
-
https://huggingface.co/raghava4u/Trading-Bot-M20/resolve/main/ml/ensemble_engine.py
- Command line
-
hf download hf://raghava4u/Trading-Bot-M20/ml/ensemble_engine.py
-
curl -L -o ensemble_engine.py https://huggingface.co/raghava4u/Trading-Bot-M20/resolve/main/ml/ensemble_engine.py
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() | |