diff --git a/crates/ml/src/cuda_pipeline/curiosity_training_kernel.cu b/crates/ml/src/cuda_pipeline/curiosity_training_kernel.cu index a8ff224da..b7e31cf40 100644 --- a/crates/ml/src/cuda_pipeline/curiosity_training_kernel.cu +++ b/crates/ml/src/cuda_pipeline/curiosity_training_kernel.cu @@ -298,42 +298,22 @@ extern "C" __global__ void curiosity_fused_zero_fwd_bwd_adam( int adam_step /* 1-based step counter */ ) { int tid = blockIdx.x * blockDim.x + threadIdx.x; - int grid_size = gridDim.x * blockDim.x; - /* ================================================================ */ - /* PHASE 1: Zero gradient buffers (grid-stride loop) */ - /* ================================================================ */ - int total_grad_elems = CUR_TOTAL_PARAMS; - /* Pack all 4 grad arrays into a single linear index space: - * [0, w1_len) -> grad_w1 - * [w1_len, w1_len+b1_len) -> grad_b1 - * [w1_len+b1_len, w1_len+b1_len+w2_len) -> grad_w2 - * [w1_len+b1_len+w2_len, total) -> grad_b2 */ + /* Gradient buffers are pre-zeroed by host-side memset_zeros on the same + * CUDA stream before this kernel launch. Stream ordering guarantees the + * memsets complete before this kernel starts. No in-kernel zeroing needed. + * + * The previous in-kernel Phase 1 zeroing had an inter-block race: + * __syncthreads() only syncs within a block, so fast blocks could start + * atomicAdd-ing gradients while slow blocks were still zeroing. */ + int w1_len = CUR_HIDDEN * CUR_INPUT; int b1_len = CUR_HIDDEN; int w2_len = CUR_OUTPUT * CUR_HIDDEN; - /* b2_len = CUR_OUTPUT (implicit) */ - - for (int i = tid; i < total_grad_elems; i += grid_size) { - if (i < w1_len) { - grad_w1[i] = 0.0f; - } else if (i < w1_len + b1_len) { - grad_b1[i - w1_len] = 0.0f; - } else if (i < w1_len + b1_len + w2_len) { - grad_w2[i - w1_len - b1_len] = 0.0f; - } else { - grad_b2[i - w1_len - b1_len - w2_len] = 0.0f; - } - } - - /* Ensure all gradient buffers are zeroed before any thread starts - * the forward/backward pass. __syncthreads() synchronizes within - * the block; __threadfence() ensures global memory visibility. */ - __threadfence(); - __syncthreads(); + int total_grad_elems = CUR_TOTAL_PARAMS; /* ================================================================ */ - /* PHASE 2: Forward + backward pass (one sample per thread) */ + /* PHASE 1: Forward + backward pass (one sample per thread) */ /* ================================================================ */ if (tid < N) { const float* state = states + tid * state_dim; diff --git a/crates/ml/src/cuda_pipeline/gpu_curiosity_trainer.rs b/crates/ml/src/cuda_pipeline/gpu_curiosity_trainer.rs index e93761d42..213adb390 100644 --- a/crates/ml/src/cuda_pipeline/gpu_curiosity_trainer.rs +++ b/crates/ml/src/cuda_pipeline/gpu_curiosity_trainer.rs @@ -178,6 +178,20 @@ impl GpuCuriosityTrainer { state_dim: usize, max_samples: usize, ) -> Result { + // The fused curiosity kernel allocates ~3KB/thread on stack (6 arrays of + // 42-128 floats: input[45], pre_act[128], hidden[128], pred[42], d_pred[42], + // d_hidden[128]). CUDA default stack is 1024B → stack overflow → async crash + // → stream poisoned → deadlock at next synchronize(). Set to 4096B. + { + use cudarc::driver::sys::{cuCtxSetLimit, CUlimit, CUresult}; + let result = unsafe { cuCtxSetLimit(CUlimit::CU_LIMIT_STACK_SIZE, 8192) }; + if result != CUresult::CUDA_SUCCESS { + return Err(MLError::ModelError(format!( + "cuCtxSetLimit(STACK_SIZE, 8192) failed: {result:?}" + ))); + } + } + // ---- Compile and load kernels ---- let ptx_result = CURIOSITY_TRAINING_PTX.get_or_init(|| compile_curiosity_training_ptx(&stream.context())); let ptx = ptx_result @@ -365,12 +379,20 @@ impl GpuCuriosityTrainer { let step = self.step; let bs_i32 = n_train as i32; - // Zero the block-arrival counter before launch. - self.stream - .memset_zeros(&mut self.block_counter) - .map_err(|e| { - MLError::ModelError(format!("memset block_counter: {e}")) - })?; + // Zero gradient buffers + block-arrival counter before fused kernel launch. + // Gradients are accumulated via atomicAdd in the kernel — they MUST start at zero. + // Previously this was done inside the kernel (Phase 1) but had an inter-block race. + // memset_zeros is GPU-side cuMemsetD8Async, ordered on the same stream. + self.stream.memset_zeros(&mut self.grad_w1) + .map_err(|e| MLError::ModelError(format!("memset grad_w1: {e}")))?; + self.stream.memset_zeros(&mut self.grad_b1) + .map_err(|e| MLError::ModelError(format!("memset grad_b1: {e}")))?; + self.stream.memset_zeros(&mut self.grad_w2) + .map_err(|e| MLError::ModelError(format!("memset grad_w2: {e}")))?; + self.stream.memset_zeros(&mut self.grad_b2) + .map_err(|e| MLError::ModelError(format!("memset grad_b2: {e}")))?; + self.stream.memset_zeros(&mut self.block_counter) + .map_err(|e| MLError::ModelError(format!("memset block_counter: {e}")))?; // Grid covers N training samples (one thread per sample for fwd/bwd). // The last block to finish also handles Adam update for all params. diff --git a/crates/ml/src/trainers/dqn/trainer/training_loop.rs b/crates/ml/src/trainers/dqn/trainer/training_loop.rs index 505061e3e..7be94ab44 100644 --- a/crates/ml/src/trainers/dqn/trainer/training_loop.rs +++ b/crates/ml/src/trainers/dqn/trainer/training_loop.rs @@ -20,7 +20,7 @@ use anyhow::{Context, Result}; use ml_core::device::MlDevice; use ml_core::cuda_autograd::GpuTensor; use cudarc; -use cudarc::driver::{CudaSlice, CudaStream, DevicePtr, DevicePtrMut}; +use cudarc::driver::{CudaSlice, CudaStream}; use common::metrics::{questdb_sink, training_metrics}; use tracing::{debug, info, warn}; @@ -93,14 +93,19 @@ impl DQNTrainer { let epoch_start = std::time::Instant::now(); // Phase 1/1b/1c: GPU data upload + experience collector init (once) + eprintln!("[TRAIN] init_gpu_data"); self.init_gpu_data(training_data).await?; + eprintln!("[TRAIN] init_gpu_raw_buffers"); self.init_gpu_raw_buffers(training_data).await?; + eprintln!("[TRAIN] init_gpu_experience_collector"); self.init_gpu_experience_collector().await?; + eprintln!("[TRAIN] collect_gpu_experiences"); // Phase 3: GPU experience collection let gpu_experiences_collected = self.collect_gpu_experiences( training_data, ).await?; + eprintln!("[TRAIN] collected={}", gpu_experiences_collected); // Free DqnGpuData when GPU experience collector is active if gpu_experiences_collected && self.gpu_data.is_some() { @@ -1992,21 +1997,12 @@ fn cuda_slice_to_tensor_f32( let n_elems: usize = shape.iter().product(); let mut dst = stream.alloc_zeros::(n_elems) .map_err(|e| anyhow::anyhow!("alloc f32: {e}"))?; - { - let src_view = src.slice(..n_elems); - let (src_ptr, _sg) = src_view.device_ptr(stream); - let (dst_ptr, _dg) = dst.device_ptr_mut(stream); - let num_bytes = n_elems * 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| anyhow::anyhow!("DtoD f32: {e}"))?; - } - } + let src_view = src.slice(..n_elems); + stream.memcpy_dtod(&src_view, &mut dst) + .map_err(|e| anyhow::anyhow!("DtoD f32: {e}"))?; stream.synchronize() .map_err(|e| anyhow::anyhow!("stream sync: {e}"))?; - GpuTensor::new(dst, shape.to_vec()) // cpu-side shape clone + GpuTensor::new(dst, shape.to_vec()) .map_err(|e| anyhow::anyhow!("GpuTensor::new f32: {e}")) } diff --git a/crates/ml/tests/smoke_test_real_data.rs b/crates/ml/tests/smoke_test_real_data.rs index 5e980c85b..4b5dff05d 100644 --- a/crates/ml/tests/smoke_test_real_data.rs +++ b/crates/ml/tests/smoke_test_real_data.rs @@ -773,12 +773,17 @@ mod gpu_smoke { async fn smoke_e2e_dqn_training_loop() { use ml::trainers::dqn::{DQNHyperparameters, DQNTrainer}; - // Smoke test: 204K bars + batch=32 + 5 epochs ~ 75 MB GPU. 2 GB is sufficient. + eprintln!("[E2E] vram={}MB", gpu_vram_mb()); if gpu_vram_mb() > 0 && gpu_vram_mb() < 2048 { - warn!(vram = gpu_vram_mb(), min = 2048, "SKIP: insufficient GPU VRAM for E2E smoke test"); + eprintln!("[E2E] SKIP: insufficient VRAM"); return; } - let Some(ohlcv_dir) = try_data("ohlcv") else { return }; + let Some(ohlcv_dir) = try_data("ohlcv") else { + eprintln!("[E2E] SKIP: no test data at ohlcv"); + return; + }; + eprintln!("[E2E] data_dir={}", ohlcv_dir.display()); + eprintln!("[E2E] creating DQNTrainer..."); let data_dir = ohlcv_dir.to_string_lossy().to_string(); // Configure for a fast smoke run: few epochs, small batch, low warmup @@ -808,7 +813,9 @@ async fn smoke_e2e_dqn_training_loop() { hyperparams.curiosity_weight = 0.0; let checkpoint_dir = tempfile::tempdir().expect("Failed to create temp dir"); + eprintln!("[E2E] DQNTrainer::new()..."); let mut trainer = DQNTrainer::new(hyperparams).expect("Failed to create DQN trainer"); + eprintln!("[E2E] DQNTrainer created, calling train()..."); info!(data_dir, "Starting E2E DQN training smoke test");