Coverage for src / monte_neo / core / gpu_scenarios.py: 91%
94 statements
« prev ^ index » next coverage.py v7.13.1, created at 2026-01-28 16:27 +0200
« prev ^ index » next coverage.py v7.13.1, created at 2026-01-28 16:27 +0200
1"""GPU scenario backtesting helper."""
3from __future__ import annotations
5from typing import TYPE_CHECKING, Any
7import mlx.core as mx
8import numpy as np
9import pandas as pd
11from monte_neo.metrics.calculator import MetricsCalculator
12from monte_neo.monte_carlo.workers import _generate_signals_wrapper
13from monte_neo.utils.parallel import ParallelExecutor
15if TYPE_CHECKING:
16 from monte_neo.indicators.base import BaseIndicator
19def normalize_signal_array(
20 signals: Any,
21 target_len: int,
22) -> np.ndarray:
23 """Normalize signals to numpy array of target length."""
24 if target_len <= 0:
25 return np.zeros(0, dtype=np.float32)
26 if signals is None:
27 return np.zeros(target_len, dtype=np.float32)
28 if isinstance(signals, pd.DataFrame):
29 if "signal" in signals.columns:
30 arr = signals["signal"].to_numpy()
31 else:
32 arr = signals.to_numpy().reshape(-1)
33 elif isinstance(signals, pd.Series):
34 arr = signals.to_numpy()
35 else:
36 arr = np.asarray(signals)
37 arr = arr.astype(np.float32, copy=False).reshape(-1)
38 if arr.size >= target_len:
39 return arr[:target_len]
40 return np.pad(arr, (0, target_len - arr.size), "constant", constant_values=0)
43def run_scenarios_backtest(
44 indicator: BaseIndicator,
45 scenarios: list[pd.DataFrame],
46 executor: ParallelExecutor | None = None,
47 use_sl_tp: bool = False,
48 sl_pct: float = 0.0,
49 tp_pct: float = 0.0,
50) -> list[dict[str, Any]]:
51 """Run one indicator across many data scenarios on GPU."""
52 if use_sl_tp:
53 # Optimization: Use parallelized batch Numba for SL/TP on multiple scenarios
54 # This is much faster than the previous Python loop
56 # 1. Get Signals (CPU parallelized)
57 tasks = [(indicator, df) for df in scenarios]
58 if executor is None:
59 local_executor = ParallelExecutor()
60 raw_signals = local_executor.map(_generate_signals_wrapper, tasks)
61 else:
62 raw_signals = executor.map(_generate_signals_wrapper, tasks)
64 # 2. Prepare Data and Signal Matrix
65 max_len = max(len(df) for df in scenarios)
67 # We need a unified price matrix for Numba batch
68 # Since scenarios can have different prices, we pad them
69 close_matrix = np.zeros((len(scenarios), max_len), dtype=np.float64)
70 high_matrix = np.zeros((len(scenarios), max_len), dtype=np.float64)
71 low_matrix = np.zeros((len(scenarios), max_len), dtype=np.float64)
72 signal_matrix = np.zeros((len(scenarios), max_len), dtype=np.int32)
74 for i, df in enumerate(scenarios):
75 l = len(df)
76 close_matrix[i, :l] = df["close"].to_numpy().astype(np.float64)
77 high_matrix[i, :l] = df["high"].to_numpy().astype(np.float64)
78 low_matrix[i, :l] = df["low"].to_numpy().astype(np.float64)
79 signal_matrix[i, :l] = normalize_signal_array(raw_signals[i], l).astype(np.int32)
81 # 3. Run Batch Calculation (Multi-scenario version)
82 # We use calculate_batch_multi_price_fast because each scenario has its own prices
83 batch_metrics = MetricsCalculator.calculate_batch_multi_price_fast(
84 close_matrix,
85 high_matrix,
86 low_matrix,
87 signal_matrix,
88 use_sl_tp,
89 sl_pct,
90 tp_pct
91 )
93 results = []
94 for i in range(len(scenarios)):
95 total_return = float(batch_metrics[i, 0])
96 max_dd = float(batch_metrics[i, 1])
97 pf = float(batch_metrics[i, 2])
98 trade_count = int(batch_metrics[i, 3])
100 results.append({
101 "total_return": total_return,
102 "max_drawdown": max_dd,
103 "profit_factor": pf,
104 "passed": bool(total_return > 0.0 and max_dd < 0.2),
105 "metrics": {
106 "total_return": total_return,
107 "max_drawdown": max_dd,
108 "profit_factor": pf,
109 "trade_count": trade_count,
110 }
111 })
112 return results
114 # 1. Prepare Returns Matrix (S_scenarios x T_bars)
115 # Assuming OHLCV format, we pre-calculate returns for all scenarios
116 returns_list = []
117 for df in scenarios:
118 rets = (df["close"].values[1:] / df["close"].values[:-1]) - 1
119 returns_list.append(rets.astype(np.float32))
121 # Handle variable lengths by padding with 0
122 if not returns_list:
123 return []
125 max_len = max(len(r) for r in returns_list)
126 padded_returns = []
128 for r in returns_list:
129 pad_width = max_len - len(r)
130 if pad_width > 0:
131 padded_returns.append(np.pad(r, (0, pad_width), "constant", constant_values=0))
132 else:
133 padded_returns.append(r)
135 # Matrix: (S, T-1)
136 returns_matrix = mx.array(np.stack(padded_returns))
138 # 2. Get Signals (CPU parallelized)
139 # We use ParallelExecutor to speed up signal generation for many scenarios
140 signal_list = []
142 # Prepare args for parallel execution
143 # (indicator is pickleable usually)
144 tasks = [(indicator, df) for df in scenarios]
146 # Use ProcessPoolExecutor for CPU-bound signal generation
147 # Adjust n_workers as needed, default is usually fine
148 # NOTE: Parallelizing might have overhead for small DataFrames.
149 # But for 1000 scenarios, it should help.
151 # If scenarios are few, do serial
152 if len(scenarios) < 50:
153 raw_signals = []
154 for df in scenarios:
155 try:
156 raw_signals.append(indicator.generate_signals(df))
157 except Exception:
158 raw_signals.append(None)
159 else:
160 # Parallel execution
161 if executor is None:
162 local_executor = ParallelExecutor()
163 # Map returns results in order
164 raw_signals = local_executor.map(_generate_signals_wrapper, tasks)
165 else:
166 # Use shared executor
167 raw_signals = executor.map(_generate_signals_wrapper, tasks)
169 for sigs in raw_signals:
170 signal_list.append(normalize_signal_array(sigs, max_len))
172 # Matrix: (S, T-1)
173 signal_matrix_np = np.stack(signal_list)
174 signal_matrix_mx: Any = mx.array(signal_matrix_np.astype(np.int32))
176 # 3. Massive GPU calc
177 strat_returns = signal_matrix_mx * returns_matrix
179 # Vectorized metrics
180 equity_curves = mx.exp(
181 mx.cumsum(mx.log1p(mx.clip(strat_returns, -0.9, 10.0)), axis=1)
182 )
184 final_rets = np.array(equity_curves[:, -1])
186 # Max DD
187 # running_max = mx.maximum.accumulate(equity_curves, axis=1) # type: ignore
188 running_max = mx.cummax(equity_curves, axis=1)
189 max_dds = np.array(mx.max((running_max - equity_curves) / running_max, axis=1))
191 # Profit Factor
192 wins = mx.where(strat_returns > 0, strat_returns, 0)
193 losses = mx.where(strat_returns < 0, strat_returns, 0)
194 gross_profit = mx.sum(wins, axis=1)
195 gross_loss = mx.abs(mx.sum(losses, axis=1))
196 profit_factor = np.array(mx.where(gross_loss > 0, gross_profit / gross_loss, 100.0))
198 results = []
199 for i in range(len(scenarios)):
200 results.append(
201 {
202 "total_return": float(final_rets[i]) - 1.0,
203 "max_drawdown": float(max_dds[i]),
204 "profit_factor": float(profit_factor[i]),
205 "passed": bool(final_rets[i] > 1.0 and max_dds[i] < 0.2),
206 "metrics": {
207 "total_return": float(final_rets[i]) - 1.0,
208 "max_drawdown": float(max_dds[i]),
209 "profit_factor": float(profit_factor[i]),
210 },
211 }
212 )
213 return results