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

1"""GPU scenario backtesting helper.""" 

2 

3from __future__ import annotations 

4 

5from typing import TYPE_CHECKING, Any 

6 

7import mlx.core as mx 

8import numpy as np 

9import pandas as pd 

10 

11from monte_neo.metrics.calculator import MetricsCalculator 

12from monte_neo.monte_carlo.workers import _generate_signals_wrapper 

13from monte_neo.utils.parallel import ParallelExecutor 

14 

15if TYPE_CHECKING: 

16 from monte_neo.indicators.base import BaseIndicator 

17 

18 

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) 

41 

42 

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 

55 

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) 

63 

64 # 2. Prepare Data and Signal Matrix 

65 max_len = max(len(df) for df in scenarios) 

66 

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) 

73 

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) 

80 

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 ) 

92 

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]) 

99 

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 

113 

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

120 

121 # Handle variable lengths by padding with 0 

122 if not returns_list: 

123 return [] 

124 

125 max_len = max(len(r) for r in returns_list) 

126 padded_returns = [] 

127 

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) 

134 

135 # Matrix: (S, T-1) 

136 returns_matrix = mx.array(np.stack(padded_returns)) 

137 

138 # 2. Get Signals (CPU parallelized) 

139 # We use ParallelExecutor to speed up signal generation for many scenarios 

140 signal_list = [] 

141 

142 # Prepare args for parallel execution 

143 # (indicator is pickleable usually) 

144 tasks = [(indicator, df) for df in scenarios] 

145 

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. 

150 

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) 

168 

169 for sigs in raw_signals: 

170 signal_list.append(normalize_signal_array(sigs, max_len)) 

171 

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

175 

176 # 3. Massive GPU calc 

177 strat_returns = signal_matrix_mx * returns_matrix 

178 

179 # Vectorized metrics 

180 equity_curves = mx.exp( 

181 mx.cumsum(mx.log1p(mx.clip(strat_returns, -0.9, 10.0)), axis=1) 

182 ) 

183 

184 final_rets = np.array(equity_curves[:, -1]) 

185 

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

190 

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

197 

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