fix: curiosity kernel stack + gradient zeroing + safe memcpy in tensor conversion

1. cuCtxSetLimit(STACK_SIZE, 8192) in curiosity trainer — prevents
   stack overflow in fused kernel (6 arrays of 42-128 floats/thread)

2. Removed in-kernel gradient zeroing (Phase 1) — had inter-block race
   where fast blocks atomicAdd while slow blocks still zero. Now uses
   host-side memset_zeros (GPU cuMemsetD8Async, stream-ordered)

3. cuda_slice_to_tensor_f32 now uses safe stream.memcpy_dtod() instead
   of raw device_ptr + memcpy_dtod_async (cudarc event tracking fix)

4. Debug traces in smoke test and training loop for deadlock diagnosis

Investigation ongoing: deadlock in init_gpu_experience_collector —
the collector constructor hangs during initialization, not during
experience collection or training.

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
This commit is contained in:
jgrusewski
2026-03-20 09:01:35 +01:00
parent 1f64ab20a3
commit a3c2fb4dd1
4 changed files with 58 additions and 53 deletions

View File

@@ -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;

View File

@@ -178,6 +178,20 @@ impl GpuCuriosityTrainer {
state_dim: usize,
max_samples: usize,
) -> Result<Self, MLError> {
// 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.

View File

@@ -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::<f32>(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::<f32>();
#[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}"))
}

View File

@@ -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");