revert: dead code removal caused NaN at step 22 — needs investigation

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
This commit is contained in:
jgrusewski
2026-04-03 02:09:40 +02:00
parent ca25b5102c
commit 8b3c209704
7 changed files with 1820 additions and 115 deletions

View File

@@ -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<f64>)],
) -> 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<f64>)],
) -> 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<f64>)> {
(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());

View File

@@ -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<f64>` targets to f32,
/// flattens into contiguous arrays, and uploads once via `clone_htod`.
pub fn upload(
data: &[([f64; 42], Vec<f64>)],
stream: &Arc<CudaStream>,
) -> Result<Self, MLError> {
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<f32>` (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<f64>)],
stream: &Arc<CudaStream>,
) -> Result<DqnGpuData, MLError> {
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::<f32>()
@@ -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<f64>)> = (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<f32> = 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<f64>)> = 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<f64>)> = 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<f64>)> = 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<f64>)> = (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<f32> = 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<f64>)> = 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<f64>)> = (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<f64>)> = (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<f64>)> = (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<f32> = 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<f64>)> = (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<f64>)> = vec![];
assert!(pool.upload_dqn(&data, &stream).is_err());
}
#[test]

View File

@@ -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(())
}

View File

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

View File

@@ -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<String>`
///
/// # Returns
///
/// Training metrics (loss, accuracy, gradient norms, Q-values)
pub async fn train<F>(
&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>) -> ([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<F>(
&mut self,
training_data: Vec<(FeatureVector, Vec<f64>)>,
val_data: Vec<(FeatureVector, Vec<f64>)>,
checkpoint_callback: F,
) -> Result<TrainingMetrics>
where
F: FnMut(usize, Vec<u8>, bool) -> Result<String> + 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<F>(
&mut self,
training_data: &[(FeatureVector, Vec<f64>)],
val_data: Vec<(FeatureVector, Vec<f64>)>,
checkpoint_callback: F,
) -> Result<TrainingMetrics>
where
F: FnMut(usize, Vec<u8>, bool) -> Result<String> + 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<F>(
&mut self,
training_data: &[(FeatureVector, Vec<f64>)],
mut checkpoint_callback: F,
) -> Result<TrainingMetrics>
where
F: FnMut(usize, Vec<u8>, bool) -> Result<String> + 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<f64> = 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<f64> = 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<f64> = Vec::with_capacity(num_folds);
let mut fold_sharpes: Vec<f64> = 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<f64>)> =
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<f64> = 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::<f64>() / n;
let mean_s = fold_sharpes.iter().sum::<f64>() / 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

View File

@@ -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<f64>)> = 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!(

File diff suppressed because it is too large Load Diff