From 8b3c209704d205ec1402b8bae2492108e4ece54f Mon Sep 17 00:00:00 2001 From: jgrusewski Date: Fri, 3 Apr 2026 02:09:40 +0200 Subject: [PATCH] =?UTF-8?q?revert:=20dead=20code=20removal=20caused=20NaN?= =?UTF-8?q?=20at=20step=2022=20=E2=80=94=20needs=20investigation?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Opus 4.6 (1M context) --- crates/ml/src/cuda_pipeline/double_buffer.rs | 58 +- crates/ml/src/cuda_pipeline/mod.rs | 229 +++- .../src/trainers/dqn/smoke_tests/hyperopt.rs | 126 +- .../dqn/smoke_tests/training_stability.rs | 22 + crates/ml/src/trainers/dqn/trainer/mod.rs | 367 +++++- crates/ml/src/trainers/dqn/trainer/tests.rs | 4 +- .../src/trainers/dqn/trainer/training_loop.rs | 1129 ++++++++++++++++- 7 files changed, 1820 insertions(+), 115 deletions(-) diff --git a/crates/ml/src/cuda_pipeline/double_buffer.rs b/crates/ml/src/cuda_pipeline/double_buffer.rs index f6c3d15ea..9759a7a2d 100644 --- a/crates/ml/src/cuda_pipeline/double_buffer.rs +++ b/crates/ml/src/cuda_pipeline/double_buffer.rs @@ -72,11 +72,9 @@ impl DoubleBufferedLoader { /// Upload initial data into the active slot. pub fn upload_initial( &mut self, - features: &[[f64; 42]], - targets: &[[f64; 4]], - ofi: &[[f64; 8]], + data: &[([f64; 42], Vec)], ) -> Result<(), MLError> { - let gpu_data = DqnGpuData::upload_slices(features, targets, ofi, &self.stream)?; + let gpu_data = DqnGpuData::upload(data, &self.stream)?; info!( "DoubleBuffer: initial upload -- {} bars, {:.1} MB VRAM", gpu_data.num_bars, @@ -94,11 +92,9 @@ impl DoubleBufferedLoader { /// [`swap()`](Self::swap), which calls it internally) completes. pub fn upload_to_staging( &mut self, - features: &[[f64; 42]], - targets: &[[f64; 4]], - ofi: &[[f64; 8]], + data: &[([f64; 42], Vec)], ) -> Result<(), MLError> { - let gpu_data = DqnGpuData::upload_slices(features, targets, ofi, &self.stream)?; + let gpu_data = DqnGpuData::upload(data, &self.stream)?; info!( "DoubleBuffer: staging upload -- {} bars, {:.1} MB VRAM", gpu_data.num_bars, @@ -221,14 +217,14 @@ mod tests { context.default_stream() } - fn make_features(n: usize) -> Vec<[f64; 42]> { - (0..n).map(|i| [i as f64 * 0.01; 42]).collect() - } - fn make_targets(n: usize) -> Vec<[f64; 4]> { - (0..n).map(|i| [100.0 + i as f64, 101.0, 100.5, 101.5]).collect() - } - fn make_ofi(n: usize) -> Vec<[f64; 8]> { - vec![[0.0; 8]; n] + fn make_data(n: usize) -> Vec<([f64; 42], Vec)> { + (0..n) + .map(|i| { + let features = [i as f64 * 0.01; 42]; + let targets = vec![100.0 + i as f64, 101.0, 100.5, 101.5]; + (features, targets) + }) + .collect() } #[test] @@ -236,7 +232,7 @@ mod tests { let mut loader = DoubleBufferedLoader::new(cuda_stream()); assert!(loader.active().is_none()); - loader.upload_initial(&make_features(50), &make_targets(50), &make_ofi(50)).ok(); + loader.upload_initial(&make_data(50)).ok(); // On CI without CUDA this will fail; on GPU it works } @@ -244,10 +240,10 @@ mod tests { fn test_staging_and_swap() { let stream = cuda_stream(); let mut loader = DoubleBufferedLoader::new(stream); - if loader.upload_initial(&make_features(100), &make_targets(100), &make_ofi(100)).is_err() { + if loader.upload_initial(&make_data(100)).is_err() { return; // skip on non-CUDA } - if loader.upload_to_staging(&make_features(200), &make_targets(200), &make_ofi(200)).is_err() { + if loader.upload_to_staging(&make_data(200)).is_err() { return; } @@ -262,7 +258,7 @@ mod tests { #[test] fn test_swap_without_staging_fails() { let mut loader = DoubleBufferedLoader::new(cuda_stream()); - if loader.upload_initial(&make_features(10), &make_targets(10), &make_ofi(10)).is_err() { + if loader.upload_initial(&make_data(10)).is_err() { return; } assert!(loader.swap().is_err()); @@ -271,13 +267,13 @@ mod tests { #[test] fn test_vram_tracking() { let mut loader = DoubleBufferedLoader::new(cuda_stream()); - if loader.upload_initial(&make_features(100), &make_targets(100), &make_ofi(100)).is_err() { + if loader.upload_initial(&make_data(100)).is_err() { return; } let single = loader.total_vram_bytes(); assert!(single > 0); - if loader.upload_to_staging(&make_features(100), &make_targets(100), &make_ofi(100)).is_err() { + if loader.upload_to_staging(&make_data(100)).is_err() { return; } let double = loader.total_vram_bytes(); @@ -287,10 +283,10 @@ mod tests { #[test] fn test_double_buffer_sync_staging() { let mut loader = DoubleBufferedLoader::new(cuda_stream()); - if loader.upload_initial(&make_features(50), &make_targets(50), &make_ofi(50)).is_err() { + if loader.upload_initial(&make_data(50)).is_err() { return; } - if loader.upload_to_staging(&make_features(100), &make_targets(100), &make_ofi(100)).is_err() { + if loader.upload_to_staging(&make_data(100)).is_err() { return; } @@ -305,10 +301,10 @@ mod tests { #[test] fn test_double_buffer_sync_staging_idempotent() { let mut loader = DoubleBufferedLoader::new(cuda_stream()); - if loader.upload_initial(&make_features(50), &make_targets(50), &make_ofi(50)).is_err() { + if loader.upload_initial(&make_data(50)).is_err() { return; } - if loader.upload_to_staging(&make_features(100), &make_targets(100), &make_ofi(100)).is_err() { + if loader.upload_to_staging(&make_data(100)).is_err() { return; } @@ -329,10 +325,10 @@ mod tests { #[test] fn test_double_buffer_swap_syncs_automatically() { let mut loader = DoubleBufferedLoader::new(cuda_stream()); - if loader.upload_initial(&make_features(50), &make_targets(50), &make_ofi(50)).is_err() { + if loader.upload_initial(&make_data(50)).is_err() { return; } - if loader.upload_to_staging(&make_features(100), &make_targets(100), &make_ofi(100)).is_err() { + if loader.upload_to_staging(&make_data(100)).is_err() { return; } @@ -354,12 +350,12 @@ mod tests { #[test] fn test_double_buffer_upload_resets_synced() { let mut loader = DoubleBufferedLoader::new(cuda_stream()); - if loader.upload_initial(&make_features(50), &make_targets(50), &make_ofi(50)).is_err() { + if loader.upload_initial(&make_data(50)).is_err() { return; } // First staging upload - if loader.upload_to_staging(&make_features(100), &make_targets(100), &make_ofi(100)).is_err() { + if loader.upload_to_staging(&make_data(100)).is_err() { return; } assert!(!loader.is_staging_synced()); @@ -368,7 +364,7 @@ mod tests { assert!(loader.is_staging_synced()); // Second staging upload resets the flag - if loader.upload_to_staging(&make_features(200), &make_targets(200), &make_ofi(200)).is_err() { + if loader.upload_to_staging(&make_data(200)).is_err() { return; } assert!(!loader.is_staging_synced()); diff --git a/crates/ml/src/cuda_pipeline/mod.rs b/crates/ml/src/cuda_pipeline/mod.rs index 4a9aa27ac..3f38271ca 100644 --- a/crates/ml/src/cuda_pipeline/mod.rs +++ b/crates/ml/src/cuda_pipeline/mod.rs @@ -246,6 +246,58 @@ impl std::fmt::Debug for DqnGpuData { } impl DqnGpuData { + /// Upload DQN training data to GPU as f32 CudaSlice buffers. + /// + /// Converts `[f64; 42]` features and `Vec` targets to f32, + /// flattens into contiguous arrays, and uploads once via `clone_htod`. + pub fn upload( + data: &[([f64; 42], Vec)], + stream: &Arc, + ) -> Result { + let num_bars = data.len(); + if num_bars == 0 { + return Err(MLError::ModelError("Empty training data".to_owned())); + } + + let feature_dim = 42; + let target_dim = 4; + + let estimated_bytes = estimate_vram_bytes(num_bars * (feature_dim + target_dim)); + if estimated_bytes > MAX_UPLOAD_BYTES { + return Err(MLError::ModelError(format!( + "Training data too large for GPU upload: {:.1} GB > 2.0 GB limit ({} bars)", + estimated_bytes as f64 / 1_073_741_824.0, + num_bars, + ))); + } + + let mut flat_features = Vec::with_capacity(num_bars * feature_dim); + for (features, _) in data { + for &v in features.iter() { + flat_features.push(v as f32); + } + } + + let mut flat_targets = Vec::with_capacity(num_bars * target_dim); + for (_, targets) in data { + for i in 0..target_dim { + flat_targets.push(targets.get(i).copied().unwrap_or(0.0) as f32); + } + } + + let features = clone_htod_f32_to_bf16(stream, &flat_features)?; + let targets = clone_htod_f32_to_bf16(stream, &flat_targets)?; + + Ok(Self { + features, + targets, + ofi_features: None, + num_bars, + feature_dim, + aligned_state_dim: None, + }) + } + /// Upload DQN training data to GPU from separate contiguous arrays. /// /// Accepts features, targets, and OFI as separate slices matching the fxcache @@ -585,8 +637,11 @@ impl DqnGpuData { /// /// # Usage /// ```ignore -/// let pool = GpuBufferPool::new(300_000, 51, 4); -/// // staging_bytes() returns total pre-allocated CPU memory +/// let mut pool = GpuBufferPool::new(300_000, 51, 4); +/// for fold in folds { +/// let gpu_data = pool.upload_dqn(&fold_data, &device)?; +/// // ... train on gpu_data ... +/// } /// ``` #[derive(Debug)] pub struct GpuBufferPool { @@ -614,6 +669,58 @@ impl GpuBufferPool { } } + /// Upload DQN data to GPU, reusing the pre-allocated staging buffers. + /// + /// Data is copied into the staging `Vec` (no heap alloc if within + /// `max_bars`), then uploaded to GPU in a single `Tensor::from_vec` call. + pub fn upload_dqn( + &mut self, + data: &[([f64; 42], Vec)], + stream: &Arc, + ) -> Result { + let num_bars = data.len(); + if num_bars == 0 { + return Err(MLError::ModelError("Empty training data".to_owned())); + } + + let feat_len = num_bars * self.feature_dim; + let targ_len = num_bars * self.target_dim; + + // Grow staging buffers only if the fold exceeds the pre-allocated size. + if feat_len > self.feature_buf.len() { + self.feature_buf.resize(feat_len, 0.0); + } + if targ_len > self.target_buf.len() { + self.target_buf.resize(targ_len, 0.0); + } + + // Fill staging buffers (zero-alloc copy). + for (i, (features, targets)) in data.iter().enumerate() { + let f_start = i * self.feature_dim; + for (j, &v) in features.iter().enumerate() { + self.feature_buf[f_start + j] = v as f32; + } + + let t_start = i * self.target_dim; + for j in 0..self.target_dim { + self.target_buf[t_start + j] = targets.get(j).copied().unwrap_or(0.0) as f32; + } + } + + // Upload the used slice to GPU via f32→bf16 conversion + clone_htod. + let features = clone_htod_f32_to_bf16(stream, &self.feature_buf[..feat_len])?; + let targets = clone_htod_f32_to_bf16(stream, &self.target_buf[..targ_len])?; + + Ok(DqnGpuData { + features, + targets, + ofi_features: None, + num_bars, + feature_dim: self.feature_dim, + aligned_state_dim: None, + }) + } + /// Total bytes of pre-allocated CPU staging memory. pub fn staging_bytes(&self) -> usize { (self.feature_buf.len() + self.target_buf.len()) * std::mem::size_of::() @@ -721,17 +828,17 @@ mod tests { } #[test] - fn test_dqn_gpu_data_upload_slices() { + fn test_dqn_gpu_data_upload() { let stream = cuda_stream(); - let features: Vec<[f64; 42]> = (0..100) - .map(|i| [i as f64 * 0.01; 42]) + let data: Vec<([f64; 42], Vec)> = (0..100) + .map(|i| { + let features = [i as f64 * 0.01; 42]; + let targets = vec![100.0 + i as f64, 101.0 + i as f64, 100.5, 101.5]; + (features, targets) + }) .collect(); - let targets: Vec<[f64; 4]> = (0..100) - .map(|i| [100.0 + i as f64, 101.0 + i as f64, 100.5, 101.5]) - .collect(); - let ofi: Vec<[f64; 8]> = vec![[0.0; 8]; 100]; - let gpu_data = DqnGpuData::upload_slices(&features, &targets, &ofi, &stream).expect("upload"); + let gpu_data = DqnGpuData::upload(&data, &stream).expect("upload"); assert_eq!(gpu_data.num_bars, 100); assert_eq!(gpu_data.feature_dim, 42); @@ -744,7 +851,7 @@ mod tests { // Download target values and verify on CPU (bf16 -> f32 conversion) let targets_gpu = gpu_data.bar_target_values(0, &stream).expect("targets"); let mut targets_bf16 = vec![half::bf16::ZERO; 4]; - stream.memcpy_dtoh(&targets_gpu, &mut targets_bf16).expect("DtoH"); + stream.memcpy_dtoh(&targets_gpu, &mut targets_bf16).expect("DtoH"); // test readback let targets_host: Vec = targets_bf16.iter().map(|v| v.to_f32()).collect(); assert!((targets_host.first().copied().unwrap_or(0.0) - 100.0).abs() < 0.5); } @@ -752,11 +859,11 @@ mod tests { #[test] fn test_dqn_build_state_tensor() { let stream = cuda_stream(); - let features = vec![[1.0_f64; 42]]; - let targets = vec![[100.0_f64, 101.0, 100.5, 101.5]]; - let ofi = vec![[0.0_f64; 8]]; + let data: Vec<([f64; 42], Vec)> = vec![ + ([1.0; 42], vec![100.0, 101.0, 100.5, 101.5]), + ]; - let gpu_data = DqnGpuData::upload_slices(&features, &targets, &ofi, &stream).expect("upload"); + let gpu_data = DqnGpuData::upload(&data, &stream).expect("upload"); let portfolio = [0.95_f32, 0.5, 0.0001]; let state = gpu_data.build_state_tensor(0, &portfolio, &stream).expect("state"); assert_eq!(state.len(), 45); // 42 market + 3 portfolio @@ -838,10 +945,8 @@ mod tests { #[test] fn test_dqn_gpu_data_empty() { let stream = cuda_stream(); - let features: Vec<[f64; 42]> = vec![]; - let targets: Vec<[f64; 4]> = vec![]; - let ofi: Vec<[f64; 8]> = vec![]; - let result = DqnGpuData::upload_slices(&features, &targets, &ofi, &stream); + let data: Vec<([f64; 42], Vec)> = vec![]; + let result = DqnGpuData::upload(&data, &stream); assert!(result.is_err()); } @@ -864,10 +969,10 @@ mod tests { #[test] fn test_dqn_boundary_access() { let stream = cuda_stream(); - let features = vec![[0.5_f64; 42]]; - let targets = vec![[100.0_f64, 101.0, 100.5, 101.5]]; - let ofi = vec![[0.0_f64; 8]]; - let gpu_data = DqnGpuData::upload_slices(&features, &targets, &ofi, &stream).expect("upload"); + let data: Vec<([f64; 42], Vec)> = vec![ + ([0.5; 42], vec![100.0, 101.0, 100.5, 101.5]), + ]; + let gpu_data = DqnGpuData::upload(&data, &stream).expect("upload"); assert!(gpu_data.bar_features(0, &stream).is_ok()); assert!(gpu_data.bar_features(1, &stream).is_err()); @@ -879,13 +984,15 @@ mod tests { #[test] fn test_dqn_build_batch_states() { let stream = cuda_stream(); - let features: Vec<[f64; 42]> = (0..100) - .map(|i| [i as f64 * 0.01; 42]) + let data: Vec<([f64; 42], Vec)> = (0..100) + .map(|i| { + let features = [i as f64 * 0.01; 42]; + let targets = vec![100.0, 101.0, 100.5, 101.5]; + (features, targets) + }) .collect(); - let targets: Vec<[f64; 4]> = vec![[100.0, 101.0, 100.5, 101.5]; 100]; - let ofi: Vec<[f64; 8]> = vec![[0.0; 8]; 100]; - let gpu_data = DqnGpuData::upload_slices(&features, &targets, &ofi, &stream).expect("upload"); + let gpu_data = DqnGpuData::upload(&data, &stream).expect("upload"); let portfolio = [0.95_f32, 0.5, 0.0001]; // Build batch of 32 states starting at bar 10 @@ -894,7 +1001,7 @@ mod tests { // Download and verify portfolio features in row 0 (bf16 -> f32 conversion) let mut host_bf16 = vec![half::bf16::ZERO; 32 * 45]; - stream.memcpy_dtoh(&batch, &mut host_bf16).expect("DtoH"); + stream.memcpy_dtoh(&batch, &mut host_bf16).expect("DtoH"); // test readback let host: Vec = host_bf16.iter().map(|v| v.to_f32()).collect(); // Row 0: features[42..45] should be portfolio [0.95, 0.5, 0.0001] @@ -915,10 +1022,11 @@ mod tests { #[test] fn test_dqn_build_batch_states_clamping() { let stream = cuda_stream(); - let features = vec![[1.0_f64; 42], [2.0_f64; 42]]; - let targets = vec![[100.0_f64, 101.0, 100.5, 101.5]; 2]; - let ofi = vec![[0.0_f64; 8]; 2]; - let gpu_data = DqnGpuData::upload_slices(&features, &targets, &ofi, &stream).expect("upload"); + let data: Vec<([f64; 42], Vec)> = vec![ + ([1.0; 42], vec![100.0, 101.0, 100.5, 101.5]), + ([2.0; 42], vec![100.0, 101.0, 100.5, 101.5]), + ]; + let gpu_data = DqnGpuData::upload(&data, &stream).expect("upload"); let portfolio = [1.0_f32, 0.0, 0.0001]; // Request more bars than available -> clamped to 2 @@ -933,9 +1041,62 @@ mod tests { #[test] fn test_gpu_buffer_pool_basic() { let stream = cuda_stream(); - let pool = GpuBufferPool::new(1000, 42, 4); + let mut pool = GpuBufferPool::new(1000, 42, 4); assert_eq!(pool.max_bars, 1000); assert_eq!(pool.staging_bytes(), (1000 * 42 + 1000 * 4) * 4); + + let data: Vec<([f64; 42], Vec)> = (0..50) + .map(|i| ([i as f64 * 0.01; 42], vec![100.0, 101.0, 100.5, 101.5])) + .collect(); + let gpu_data = pool.upload_dqn(&data, &stream).expect("upload"); + assert_eq!(gpu_data.num_bars, 50); + assert_eq!(gpu_data.feature_dim, 42); + } + + #[test] + fn test_gpu_buffer_pool_reuse() { + let stream = cuda_stream(); + let mut pool = GpuBufferPool::new(500, 42, 4); + + let data1: Vec<([f64; 42], Vec)> = (0..100) + .map(|i| ([i as f64 * 0.01; 42], vec![100.0, 101.0, 100.5, 101.5])) + .collect(); + let g1 = pool.upload_dqn(&data1, &stream).expect("upload1"); + assert_eq!(g1.num_bars, 100); + + let data2: Vec<([f64; 42], Vec)> = (0..200) + .map(|i| ([i as f64 * 0.02; 42], vec![200.0, 201.0, 200.5, 201.5])) + .collect(); + let g2 = pool.upload_dqn(&data2, &stream).expect("upload2"); + assert_eq!(g2.num_bars, 200); + + // Download target values and verify (bf16 -> f32 conversion) + let t = g2.bar_target_values(0, &stream).expect("targets"); + let mut host_bf16 = vec![half::bf16::ZERO; 4]; + stream.memcpy_dtoh(&t, &mut host_bf16).expect("DtoH"); // test readback + let host: Vec = host_bf16.iter().map(|v| v.to_f32()).collect(); + assert!((host[0] - 200.0).abs() < 1.0); + } + + #[test] + fn test_gpu_buffer_pool_grow() { + let stream = cuda_stream(); + let mut pool = GpuBufferPool::new(10, 42, 4); + + let data: Vec<([f64; 42], Vec)> = (0..50) + .map(|i| ([i as f64; 42], vec![1.0, 2.0, 3.0, 4.0])) + .collect(); + let gpu_data = pool.upload_dqn(&data, &stream).expect("upload"); + assert_eq!(gpu_data.num_bars, 50); + assert!(pool.feature_buf.len() >= 50 * 42); + } + + #[test] + fn test_gpu_buffer_pool_empty() { + let stream = cuda_stream(); + let mut pool = GpuBufferPool::new(100, 42, 4); + let data: Vec<([f64; 42], Vec)> = vec![]; + assert!(pool.upload_dqn(&data, &stream).is_err()); } #[test] diff --git a/crates/ml/src/trainers/dqn/smoke_tests/hyperopt.rs b/crates/ml/src/trainers/dqn/smoke_tests/hyperopt.rs index 7b774adfb..33afe5d59 100644 --- a/crates/ml/src/trainers/dqn/smoke_tests/hyperopt.rs +++ b/crates/ml/src/trainers/dqn/smoke_tests/hyperopt.rs @@ -1,25 +1,27 @@ //! Hyperopt training path smoke tests. //! -//! These tests exercise the fxcache-based training path used by the hyperopt -//! adapter: `init_from_fxcache` + `train_fold_from_slices`. -//! -//! The old `train_with_preloaded_data` and `train_with_shared_data` entry -//! points have been removed. The `trainer.train(&data_dir, ...)` convenience -//! wrapper now delegates through the slice-based loop internally. +//! These tests exercise the `train_with_preloaded_data` and `train_with_shared_data` +//! entry points which are used by the hyperopt adapter. Unlike `trainer.train()`, +//! these paths bypass GPU walk-forward and call `train_with_data_full_loop` directly. //! //! Key assertions: //! - Training completes without hang or NaN //! - Loss and gradient norms are finite +//! - Both preloaded (owned Vec) and shared (&slice) paths produce equivalent results use super::helpers::*; -/// Hyperopt path: load data from disk via `trainer.train(...)`, which now -/// uses `init_from_fxcache` + `train_fold_from_slices` internally. +/// Hyperopt preloaded path: load data once, pass owned Vecs to trainer. +/// +/// This is the primary hyperopt path — data is loaded once at trial start +/// and passed to each trial's trainer. Exercises train_with_data_full_loop +/// directly, bypassing walk-forward. #[test] #[ignore] // Loads real training data + GPU -fn test_hyperopt_train_via_data_dir() -> anyhow::Result<()> { - let data_dir = test_data_dir() - .expect("FOXHUNT_TEST_DATA or test_data/ must exist"); +fn test_hyperopt_preloaded_data() -> anyhow::Result<()> { + let (train, val) = load_smoke_data()?; + assert!(!train.is_empty(), "Training data must be non-empty"); + assert!(!val.is_empty(), "Validation data must be non-empty"); let p = smoke_params(); let mut trainer = smoke_trainer_with(p)?; @@ -27,24 +29,112 @@ fn test_hyperopt_train_via_data_dir() -> anyhow::Result<()> { .enable_all() .build()?; - let metrics = rt.block_on(trainer.train(&data_dir, |_epoch, _bytes, _best| { - Ok("skip".to_owned()) - }))?; + let metrics = rt.block_on(trainer.train_with_preloaded_data( + train, val, |_epoch, _bytes, _best| Ok("skip".to_owned()), + ))?; assert!( metrics.epochs_trained >= 2, - "Hyperopt train: only trained {} epochs (expected >=2)", + "Hyperopt preloaded: only trained {} epochs (expected >=2)", metrics.epochs_trained ); - assert_finite(metrics.loss, "hyperopt loss"); + assert_finite(metrics.loss, "hyperopt_preloaded loss"); let grad_norm = metrics.additional_metrics.get("avg_gradient_norm") .copied().unwrap_or(0.0); - assert!(grad_norm > 0.0, "Hyperopt: grad_norm={grad_norm} -- model must be learning"); + assert!(grad_norm > 0.0, "Hyperopt preloaded: grad_norm={grad_norm} — model must be learning"); tracing::info!( - "Hyperopt: loss={:.6}, grad_norm={:.4}, epochs={}", + "Hyperopt preloaded: loss={:.6}, grad_norm={:.4}, epochs={}", metrics.loss, grad_norm, metrics.epochs_trained ); Ok(()) } + +/// Hyperopt shared path: load data once, pass &slice to trainer (zero-copy). +/// +/// Production hyperopt uses Arc-wrapped data to avoid ~150 MB deep clone +/// per trial. This test validates the zero-copy borrow path. +#[test] +#[ignore] // Loads real training data + GPU +fn test_hyperopt_shared_data() -> anyhow::Result<()> { + let (train, val) = load_smoke_data()?; + assert!(!train.is_empty(), "Training data must be non-empty"); + assert!(!val.is_empty(), "Validation data must be non-empty"); + + let p = smoke_params(); + let mut trainer = smoke_trainer_with(p)?; + let rt = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build()?; + + let metrics = rt.block_on(trainer.train_with_shared_data( + &train, val, |_epoch, _bytes, _best| Ok("skip".to_owned()), + ))?; + + assert!( + metrics.epochs_trained >= 2, + "Hyperopt shared: only trained {} epochs (expected >=2)", + metrics.epochs_trained + ); + assert_finite(metrics.loss, "hyperopt_shared loss"); + + let grad_norm = metrics.additional_metrics.get("avg_gradient_norm") + .copied().unwrap_or(0.0); + assert!(grad_norm > 0.0, "Hyperopt shared: grad_norm={grad_norm} — model must be learning"); + + tracing::info!( + "Hyperopt shared: loss={:.6}, grad_norm={:.4}, epochs={}", + metrics.loss, grad_norm, metrics.epochs_trained + ); + Ok(()) +} + +/// Hyperopt consistency: preloaded and shared paths produce similar loss. +/// +/// Both paths should call train_with_data_full_loop with the same data. +/// Due to GPU RNG (epsilon-greedy, noisy nets, domain rand), results won't +/// be identical, but loss magnitude should be in the same ballpark. +#[test] +#[ignore] // Loads real training data + GPU — runs 2 training loops +fn test_hyperopt_paths_consistent() -> anyhow::Result<()> { + let (train, val) = load_smoke_data()?; + + let rt = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build()?; + + // Run preloaded path + let mut trainer1 = smoke_trainer_with(smoke_params())?; + let m1 = rt.block_on(trainer1.train_with_preloaded_data( + train.clone(), val.clone(), |_e, _b, _best| Ok("skip".to_owned()), + ))?; + + // Run shared path + let mut trainer2 = smoke_trainer_with(smoke_params())?; + let m2 = rt.block_on(trainer2.train_with_shared_data( + &train, val, |_e, _b, _best| Ok("skip".to_owned()), + ))?; + + // Both must complete + assert!(m1.epochs_trained >= 2, "Preloaded path failed: {} epochs", m1.epochs_trained); + assert!(m2.epochs_trained >= 2, "Shared path failed: {} epochs", m2.epochs_trained); + + // Losses should be same order of magnitude (within 10x) + let ratio = if m1.loss.abs() > 1e-10 && m2.loss.abs() > 1e-10 { + (m1.loss / m2.loss).abs() + } else { + 1.0 // both near-zero is fine + }; + assert!( + ratio > 0.1 && ratio < 10.0, + "Hyperopt paths diverged: preloaded={:.6}, shared={:.6}, ratio={:.2}", + m1.loss, m2.loss, ratio + ); + + tracing::info!( + "Hyperopt consistency: preloaded={:.6}, shared={:.6}, ratio={:.2}", + m1.loss, m2.loss, ratio + ); + Ok(()) +} 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 27a4ec5c8..b02e9a3f6 100644 --- a/crates/ml/src/trainers/dqn/smoke_tests/training_stability.rs +++ b/crates/ml/src/trainers/dqn/smoke_tests/training_stability.rs @@ -397,6 +397,28 @@ fn test_50_epoch_convergence() -> anyhow::Result<()> { Ok(()) } +/// CUDA training auto-initializes the GPU experience collector. +/// (Previously this test verified rejection when collector was missing, +/// but the collector now initializes automatically during training.) +#[test] +#[ignore] // Loads real training data — run via nightly CI or manual trigger +fn test_gpu_collector_auto_initializes() -> anyhow::Result<()> { + let data_dir = test_data_dir() + .expect("FOXHUNT_TEST_DATA or test_data/ must exist"); + let params = smoke_params(); + let mut trainer = smoke_trainer_with(params)?; + let rt = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build()?; + let result = rt.block_on(trainer.train(&data_dir, |_epoch, _bytes, _is_best| { + Ok(String::new()) + })); + assert!(result.is_ok(), "Training must succeed with auto-initialized GPU collector: {:?}", result.err()); + drop(trainer); + drop(rt); + Ok(()) +} + /// Validate the zero-copy fxcache training path end-to-end. /// /// This is the EXACT code path that H100 production uses: diff --git a/crates/ml/src/trainers/dqn/trainer/mod.rs b/crates/ml/src/trainers/dqn/trainer/mod.rs index c5d56db4d..a86d659a3 100644 --- a/crates/ml/src/trainers/dqn/trainer/mod.rs +++ b/crates/ml/src/trainers/dqn/trainer/mod.rs @@ -418,11 +418,16 @@ impl DQNTrainer { self } - /// Train DQN by loading data from DBN files, then using the slice-based loop. + /// Train DQN on market data from DBN files /// - /// This is a convenience wrapper for smoke tests and the legacy hyperopt - /// adapter. Production training uses `init_from_fxcache` + - /// `train_fold_from_slices` directly. + /// # Arguments + /// + /// * `dbn_data_dir` - Directory containing DBN files (e.g., "test_data/real/databento/ml_training/") + /// * `checkpoint_callback` - Callback for saving checkpoints (epoch, model_data, is_final) -> `Result` + /// + /// # Returns + /// + /// Training metrics (loss, accuracy, gradient norms, Q-values) pub async fn train( &mut self, dbn_data_dir: &str, @@ -436,6 +441,7 @@ impl DQNTrainer { self.hyperparams.epochs, self.hyperparams.batch_size ); + // Load market data from DBN files (ALL data for walk-forward or single-pass) let (training_data, val_data) = self.load_training_data(dbn_data_dir).await?; info!( @@ -444,36 +450,81 @@ impl DQNTrainer { val_data.len() ); - // Convert (FeatureVector, Vec) -> ([f64; 42], [f64; 4]) - let features: Vec<[f64; 42]> = training_data.iter().map(|(f, _)| *f).collect(); - let targets: Vec<[f64; 4]> = training_data.iter().map(|(_, t)| { - let mut arr = [0.0_f64; 4]; - for (i, v) in t.iter().take(4).enumerate() { - arr[i] = *v; - } - arr - }).collect(); - // Clone OFI to break borrow on self (init_from_fxcache needs &mut self) - let ofi_owned: Vec<[f64; 8]> = self.ofi_features - .as_deref() - .map(|s| s.to_vec()) - .unwrap_or_default(); + // GPU walk-forward: upload ALL data to GPU, run expanding-window folds (always active when CUDA available) + if self.cuda_stream.is_some() { + // Merge train+val into a single dataset for walk-forward splitting + let mut all_data = training_data; + all_data.extend(val_data); + info!( + "GPU walk-forward enabled: {} total bars, uploading to VRAM", + all_data.len(), + ); + return self.train_walk_forward(&all_data, checkpoint_callback).await; + } - let n = features.len(); - let n_val = val_data.len(); + // Standard single-pass training + self.ofi_val_offset = training_data.len(); + self.val_data = val_data; + self.val_features_gpu = None; + self.val_closes_gpu = None; + self.val_ofi_gpu = None; + self.train_with_data_full_loop(&training_data, checkpoint_callback) + .await + } - // Upload to GPU - self.init_from_fxcache(&features, &targets, &ofi_owned).await?; - self.set_training_range(0, n, n, n + n_val); + /// Train with preloaded data (skips disk I/O and feature extraction). + /// + /// Accepts pre-split training and validation data that was loaded once and + /// cached across hyperopt trials. This avoids re-reading 36 `.dbn.zst` files + /// and re-extracting 42 features on every trial, eliminating minutes of GPU + /// idle time at each trial boundary. + /// + /// # Arguments + /// + /// * `training_data` - Pre-extracted (features, targets) for training split + /// * `val_data` - Pre-extracted (features, targets) for validation split + /// * `checkpoint_callback` - Checkpoint save callback + /// + /// # Returns + /// + /// Training metrics from the completed run + /// + /// # Errors + /// + /// Returns error if the training loop fails + pub async fn train_with_preloaded_data( + &mut self, + training_data: Vec<(FeatureVector, Vec)>, + val_data: Vec<(FeatureVector, Vec)>, + checkpoint_callback: F, + ) -> Result + where + F: FnMut(usize, Vec, bool) -> Result + Send, + { + // Clear stale CUDA context errors from previous training runs. + // cudarc stores errors from CUDA Graph capture (cuStreamWaitEvent on + // disabled events). These persist across DQNTrainer instances and + // cause bind_to_thread() to fail on the next run. + if let MlDevice::Cuda { ref context, .. } = self.device { + let _ = context.check_err(); + } - // Store validation data - self.ofi_val_offset = n; + info!( + "Starting DQN training with preloaded data: {} train, {} val samples", + training_data.len(), + val_data.len() + ); + + // Store validation data for loss computation + self.ofi_val_offset = training_data.len(); self.val_data = val_data; self.val_features_gpu = None; self.val_closes_gpu = None; self.val_ofi_gpu = None; - self.train_fold_from_slices(&features, &targets, checkpoint_callback).await + // Use the common training loop (Wave 12 Group 3 refactor) + self.train_with_data_full_loop(&training_data, checkpoint_callback) + .await } /// Train one fold from contiguous feature/target slices. @@ -513,6 +564,270 @@ impl DQNTrainer { .await } + /// Train with shared preloaded data (zero-copy for hyperopt). + /// + /// Same as [`train_with_preloaded_data`] but accepts `Arc`-wrapped data, + /// avoiding a ~150 MB deep clone per hyperopt trial. + pub async fn train_with_shared_data( + &mut self, + training_data: &[(FeatureVector, Vec)], + val_data: Vec<(FeatureVector, Vec)>, + checkpoint_callback: F, + ) -> Result + where + F: FnMut(usize, Vec, bool) -> Result + Send, + { + info!( + "Starting DQN training with shared data: {} train, {} val samples", + training_data.len(), + val_data.len() + ); + + self.ofi_val_offset = training_data.len(); + self.val_data = val_data; + self.val_features_gpu = None; + self.val_closes_gpu = None; + self.val_ofi_gpu = None; + self.train_with_data_full_loop(training_data, checkpoint_callback) + .await + } + + /// Train with GPU-resident walk-forward cross-validation. + /// + /// Uploads the ENTIRE dataset to GPU VRAM once, then runs expanding-window + /// walk-forward: each fold trains on [0..T], validates on [T..V], tests on + /// [V..E]. Fold transitions are zero-copy (index range changes only). + /// + /// Returns the metrics from the LAST fold (most data, most representative). + pub async fn train_walk_forward( + &mut self, + training_data: &[(FeatureVector, Vec)], + mut checkpoint_callback: F, + ) -> Result + where + F: FnMut(usize, Vec, bool) -> Result + Send, + { + use crate::cuda_pipeline::gpu_walk_forward::{GpuWalkForwardConfig, GpuWalkForwardData}; + + let wf_config = GpuWalkForwardConfig { + initial_train_fraction: self.hyperparams.wf_initial_train_fraction, + val_fraction: self.hyperparams.wf_val_fraction, + test_fraction: self.hyperparams.wf_test_fraction, + step_fraction: self.hyperparams.wf_step_fraction, + }; + + // Upload ALL data to GPU once + let wf_stream = self.cuda_stream.as_ref() + .ok_or_else(|| anyhow::anyhow!("CUDA stream required for walk-forward upload"))?; + let wf_data = GpuWalkForwardData::upload( + training_data, + self.ofi_features.as_deref(), + &wf_config, + wf_stream, + ).map_err(|e| anyhow::anyhow!("GPU walk-forward upload: {e}"))?; + + let num_folds = wf_data.num_folds(); + if num_folds == 0 { + return Err(anyhow::anyhow!( + "Insufficient data for walk-forward: {} bars, need at least {} for one fold", + training_data.len(), + ((wf_config.initial_train_fraction + wf_config.val_fraction + wf_config.test_fraction) * training_data.len() as f64) as usize, + )); + } + + info!( + "GPU walk-forward: {} folds, {:.1} MB VRAM, {} total bars", + num_folds, wf_data.vram_bytes as f64 / 1_048_576.0, wf_data.total_bars, + ); + + // Log stratification quality summary: compute max regime deviation across folds + { + let global_features_flat: Vec = training_data.iter() + .flat_map(|(features, _)| features.iter().copied()) + .collect(); + let (g_t, g_r, g_v) = + ml_dqn::experience::regime_distribution_flat_f64(&global_features_flat, 42); + + let mut max_dev_across_folds = 0.0_f64; + for fold in &wf_data.folds { + let val_flat: Vec = training_data[fold.val_start..fold.val_end] + .iter() + .flat_map(|(features, _)| features.iter().copied()) + .collect(); + let (v_t, v_r, v_v) = + ml_dqn::experience::regime_distribution_flat_f64(&val_flat, 42); + let dev = (v_t - g_t).abs().max((v_r - g_r).abs()).max((v_v - g_v).abs()); + max_dev_across_folds = max_dev_across_folds.max(dev); + } + + info!( + "Regime stratification: global T={:.1}% R={:.1}% V={:.1}%, max fold deviation={:.1}pp", + g_t, g_r, g_v, max_dev_across_folds, + ); + } + + // Store GPU walk-forward data and cudarc buffers for the experience collector + self.features_raw_cuda = Some(wf_data.features); + self.targets_raw_cuda = Some(wf_data.targets); + + let mut last_metrics = TrainingMetrics::new(); + + // Task 14: Collect per-fold data for post-loop R² gap analysis + let mut fold_volatile_pcts: Vec = Vec::with_capacity(num_folds); + let mut fold_sharpes: Vec = Vec::with_capacity(num_folds); + + for fold_idx in 0..num_folds { + let fold = wf_data.folds.get(fold_idx).ok_or_else(|| { + anyhow::anyhow!("Fold {fold_idx} out of range") + })?; + + info!( + "=== Walk-Forward Fold {}/{} === train: {} bars, val: {} bars, test: {} bars", + fold_idx + 1, num_folds, fold.train_len(), fold.val_len(), fold.test_len(), + ); + + // Split training_data into fold's train and val slices (for CPU-side data) + let fold_train = &training_data[fold.train_start..fold.train_end]; + let fold_val: Vec<(FeatureVector, Vec)> = + training_data[fold.val_start..fold.val_end].to_vec(); // cpu-side fold split + + // Store validation data for this fold + self.ofi_val_offset = fold.train_end; + self.val_data = fold_val; + + // Reset training state for new fold + self.gpu_data = None; // Force re-upload via DqnGpuData for the fold's range + self.best_sharpe = f64::NEG_INFINITY; + self.best_val_loss = f64::INFINITY; + self.loss_history.clear(); + self.q_value_history.clear(); + self.val_loss_history.clear(); + self.sharpe_history.clear(); + + // Compute regime distribution for this fold's validation data + let val_features_flat: Vec = self.val_data.iter() + .flat_map(|(features, _)| features.iter().copied()) + .collect(); + let (trending_pct, ranging_pct, volatile_pct) = + ml_dqn::experience::regime_distribution_flat_f64(&val_features_flat, 42); + + // ── Task 15: Per-fold adaptive regime replay decay ────────────── + // More concentrated regime → lower decay (more aggressive bias toward that regime) + // Balanced regime (33/33/33) → higher decay (less bias needed) + let max_regime_pct = trending_pct.max(ranging_pct).max(volatile_pct); + let adaptive_decay = self.hyperparams.regime_replay_decay + * (1.0 - max_regime_pct / 100.0); + let adaptive_decay = adaptive_decay.clamp(0.05, 0.95); + + // Determine dominant regime for sample_regime_biased() current_regime parameter + let dominant_regime: u8 = if trending_pct >= ranging_pct && trending_pct >= volatile_pct { + 0 // Trending + } else if volatile_pct >= ranging_pct { + 2 // Volatile + } else { + 1 // Ranging + }; + + self.regime_replay_decay_override = adaptive_decay as f32; + self.fold_dominant_regime = dominant_regime; + + info!( + "Fold {}/{} adaptive replay decay: {:.3} (base={:.3}, max_regime={:.1}%, dominant={})", + fold_idx + 1, num_folds, adaptive_decay, self.hyperparams.regime_replay_decay, + max_regime_pct, + match dominant_regime { 0 => "Trending", 2 => "Volatile", _ => "Ranging" }, + ); + + // Run training loop on this fold's data + last_metrics = self + .train_with_data_full_loop(fold_train, &mut checkpoint_callback) + .await?; + + // Store regime distribution in metrics + last_metrics.add_metric("regime_trending_pct", trending_pct); + last_metrics.add_metric("regime_ranging_pct", ranging_pct); + last_metrics.add_metric("regime_volatile_pct", volatile_pct); + + // ── Task 14: Regime-normalized Sharpe ratio ───────────────────── + let raw_sharpe = self.best_sharpe; + // Regime difficulty weight: volatile = harder (scale up), trending = easier (scale down) + let regime_difficulty = 1.0 + 0.5 * (volatile_pct / 100.0) - 0.3 * (trending_pct / 100.0); + let normalized_sharpe = raw_sharpe * regime_difficulty; + + last_metrics.add_metric("regime_normalized_sharpe", normalized_sharpe); + last_metrics.add_metric("regime_difficulty", regime_difficulty); + + info!( + "Fold {}/{} Sharpe: raw={:.4}, regime-normalized={:.4} (difficulty={:.3})", + fold_idx + 1, num_folds, raw_sharpe, normalized_sharpe, regime_difficulty, + ); + + // Collect per-fold data for post-loop R² analysis + fold_volatile_pcts.push(volatile_pct); + fold_sharpes.push(raw_sharpe); + + info!( + "Fold {}/{} complete: loss={:.6}, epochs={}, best_sharpe={:.4}", + fold_idx + 1, num_folds, + last_metrics.loss, + last_metrics.epochs_trained, + self.best_sharpe, + ); + info!( + "Fold {}/{} regime distribution: Trending={:.1}% Ranging={:.1}% Volatile={:.1}%", + fold_idx + 1, num_folds, + trending_pct, ranging_pct, volatile_pct, + ); + } + + // Reset regime replay decay override after walk-forward completes + self.regime_replay_decay_override = 1.0; + self.fold_dominant_regime = 1; + + // ── Task 14: IS→OOS gap analysis (R² between volatile% and Sharpe) ── + let n = fold_volatile_pcts.len() as f64; + if n >= 3.0 { + let mean_v = fold_volatile_pcts.iter().sum::() / n; + let mean_s = fold_sharpes.iter().sum::() / n; + let mut cov = 0.0; + let mut var_v = 0.0; + let mut var_s = 0.0; + for i in 0..fold_volatile_pcts.len() { + let dv = fold_volatile_pcts[i] - mean_v; + let ds = fold_sharpes[i] - mean_s; + cov += dv * ds; + var_v += dv * dv; + var_s += ds * ds; + } + let r = if var_v > 1e-10 && var_s > 1e-10 { + cov / (var_v.sqrt() * var_s.sqrt()) + } else { + 0.0 + }; + let r_squared = r * r; + + let verdict = if r_squared > 0.7 { + "regime-driven (data problem \u{2014} regime distribution varies across folds)" + } else if r_squared < 0.3 { + "model-driven (architecture/hyperparams \u{2014} model generalizes poorly)" + } else { + "mixed (both regime and model contribute)" + }; + + info!( + "IS\u{2192}OOS gap analysis: R\u{00B2}={:.3} between volatile% and Sharpe \u{2192} {}", + r_squared, verdict, + ); + + last_metrics.add_metric("gap_r_squared", r_squared); + } + + // Clean up GPU walk-forward buffers (features/targets already stored in self) + self.gpu_walk_forward = None; + + Ok(last_metrics) + } + /// Calculate adaptive bounds with margin diff --git a/crates/ml/src/trainers/dqn/trainer/tests.rs b/crates/ml/src/trainers/dqn/trainer/tests.rs index 28f5ac9ef..351aea685 100644 --- a/crates/ml/src/trainers/dqn/trainer/tests.rs +++ b/crates/ml/src/trainers/dqn/trainer/tests.rs @@ -436,11 +436,11 @@ async fn test_train_with_empty_data_completes_gracefully() { params.buffer_size = 1024; // MIN_GPU_CAPACITY — GPU PER mandatory let device = MlDevice::new_cuda(0).expect("CUDA device required"); let mut trainer = DQNTrainer::new_with_device(params, device).unwrap(); - let empty_data: Vec<([f64; 42], [f64; 4])> = vec![]; + let empty_data: Vec<(FeatureVector, Vec)> = vec![]; let checkpoint_callback = |_, _, _| Ok(String::new()); let result = trainer - .train_with_data_full_loop_slices(&empty_data, checkpoint_callback) + .train_with_data_full_loop(&empty_data, checkpoint_callback) .await; assert!( diff --git a/crates/ml/src/trainers/dqn/trainer/training_loop.rs b/crates/ml/src/trainers/dqn/trainer/training_loop.rs index 1c59eabd0..832d13068 100644 --- a/crates/ml/src/trainers/dqn/trainer/training_loop.rs +++ b/crates/ml/src/trainers/dqn/trainer/training_loop.rs @@ -1,12 +1,13 @@ #![allow(unsafe_code)] // Required for CudaSlice reinterpret-cast in GPU PER insert path -//! DQN main training loop — `train_with_data_full_loop_slices`. +//! DQN main training loop — `train_with_data_full_loop`. //! //! The main loop delegates to focused helper methods on `DQNTrainer`: //! - `log_training_config`: one-time Rainbow/target-update logging -//! - `init_gpu_raw_buffers_from_slices`: upload raw cudarc targets/features (Phase 1b) +//! - `init_gpu_data`: upload training data to GPU (Phase 1) +//! - `init_gpu_raw_buffers`: upload raw cudarc targets/features + portfolio sim (Phase 1b) //! - `init_gpu_experience_collector`: build the zero-roundtrip CUDA collector (Phase 1c) -//! - `collect_gpu_experiences_slices`: run the GPU experience collection kernel (Phase 3) -//! - `run_training_steps_slices`: batched training from replay buffer with guard kernel +//! - `collect_gpu_experiences`: run the GPU experience collection kernel (Phase 3) +//! - `run_training_steps`: batched training from replay buffer with guard kernel //! - `process_epoch_boundary`: single readback + safety checks at epoch end //! - `sync_gpu_weights`: push updated weights to GPU collector //! - `refresh_stale_per_priorities`: M2 PER staleness refresh @@ -20,11 +21,13 @@ use cudarc::driver::CudaSlice; use common::metrics::{questdb_sink, training_metrics}; use tracing::{debug, info, warn}; +use crate::cuda_pipeline::DqnGpuData; use crate::dqn::logging::{log_epoch_start, log_epoch_end, log_training_progress}; use crate::dqn::target_update::convergence_half_life; use crate::evaluation::metrics::calculate_var_cvar; use crate::trainers::TargetUpdateMode; use crate::TrainingMetrics; +use crate::features::extraction::FeatureVector; use super::super::config::DQNAgentType; use super::super::financials::compute_epoch_financials; use super::super::monitoring::TrainingMonitor; @@ -50,6 +53,470 @@ pub(crate) struct EpochLogOutput { } impl DQNTrainer { + // ═══════════════════════════════════════════════════════════════════════ + // Main training loop — orchestrates helpers + // ═══════════════════════════════════════════════════════════════════════ + + pub(crate) async fn train_with_data_full_loop( + &mut self, + training_data: &[(FeatureVector, Vec)], + mut checkpoint_callback: F, + ) -> Result + where + F: FnMut(usize, Vec, bool) -> Result + Send, + { + let start_time = std::time::Instant::now(); + let mut total_loss = 0.0; + let mut total_q_value = 0.0; + let mut total_gradient_norm = 0.0; + let mut total_reward = 0.0; + let mut total_action_counts = [0_usize; 9]; + let mut total_factored_action_counts = [0_usize; 81]; + + self.log_training_config().await; + + // Decision Transformer pre-training on offline trajectories + if self.hyperparams.dt_pretrain_epochs > 0 { + info!( + epochs = self.hyperparams.dt_pretrain_epochs, + context_len = self.hyperparams.dt_context_len, + embed_dim = self.hyperparams.dt_embed_dim, + "Starting Decision Transformer pre-training" + ); + + // Ensure GPU data is uploaded before DT pre-training + self.init_gpu_raw_buffers(training_data).await?; + + if let Some(ref stream) = self.cuda_stream { + let stream = Arc::clone(stream); + + // Build DT config from hyperparameters. + // DT state_dim = 42 (raw market features), NOT agent.get_state_dim() + // which includes portfolio dims appended by the experience collector. + // The raw features_raw_cuda buffer is [num_bars, 42]. + let dt_state_dim: usize = 42; + + let dt_config = crate::cuda_pipeline::decision_transformer::DecisionTransformerConfig { + state_dim: dt_state_dim, + num_actions: 9, // DT uses branch_0 exposure actions only + embed_dim: self.hyperparams.dt_embed_dim, + num_layers: self.hyperparams.dt_num_layers, + num_heads: 4, // standard default + context_len: self.hyperparams.dt_context_len, + dropout: 0.1, + batch_size: self.hyperparams.batch_size.min(256), + }; + + let mut dt = crate::cuda_pipeline::decision_transformer::DecisionTransformer::new( + stream, dt_config, + ).map_err(|e| anyhow::anyhow!("DT init: {e}"))?; + + // Build trajectories from GPU-resident features and targets + let num_bars = training_data.len(); + if let (Some(ref features_gpu), Some(ref targets_gpu)) = + (&self.features_raw_cuda, &self.targets_raw_cuda) + { + let (trajectories, target_actions, num_batches) = dt + .build_dt_trajectories( + features_gpu, + targets_gpu, + num_bars, + self.hyperparams.gamma, + ) + .map_err(|e| anyhow::anyhow!("DT trajectory build: {e}"))?; + + if num_batches == 0 { + warn!( + num_bars, + context_len = self.hyperparams.dt_context_len, + batch_size = dt.config().batch_size, + "DT pre-training: not enough data for a full batch, skipping" + ); + } else { + let batch_size = dt.config().batch_size; + let context_len = dt.config().context_len; + let input_dim = dt.config().state_dim + 2; + + for epoch in 0..self.hyperparams.dt_pretrain_epochs { + let mut epoch_loss = 0.0_f32; + let mut batch_count = 0_usize; + + for batch_idx in 0..num_batches { + // Offset into the trajectory/target buffers for this batch. + // trajectories: [num_episodes, T, input_dim] + // Each batch is batch_size consecutive episodes. + let ep_offset = batch_idx * batch_size; + let traj_elem_offset = ep_offset * context_len * input_dim; + let act_elem_offset = ep_offset * context_len; + + let loss = dt.pretrain_step( + &trajectories, + &target_actions, + batch_size, + traj_elem_offset, + act_elem_offset, + ).map_err(|e| anyhow::anyhow!("DT pretrain_step: {e}"))?; + + epoch_loss += loss; + batch_count += 1; + } + + let avg_loss = if batch_count > 0 { + epoch_loss / batch_count as f32 + } else { + 0.0 + }; + + info!( + epoch = epoch + 1, + total_epochs = self.hyperparams.dt_pretrain_epochs, + avg_loss = format!("{avg_loss:.4}"), + batches = batch_count, + "DT pre-training" + ); + + training_metrics::set_epoch( + "dt_pretrain", + "loss", + avg_loss as f64, + ); + } + + info!( + epochs = self.hyperparams.dt_pretrain_epochs, + "DT pre-training complete" + ); + } + } else { + warn!("DT pre-training: GPU features/targets not available, skipping"); + } + } else { + warn!("DT pre-training: no CUDA stream available, skipping"); + } + } + + // Training loop + for epoch in 0..self.hyperparams.epochs { + self.current_epoch = epoch; + self.reset_epoch_state(epoch); + + // C51 warmup: linearly ramp c51_alpha from 0→1 over warmup epochs. + // MSE provides 10-100x stronger gradient when per-bar returns are ~0.001 + // and C51 cross-entropy between nearly-identical distributions produces + // near-zero gradient for the value mean. Smooth blending avoids the + // gradient discontinuity of a hard MSE→C51 switch. + { + let alpha = if self.hyperparams.c51_warmup_epochs > 0 { + (epoch as f32 / self.hyperparams.c51_warmup_epochs as f32).min(1.0) + } else { + 1.0 + }; + if let Some(ref mut fused) = self.fused_ctx { + fused.set_c51_alpha(alpha); + } + } + + // Q-gap warmup: ramp from 0 to configured threshold over first 5 epochs. + // In early training, Q-values are near-random → Q-gaps are tiny → static + // threshold forces ALL greedy actions to Flat, creating a "only learn Flat" loop. + // Linear ramp lets exploration happen in early epochs. + if self.hyperparams.q_gap_threshold > 0.0 { + let warmup_epochs = 5.0_f32; + let ramp = (epoch as f32 / warmup_epochs).min(1.0); + let ramped_threshold = self.hyperparams.q_gap_threshold as f32 * ramp; + if let Some(ref mut selector) = self.gpu_action_selector { + selector.set_q_gap_threshold(ramped_threshold); + } + } + + log_epoch_start(epoch + 1, self.hyperparams.epochs, self.hyperparams.learning_rate); + training_metrics::set_epoch("dqn", "current", (epoch + 1) as f64); + + let mut monitor = TrainingMonitor::new(epoch + 1); + + // BUG #15 FIX: Portfolio compounds across epochs (see detailed comment in reset_epoch_state) + self.portfolio_tracker.reset_drawdown_tracking(); + + let epoch_start = std::time::Instant::now(); + + // ── Phase 1: GPU data upload + experience collector init ── + let phase1_start = std::time::Instant::now(); + self.init_gpu_data(training_data).await?; + self.init_gpu_raw_buffers(training_data).await?; + self.init_gpu_experience_collector().await?; + let phase1_ms = phase1_start.elapsed().as_secs_f64() * 1000.0; + + // ── #23 Causal feature masking: GPU-native via threshold on per-feature RNG ── + // Instead of CPU Fisher-Yates shuffle, each feature independently gets a + // random draw. Features where random < mask_fraction are zeroed. + // This is a slightly different distribution (binomial vs exact count) + // but equally effective for regularization and fully GPU-native. + { + let frac = self.hyperparams.feature_mask_fraction; + if frac > 0.0 && frac < 1.0 { + // Pass fraction to the state_gather kernel which handles masking + // via per-feature RNG draw in the feature_noise section. + // The mask is now probabilistic (each feature independently masked + // with probability frac) rather than exact count. + // Set epoch_feature_mask = None to disable the old upload path. + // The state_gather kernel uses feature_noise_scale > 0 OR the + // existing feature_mask upload. Since we want GPU-native masking, + // we'll generate the mask in the state_gather kernel directly + // using the RNG that's already there. + self.epoch_feature_mask = None; + // The feature_mask_fraction is already passed via ExperienceCollectorConfig + // and consumed by the state_gather kernel's feature_mask logic. + } else { + self.epoch_feature_mask = None; + } + } + + // ── #13 Vol normalization: compute realized vol from training data ── + // #13 Vol normalization: always active (one production path) + self.epoch_vol_normalizer = if training_data.len() > 20 { + let returns: Vec = training_data.iter() + .map(|(fv, _)| fv[3]) + .collect(); + let n = returns.len() as f64; + let mean = returns.iter().sum::() / n; + let var = returns.iter().map(|r| (r - mean).powi(2)).sum::() / (n - 1.0); + var.sqrt().max(1e-8) as f32 + } else { + 0.0 + }; + + // ── Phase 2: GPU experience collection ── + eprintln!("[DEBUG] Phase 2: starting GPU experience collection (epoch {})", epoch); + let phase2_start = std::time::Instant::now(); + let gpu_experiences_collected = self.collect_gpu_experiences( + training_data, + ).await?; + let phase2_ms = phase2_start.elapsed().as_secs_f64() * 1000.0; + eprintln!("[DEBUG] Phase 2: done in {:.1}ms, collected={}", phase2_ms, gpu_experiences_collected); + + // CUDA builds: GPU experience collector is MANDATORY + if !gpu_experiences_collected { + return Err(anyhow::anyhow!( + "GPU experience collector MUST be active for CUDA training. \ + No CPU fallback path exists in CUDA builds. \ + or check GPU collector initialization errors above." + )); + } + + // ── Periodic shrink-and-perturb: kill memorized weights ── + let sp_interval = self.hyperparams.shrink_perturb_interval; + if sp_interval > 0 && epoch > 0 && epoch % sp_interval == 0 { + if let Some(ref mut fused) = self.fused_ctx { + let alpha = self.hyperparams.shrink_perturb_alpha as f32; + let sigma = self.hyperparams.shrink_perturb_sigma as f32; + match fused.shrink_and_perturb(alpha, sigma) { + Ok(()) => info!(epoch, alpha, sigma, "Periodic shrink-and-perturb applied"), + Err(e) => warn!(epoch, "Shrink-and-perturb failed (non-fatal): {e}"), + } + } + } + + // ── Phase 3: Batched training from replay buffer ── + eprintln!("[DEBUG] Phase 3: starting training steps (epoch {})", epoch); + let phase3_start = std::time::Instant::now(); + let train_step_count = self.run_training_steps(training_data).await?; + let phase3_ms = phase3_start.elapsed().as_secs_f64() * 1000.0; + eprintln!("[DEBUG] Phase 3: done in {:.1}ms, steps={}", phase3_ms, train_step_count); + + // ── Phase 4: Epoch boundary + validation ── + let phase4_start = std::time::Instant::now(); + let boundary = if train_step_count > 0 { + let b = self.process_epoch_boundary(epoch, train_step_count, &mut monitor).await?; + Some(b) + } else { + None + }; + + // Sync GPU weights after training + self.sync_gpu_weights().await?; + + // ── Three-Phase Training Pipeline ────────────────────────────── + // + // Phase 1 (Behavioral Cloning): epochs 0..c51_warmup_epochs + // C51 alpha ramps 0→1 (MSE→C51 blend). Expert demos active. + // Already implemented via c51_alpha ramp + expert_demo_ratio above. + // + // Phase 2 (Full-Stack Online RL): c51_warmup_epochs..phase3_start + // All features active. Expert ratio decays to 0 via ExpertDemoGenerator. + // Already implemented via expert_demo_decay_epochs linear decay. + // + // Phase 3 (Refinement): last 20% of epochs + // Pure C51. No expert demos. Shrink-and-perturb for plasticity. + // Forces the network to consolidate learned policy without expert crutch. + + let total_epochs = self.hyperparams.epochs; + let phase3_start = total_epochs * 80 / 100; + let in_phase3 = epoch >= phase3_start && total_epochs > 4; + + if in_phase3 { + // Phase 3: Force expert demos off (pure RL refinement) + if let Some(ref mut c) = self.gpu_experience_collector { + c.set_expert_ratio(0.0); + } + + // Phase 3: Shrink-and-perturb at the phase boundary for plasticity + if epoch == phase3_start { + if let Some(ref mut fused) = self.fused_ctx { + let alpha = 0.9_f32; // 90% old weights, 10% noise + let sigma = 0.01_f32; // noise scale + if let Err(e) = fused.shrink_and_perturb(alpha, sigma) { + tracing::warn!("Shrink-and-Perturb failed at Phase 3 start (non-fatal): {e}"); + } else { + tracing::info!( + epoch, + phase3_start, + alpha, + sigma, + "Phase 3 (Refinement): Shrink-and-Perturb applied at phase boundary" + ); + } + } + } + } + + // Shrink-and-Perturb for plasticity maintenance (Phases 1-2). + // Applied every 25% of total epochs (at epoch boundaries only). + // Prevents the network from losing its ability to learn new patterns + // across market regime shifts during long training runs. + // In Phase 3, the boundary perturb above handles it — skip periodic. + if !in_phase3 && self.hyperparams.epochs > 4 { + let interval = (self.hyperparams.epochs / 4).max(1); + if epoch > 0 && epoch % interval == 0 { + if let Some(ref mut fused) = self.fused_ctx { + let alpha = 0.9_f32; // 90% old weights, 10% noise + let sigma = 0.01_f32; // noise scale + if let Err(e) = fused.shrink_and_perturb(alpha, sigma) { + tracing::warn!("Shrink-and-Perturb failed (non-fatal): {e}"); + } else { + tracing::info!( + epoch, + alpha, + sigma, + "Shrink-and-Perturb applied (plasticity maintenance)" + ); + } + } + } + } + + // Compute CVaR scales from IQN head for next epoch's experience collection. + // C51 picks direction, IQN sizes the risk via CVaR at 5th percentile. + if let Some(ref mut fused) = self.fused_ctx { + let bs = fused.batch_size(); + // Clone the actions CudaSlice reference to avoid borrow conflict + // (actions_buf borrows immutably, compute_cvar borrows mutably) + let actions_clone = { + let actions = fused.actions_buf(); + actions.clone() + }; + let cvar_ptr = fused.compute_cvar_device_ptr(&actions_clone, bs, 0.05) + .unwrap_or(0); + if let Some(ref mut collector) = self.gpu_experience_collector { + collector.set_cvar_scales(cvar_ptr); + } + } + + // Flush GPU-accumulated max priority + if train_step_count > 0 { + let agent = self.agent.read().await; + if let Err(e) = agent.flush_max_priority() { + debug!("GPU max_priority flush failed (non-fatal): {}", e); + } + } + let phase4_ms = phase4_start.elapsed().as_secs_f64() * 1000.0; + + let epoch_duration = epoch_start.elapsed(); + + // ── Epoch phase breakdown ── + info!( + "Epoch {}/{} phase breakdown: init={:.0}ms experience={:.0}ms training={:.0}ms validation={:.0}ms total={:.0}ms", + epoch + 1, self.hyperparams.epochs, + phase1_ms, phase2_ms, phase3_ms, phase4_ms, + epoch_duration.as_secs_f64() * 1000.0, + ); + + // Calculate and log epoch metrics + let log_output = self.log_epoch_metrics_and_financials( + epoch, + train_step_count, + &boundary, + &mut monitor, + epoch_duration, + &mut total_action_counts, + &mut total_factored_action_counts, + ).await?; + + total_loss += log_output.avg_loss; + total_q_value += log_output.avg_q_value; + total_gradient_norm += log_output.avg_grad_norm; + + let epoch_avg_reward = if !monitor.reward_history.is_empty() { + monitor.reward_history.iter().sum::() / monitor.reward_history.len() as f32 + } else { + 0.0 + }; + total_reward += epoch_avg_reward as f64; + + // Checkpoints + early stopping (returns Err to signal hyperopt on stop) + self.handle_epoch_checkpoints_and_early_stopping( + epoch, + train_step_count, + &log_output, + &mut checkpoint_callback, + ).await?; + + // WAVE 13-A2: Save periodic checkpoint every N epochs + if (epoch + 1) % self.hyperparams.checkpoint_frequency == 0 { + info!( + "Saving periodic checkpoint at epoch {}/{}", + epoch + 1, self.hyperparams.epochs + ); + let checkpoint_data = self.serialize_model().await?; + let checkpoint_size = checkpoint_data.len(); + let checkpoint_path = checkpoint_callback(epoch + 1, checkpoint_data, false) + .context("Failed to save periodic checkpoint")?; + info!( + "Periodic checkpoint saved: {} ({} bytes)", + checkpoint_path, checkpoint_size + ); + } + } + + let training_duration = start_time.elapsed(); + + let metrics = self.create_final_metrics( + total_loss, total_q_value, total_gradient_norm, total_reward, + self.hyperparams.epochs, training_duration, false, + total_action_counts, total_factored_action_counts, + ).await?; + + { + let mut stored_metrics = self.metrics.write().await; + *stored_metrics = metrics.clone(); + } + + info!( + "Training completed in {:.2}s: final_loss={:.6}, avg_q_value={:.4}", + training_duration.as_secs_f64(), + metrics.loss, + metrics.additional_metrics.get("avg_q_value").unwrap_or(&0.0) + ); + + info!("Best model summary:"); + info!( + " Best Sharpe: {:.4} at epoch {} (val_loss={:.6})", + self.best_sharpe, self.best_epoch, self.best_val_loss + ); + info!(" Best model checkpoint: best_model.safetensors"); + + Ok(metrics) + } + // ═══════════════════════════════════════════════════════════════════════ // Zero-alloc training loop — accepts fixed-size arrays, no Vec // ═══════════════════════════════════════════════════════════════════════ @@ -542,6 +1009,127 @@ impl DQNTrainer { } } + // ═══════════════════════════════════════════════════════════════════════ + // Helper: Phase 1 — GPU data upload (once) + // ═══════════════════════════════════════════════════════════════════════ + + pub(crate) async fn init_gpu_data( + &mut self, + training_data: &[(FeatureVector, Vec)], + ) -> Result<()> { + if self.gpu_data.is_some() { + return Ok(()); + } + + // Check double-buffer: skip upload if active slot already populated + let skip_upload = self.double_buffer.as_ref().is_some_and(|db| db.active().is_some()); + + if skip_upload { + info!("DoubleBuffer: active slot populated, skipping re-upload"); + return Ok(()); + } + + let upload_result = if let Some(ref mut pool) = self.buffer_pool { + info!("GpuBufferPool: reusing pre-allocated staging buffers for {} bars", training_data.len()); + pool.upload_dqn(training_data, self.cuda_stream.as_ref().ok_or_else(|| anyhow::anyhow!("CUDA stream required for upload"))?) + } else { + DqnGpuData::upload(training_data, self.cuda_stream.as_ref().ok_or_else(|| anyhow::anyhow!("CUDA stream required for DqnGpuData upload"))?) + }; + + match upload_result { + Ok(mut gpu_data) => { + info!("GPU data pre-uploaded: {} bars x {} features ({:.1} MB)", + gpu_data.num_bars, + gpu_data.feature_dim, + (gpu_data.num_bars * (42 + 4) * 4) as f64 / 1_048_576.0 + ); + // Upload OFI features to GPU if available + if let Some(ref ofi) = self.ofi_features { + match gpu_data.upload_ofi(ofi, self.cuda_stream.as_ref().ok_or_else(|| anyhow::anyhow!("CUDA stream for OFI upload"))?) { + Ok(()) => info!("GPU OFI features uploaded: {} bars x 8 dims", ofi.len()), + Err(e) => { + return Err(anyhow::anyhow!("GPU OFI upload FAILED (no CPU fallback): {e}")); + } + } + } + // Set tensor core alignment so build_*_states pads output + let ofi_enabled = self.ofi_features.is_some(); + let raw_dim = if ofi_enabled { 53 } else { 45 }; + let aligned_dim = (raw_dim + 7) & !7; + gpu_data.set_aligned_state_dim(aligned_dim); + self.gpu_data = Some(gpu_data); + } + Err(e) => { + return Err(anyhow::anyhow!("GPU data pre-upload FAILED (no CPU fallback): {e}")); + } + } + Ok(()) + } + + // ═══════════════════════════════════════════════════════════════════════ + // Helper: Phase 1b — raw cudarc targets + features + portfolio sim (once) + // ═══════════════════════════════════════════════════════════════════════ + + pub(crate) async fn init_gpu_raw_buffers( + &mut self, + training_data: &[(FeatureVector, Vec)], + ) -> Result<()> { + if self.targets_raw_cuda.is_some() { + return Ok(()); + } + + let stream = match self.cuda_stream { + Some(ref s) => Arc::clone(s), + None => return Ok(()), + }; + + let target_dim = 4; + let feature_dim = 42; + let num_bars = training_data.len(); + + // Build flat targets (same as DqnGpuData::upload but for cudarc) + let mut flat_targets = Vec::with_capacity(num_bars * target_dim); + for (_, targets) in training_data { + for i in 0..target_dim { + flat_targets.push(targets.get(i).copied().unwrap_or(0.0) as f32); + } + } + + match crate::cuda_pipeline::clone_htod_f32_to_bf16(&stream, &flat_targets) { + Ok(buf) => { + info!("CUDA targets_raw uploaded: {} bars x 4 ({:.1} KB)", + num_bars, (num_bars * target_dim * 4) as f64 / 1024.0); + self.targets_raw_cuda = Some(buf); + } + Err(e) => { + return Err(anyhow::anyhow!("CUDA targets_raw upload FAILED (no CPU fallback): {e}")); + } + } + + // Build flat features [num_bars * 42] for GPU experience kernel + if self.features_raw_cuda.is_none() { + let mut flat_features = Vec::with_capacity(num_bars * feature_dim); + for (features, _) in training_data { + for &v in features.iter() { + flat_features.push(v as f32); + } + } + match crate::cuda_pipeline::clone_htod_f32_to_bf16(&stream, &flat_features) { + Ok(buf) => { + info!("CUDA features_raw uploaded: {} bars x {} ({:.1} KB)", + num_bars, feature_dim, (num_bars * feature_dim * 4) as f64 / 1024.0); + self.features_raw_cuda = Some(buf); + } + Err(e) => { + return Err(anyhow::anyhow!("CUDA features_raw upload FAILED (no CPU fallback): {e}")); + } + } + } + + // GpuPortfolioSimulator removed — backtest evaluation uses GpuBacktestEvaluator. + + Ok(()) + } /// Upload raw features + targets from contiguous fxcache slices. /// No tuple unpacking — flat_map directly from &[[f64; N]] arrays. @@ -763,6 +1351,346 @@ impl DQNTrainer { self.hyperparams.gpu_n_episodes.max(32) } + // ═══════════════════════════════════════════════════════════════════════ + // Helper: Phase 3 — GPU experience collection + // ═══════════════════════════════════════════════════════════════════════ + + pub(crate) async fn collect_gpu_experiences( + &mut self, + training_data: &[(FeatureVector, Vec)], + ) -> Result { + // Feature normalization stats calculation epoch (kept for API compat) + let _stats_collection_epochs = { + let ratio_based = (self.hyperparams.epochs as f32 * self.hyperparams.feature_stats_collection_ratio) as usize; + let capped = match self.hyperparams.max_feature_stats_epochs { + Some(max_epochs) => ratio_based.min(max_epochs), + None => ratio_based, + }; + capped.max(1) + }; + + // Set expert demo ratio BEFORE borrowing collector (avoids borrow conflict) + if self.hyperparams.expert_demo_ratio > 0.0 { + use crate::trainers::dqn::expert_demos::ExpertDemoGenerator; + let effective = ExpertDemoGenerator::effective_ratio( + self.hyperparams.expert_demo_ratio, + self.hyperparams.expert_demo_decay_epochs, + self.current_epoch, + ) as f32; + if let Some(ref mut c) = self.gpu_experience_collector { + c.set_expert_ratio(effective); + } + } + + // n_episodes/stride/dr computed later — defer GPU setup to after those computations. + // (See "ALL GPU experience collector setup" block below) + + let ( + Some(ref mut collector), + Some(ref features_buf), + Some(ref targets_buf), + ) = ( + &mut self.gpu_experience_collector, + &self.features_raw_cuda, + &self.targets_raw_cuda, + ) else { + return Ok(false); + }; + + use crate::cuda_pipeline::gpu_experience_collector::ExperienceCollectorConfig; + + let n_episodes = self.hyperparams.gpu_n_episodes.max(32) as i32; + + // Domain randomization: all randomization is GPU-native (zero CPU RNG) + // Domain randomization, mirror, causal, vaccine: ALL always active (one production path) + let timesteps = self.hyperparams.gpu_timesteps_per_episode.min(1000) as i32; + let total_bars = training_data.len() as i32; + let usable_bars = (total_bars - timesteps).max(1); + let stride = (usable_bars / n_episodes).max(1); + + // GPU-native domain randomization — use already-borrowed `collector` + collector.generate_episode_starts_gpu( + n_episodes as usize, stride, usable_bars, + ).map_err(|e| anyhow::anyhow!("GPU episode starts: {e}"))?; + collector.generate_sim_params_gpu( + n_episodes as usize, + self.hyperparams.transaction_cost_multiplier as f32, + self.hyperparams.avg_spread as f32, + self.hyperparams.fill_ioc_fill_prob as f32, + self.hyperparams.fill_limit_fill_min as f32, + self.hyperparams.fill_limit_fill_max as f32, + ).map_err(|e| anyhow::anyhow!("GPU sim params: {e}"))?; + + // Episode starts for curriculum filtering (still need host copy for ADX filter) + let mut episode_starts: Vec = (0..n_episodes) + .map(|i| { + let base = (i * stride).rem_euclid(usable_bars); + base + }) + .collect(); + + // ── Curriculum learning: filter episode starts by difficulty ────── + // Phase 0 (Easy): only start episodes at trending bars (ADX > 30) + // Phase 1 (Mixed): all starts, but duplicate trending ones (2x weight) + // Phase 2 (Full): all starts, equal weight (default) + let curriculum = self.hyperparams.curriculum_phase; + if curriculum < 2 && !training_data.is_empty() { + let adx_idx = 40; // ADX at market feature index 40 + let adx_threshold = 30.0_f64; + let original = episode_starts.clone(); + episode_starts.clear(); + + for &start in &original { + let bar = (start as usize).min(training_data.len().saturating_sub(1)); + let features = &training_data[bar].0; + let adx = if features.len() > adx_idx { features[adx_idx] } else { 0.0 }; + + match curriculum { + 0 => { + // Easy: only trending bars + if adx > adx_threshold { + episode_starts.push(start); + } + } + 1 => { + // Mixed: all bars, trending bars duplicated + episode_starts.push(start); + if adx > adx_threshold { + episode_starts.push(start); + } + } + _ => episode_starts.push(start), + } + } + + // Ensure at least 16 episodes (fall back to Full if too few) + if episode_starts.len() < 16 { + tracing::warn!( + curriculum, + filtered = episode_starts.len(), + original = original.len(), + "Curriculum phase filtered too aggressively, falling back to Full" + ); + episode_starts = original; + } else { + tracing::debug!( + curriculum, + filtered = episode_starts.len(), + original = original.len(), + "Curriculum learning: episode starts filtered by ADX difficulty" + ); + } + } + + // Reset per-episode state before each epoch + if let Err(e) = collector.reset_episodes( + self.hyperparams.initial_capital as f32, + self.hyperparams.avg_spread as f32, + self.hyperparams.cash_reserve_percent as f32, + ) { + return Err(anyhow::anyhow!("GPU episode reset FAILED (no CPU fallback): {e}")); + } + + let agent = self.agent.read().await; + let epsilon = agent.get_epsilon(); + drop(agent); + + // Expert ratio already set above (before collector borrow). + + // Adversarial regime: harsh multipliers to stress-test the policy + let adversarial = self.adversarial_active; + if adversarial { + info!("Adversarial regime ACTIVE this epoch: 3x spread, 2x tx_cost, 0.5x fill"); + } + + let config = ExperienceCollectorConfig { + n_episodes, + timesteps_per_episode: timesteps, + total_bars, + epsilon, + gamma: self.hyperparams.gamma as f32, + max_position: self.max_position as f32, + enable_action_masking: self.enable_action_masking, + curiosity_scale: 1.0, // always active + loss_aversion: self.hyperparams.loss_aversion as f32, + // Base values only — per-episode randomization is GPU-native (domain_rand_sim_params kernel) + tx_cost_multiplier: { + let base = self.hyperparams.transaction_cost_multiplier as f32; + if adversarial { base * 2.0 } else { base } + }, + count_bonus_coefficient: self.hyperparams.count_bonus_coefficient + .unwrap_or(0.0) as f32, + q_clip_min: self.hyperparams.q_clip_min as f32, + q_clip_max: self.hyperparams.q_clip_max as f32, + huber_kappa: if true { + self.hyperparams.huber_delta as f32 + } else { + 0.0 + }, + use_noisy_nets: true, + noisy_sigma_init: self.hyperparams.noisy_sigma_init as f32, + use_distributional: true, + num_atoms: self.hyperparams.num_atoms as i32, + v_min: self.hyperparams.v_min as f32, + v_max: self.hyperparams.v_max as f32, + fill_median_spread: { + let base = self.hyperparams.avg_spread as f32; + if adversarial { base * 3.0 } else { base } + }, + fill_median_vol: self.median_vol as f32, + fill_ioc_fill_prob: { + let base = self.hyperparams.fill_ioc_fill_prob as f32; + if adversarial { base * 0.5 } else { base } + }, + fill_limit_fill_min: self.hyperparams.fill_limit_fill_min as f32, + fill_limit_fill_max: self.hyperparams.fill_limit_fill_max as f32, + fill_spread_cost_frac: self.hyperparams.fill_spread_cost_frac as f32, + fill_spread_capture_frac: self.hyperparams.fill_spread_capture_frac as f32, + fill_simulation_enabled: self.median_vol > 0.0, + // Q-gap cold-start fix: ramp threshold from 0 to target over 10 epochs. + // At epoch 0, Q-values are random — forcing Q-gap > 0.1 prevents ALL + // trades (except 10% epsilon), starving the model of learning signal. + // Linear ramp: epoch 0→0.0, epoch 5→0.05, epoch 10→full threshold. + q_gap_threshold: { + let target = self.hyperparams.q_gap_threshold as f32; + let epoch = self.current_epoch.min(10) as f32; + target * (epoch / 10.0) + }, + dsr_eta: self.hyperparams.dsr_eta as f32, + n_steps: self.hyperparams.n_steps as i32, + min_hold_bars: self.hyperparams.min_hold_bars as i32, + // spread_cost = tick_size * multiplier * fraction + // Matches backtest_env_kernel's spread_cost from GpuBacktestConfig + spread_cost: { + let base = (self.hyperparams.tick_size * self.hyperparams.contract_multiplier + * self.hyperparams.fill_spread_cost_frac) as f32; + if adversarial { base * 3.0 } else { base } + }, + contract_multiplier: self.hyperparams.contract_multiplier as f32, + margin_pct: self.hyperparams.margin_pct as f32, + dd_threshold: self.hyperparams.dd_threshold as f32, + w_dd: self.hyperparams.w_dd as f32, + beta_penalty: self.hyperparams.beta_penalty_strength as f32, + // Gems & Pearls generalization params + time_reversal_mod: self.hyperparams.time_reversal_mod as i32, + regret_blend: self.hyperparams.regret_blend as f32, + mirror_active: self.current_epoch % 2 == 1, // Mirror every other epoch (always active) + position_entropy_weight: self.hyperparams.position_entropy_weight as f32, + trade_clustering_penalty: self.hyperparams.trade_clustering_penalty as f32, + feature_noise_scale: self.hyperparams.feature_noise_scale as f32, + vol_normalizer: self.epoch_vol_normalizer, + feature_mask: self.epoch_feature_mask.clone(), + ..Default::default() + }; + + // #33 Saboteur state already set above (before collector borrow) + + // Zero-roundtrip GPU path — GPU PER is always active in CUDA builds + let gpu_batch = collector.collect_experiences_gpu( + features_buf, targets_buf, &episode_starts, &config, + ).map_err(|e| anyhow::anyhow!( + "GPU zero-roundtrip collection FAILED (no CPU fallback): {e}" + ))?; + + let count = gpu_batch.n_episodes * gpu_batch.timesteps * 2; // counterfactual doubles experiences + + // Record CUDA event after experience collection completes on the forked stream. + // This replaces stream.synchronize() — the event is checked lazily at the start + // of run_training_steps, allowing CPU-side training setup (lock acquisition, + // variable initialization, shrink-and-perturb) to overlap with the tail end + // of experience collection GPU kernels. + // The PER insert (insert_batch_bf16) uses the same forked stream, so in-stream + // ordering guarantees data visibility without cross-stream synchronization. + if let Some(ref stream) = self.cuda_stream { + self.experience_done_event = Some(stream.record_event(None) + .map_err(|e| { eprintln!("!!! EVENT RECORD FAILED: {e}"); anyhow::anyhow!("experience event record: {e}") })?); + } + + if let Some(ref mut mon) = self.gpu_monitoring { + if let Err(e) = mon.reduce(collector.rewards_gpu(), collector.actions_gpu(), count) { + debug!("GPU monitoring reduce failed: {e}"); + } + // Catch async monitoring kernel crashes before they poison the stream + if let Some(ref stream) = self.cuda_stream { + if stream.synchronize().is_err() { + warn!("GPU monitoring kernel crashed — disabling monitoring"); + self.gpu_monitoring = None; + } + } + } + + if count > 0 { + if let Err(e) = collector.train_curiosity_gpu( + gpu_batch.n_episodes, gpu_batch.timesteps, + ) { + debug!("GPU curiosity training failed (non-fatal): {e}"); + } + // Record event after curiosity kernel to detect async crashes. + // The event is checked immediately (non-blocking poll first, then sync + // if still pending). This catches poisoned-stream errors from the + // curiosity kernel before they propagate to the PER insert path. + if let Some(ref stream) = self.cuda_stream { + match stream.record_event(None) { + Ok(event) => { + if !event.is_complete() { + if let Err(e) = event.synchronize() { + warn!("Curiosity kernel async error detected: {e} — disabling curiosity for this run"); + self.curiosity_module = None; + collector.disable_curiosity(); + } + } + } + Err(e) => { + warn!("Curiosity event record failed: {e} — disabling curiosity for this run"); + self.curiosity_module = None; + collector.disable_curiosity(); + } + } + } + } + + if count > 0 { + // Direct CudaSlice insertion — zero Candle, zero CPU roundtrip. + let total = gpu_batch.n_episodes * gpu_batch.timesteps * 2; // counterfactual doubles experiences + + // actions: CudaSlice → CudaSlice (safe reinterpret for values [0..44]) + let actions_u32: &CudaSlice = unsafe { + &*(&gpu_batch.actions as *const CudaSlice as *const CudaSlice) + }; + + // rewards/dones: CudaSlice (kernel outputs float directly, no bf16 NaN) + let agent = self.agent.read().await; + let mut gpu_buf = agent.memory().as_gpu_buffer() + .ok_or_else(|| anyhow::anyhow!("GPU PER buffer required"))?; + gpu_buf.gpu.insert_batch_bf16( + &gpu_batch.states, + &gpu_batch.next_states, + actions_u32, + &gpu_batch.rewards, + &gpu_batch.dones, + total, + ).map_err(|e| anyhow::anyhow!("GPU PER insert_batch: {e}"))?; + } + + // Sync GPU DSR EMA state back to CPU for logging/checkpointing. + // The GPU kernel's dsr_step() maintains its own accumulators in epoch_state[5..6]. + // After experience collection, read them back so the CPU RewardFunction stays current. + match collector.read_epoch_dsr_state() { + Ok((dsr_a, dsr_b)) => { + self.reward_fn.sync_dsr_from_gpu(dsr_a as f64, dsr_b as f64); + debug!( + "DSR GPU->CPU sync: ema_return={:.6}, ema_return_sq={:.6}", + dsr_a, dsr_b + ); + } + Err(e) => { + debug!("DSR epoch_state readback failed (non-fatal): {e}"); + } + } + + Ok(true) + } + // ═══════════════════════════════════════════════════════════════════════ // Helper: Phase 3 — GPU experience collection (zero-alloc slice variant) // ═══════════════════════════════════════════════════════════════════════ @@ -1087,6 +2015,199 @@ impl DQNTrainer { Ok(true) } + // ═══════════════════════════════════════════════════════════════════════ + // Helper: Batched training steps with guard kernel + // ═══════════════════════════════════════════════════════════════════════ + + /// Runs all training steps for one epoch, returning the step count. + pub(crate) async fn run_training_steps( + &mut self, + training_data: &[(FeatureVector, Vec)], + ) -> Result { + let batch_size = self.hyperparams.batch_size; + let num_training_steps = if self.can_train().await? { + (training_data.len() / batch_size).max(1) + } else { + 0 + }; + + let mut train_step_count = 0; + + // Lazy-init fused CUDA training context. + // Recreate if batch_size changed (OOM recovery). + if num_training_steps > 0 { + let needs_init = match &self.fused_ctx { + None => self.device.is_cuda(), + Some(ctx) => ctx.batch_size() != self.current_batch_size, + }; + if needs_init { + if let Some(ref stream) = self.cuda_stream { + if self.fused_ctx.is_some() { + tracing::info!("Fused CUDA context: batch_size changed, recreating"); + self.fused_ctx = None; + } + let agent = self.agent.read().await; + match super::super::fused_training::FusedTrainingCtx::new( + &self.device, &*agent, &self.hyperparams, self.current_batch_size, std::sync::Arc::clone(stream), + ) { + Ok(ctx) => { self.fused_ctx = Some(ctx); } + Err(e) => { tracing::error!("Fused CUDA context init failed: {e}"); } + } + } + } + } + + // 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::from_stream(Arc::clone(self.cuda_stream.as_ref().ok_or_else(|| anyhow::anyhow!("CUDA stream for training guard"))?)) { + Ok(g) => { + info!("GPU training guard initialized (epoch loop)"); + self.training_guard = Some(g); + } + Err(e) => return Err(anyhow::anyhow!("GPU training guard init: {e}")), + } + } + if let Some(ref mut guard) = self.training_guard { + guard.reset_accumulators() + .map_err(|e| anyhow::anyhow!("guard reset: {e}"))?; + } + } + + let guard_collapse_thresh = + self.hyperparams.learning_rate as f32 + * self.hyperparams.gradient_collapse_multiplier as f32; + let guard_past_warmup = { + let ws = (self.collapse_warmup_buffer_size as f64 * 0.2) as u64; + self.gradient_logging_step as u64 > ws + }; + + let mut sample_total_us = 0_u64; + let mut fused_total_us = 0_u64; + let mut guard_total_us = 0_u64; + + // Wait for experience collection GPU kernels to complete before sampling. + if let Some(event) = self.experience_done_event.take() { + if !event.is_complete() { + event.synchronize() + .map_err(|e| anyhow::anyhow!("experience event wait: {e}"))?; + } + } + + + // Sample-train loop: one batch at a time, zero prefetch. + // GPU PER sampling (seg_tree_sample kernel) + GPU training (graph replay) + // interleave naturally. No CPU-side Vec accumulation. + for _step in 0..num_training_steps { + // Sample one batch (GPU-native PER for GpuPrioritized) + let sample_start = std::time::Instant::now(); + let (batch, vaccine_batch) = { + let agent = self.agent.read().await; + let buffer = agent.memory(); + if !buffer.can_sample(self.current_batch_size) { + break; + } + let b = buffer.sample(self.current_batch_size) + .map_err(|e| anyhow::anyhow!("PER sample: {e}"))?; + let vb = buffer.sample(self.current_batch_size).ok(); + (b, vb) + }; + sample_total_us += sample_start.elapsed().as_micros() as u64; + + // Train one step + { + let mut agent = self.agent.write().await; + + let fused_start = std::time::Instant::now(); + let _gpu_result = if let Some(ref mut fused) = self.fused_ctx { + if let Some(vb) = vaccine_batch { + fused.pending_vaccine_batch = vb.gpu_batch; + } + let result = fused.run_full_step(&batch, &mut *agent, &self.device) + .map_err(|e| { eprintln!("!!! FUSED STEP ERROR: {:#}", e); e }) + .context("Fused CUDA training step failed")?; + result + } else { + unreachable!("Fused CUDA training is the only production path") + }; + + fused_total_us += fused_start.elapsed().as_micros() as u64; + + let guard_start = std::time::Instant::now(); + if let Some(ref mut guard) = self.training_guard { + // Read loss/grad directly from fused trainer's GPU buffers. + // GpuTrainResult returns hardcoded zeros per-step (no sync). + let (loss_raw, grad_raw) = if let Some(ref fused) = self.fused_ctx { + (fused.loss_gpu_buf().raw_ptr(), fused.grad_norm_gpu_buf().raw_ptr()) + } else { + let ls = _gpu_result.loss_cuda_slice() + .map_err(|e| anyhow::anyhow!("guard loss: {e}"))?; + let gs = _gpu_result.grad_norm_cuda_slice() + .map_err(|e| anyhow::anyhow!("guard grad: {e}"))?; + (ls.raw_ptr(), gs.raw_ptr()) + }; + let gr = guard.check_and_accumulate( + loss_raw, + grad_raw, + 1e6_f32, + guard_collapse_thresh, + !guard_past_warmup, + ).map_err(|e| anyhow::anyhow!("guard check: {e}"))?; + if gr.halt_nan { + return Err(anyhow::anyhow!( + "NaN/Inf at step {}: loss={}, grad={}", + train_step_count, gr.raw_loss, gr.raw_grad_norm + )); + } + if gr.halt_grad_collapse { + agent.check_gradient_collapse(gr.raw_grad_norm).map_err(|e| { + tracing::info!("Early stopping (gradient collapse): {}", e); + anyhow::anyhow!("Early stopping: {}", e) + })?; + } + } + guard_total_us += guard_start.elapsed().as_micros() as u64; + + // Q-value stats: reduce from training batch every 50 steps + if train_step_count % 50 == 0 { + if let Some(ref mut fused) = self.fused_ctx { + if let Ok(stats) = fused.reduce_current_q_stats() { + self.cached_avg_q = stats.q_mean as f64; + if stats.q_min < self.epoch_q_min { self.epoch_q_min = stats.q_min; } + if stats.q_max > self.epoch_q_max { self.epoch_q_max = stats.q_max; } + } + } + } + train_step_count += 1; + self.gradient_logging_step += 1; + } + } + + // Flush the last async readbacks before epoch-end metrics + if let Some(ref mut fused) = self.fused_ctx { + let _ = fused.flush_readback(); + if let Ok(stats) = fused.flush_q_stats_readback() { + self.cached_avg_q = stats.q_mean as f64; + if stats.q_min < self.epoch_q_min { self.epoch_q_min = stats.q_min; } + if stats.q_max > self.epoch_q_max { self.epoch_q_max = stats.q_max; } + } + } + + if train_step_count > 0 { + info!( + "Training step breakdown ({} steps): sample={:.0}ms fused={:.0}ms guard={:.0}ms (per-step: sample={:.1}ms fused={:.1}ms guard={:.1}ms)", + train_step_count, + sample_total_us as f64 / 1000.0, + fused_total_us as f64 / 1000.0, + guard_total_us as f64 / 1000.0, + sample_total_us as f64 / 1000.0 / train_step_count as f64, + fused_total_us as f64 / 1000.0 / train_step_count as f64, + guard_total_us as f64 / 1000.0 / train_step_count as f64, + ); + } + Ok(train_step_count) + } + /// Like `run_training_steps` but accepts `num_bars` directly instead of /// extracting `.len()` from a `&[(FeatureVector, Vec)]` slice. /// This avoids requiring a `Vec` allocation per bar.