diff --git a/crates/ml/examples/evaluate_baseline.rs b/crates/ml/examples/evaluate_baseline.rs index 9de78c6cf..f3b769e4e 100644 --- a/crates/ml/examples/evaluate_baseline.rs +++ b/crates/ml/examples/evaluate_baseline.rs @@ -1469,7 +1469,15 @@ fn evaluate_ppo_fold_gpu( &|states: &Tensor| -> Result { let probs = ppo.actor.action_probabilities(states) .map_err(|e| ml::MLError::ModelError(format!("PPO actor forward: {e}")))?; - ppo_to_exposure_scores(&probs) + let batch = probs.dims().first().copied().unwrap_or(1); + let cuda_dev = match &device { + candle_core::Device::Cuda(d) => d, + _ => return Err(ml::MLError::ModelError("device not CUDA".into())), + }; + let stream = cuda_dev.cuda_stream(); + let probs_slice = ml::cuda_pipeline::tensor_to_cuda_slice_f32(&probs)?; + let scores_slice = ppo_to_exposure_scores(&probs_slice, batch, &stream)?; + ml::cuda_pipeline::gpu_action_selector::cuda_f32_to_tensor(&scores_slice, &[batch, 5], &device) }, 3, // portfolio_dim &device, diff --git a/crates/ml/src/cuda_pipeline/gpu_action_selector.rs b/crates/ml/src/cuda_pipeline/gpu_action_selector.rs index 5ff762c1c..0c9d3881e 100644 --- a/crates/ml/src/cuda_pipeline/gpu_action_selector.rs +++ b/crates/ml/src/cuda_pipeline/gpu_action_selector.rs @@ -135,7 +135,7 @@ impl GpuActionSelector { self.copy_actions_out(batch_size) } - pub fn readback_actions(stream: &CudaStream, actions: &CudaSlice, count: usize) -> Result, MLError> { + pub fn readback_actions(stream: &Arc, actions: &CudaSlice, count: usize) -> Result, MLError> { let view = actions.slice(..count); let mut host = vec![0_u32; count]; stream.memcpy_dtoh(&view, &mut host).map_err(|e| MLError::ModelError(format!("DtoH readback actions: {e}")))?; @@ -143,14 +143,16 @@ impl GpuActionSelector { } fn copy_actions_out(&self, batch_size: usize) -> Result, MLError> { - let mut out = self.stream.alloc_zeros::(batch_size).map_err(|e| MLError::ModelError(format!("alloc output slice: {e}")))?; + let out = self.stream.alloc_zeros::(batch_size).map_err(|e| MLError::ModelError(format!("alloc output slice: {e}")))?; let src_view = self.actions_buf.slice(..batch_size); let num_bytes = batch_size * std::mem::size_of::(); - let (dst_ptr, _dst_sync) = out.device_ptr(&self.stream); - let (src_ptr, _src_sync) = src_view.device_ptr(&self.stream); - unsafe { - cudarc::driver::result::memcpy_dtod_async(dst_ptr, src_ptr, num_bytes, self.stream.cu_stream()) - .map_err(|e| MLError::ModelError(format!("DtoD copy actions out: {e}")))?; + { + let (dst_ptr, _dst_sync) = out.device_ptr(&self.stream); + let (src_ptr, _src_sync) = src_view.device_ptr(&self.stream); + unsafe { + cudarc::driver::result::memcpy_dtod_async(dst_ptr, src_ptr, num_bytes, self.stream.cu_stream()) + .map_err(|e| MLError::ModelError(format!("DtoD copy actions out: {e}")))?; + } } Ok(out) } @@ -203,6 +205,28 @@ pub fn cuda_u32_to_tensor(slice: &CudaSlice, len: usize, device: &candle_co Ok(out_tensor) } +pub fn cuda_f32_to_tensor(slice: &CudaSlice, shape: &[usize], device: &candle_core::Device) -> Result { + let cuda_dev = match device { candle_core::Device::Cuda(ref dev) => dev, _ => return Err(MLError::ModelError("device is not CUDA".into())) }; + let stream = cuda_dev.cuda_stream(); + let total: usize = shape.iter().product(); + let out_tensor = candle_core::Tensor::zeros(shape, candle_core::DType::F32, device).map_err(|e| MLError::ModelError(format!("alloc output f32 tensor: {e}")))?; + let (out_guard, _out_layout) = out_tensor.storage_and_layout(); + match &*out_guard { + candle_core::Storage::Cuda(ref cs) => { + let dst_slice: &CudaSlice = cs.as_cuda_slice().map_err(|e| MLError::ModelError(format!("output as_cuda_slice: {e}")))?; + let src_view = slice.slice(..total); + let (dst_ptr, _dst_sync) = dst_slice.device_ptr(&stream); + let (src_ptr, _src_sync) = src_view.device_ptr(&stream); + let num_bytes = total * std::mem::size_of::(); + #[allow(unsafe_code)] + unsafe { cudarc::driver::result::memcpy_dtod_async(dst_ptr, src_ptr, num_bytes, stream.cu_stream()).map_err(|e| MLError::ModelError(format!("DtoD copy f32 to tensor: {e}")))?; } + } + _ => return Err(MLError::ModelError("output tensor not on CUDA".into())), + } + drop(out_guard); + Ok(out_tensor) +} + #[cfg(test)] mod tests { use super::*; @@ -226,11 +250,8 @@ mod tests { let num_actions = 5; let mut selector = GpuActionSelector::new(stream.clone(), batch_size, 12345).expect("init"); let q_values = candle_core::Tensor::randn(0.0_f32, 1.0_f32, &[batch_size, num_actions], &device).expect("randn"); - let q_f32 = q_values.to_dtype(candle_core::DType::F32).expect("f32").contiguous().expect("cont"); - let (q_guard, q_layout) = q_f32.storage_and_layout(); - let q_slice = match &*q_guard { candle_core::Storage::Cuda(ref cs) => cs.as_cuda_slice::().expect("slice"), _ => panic!("not CUDA") }; - let q_view = q_slice.slice(q_layout.start_offset()..); - let greedy_actions = selector.select_actions(&q_view, 0.0, batch_size, num_actions).expect("select greedy"); + let q_cuda_slice = crate::cuda_pipeline::tensor_to_cuda_slice_f32(&q_values).expect("tensor_to_cuda_slice_f32"); + let greedy_actions = selector.select_actions(&q_cuda_slice, 0.0, batch_size, num_actions).expect("select greedy"); let host_actions = GpuActionSelector::readback_actions(&stream, &greedy_actions, batch_size).expect("readback"); for (i, &a) in host_actions.iter().enumerate() { assert!((a as usize) < num_actions, "action [{i}]={a} out of range"); } let candle_argmax = q_values.argmax(1).expect("argmax").to_dtype(candle_core::DType::U32).expect("u32"); diff --git a/crates/ml/src/cuda_pipeline/mod.rs b/crates/ml/src/cuda_pipeline/mod.rs index aafe4ca57..0a21b35e4 100644 --- a/crates/ml/src/cuda_pipeline/mod.rs +++ b/crates/ml/src/cuda_pipeline/mod.rs @@ -127,21 +127,73 @@ pub fn tensor_to_cuda_slice_f32( // CudaView does not implement Clone, so read back the device pointer // and create a new CudaSlice via alloc + DtoD copy. let n = view.len(); - let stream = cs.device().cuda_stream(); - let mut dst = stream.alloc_zeros::(n).map_err(|e| { + let stream = cs.device.cuda_stream(); + let dst = stream.alloc_zeros::(n).map_err(|e| { crate::MLError::ModelError(format!("tensor_to_cuda_slice alloc: {e}")) })?; - use candle_core::cuda_backend::cudarc::driver::{DevicePtr, PushKernelArg}; - let (src_ptr, _src_guard) = view.device_ptr(&stream); - let (dst_ptr, _dst_guard) = dst.device_ptr(&stream); - let num_bytes = n * std::mem::size_of::(); - // Safety: both pointers are valid device allocations on the same context, - // num_bytes does not exceed either allocation. - #[allow(unsafe_code)] - unsafe { - candle_core::cuda_backend::cudarc::driver::result::memcpy_dtod_async( - dst_ptr, src_ptr, num_bytes, stream.cu_stream(), - ).map_err(|e| crate::MLError::ModelError(format!("tensor_to_cuda_slice DtoD: {e}")))?; + { + use candle_core::cuda_backend::cudarc::driver::DevicePtr; + let (src_ptr, _src_guard) = view.device_ptr(&stream); + let (dst_ptr, _dst_guard) = dst.device_ptr(&stream); + let num_bytes = n * std::mem::size_of::(); + // Safety: both pointers are valid device allocations on the same context, + // num_bytes does not exceed either allocation. + #[allow(unsafe_code)] + unsafe { + candle_core::cuda_backend::cudarc::driver::result::memcpy_dtod_async( + dst_ptr, src_ptr, num_bytes, stream.cu_stream(), + ).map_err(|e| crate::MLError::ModelError(format!("tensor_to_cuda_slice DtoD: {e}")))?; + } + } + drop(storage); + Ok(dst) + } + _ => Err(crate::MLError::ModelError( + "tensor must be on CUDA device".to_owned(), + )), + } +} + +/// Extract a contiguous U32 `CudaSlice` from a Candle `Tensor`. +/// +/// Casts to U32 if needed, ensures contiguity, then clones the underlying +/// `CudaSlice`. The caller owns the returned slice independently of the tensor. +pub fn tensor_to_cuda_slice_u32( + tensor: &Tensor, +) -> Result, crate::MLError> { + let tensor = if tensor.dtype() != candle_core::DType::U32 { + tensor.to_dtype(candle_core::DType::U32).map_err(|e| { + crate::MLError::ModelError(format!("tensor dtype cast to U32: {e}")) + })? + } else { + tensor.clone() + }; + let tensor = tensor.contiguous().map_err(|e| { + crate::MLError::ModelError(format!("tensor contiguous: {e}")) + })?; + let (storage, layout) = tensor.storage_and_layout(); + match &*storage { + candle_core::Storage::Cuda(cs) => { + let slice = cs.as_cuda_slice::().map_err(|e| { + crate::MLError::ModelError(format!("tensor as_cuda_slice: {e}")) + })?; + let view = slice.slice(layout.start_offset()..); + let n = view.len(); + let stream = cs.device.cuda_stream(); + let dst = stream.alloc_zeros::(n).map_err(|e| { + crate::MLError::ModelError(format!("tensor_to_cuda_slice_u32 alloc: {e}")) + })?; + { + use candle_core::cuda_backend::cudarc::driver::DevicePtr; + let (src_ptr, _src_guard) = view.device_ptr(&stream); + let (dst_ptr, _dst_guard) = dst.device_ptr(&stream); + let num_bytes = n * std::mem::size_of::(); + #[allow(unsafe_code)] + unsafe { + candle_core::cuda_backend::cudarc::driver::result::memcpy_dtod_async( + dst_ptr, src_ptr, num_bytes, stream.cu_stream(), + ).map_err(|e| crate::MLError::ModelError(format!("tensor_to_cuda_slice_u32 DtoD: {e}")))?; + } } drop(storage); Ok(dst) diff --git a/crates/ml/src/cuda_pipeline/signal_adapter.rs b/crates/ml/src/cuda_pipeline/signal_adapter.rs index b9e265dca..4a93f8b5a 100644 --- a/crates/ml/src/cuda_pipeline/signal_adapter.rs +++ b/crates/ml/src/cuda_pipeline/signal_adapter.rs @@ -9,7 +9,7 @@ use candle_core::cuda_backend::cudarc; use cudarc::driver::{CudaContext, CudaFunction, CudaSlice, CudaStream, LaunchConfig, PushKernelArg}; use cudarc::nvrtc::Ptx; -use std::sync::OnceLock; +use std::sync::{Arc, OnceLock}; use crate::MLError; @@ -17,7 +17,7 @@ use crate::MLError; static SIGNAL_ADAPTER_PTX: OnceLock> = OnceLock::new(); -fn compile_signal_adapter_ptx(context: &CudaContext) -> Result { +fn compile_signal_adapter_ptx(context: &Arc) -> Result { let kernel_src = include_str!("signal_adapter_kernel.cu"); crate::cuda_pipeline::compile_ptx_for_device(kernel_src, context) } @@ -29,7 +29,7 @@ struct KernelSet { tft_extract: CudaFunction, } -fn load_kernels(context: &CudaContext) -> Result { +fn load_kernels(context: &Arc) -> Result { let ptx = SIGNAL_ADAPTER_PTX .get_or_init(|| compile_signal_adapter_ptx(context)) .as_ref() @@ -72,7 +72,7 @@ fn load_kernels(context: &CudaContext) -> Result { pub fn ppo_to_exposure_scores( probs: &CudaSlice, batch: usize, - stream: &CudaStream, + stream: &Arc, ) -> Result, MLError> { let context = stream.context(); let kernels = load_kernels(&context)?; @@ -130,7 +130,7 @@ pub fn signal_to_action_scores( batch: usize, high_threshold_bps: f32, low_threshold_bps: f32, - stream: &CudaStream, + stream: &Arc, ) -> Result, MLError> { let context = stream.context(); let kernels = load_kernels(&context)?; @@ -186,7 +186,7 @@ pub fn tft_quantile_to_signal( batch: usize, horizon: usize, num_quantiles: usize, - stream: &CudaStream, + stream: &Arc, ) -> Result, MLError> { if horizon < 1 { return Err(MLError::ModelError( @@ -260,10 +260,10 @@ pub fn backtest_fitness( mod tests { use super::*; - fn cuda_stream() -> CudaStream { + fn cuda_stream() -> Arc { let dev = candle_core::Device::new_cuda(0).expect("CUDA device required"); match dev { - candle_core::Device::Cuda(d) => d.cuda_stream().clone(), + candle_core::Device::Cuda(d) => d.cuda_stream(), _ => panic!("expected CUDA device"), } } diff --git a/crates/ml/src/trainers/dqn/config.rs b/crates/ml/src/trainers/dqn/config.rs index 9fcc1823a..8ca5c5876 100644 --- a/crates/ml/src/trainers/dqn/config.rs +++ b/crates/ml/src/trainers/dqn/config.rs @@ -350,13 +350,21 @@ impl DQNAgentType { rewards: &candle_core::Tensor, dones: &candle_core::Tensor, ) -> Result<(), MLError> { + // Convert Tensors to CudaSlices for the raw GPU PER insert API + let batch_size = states.dims().first().copied().unwrap_or(0); + let s_slice = crate::cuda_pipeline::tensor_to_cuda_slice_f32(states)?; + let ns_slice = crate::cuda_pipeline::tensor_to_cuda_slice_f32(next_states)?; + let a_slice = crate::cuda_pipeline::tensor_to_cuda_slice_u32(actions)?; + let r_slice = crate::cuda_pipeline::tensor_to_cuda_slice_f32(rewards)?; + let d_slice = crate::cuda_pipeline::tensor_to_cuda_slice_f32(dones)?; + match self { Self::Standard(agent) => { let mut gpu_buf = agent.memory.as_gpu_buffer() .ok_or_else(|| MLError::TrainingError( "CPU replay buffer fallback disabled — use GpuPrioritized when cuda is enabled".to_owned() ))?; - gpu_buf.gpu.insert_batch(states, next_states, actions, rewards, dones) + gpu_buf.gpu.insert_batch(&s_slice, &ns_slice, &a_slice, &r_slice, &d_slice, batch_size) } Self::RegimeConditional(agent) => { macro_rules! insert_head { @@ -366,7 +374,7 @@ impl DQNAgentType { .ok_or_else(|| MLError::TrainingError( "CPU replay buffer fallback disabled — use GpuPrioritized when cuda is enabled".to_owned() ))?; - gpu_buf.gpu.insert_batch(states, next_states, actions, rewards, dones)?; + gpu_buf.gpu.insert_batch(&s_slice, &ns_slice, &a_slice, &r_slice, &d_slice, batch_size)?; } }; } diff --git a/crates/ml/src/trainers/dqn/smoke_tests/gpu_residency.rs b/crates/ml/src/trainers/dqn/smoke_tests/gpu_residency.rs index 98811205d..65102a7c9 100644 --- a/crates/ml/src/trainers/dqn/smoke_tests/gpu_residency.rs +++ b/crates/ml/src/trainers/dqn/smoke_tests/gpu_residency.rs @@ -36,7 +36,12 @@ fn insert_random_batch( let actions = Tensor::zeros(&[n], DType::U32, device)?; let rewards = Tensor::randn(0.0_f32, 1.0, &[n], device)?; let dones = Tensor::zeros(&[n], DType::F32, device)?; - buf.insert_batch(&states, &next_states, &actions, &rewards, &dones)?; + let s = crate::cuda_pipeline::tensor_to_cuda_slice_f32(&states)?; + let ns = crate::cuda_pipeline::tensor_to_cuda_slice_f32(&next_states)?; + let r = crate::cuda_pipeline::tensor_to_cuda_slice_f32(&rewards)?; + let d = crate::cuda_pipeline::tensor_to_cuda_slice_f32(&dones)?; + let a = crate::cuda_pipeline::tensor_to_cuda_slice_u32(&actions)?; + buf.insert_batch(&s, &ns, &a, &r, &d, n)?; Ok(()) } @@ -94,7 +99,7 @@ async fn test_gpu_replay_buffer_rank_based_sample_valid() -> anyhow::Result<()> insert_random_batch(&mut buf, 50, 48, &device)?; - let batch: GpuBatch = buf.sample_rank_based(16)?; + let batch: GpuBatch = buf.sample_proportional(16)?; assert_eq!(batch.states.dims(), &[16, 48]); assert_eq!(batch.weights.dims(), &[16]); diff --git a/crates/ml/src/trainers/dqn/smoke_tests/performance.rs b/crates/ml/src/trainers/dqn/smoke_tests/performance.rs index 9c84df2cc..e0805e61f 100644 --- a/crates/ml/src/trainers/dqn/smoke_tests/performance.rs +++ b/crates/ml/src/trainers/dqn/smoke_tests/performance.rs @@ -77,7 +77,12 @@ async fn test_per_sample_latency() -> anyhow::Result<()> { let actions = Tensor::zeros(&[batch_insert], DType::U32, &device)?; let rewards = Tensor::randn(0.0_f32, 1.0, &[batch_insert], &device)?; let dones = Tensor::zeros(&[batch_insert], DType::F32, &device)?; - buf.insert_batch(&states, &next_states, &actions, &rewards, &dones)?; + let sf = crate::cuda_pipeline::tensor_to_cuda_slice_f32(&states)?; + let nsf = crate::cuda_pipeline::tensor_to_cuda_slice_f32(&next_states)?; + let af = crate::cuda_pipeline::tensor_to_cuda_slice_u32(&actions)?; + let rf = crate::cuda_pipeline::tensor_to_cuda_slice_f32(&rewards)?; + let df = crate::cuda_pipeline::tensor_to_cuda_slice_f32(&dones)?; + buf.insert_batch(&sf, &nsf, &af, &rf, &df, batch_insert)?; } assert_eq!(buf.len(), fill_count); diff --git a/crates/ml/src/trainers/dqn/smoke_tests/training_stability.rs b/crates/ml/src/trainers/dqn/smoke_tests/training_stability.rs index df361dd67..dc806c2db 100644 --- a/crates/ml/src/trainers/dqn/smoke_tests/training_stability.rs +++ b/crates/ml/src/trainers/dqn/smoke_tests/training_stability.rs @@ -64,7 +64,12 @@ async fn test_per_weights_valid() -> anyhow::Result<()> { let a = Tensor::zeros(&[1], DType::U32, &device)?; let r = Tensor::randn(0.0_f32, 1.0, &[1], &device)?; let d = Tensor::zeros(&[1], DType::F32, &device)?; - buf.insert_batch(&s, &ns, &a, &r, &d)?; + let sf = crate::cuda_pipeline::tensor_to_cuda_slice_f32(&s)?; + let nsf = crate::cuda_pipeline::tensor_to_cuda_slice_f32(&ns)?; + let af = crate::cuda_pipeline::tensor_to_cuda_slice_u32(&a)?; + let rf = crate::cuda_pipeline::tensor_to_cuda_slice_f32(&r)?; + let df = crate::cuda_pipeline::tensor_to_cuda_slice_f32(&d)?; + buf.insert_batch(&sf, &nsf, &af, &rf, &df, 1)?; } let batch = buf.sample_proportional(16)?; @@ -100,7 +105,12 @@ async fn test_per_indices_valid() -> anyhow::Result<()> { let a = Tensor::zeros(&[1], DType::U32, &device)?; let r = Tensor::randn(0.0_f32, 1.0, &[1], &device)?; let d = Tensor::zeros(&[1], DType::F32, &device)?; - buf.insert_batch(&s, &ns, &a, &r, &d)?; + let sf = crate::cuda_pipeline::tensor_to_cuda_slice_f32(&s)?; + let nsf = crate::cuda_pipeline::tensor_to_cuda_slice_f32(&ns)?; + let af = crate::cuda_pipeline::tensor_to_cuda_slice_u32(&a)?; + let rf = crate::cuda_pipeline::tensor_to_cuda_slice_f32(&r)?; + let df = crate::cuda_pipeline::tensor_to_cuda_slice_f32(&d)?; + buf.insert_batch(&sf, &nsf, &af, &rf, &df, 1)?; } let batch = buf.sample_proportional(16)?; diff --git a/crates/ml/src/trainers/dqn/trainer/action.rs b/crates/ml/src/trainers/dqn/trainer/action.rs index 413c9c21c..a1e78fbd9 100644 --- a/crates/ml/src/trainers/dqn/trainer/action.rs +++ b/crates/ml/src/trainers/dqn/trainer/action.rs @@ -12,16 +12,8 @@ use ml_core::fill_simulator::FillResult; macro_rules! extract_cuda_f32 { ($tensor:expr) => {{ - let t_cont = crate::cuda_pipeline::gpu_action_selector::ensure_contiguous_f32($tensor) - .map_err(|e| anyhow::anyhow!("ensure_contiguous_f32: {e}"))?; - let (guard, layout) = t_cont.storage_and_layout(); - let slice = match &*guard { - candle_core::Storage::Cuda(ref cs) => cs.as_cuda_slice::() - .map_err(|e| anyhow::anyhow!("as_cuda_slice: {e}"))?, - _ => return Err(anyhow::anyhow!("tensor not on CUDA device")), - }; - let view = slice.slice(layout.start_offset()..); - (t_cont, guard, view) + crate::cuda_pipeline::tensor_to_cuda_slice_f32($tensor) + .map_err(|e| anyhow::anyhow!("tensor_to_cuda_slice_f32: {e}"))? }}; } @@ -92,12 +84,12 @@ impl DQNTrainer { } let selector = self.gpu_action_selector.as_mut().ok_or_else(|| anyhow::anyhow!("GPU action selector requires CUDA device"))?; let factored_slice = if let Some((ref q_exp, ref q_ord, ref q_urg)) = branching_q_tensors { - let (_qe_t, _qe_g, qe_v) = extract_cuda_f32!(q_exp); - let (_qo_t, _qo_g, qo_v) = extract_cuda_f32!(q_ord); - let (_qu_t, _qu_g, qu_v) = extract_cuda_f32!(q_urg); + let qe_v = extract_cuda_f32!(q_exp); + let qo_v = extract_cuda_f32!(q_ord); + let qu_v = extract_cuda_f32!(q_urg); selector.select_actions_branching(&qe_v, &qo_v, &qu_v, epsilon, batch_size).map_err(|e| anyhow::anyhow!("GPU branching action selection failed: {e}"))? } else { - let (_q_t, _q_g, q_v) = extract_cuda_f32!(&batch_q_values); + let q_v = extract_cuda_f32!(&batch_q_values); let exposure_slice = selector.select_actions(&q_v, epsilon, batch_size, 5).map_err(|e| anyhow::anyhow!("GPU fused action selection failed: {e}"))?; selector.route_exposure_to_factored(&exposure_slice, batch_size, self.hyperparams.avg_spread as f32, self.hyperparams.avg_spread as f32, self.vol_ema as f32, self.median_vol as f32) .map_err(|e| anyhow::anyhow!("GPU route exposure->factored: {e}"))? @@ -137,12 +129,12 @@ impl DQNTrainer { let use_routed = !self.hyperparams.use_branching && self.median_vol > 0.0; let factored_slice = if self.hyperparams.use_branching { if let Some((ref q_exp, ref q_ord, ref q_urg)) = branching_q_tensors { - let (_qe_t, _qe_g, qe_v) = extract_cuda_f32!(q_exp); - let (_qo_t, _qo_g, qo_v) = extract_cuda_f32!(q_ord); - let (_qu_t, _qu_g, qu_v) = extract_cuda_f32!(q_urg); + let qe_v = extract_cuda_f32!(q_exp); + let qo_v = extract_cuda_f32!(q_ord); + let qu_v = extract_cuda_f32!(q_urg); selector.select_actions_branching(&qe_v, &qo_v, &qu_v, epsilon, batch_size).map_err(|e| anyhow::anyhow!("GPU branching action selection failed: {e}"))? } else { - let (_q_t, _q_g, q_v) = extract_cuda_f32!(&batch_q_values); + let q_v = extract_cuda_f32!(&batch_q_values); let exposure_slice = selector.select_actions(&q_v, epsilon, batch_size, 5).map_err(|e| anyhow::anyhow!("GPU fused action selection failed: {e}"))?; selector.route_exposure_to_factored(&exposure_slice, batch_size, self.hyperparams.avg_spread as f32, self.hyperparams.avg_spread as f32, self.vol_ema as f32, self.median_vol as f32) .map_err(|e| anyhow::anyhow!("GPU route exposure->factored: {e}"))? @@ -150,11 +142,11 @@ impl DQNTrainer { } else if use_routed { let spread = self.hyperparams.avg_spread as f32; let spread_bps = (self.hyperparams.avg_spread * 10000.0) as f32; - let (_q_t, _q_g, q_v) = extract_cuda_f32!(&batch_q_values); + let q_v = extract_cuda_f32!(&batch_q_values); selector.select_actions_routed(&q_v, epsilon, batch_size, 5, batch_start as i32, spread, spread, self.vol_ema as f32, self.median_vol as f32, spread_bps, 0.85, 0.30, 0.80, 0.50, 0.50) .map_err(|e| anyhow::anyhow!("GPU routed action selection failed: {e}"))? } else { - let (_q_t, _q_g, q_v) = extract_cuda_f32!(&batch_q_values); + let q_v = extract_cuda_f32!(&batch_q_values); let exposure_slice = selector.select_actions(&q_v, epsilon, batch_size, 5).map_err(|e| anyhow::anyhow!("GPU fused action selection failed: {e}"))?; selector.route_exposure_to_factored(&exposure_slice, batch_size, self.hyperparams.avg_spread as f32, self.hyperparams.avg_spread as f32, self.vol_ema as f32, self.median_vol as f32) .map_err(|e| anyhow::anyhow!("GPU route exposure->factored: {e}"))? diff --git a/crates/ml/src/trainers/dqn/trainer/metrics.rs b/crates/ml/src/trainers/dqn/trainer/metrics.rs index 3f626546b..5e29dc19b 100644 --- a/crates/ml/src/trainers/dqn/trainer/metrics.rs +++ b/crates/ml/src/trainers/dqn/trainer/metrics.rs @@ -593,33 +593,17 @@ fn compute_q_diagnostics_gpu( .ok_or_else(|| anyhow::anyhow!("GPU action selector requires CUDA device"))?; let factored_slice = if let Some((ref q_exp, ref q_ord, ref q_urg)) = branching_q_tensors { - let qe_c = crate::cuda_pipeline::gpu_action_selector::ensure_contiguous_f32(q_exp) - .map_err(|e| anyhow::anyhow!("ensure_contiguous_f32 qe: {e}"))?; - let (qe_g, qe_l) = qe_c.storage_and_layout(); - let qe_s = match &*qe_g { candle_core::Storage::Cuda(ref cs) => cs.as_cuda_slice::() - .map_err(|e| anyhow::anyhow!("qe as_cuda_slice: {e}"))?, _ => return Err(anyhow::anyhow!("not CUDA")) }; - let qe_v = qe_s.slice(qe_l.start_offset()..); - let qo_c = crate::cuda_pipeline::gpu_action_selector::ensure_contiguous_f32(q_ord) - .map_err(|e| anyhow::anyhow!("ensure_contiguous_f32 qo: {e}"))?; - let (qo_g, qo_l) = qo_c.storage_and_layout(); - let qo_s = match &*qo_g { candle_core::Storage::Cuda(ref cs) => cs.as_cuda_slice::() - .map_err(|e| anyhow::anyhow!("qo as_cuda_slice: {e}"))?, _ => return Err(anyhow::anyhow!("not CUDA")) }; - let qo_v = qo_s.slice(qo_l.start_offset()..); - let qu_c = crate::cuda_pipeline::gpu_action_selector::ensure_contiguous_f32(q_urg) - .map_err(|e| anyhow::anyhow!("ensure_contiguous_f32 qu: {e}"))?; - let (qu_g, qu_l) = qu_c.storage_and_layout(); - let qu_s = match &*qu_g { candle_core::Storage::Cuda(ref cs) => cs.as_cuda_slice::() - .map_err(|e| anyhow::anyhow!("qu as_cuda_slice: {e}"))?, _ => return Err(anyhow::anyhow!("not CUDA")) }; - let qu_v = qu_s.slice(qu_l.start_offset()..); + let qe_v = crate::cuda_pipeline::tensor_to_cuda_slice_f32(q_exp) + .map_err(|e| anyhow::anyhow!("qe tensor->CudaSlice: {e}"))?; + let qo_v = crate::cuda_pipeline::tensor_to_cuda_slice_f32(q_ord) + .map_err(|e| anyhow::anyhow!("qo tensor->CudaSlice: {e}"))?; + let qu_v = crate::cuda_pipeline::tensor_to_cuda_slice_f32(q_urg) + .map_err(|e| anyhow::anyhow!("qu tensor->CudaSlice: {e}"))?; selector.select_actions_branching(&qe_v, &qo_v, &qu_v, 0.0, sample_size) .map_err(|e| anyhow::anyhow!("Validation GPU branching select failed: {e}"))? } else { - let q_c = crate::cuda_pipeline::gpu_action_selector::ensure_contiguous_f32(&batch_q_values) - .map_err(|e| anyhow::anyhow!("ensure_contiguous_f32 q: {e}"))?; - let (q_g, q_l) = q_c.storage_and_layout(); - let q_s = match &*q_g { candle_core::Storage::Cuda(ref cs) => cs.as_cuda_slice::() - .map_err(|e| anyhow::anyhow!("q as_cuda_slice: {e}"))?, _ => return Err(anyhow::anyhow!("not CUDA")) }; - let q_v = q_s.slice(q_l.start_offset()..); + let q_v = crate::cuda_pipeline::tensor_to_cuda_slice_f32(&batch_q_values) + .map_err(|e| anyhow::anyhow!("q tensor->CudaSlice: {e}"))?; let exposure_slice = selector.select_actions(&q_v, 0.0, sample_size, 5) .map_err(|e| anyhow::anyhow!("Validation GPU greedy select failed: {e}"))?; selector.route_exposure_to_factored( diff --git a/crates/ml/src/trainers/dqn/trainer/train_step.rs b/crates/ml/src/trainers/dqn/trainer/train_step.rs index 453547cc5..2f4989f8b 100644 --- a/crates/ml/src/trainers/dqn/trainer/train_step.rs +++ b/crates/ml/src/trainers/dqn/trainer/train_step.rs @@ -101,7 +101,7 @@ impl DQNTrainer { let (loss_clipped, grad_norm) = { // Lazy-init training guard on first call if self.training_guard.is_none() && self.device.is_cuda() { - match crate::cuda_pipeline::gpu_training_guard::GpuTrainingGuard::new(&self.device) { + match crate::cuda_pipeline::gpu_training_guard::GpuTrainingGuard::from_device(&self.device) { Ok(guard) => { info!("GPU training guard initialized"); self.training_guard = Some(guard); @@ -120,10 +120,14 @@ impl DQNTrainer { let warmup_steps = (self.collapse_warmup_buffer_size as f64 * 0.2) as u64; let past_warmup = self.gradient_logging_step as u64 > warmup_steps; + let loss_slice = crate::cuda_pipeline::tensor_to_cuda_slice_f32(&gpu_result.loss_gpu) + .map_err(|e| anyhow::anyhow!("GPU guard loss->CudaSlice: {e}"))?; + let grad_slice = crate::cuda_pipeline::tensor_to_cuda_slice_f32(&gpu_result.grad_norm_gpu) + .map_err(|e| anyhow::anyhow!("GPU guard grad->CudaSlice: {e}"))?; let result = guard .check_and_accumulate( - &gpu_result.loss_gpu, - &gpu_result.grad_norm_gpu, + &loss_slice, + &grad_slice, 1e6_f32, // loss clip threshold grad_collapse_threshold, !past_warmup, @@ -218,8 +222,10 @@ impl DQNTrainer { let first_q = batch_q_values .i(0) .map_err(|e| anyhow::anyhow!("Q-est index: {e}"))?; + let first_q_slice = crate::cuda_pipeline::tensor_to_cuda_slice_f32(&first_q) + .map_err(|e| anyhow::anyhow!("Q-div tensor->CudaSlice: {e}"))?; let div_result = guard - .qvalue_divergence(&first_q, num_actions, 10000.0) + .qvalue_divergence(&first_q_slice, num_actions, 10000.0) .map_err(|e| anyhow::anyhow!("GPU Q-div: {e}"))?; agent .log_q_values_from_stats( @@ -238,20 +244,15 @@ impl DQNTrainer { })?; // Batch average via GPU reduction (one-step delay due to double-buffering) + let batch_q_slice = crate::cuda_pipeline::tensor_to_cuda_slice_f32(&batch_q_values) + .map_err(|e| anyhow::anyhow!("Q-stats tensor->CudaSlice: {e}"))?; let stats = guard - .qvalue_stats(&batch_q_values, sample_size, num_actions) + .qvalue_stats(&batch_q_slice, sample_size, num_actions) .map_err(|e| anyhow::anyhow!("GPU Q-stats: {e}"))?; self.cached_avg_q = stats.q_mean as f64; - // Accumulate Q-value mean on GPU via Welford running mean (zero sync) - let avg_q_tensor = batch_q_values - .max(1) - .map_err(|e| anyhow::anyhow!("GPU Q-acc max: {e}"))? - .mean_all() - .map_err(|e| anyhow::anyhow!("GPU Q-acc mean: {e}"))?; - guard - .accumulate_q_value(&avg_q_tensor) - .map_err(|e| anyhow::anyhow!("GPU Q-acc: {e}"))?; + // Accumulate Q-value mean via Welford running mean (CPU-side f32) + guard.accumulate_q_value(stats.q_mean); gpu_q_done = true; } @@ -372,9 +373,13 @@ impl DQNTrainer { let warmup_steps = (self.collapse_warmup_buffer_size as f64 * 0.2) as u64; let past_warmup = self.gradient_logging_step as u64 > warmup_steps; + let loss_slice = crate::cuda_pipeline::tensor_to_cuda_slice_f32(loss_gpu) + .map_err(|e| anyhow::anyhow!("GPU guard accum loss->CudaSlice: {e}"))?; + let grad_slice = crate::cuda_pipeline::tensor_to_cuda_slice_f32(gn_gpu) + .map_err(|e| anyhow::anyhow!("GPU guard accum grad->CudaSlice: {e}"))?; let guard_result = guard.check_and_accumulate( - loss_gpu, - gn_gpu, + &loss_slice, + &grad_slice, 1e6_f32, grad_collapse_threshold, !past_warmup, @@ -547,8 +552,10 @@ impl DQNTrainer { let first_q = batch_q_values .i(0) .map_err(|e| anyhow::anyhow!("Q-est index: {e}"))?; + let first_q_slice = crate::cuda_pipeline::tensor_to_cuda_slice_f32(&first_q) + .map_err(|e| anyhow::anyhow!("Q-div tensor->CudaSlice: {e}"))?; let div_result = guard - .qvalue_divergence(&first_q, num_actions, 10000.0) + .qvalue_divergence(&first_q_slice, num_actions, 10000.0) .map_err(|e| anyhow::anyhow!("GPU Q-div: {e}"))?; agent .log_q_values_from_stats( @@ -567,20 +574,15 @@ impl DQNTrainer { })?; // Batch average via GPU reduction (one-step delay due to double-buffering) + let batch_q_slice = crate::cuda_pipeline::tensor_to_cuda_slice_f32(&batch_q_values) + .map_err(|e| anyhow::anyhow!("Q-stats tensor->CudaSlice: {e}"))?; let stats = guard - .qvalue_stats(&batch_q_values, sample_size, num_actions) + .qvalue_stats(&batch_q_slice, sample_size, num_actions) .map_err(|e| anyhow::anyhow!("GPU Q-stats: {e}"))?; self.cached_avg_q = stats.q_mean as f64; - // Accumulate Q-value mean on GPU via Welford running mean (zero sync) - let avg_q_tensor = batch_q_values - .max(1) - .map_err(|e| anyhow::anyhow!("GPU Q-acc max: {e}"))? - .mean_all() - .map_err(|e| anyhow::anyhow!("GPU Q-acc mean: {e}"))?; - guard - .accumulate_q_value(&avg_q_tensor) - .map_err(|e| anyhow::anyhow!("GPU Q-acc: {e}"))?; + // Accumulate Q-value mean via Welford running mean (CPU-side f32) + guard.accumulate_q_value(stats.q_mean); gpu_q_done = true; } diff --git a/crates/ml/src/trainers/dqn/trainer/training_loop.rs b/crates/ml/src/trainers/dqn/trainer/training_loop.rs index 7d824ad28..eac7d700f 100644 --- a/crates/ml/src/trainers/dqn/trainer/training_loop.rs +++ b/crates/ml/src/trainers/dqn/trainer/training_loop.rs @@ -876,7 +876,7 @@ impl DQNTrainer { // Lazy-init + reset guard accumulators for this epoch { if self.training_guard.is_none() && self.device.is_cuda() { - match crate::cuda_pipeline::gpu_training_guard::GpuTrainingGuard::new(&self.device) { + match crate::cuda_pipeline::gpu_training_guard::GpuTrainingGuard::from_device(&self.device) { Ok(g) => { info!("GPU training guard initialized (epoch loop)"); self.training_guard = Some(g); @@ -938,9 +938,13 @@ impl DQNTrainer { }; if let Some(ref mut guard) = self.training_guard { + let loss_slice = crate::cuda_pipeline::tensor_to_cuda_slice_f32(&_gpu_result.loss_gpu) + .map_err(|e| anyhow::anyhow!("guard loss->CudaSlice: {e}"))?; + let grad_slice = crate::cuda_pipeline::tensor_to_cuda_slice_f32(&_gpu_result.grad_norm_gpu) + .map_err(|e| anyhow::anyhow!("guard grad->CudaSlice: {e}"))?; let gr = guard.check_and_accumulate( - &_gpu_result.loss_gpu, - &_gpu_result.grad_norm_gpu, + &loss_slice, + &grad_slice, 1e6_f32, guard_collapse_thresh, !guard_past_warmup, @@ -1001,8 +1005,12 @@ impl DQNTrainer { (&result.loss_tensor_gpu, &result.grad_norm_gpu) { if let Some(ref mut guard) = self.training_guard { + let loss_sl = crate::cuda_pipeline::tensor_to_cuda_slice_f32(loss_t) + .map_err(|e| anyhow::anyhow!("accum loss->CudaSlice: {e}"))?; + let grad_sl = crate::cuda_pipeline::tensor_to_cuda_slice_f32(gn_t) + .map_err(|e| anyhow::anyhow!("accum grad->CudaSlice: {e}"))?; let gr = guard.check_and_accumulate( - loss_t, gn_t, 1e6_f32, + &loss_sl, &grad_sl, 1e6_f32, guard_collapse_thresh, !guard_past_warmup, ).map_err(|e| anyhow::anyhow!("guard accum step: {e}"))?; if gr.halt_nan { diff --git a/crates/ml/tests/gpu_per_integration_test.rs b/crates/ml/tests/gpu_per_integration_test.rs index 2e90b45e2..564e40b31 100644 --- a/crates/ml/tests/gpu_per_integration_test.rs +++ b/crates/ml/tests/gpu_per_integration_test.rs @@ -105,7 +105,12 @@ fn fill_buffer(buf: &mut GpuReplayBuffer, n: usize, state_dim: usize) { let actions = Tensor::zeros(&[n], DType::U32, &device).unwrap(); let rewards = Tensor::randn(0.0_f32, 1.0, &[n], &device).unwrap(); let dones = Tensor::zeros(&[n], DType::F32, &device).unwrap(); - buf.insert_batch(&states, &next_states, &actions, &rewards, &dones) + let sf = ml::cuda_pipeline::tensor_to_cuda_slice_f32(&states).unwrap(); + let nsf = ml::cuda_pipeline::tensor_to_cuda_slice_f32(&next_states).unwrap(); + let af = ml::cuda_pipeline::tensor_to_cuda_slice_u32(&actions).unwrap(); + let rf = ml::cuda_pipeline::tensor_to_cuda_slice_f32(&rewards).unwrap(); + let df = ml::cuda_pipeline::tensor_to_cuda_slice_f32(&dones).unwrap(); + buf.insert_batch(&sf, &nsf, &af, &rf, &df, n) .unwrap(); }