fix(cuda): resolve 35 caller-boundary type mismatches after CudaSlice migration
Update all callers to match the new pure-cudarc APIs introduced by the hive agent CudaSlice migration. Key changes: - GpuTrainingGuard::new() now takes Arc<CudaStream>; callers use from_device() - check_and_accumulate/qvalue_stats/qvalue_divergence take &CudaSlice<f32> instead of &Tensor; callers convert via tensor_to_cuda_slice_f32() - accumulate_q_value takes f32 scalar, returns () (no Result) - GpuReplayBuffer::insert_batch gains batch_size arg, takes CudaSlice params - signal_adapter functions take &Arc<CudaStream> (cudarc 0.17 Arc requirement) - Add tensor_to_cuda_slice_u32() and cuda_f32_to_tensor() utility functions - Replace CudaView usage with owned CudaSlice via tensor_to_cuda_slice_f32() - Fix CudaStorage.device field access (was method call in older API) - Fix borrow-after-move in copy_actions_out via scoped DtoD copy Zero errors, zero warnings across lib + tests + examples + full workspace. Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -1469,7 +1469,15 @@ fn evaluate_ppo_fold_gpu(
|
||||
&|states: &Tensor| -> Result<Tensor, ml::MLError> {
|
||||
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,
|
||||
|
||||
@@ -135,7 +135,7 @@ impl GpuActionSelector {
|
||||
self.copy_actions_out(batch_size)
|
||||
}
|
||||
|
||||
pub fn readback_actions(stream: &CudaStream, actions: &CudaSlice<u32>, count: usize) -> Result<Vec<u32>, MLError> {
|
||||
pub fn readback_actions(stream: &Arc<CudaStream>, actions: &CudaSlice<u32>, count: usize) -> Result<Vec<u32>, 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<CudaSlice<u32>, MLError> {
|
||||
let mut out = self.stream.alloc_zeros::<u32>(batch_size).map_err(|e| MLError::ModelError(format!("alloc output slice: {e}")))?;
|
||||
let out = self.stream.alloc_zeros::<u32>(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::<u32>();
|
||||
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<u32>, len: usize, device: &candle_co
|
||||
Ok(out_tensor)
|
||||
}
|
||||
|
||||
pub fn cuda_f32_to_tensor(slice: &CudaSlice<f32>, shape: &[usize], device: &candle_core::Device) -> Result<candle_core::Tensor, MLError> {
|
||||
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<f32> = 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::<f32>();
|
||||
#[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::<f32>().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");
|
||||
|
||||
@@ -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::<f32>(n).map_err(|e| {
|
||||
let stream = cs.device.cuda_stream();
|
||||
let dst = stream.alloc_zeros::<f32>(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::<f32>();
|
||||
// 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::<f32>();
|
||||
// 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<u32>` 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<candle_core::cuda_backend::cudarc::driver::CudaSlice<u32>, 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::<u32>().map_err(|e| {
|
||||
crate::MLError::ModelError(format!("tensor as_cuda_slice<u32>: {e}"))
|
||||
})?;
|
||||
let view = slice.slice(layout.start_offset()..);
|
||||
let n = view.len();
|
||||
let stream = cs.device.cuda_stream();
|
||||
let dst = stream.alloc_zeros::<u32>(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::<u32>();
|
||||
#[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)
|
||||
|
||||
@@ -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<Result<Ptx, String>> = OnceLock::new();
|
||||
|
||||
fn compile_signal_adapter_ptx(context: &CudaContext) -> Result<Ptx, String> {
|
||||
fn compile_signal_adapter_ptx(context: &Arc<CudaContext>) -> Result<Ptx, String> {
|
||||
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<KernelSet, MLError> {
|
||||
fn load_kernels(context: &Arc<CudaContext>) -> Result<KernelSet, MLError> {
|
||||
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<KernelSet, MLError> {
|
||||
pub fn ppo_to_exposure_scores(
|
||||
probs: &CudaSlice<f32>,
|
||||
batch: usize,
|
||||
stream: &CudaStream,
|
||||
stream: &Arc<CudaStream>,
|
||||
) -> Result<CudaSlice<f32>, 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<CudaStream>,
|
||||
) -> Result<CudaSlice<f32>, 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<CudaStream>,
|
||||
) -> Result<CudaSlice<f32>, 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<CudaStream> {
|
||||
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"),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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)?;
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
@@ -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]);
|
||||
|
||||
@@ -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);
|
||||
|
||||
|
||||
@@ -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)?;
|
||||
|
||||
@@ -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::<f32>()
|
||||
.map_err(|e| anyhow::anyhow!("as_cuda_slice<f32>: {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}"))?
|
||||
|
||||
@@ -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::<f32>()
|
||||
.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::<f32>()
|
||||
.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::<f32>()
|
||||
.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::<f32>()
|
||||
.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(
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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();
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user