From fa74f9afcfef05c6b8e18c590705eddc9820ceea Mon Sep 17 00:00:00 2001 From: jgrusewski Date: Sat, 28 Feb 2026 14:41:54 +0100 Subject: [PATCH] fix(cuda): add missing gpu_portfolio.rs and experience_kernels.cu files These files were part of the Phase 2 cuda_pipeline but were untracked and not included in previous commits. Required for compilation with --features ml/cuda. Co-Authored-By: Claude Opus 4.6 --- .../src/cuda_pipeline/experience_kernels.cu | 232 ++++++++++++++++ crates/ml/src/cuda_pipeline/gpu_portfolio.rs | 254 ++++++++++++++++++ 2 files changed, 486 insertions(+) create mode 100644 crates/ml/src/cuda_pipeline/experience_kernels.cu create mode 100644 crates/ml/src/cuda_pipeline/gpu_portfolio.rs diff --git a/crates/ml/src/cuda_pipeline/experience_kernels.cu b/crates/ml/src/cuda_pipeline/experience_kernels.cu new file mode 100644 index 000000000..725f97afb --- /dev/null +++ b/crates/ml/src/cuda_pipeline/experience_kernels.cu @@ -0,0 +1,232 @@ +/** + * CUDA Kernels for DQN Experience Collection + * + * Replaces CPU-bound per-bar portfolio simulation and reward + * computation in the DQN trainer hot loop. + * + * Architecture: One thread processes batch_size bars sequentially. + * Parallelism across hyperopt trials via separate CUDA stream launches. + */ + +/** + * Map action index (0-44) to target exposure fraction. + * + * FactoredAction: 5 exposures × 3 order types × 3 urgencies = 45 + * Exposure levels: Short100=-1.0, Short50=-0.5, Flat=0.0, Long50=0.5, Long100=1.0 + * Exposure is encoded as action_index / 9 (integer division). + */ +__device__ __forceinline__ float action_to_exposure(int action_idx) { + int exposure_level = action_idx / 9; // 0-4 + switch (exposure_level) { + case 0: return -1.0f; // Short100 + case 1: return -0.5f; // Short50 + case 2: return 0.0f; // Flat + case 3: return 0.5f; // Long50 + case 4: return 1.0f; // Long100 + default: return 0.0f; + } +} + +/** + * Map action index to transaction cost rate. + * + * OrderType encoded as (action_index / 3) % 3: + * 0 = Market (0.0001), 1 = LimitMaker (0.00005), 2 = IoC (0.00015) + */ +__device__ __forceinline__ float action_to_tx_cost(int action_idx) { + int order_type = (action_idx / 3) % 3; + switch (order_type) { + case 0: return 0.0001f; // Market + case 1: return 0.00005f; // LimitMaker + case 2: return 0.00015f; // IoC + default: return 0.0001f; + } +} + +/** + * Portfolio simulation kernel — processes batch_size bars sequentially. + * + * For each bar: + * 1. Read action → target exposure + * 2. Execute trade (update cash, position) + * 3. Mark-to-market at next price + * 4. Compute normalized portfolio features + * 5. Compute raw PnL reward + * 6. Check episode boundary + * + * @param targets [total_bars, 4] — [preproc_close, preproc_next, raw_close, raw_next] + * @param actions [batch_size] — action indices (0-44) + * @param portfolio_state [8] — IN/OUT: cash, position, entry_price, initial_capital, + * spread, last_price, reserve_pct, cum_costs + * @param portfolio_out [batch_size, 3] — normalized portfolio features per bar + * @param rewards_out [batch_size] — raw PnL reward per bar + * @param done_out [batch_size] — 1 if episode boundary, 0 otherwise + * @param batch_start First bar index in targets array + * @param batch_size Number of bars in this batch + * @param max_position Maximum position size (contracts) + * @param episode_length Bars per episode for time-based termination + * @param total_bars Total training bars (for data boundary check) + */ +extern "C" __global__ void portfolio_sim_kernel( + const float* __restrict__ targets, + const int* __restrict__ actions, + float* portfolio_state, + float* portfolio_out, + float* rewards_out, + int* done_out, + int batch_start, + int batch_size, + float max_position, + int episode_length, + int total_bars +) { + // Single thread — sequential processing + if (threadIdx.x != 0 || blockIdx.x != 0) return; + + // Load portfolio state into registers + float cash = portfolio_state[0]; + float position = portfolio_state[1]; + float entry_price = portfolio_state[2]; + float initial_cap = portfolio_state[3]; + float spread = portfolio_state[4]; + float last_price = portfolio_state[5]; + float reserve_pct = portfolio_state[6]; + float cum_costs = portfolio_state[7]; + + for (int b = 0; b < batch_size; b++) { + int global_idx = batch_start + b; + int action_idx = actions[b]; + + // Read prices from targets[global_idx, :] + int t_offset = global_idx * 4; + float current_close = targets[t_offset + 0]; // preprocessed + float next_close = targets[t_offset + 1]; // preprocessed + float current_close_raw = targets[t_offset + 2]; // raw (for barrier) + float next_close_raw = targets[t_offset + 3]; // raw + + // Use raw prices if available, otherwise fall back to preprocessed + float price = (current_close_raw != 0.0f) ? current_close_raw : current_close; + if (price <= 0.0f) price = 1.0f; // Safety + + // Capture current portfolio value BEFORE trade + float current_value = cash + position * price; + float current_norm = current_value / initial_cap; + + // --- Execute action --- + float target_exposure = action_to_exposure(action_idx); + float target_position = target_exposure * max_position; + float tx_rate = action_to_tx_cost(action_idx); + + // Detect reversal (sign change) + bool is_reversal = (position > 0.0f && target_position < 0.0f) || + (position < 0.0f && target_position > 0.0f); + + if (is_reversal) { + // Phase 1: Close current position + float close_cash = position * price; // + for long, - for short + float close_cost = fabsf(position) * price * tx_rate; + cash += close_cash - close_cost; + cum_costs += close_cost; + + // Phase 2: Open opposite position (may be partial) + float reserve = (reserve_pct > 0.0f) ? current_value * (reserve_pct / 100.0f) : 0.0f; + float affordable = fmaxf(cash - reserve, 0.0f); + float max_contracts = (price > 0.0f) ? affordable / (price * (1.0f + tx_rate)) : 0.0f; + max_contracts = floorf(max_contracts); + float actual = fminf(max_contracts, fabsf(target_position)); + + if (actual > 0.0f) { + float new_pos = (target_position > 0.0f) ? actual : -actual; + float open_cost = actual * price * tx_rate; + cash -= new_pos * price + open_cost; + cum_costs += open_cost; + position = new_pos; + entry_price = price; + } else { + position = 0.0f; + entry_price = 0.0f; + } + } else { + // Non-reversal: adjust position directly + float delta = target_position - position; + if (fabsf(delta) > 0.0f) { + float trade_cost = fabsf(delta) * price * tx_rate; + cum_costs += trade_cost; + cash -= trade_cost; + + // Cash reserve check for buys + if (delta > 0.0f && reserve_pct > 0.0f) { + float pv = cash + position * price; + float reserve = pv * (reserve_pct / 100.0f); + float buy_cost = delta * price; + if (cash - buy_cost < reserve) { + float affordable = fmaxf(cash - reserve, 0.0f); + delta = fminf(delta, floorf(affordable / price)); + } + } + + if (delta > 0.0f) { + entry_price = price; + } else if (target_position == 0.0f) { + entry_price = 0.0f; + } + cash -= delta * price; + position = position + delta; + } + } + + last_price = price; + + // --- Mark-to-market at next price --- + float next_price = (next_close_raw != 0.0f) ? next_close_raw : next_close; + if (next_price <= 0.0f) next_price = price; + float next_value = cash + position * next_price; + float next_norm = next_value / initial_cap; + + // --- Normalized portfolio features --- + float max_pos_norm = (price > 0.0f) ? initial_cap / price : 1.0f; + float pos_norm = position / max_pos_norm; + + portfolio_out[b * 3 + 0] = next_norm; // normalized value + portfolio_out[b * 3 + 1] = pos_norm; // normalized position + portfolio_out[b * 3 + 2] = spread; // spread + + // --- Raw PnL reward --- + float reward = 0.0f; + if (current_norm > 0.0f) { + reward = (next_norm - current_norm) / current_norm; + } + + // Risk penalty: penalize excessive position (>80% exposure) + float abs_pos = fabsf(pos_norm); + if (abs_pos > 0.8f) { + reward -= (abs_pos - 0.8f) * 5.0f * 0.1f; // risk_weight=0.1 default + } + + rewards_out[b] = reward; + + // --- Episode boundary --- + int step = global_idx + 1; + int time_done = (step % episode_length == 0) ? 1 : 0; + int data_done = (step >= total_bars) ? 1 : 0; + done_out[b] = (time_done || data_done) ? 1 : 0; + + // Reset portfolio on episode boundary + if (done_out[b]) { + cash = initial_cap; + position = 0.0f; + entry_price = 0.0f; + cum_costs = 0.0f; + } + } + + // Write back portfolio state + portfolio_state[0] = cash; + portfolio_state[1] = position; + portfolio_state[2] = entry_price; + portfolio_state[3] = initial_cap; + portfolio_state[4] = spread; + portfolio_state[5] = last_price; + portfolio_state[6] = reserve_pct; + portfolio_state[7] = cum_costs; +} diff --git a/crates/ml/src/cuda_pipeline/gpu_portfolio.rs b/crates/ml/src/cuda_pipeline/gpu_portfolio.rs new file mode 100644 index 000000000..3be126dbf --- /dev/null +++ b/crates/ml/src/cuda_pipeline/gpu_portfolio.rs @@ -0,0 +1,254 @@ +#![allow(unsafe_code)] // Required for CUDA kernel launches + +//! GPU-accelerated portfolio simulation for DQN training. +//! +//! Compiles CUDA kernels at runtime via NVRTC and launches them +//! to replace the CPU-bound inner loop in the DQN trainer. +//! +//! Uses cudarc 0.17 API (re-exported by candle_core::cuda_backend::cudarc): +//! - CudaContext / CudaStream (not CudaDevice — that's the older API) +//! - stream.launch_builder(&func).arg(&buf).launch(config) +//! - stream.alloc_zeros / memcpy_htod / memcpy_dtoh + +use std::sync::Arc; + +use candle_core::cuda_backend::cudarc; +use cudarc::driver::{CudaContext, CudaFunction, CudaSlice, CudaStream, LaunchConfig, PushKernelArg}; +use cudarc::nvrtc::Ptx; +use tracing::{debug, info}; + +use crate::MLError; + +/// Maximum batch size for GPU portfolio simulation +const MAX_BATCH_SIZE: usize = 256; + +/// Number of f32 elements in portfolio state +const PORTFOLIO_STATE_SIZE: usize = 8; + +/// GPU-accelerated portfolio simulator for DQN training. +/// +/// Replaces per-bar `PortfolioTracker::execute_action()` + reward +/// calculation with a CUDA kernel that processes an entire batch +/// sequentially on GPU. +#[allow(missing_debug_implementations)] // CudaFunction/CudaSlice don't impl Debug +pub struct GpuPortfolioSimulator { + stream: Arc, + kernel_func: CudaFunction, + + // GPU buffers (persistent across batches) + actions_buf: CudaSlice, + portfolio_state_buf: CudaSlice, + portfolio_out_buf: CudaSlice, + rewards_out_buf: CudaSlice, + done_out_buf: CudaSlice, + + // Config + max_position: f32, + episode_length: i32, + total_bars: i32, +} + +/// Output from a portfolio simulation batch +#[derive(Debug)] +pub struct PortfolioSimResult { + /// Normalized portfolio features [batch_size, 3]: (value, position, spread) + pub portfolio_features: Vec, + /// Raw PnL rewards [batch_size] + pub rewards: Vec, + /// Episode boundary flags [batch_size]: 1 = done, 0 = continue + pub done_flags: Vec, + /// Number of bars processed + pub batch_size: usize, +} + +fn compile_experience_kernels(context: &Arc) -> Result { + let kernel_src = include_str!("experience_kernels.cu"); + let ptx: Ptx = cudarc::nvrtc::compile_ptx(kernel_src).map_err(|e| { + MLError::ModelError(format!( + "CUDA experience kernel compilation failed (is NVRTC available?): {e}" + )) + })?; + let module = context.load_module(ptx).map_err(|e| { + MLError::ModelError(format!("Failed to load experience kernel module: {e}")) + })?; + module.load_function("portfolio_sim_kernel").map_err(|e| { + MLError::ModelError(format!("Failed to load portfolio_sim_kernel: {e}")) + }) +} + +impl GpuPortfolioSimulator { + /// Create a new GPU portfolio simulator. + /// + /// Compiles CUDA kernels via NVRTC and allocates persistent GPU buffers. + pub fn new( + stream: Arc, + initial_capital: f32, + avg_spread: f32, + cash_reserve_pct: f32, + max_position: f32, + episode_length: usize, + total_bars: usize, + ) -> Result { + // Compile and load CUDA kernel + let kernel_func = compile_experience_kernels(stream.context())?; + info!("GPU portfolio simulator: CUDA kernels compiled and loaded"); + + // Allocate persistent buffers + let actions_buf = stream + .alloc_zeros::(MAX_BATCH_SIZE) + .map_err(|e| MLError::ModelError(format!("Failed to alloc actions buffer: {e}")))?; + let mut portfolio_state_buf = stream + .alloc_zeros::(PORTFOLIO_STATE_SIZE) + .map_err(|e| MLError::ModelError(format!("Failed to alloc portfolio state: {e}")))?; + let portfolio_out_buf = stream + .alloc_zeros::(MAX_BATCH_SIZE * 3) + .map_err(|e| MLError::ModelError(format!("Failed to alloc portfolio output: {e}")))?; + let rewards_out_buf = stream + .alloc_zeros::(MAX_BATCH_SIZE) + .map_err(|e| MLError::ModelError(format!("Failed to alloc rewards output: {e}")))?; + let done_out_buf = stream + .alloc_zeros::(MAX_BATCH_SIZE) + .map_err(|e| MLError::ModelError(format!("Failed to alloc done output: {e}")))?; + + // Initialize portfolio state: [cash, position, entry_price, initial_cap, spread, last_price, reserve_pct, cum_costs] + let init_state = [ + initial_capital, + 0.0_f32, + 0.0, + initial_capital, + avg_spread, + 0.0, + cash_reserve_pct, + 0.0, + ]; + stream + .memcpy_htod(&init_state, &mut portfolio_state_buf) + .map_err(|e| MLError::ModelError(format!("Failed to upload portfolio state: {e}")))?; + + debug!( + "GPU portfolio sim initialized: capital={}, spread={}, reserve={}%, max_pos={}, episode_len={}, total_bars={}", + initial_capital, avg_spread, cash_reserve_pct, max_position, episode_length, total_bars + ); + + Ok(Self { + stream, + kernel_func, + actions_buf, + portfolio_state_buf, + portfolio_out_buf, + rewards_out_buf, + done_out_buf, + max_position, + episode_length: episode_length as i32, + total_bars: total_bars as i32, + }) + } + + /// Run portfolio simulation for a batch of bars. + /// + /// # Arguments + /// * `targets_buf` - Pre-uploaded targets tensor buffer [total_bars, 4] + /// * `actions` - Action indices (0-44) for this batch + /// * `batch_start` - Global bar index of first bar in batch + pub fn simulate_batch( + &mut self, + targets_buf: &CudaSlice, + actions: &[i32], + batch_start: usize, + ) -> Result { + let batch_size = actions.len(); + if batch_size == 0 { + return Ok(PortfolioSimResult { + portfolio_features: vec![], + rewards: vec![], + done_flags: vec![], + batch_size: 0, + }); + } + if batch_size > MAX_BATCH_SIZE { + return Err(MLError::ModelError(format!( + "Batch size {} exceeds max {}", + batch_size, MAX_BATCH_SIZE + ))); + } + + // Upload actions to GPU + self.stream + .memcpy_htod(actions, &mut self.actions_buf) + .map_err(|e| MLError::ModelError(format!("Failed to upload actions: {e}")))?; + + let config = LaunchConfig { + grid_dim: (1, 1, 1), + block_dim: (1, 1, 1), + shared_mem_bytes: 0, + }; + + let batch_start_i32 = batch_start as i32; + let batch_size_i32 = batch_size as i32; + + // Safety: kernel parameters match the CUDA function signature exactly. + // All buffers are allocated with sufficient size (MAX_BATCH_SIZE). + // targets_buf is pre-uploaded with shape [total_bars, 4]. + unsafe { + self.stream + .launch_builder(&self.kernel_func) + .arg(targets_buf) + .arg(&self.actions_buf) + .arg(&mut self.portfolio_state_buf) + .arg(&mut self.portfolio_out_buf) + .arg(&mut self.rewards_out_buf) + .arg(&mut self.done_out_buf) + .arg(&batch_start_i32) + .arg(&batch_size_i32) + .arg(&self.max_position) + .arg(&self.episode_length) + .arg(&self.total_bars) + .launch(config) + .map_err(|e| MLError::ModelError(format!("Kernel launch failed: {e}")))?; + } + + // Download results (only batch_size elements, not full buffer) + let portfolio_view = self.portfolio_out_buf.slice(..batch_size * 3); + let rewards_view = self.rewards_out_buf.slice(..batch_size); + let done_view = self.done_out_buf.slice(..batch_size); + + let mut portfolio_features = vec![0.0_f32; batch_size * 3]; + let mut rewards = vec![0.0_f32; batch_size]; + let mut done_flags = vec![0_i32; batch_size]; + + self.stream + .memcpy_dtoh(&portfolio_view, &mut portfolio_features) + .map_err(|e| MLError::ModelError(format!("Failed to download portfolio features: {e}")))?; + self.stream + .memcpy_dtoh(&rewards_view, &mut rewards) + .map_err(|e| MLError::ModelError(format!("Failed to download rewards: {e}")))?; + self.stream + .memcpy_dtoh(&done_view, &mut done_flags) + .map_err(|e| MLError::ModelError(format!("Failed to download done flags: {e}")))?; + + Ok(PortfolioSimResult { + portfolio_features, + rewards, + done_flags, + batch_size, + }) + } + + /// Reset portfolio state to initial capital. + pub fn reset(&mut self, initial_capital: f32, avg_spread: f32, cash_reserve_pct: f32) -> Result<(), MLError> { + let init_state = [initial_capital, 0.0, 0.0, initial_capital, avg_spread, 0.0, cash_reserve_pct, 0.0_f32]; + self.stream + .memcpy_htod(&init_state, &mut self.portfolio_state_buf) + .map_err(|e| MLError::ModelError(format!("Failed to reset portfolio state: {e}")))?; + Ok(()) + } + + /// Get current portfolio state from GPU (for diagnostics). + pub fn get_portfolio_state(&self) -> Result<[f32; PORTFOLIO_STATE_SIZE], MLError> { + let mut state = [0.0_f32; PORTFOLIO_STATE_SIZE]; + self.stream + .memcpy_dtoh(&self.portfolio_state_buf, &mut state) + .map_err(|e| MLError::ModelError(format!("Failed to download portfolio state: {e}")))?; + Ok(state) + } +}