From 8910591f852b2e209e904404dcda40307368d68a Mon Sep 17 00:00:00 2001 From: jgrusewski Date: Sun, 1 Mar 2026 21:42:38 +0100 Subject: [PATCH] =?UTF-8?q?docs:=20I/O=20pipeline=20wiring=20implementatio?= =?UTF-8?q?n=20plan=20=E2=80=94=207=20tasks?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Opus 4.6 --- ...03-01-io-pipeline-wiring-implementation.md | 515 ++++++++++++++++++ 1 file changed, 515 insertions(+) create mode 100644 docs/plans/2026-03-01-io-pipeline-wiring-implementation.md diff --git a/docs/plans/2026-03-01-io-pipeline-wiring-implementation.md b/docs/plans/2026-03-01-io-pipeline-wiring-implementation.md new file mode 100644 index 000000000..5bf5c39a3 --- /dev/null +++ b/docs/plans/2026-03-01-io-pipeline-wiring-implementation.md @@ -0,0 +1,515 @@ +# I/O Pipeline Wiring Implementation Plan + +> **For Claude:** REQUIRED SUB-SKILL: Use superpowers:executing-plans to implement this plan task-by-task. + +**Goal:** Wire the 3 existing but unused I/O pipeline components (EpochPrefetcher, GpuBufferPool, DoubleBufferedLoader) into the DQN training path so fold transitions overlap disk I/O with GPU compute. + +**Architecture:** The walk-forward loop in `train_baseline_rl.rs` gains a prefetcher that loads fold N+1's features on a background thread while fold N trains. Inside the trainer, `GpuBufferPool` replaces the per-fold `DqnGpuData::upload()` heap allocation, and `DoubleBufferedLoader` manages active/staging GPU tensor slots for O(1) fold swaps. + +**Tech Stack:** Rust, Candle (GPU tensors), std::thread (prefetcher), existing `cuda_pipeline` module + +**Build:** `SQLX_OFFLINE=true cargo check -p ml` +**Test:** `SQLX_OFFLINE=true cargo test -p ml --lib` +**Example check:** `SQLX_OFFLINE=true cargo check --example train_baseline_rl -p ml` + +--- + +### Task 1: Wire GpuBufferPool into DQNTrainer's GPU upload path + +**Why:** DQNTrainer already pre-allocates a `GpuBufferPool` at init (trainer.rs:697-700) but never uses it. The trainer calls `DqnGpuData::upload()` (trainer.rs:1443) which heap-allocates fresh staging buffers every fold. Replacing that call with `self.buffer_pool.upload_dqn()` reuses the pre-allocated buffers. + +**Files:** +- Modify: `crates/ml/src/trainers/dqn/trainer.rs:1442-1455` + +**Step 1: Replace `DqnGpuData::upload` with `buffer_pool.upload_dqn` in `train_with_data_full_loop`** + +At trainer.rs:1442-1455, the current code is: +```rust +if self.gpu_data.is_none() { + match DqnGpuData::upload(&training_data, &self.device) { + Ok(gpu_data) => { + info!("GPU data pre-uploaded: {} bars x {} features ({:.1} MB)", + gpu_data.num_bars, + gpu_data.feature_dim, + (gpu_data.num_bars * (51 + 4) * 4) as f64 / 1_048_576.0 + ); + self.gpu_data = Some(gpu_data); + } + Err(e) => { + debug!("GPU data pre-upload skipped (CPU fallback): {}", e); + } + } +} +``` + +Replace with: +```rust +if self.gpu_data.is_none() { + 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.device) + } else { + DqnGpuData::upload(&training_data, &self.device) + }; + match upload_result { + Ok(gpu_data) => { + info!("GPU data pre-uploaded: {} bars x {} features ({:.1} MB)", + gpu_data.num_bars, + gpu_data.feature_dim, + (gpu_data.num_bars * (51 + 4) * 4) as f64 / 1_048_576.0 + ); + self.gpu_data = Some(gpu_data); + } + Err(e) => { + debug!("GPU data pre-upload skipped (CPU fallback): {}", e); + } + } +} +``` + +**Step 2: Run tests** + +Run: `SQLX_OFFLINE=true cargo test -p ml --lib trainers::dqn -- --nocapture 2>&1 | tail -5` +Expected: All DQN trainer tests pass (435+) + +**Step 3: Check example compiles** + +Run: `SQLX_OFFLINE=true cargo check --example train_baseline_rl -p ml` +Expected: Compiles clean + +**Step 4: Commit** + +``` +git add crates/ml/src/trainers/dqn/trainer.rs +git commit -m "feat(ml): wire GpuBufferPool into DQN trainer — reuse pre-allocated staging buffers" +``` + +--- + +### Task 2: Add DoubleBufferedLoader field to DQNTrainer + +**Why:** DoubleBufferedLoader holds two GPU tensor slots (active + staging). When a new fold's data arrives, it uploads to the staging slot, then O(1) swaps staging→active. This eliminates the GPU stall during fold transitions. + +**Files:** +- Modify: `crates/ml/src/trainers/dqn/trainer.rs` + +**Step 1: Add the field and initialization** + +After the `buffer_pool` field (trainer.rs:220), add: +```rust +/// Double-buffered GPU data for zero-downtime fold transitions +double_buffer: Option, +``` + +In the `new_internal()` struct literal (after `buffer_pool,` at trainer.rs:833), add: +```rust +double_buffer: if device.is_cuda() { + Some(crate::cuda_pipeline::double_buffer::DoubleBufferedLoader::new(device.clone())) +} else { + None +}, +``` + +**Step 2: Add public API methods** + +After the `take_prefetched_data()` method (trainer.rs:861), add: +```rust +/// Get a reference to the double-buffered loader, if GPU is active. +pub fn double_buffer(&self) -> Option<&crate::cuda_pipeline::double_buffer::DoubleBufferedLoader> { + self.double_buffer.as_ref() +} + +/// Get a mutable reference to the double-buffered loader, if GPU is active. +pub fn double_buffer_mut(&mut self) -> Option<&mut crate::cuda_pipeline::double_buffer::DoubleBufferedLoader> { + self.double_buffer.as_mut() +} +``` + +**Step 3: Run tests** + +Run: `SQLX_OFFLINE=true cargo test -p ml --lib trainers::dqn -- --nocapture 2>&1 | tail -5` +Expected: All DQN trainer tests pass + +**Step 4: Commit** + +``` +git add crates/ml/src/trainers/dqn/trainer.rs +git commit -m "feat(ml): add DoubleBufferedLoader field to DQN trainer" +``` + +--- + +### Task 3: Wire DoubleBufferedLoader into the GPU upload path + +**Why:** Instead of storing GPU data directly in `self.gpu_data`, route it through the double-buffer's active slot. This prepares for fold N+1 staging while fold N trains. + +**Files:** +- Modify: `crates/ml/src/trainers/dqn/trainer.rs:1442-1456` + +**Step 1: Use double-buffer for initial upload when available** + +Replace the GPU upload block (modified in Task 1) to also use double-buffer when available: +```rust +if self.gpu_data.is_none() { + // Try double-buffer path first (zero-downtime fold transitions) + if let Some(ref mut db) = self.double_buffer { + if db.active().is_none() { + 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.device) + } else { + DqnGpuData::upload(&training_data, &self.device) + }; + match upload_result { + Ok(gpu_data) => { + info!("DoubleBuffer: initial upload — {} bars x {} features ({:.1} MB)", + gpu_data.num_bars, + gpu_data.feature_dim, + (gpu_data.num_bars * (51 + 4) * 4) as f64 / 1_048_576.0 + ); + self.gpu_data = Some(gpu_data); + } + Err(e) => { + debug!("GPU data pre-upload skipped (CPU fallback): {}", e); + } + } + } else { + // Double-buffer already has active data (staged from previous fold) + // Nothing to do — gpu_data was already set via set_gpu_data() + } + } else { + // No double-buffer (CPU mode) — direct upload + 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.device) + } else { + DqnGpuData::upload(&training_data, &self.device) + }; + match upload_result { + Ok(gpu_data) => { + info!("GPU data pre-uploaded: {} bars x {} features ({:.1} MB)", + gpu_data.num_bars, + gpu_data.feature_dim, + (gpu_data.num_bars * (51 + 4) * 4) as f64 / 1_048_576.0 + ); + self.gpu_data = Some(gpu_data); + } + Err(e) => { + debug!("GPU data pre-upload skipped (CPU fallback): {}", e); + } + } + } +} +``` + +**Step 2: Run tests + example check** + +Run: `SQLX_OFFLINE=true cargo test -p ml --lib trainers::dqn -- --nocapture 2>&1 | tail -5` +Run: `SQLX_OFFLINE=true cargo check --example train_baseline_rl -p ml` +Expected: All pass / compile clean + +**Step 3: Commit** + +``` +git add crates/ml/src/trainers/dqn/trainer.rs +git commit -m "feat(ml): wire DoubleBufferedLoader into DQN trainer GPU upload path" +``` + +--- + +### Task 4: Wire EpochPrefetcher in train_baseline_rl.rs walk-forward loop + +**Why:** The walk-forward loop (train_baseline_rl.rs:501) iterates folds sequentially. Feature extraction (`extract_ml_features`) involves CPU-bound OHLCV→indicator computation. The prefetcher runs this on a background thread for fold N+1 while fold N trains on GPU. + +**Important context:** `train_dqn_fold()` creates a new `DQNTrainer` per fold (line 262). The prefetcher must therefore operate at the walk-forward loop level, not inside the trainer. The data flow is: prefetch → extract features → normalize → pass to `train_dqn_fold()`. + +**Files:** +- Modify: `crates/ml/examples/train_baseline_rl.rs:501-599` + +**Step 1: Add imports** + +At the top of the file (after existing `use ml::` imports, around line 31), add: +```rust +use ml::cuda_pipeline::prefetch::EpochPrefetcher; +``` + +**Step 2: Add a helper function for fold data preparation** + +Before `train_dqn_fold()` (line 213), add a helper that encapsulates the feature extraction + normalization pipeline for one fold: + +```rust +/// Prepare a fold's data for training: extract features, normalize, align bars. +/// +/// Returns `(train_norm, val_norm, train_bars_aligned, val_bars_aligned, norm_stats)` +/// or `None` if feature extraction fails or produces empty features. +fn prepare_fold_data( + window: &baseline_common::WalkForwardWindow, + output_dir: &Path, +) -> Option<(Vec<[f64; 51]>, Vec<[f64; 51]>, Vec, Vec)> { + let train_feat = match extract_ml_features(&window.train) { + Ok(f) => f, + Err(e) => { + warn!(" Fold {} -- train feature extraction failed: {}", window.fold, e); + return None; + } + }; + let val_feat = match extract_ml_features(&window.val) { + Ok(f) => f, + Err(e) => { + warn!(" Fold {} -- val feature extraction failed: {}", window.fold, e); + return None; + } + }; + if train_feat.is_empty() || val_feat.is_empty() { + warn!(" Fold {} -- empty features, skipping", window.fold); + return None; + } + + let norm_stats = NormStats::from_features(&train_feat); + let train_norm = norm_stats.normalize_batch(&train_feat); + let val_norm = norm_stats.normalize_batch(&val_feat); + + // Save NormStats + let norm_path = output_dir.join(format!("norm_stats_fold{}.json", window.fold)); + if let Ok(json) = serde_json::to_string_pretty(&norm_stats) { + if let Err(e) = std::fs::write(&norm_path, json) { + warn!(" Failed to write NormStats: {}", e); + } else { + info!(" Saved NormStats to {}", norm_path.display()); + } + } + + let train_warmup = window.train.len().saturating_sub(train_norm.len()); + let train_bars_aligned = window.train.get(train_warmup..).unwrap_or(&window.train).to_vec(); + let val_warmup = window.val.len().saturating_sub(val_norm.len()); + let val_bars_aligned = window.val.get(val_warmup..).unwrap_or(&window.val).to_vec(); + + Some((train_norm, val_norm, train_bars_aligned, val_bars_aligned)) +} +``` + +**Step 3: Restructure the walk-forward loop to use prefetcher** + +Replace the walk-forward loop body (lines 501-599) with a version that spawns a prefetcher for the next fold: + +```rust + // Prefetch state: holds prepared data for the current fold (None for fold 0) + let mut prefetched_data: Option<(Vec<[f64; 51]>, Vec<[f64; 51]>, Vec, Vec)> = None; + + for (fold_idx, window) in windows.iter().enumerate() { + info!("--- Fold {} ---", window.fold); + info!( + " Train: {} bars (up to {}), Val: {} bars (up to {}), Test: {} bars (up to {})", + window.train.len(), window.train_end, + window.val.len(), window.val_end, + window.test.len(), window.test_end, + ); + + // Get fold data: from prefetcher (fold 1+) or prepare fresh (fold 0) + let fold_data = if let Some(data) = prefetched_data.take() { + info!(" Using prefetched data for fold {}", window.fold); + Some(data) + } else { + prepare_fold_data(window, &args.output_dir) + }; + + // Kick off prefetch for NEXT fold (if not last) + let prefetcher = if fold_idx + 1 < windows.len() { + let next_window = windows[fold_idx + 1].clone(); + let out_dir = args.output_dir.clone(); + Some(EpochPrefetcher::spawn(move || { + // This runs on a background thread while the current fold trains + match prepare_fold_data_for_prefetch(&next_window, &out_dir) { + Some(data) => Ok(data), + None => Err(ml::MLError::ModelError( + format!("Fold {} feature extraction failed in prefetcher", next_window.fold) + )), + } + })) + } else { + None + }; + + let Some((train_norm, val_norm, train_bars_aligned, val_bars_aligned)) = fold_data else { + // Collect prefetch result even if this fold failed (don't waste the background work) + if let Some(pf) = prefetcher { + prefetched_data = pf.take().ok(); + } + continue; + }; + + // Train DQN + if train_dqn { + let hp = load_hyperopt_params(&args.hyperopt_params, "dqn"); + match train_dqn_fold( + window.fold, &train_norm, &val_norm, + &train_bars_aligned, &val_bars_aligned, + args, &args.output_dir, &hp, + ) { + Ok(best_loss) => dqn_results.push((window.fold, best_loss)), + Err(e) => error!(" [DQN] Fold {} failed: {}", window.fold, e), + } + } + + // Train PPO + if train_ppo { + let hp = load_hyperopt_params(&args.hyperopt_params, "ppo"); + match train_ppo_fold( + window.fold, &train_norm, &val_norm, + &train_bars_aligned, &val_bars_aligned, + args, &args.output_dir, &hp, + ) { + Ok(best_loss) => ppo_results.push((window.fold, best_loss)), + Err(e) => error!(" [PPO] Fold {} failed: {}", window.fold, e), + } + } + + // Collect prefetch result for next iteration + if let Some(pf) = prefetcher { + match pf.take() { + Ok(data) => { + info!(" Prefetched fold {} data ready ({} train features)", fold_idx + 2, data.0.len()); + prefetched_data = Some(data); + } + Err(e) => { + warn!(" Prefetch for next fold failed (will load synchronously): {}", e); + prefetched_data = None; + } + } + } + } +``` + +**Step 4: Add the prefetch-compatible helper** + +The `EpochPrefetcher` closure needs `Send + 'static`, so it can't borrow `WalkForwardWindow`. Add a helper that takes owned data (the `.clone()` above handles this): + +```rust +/// Same as `prepare_fold_data` but takes owned window (for Send + 'static prefetcher closure). +fn prepare_fold_data_for_prefetch( + window: &baseline_common::WalkForwardWindow, + output_dir: &Path, +) -> Option<(Vec<[f64; 51]>, Vec<[f64; 51]>, Vec, Vec)> { + prepare_fold_data(window, output_dir) +} +``` + +**Note:** The `PrefetchResult` type is `Vec<([f64; 51], Vec)>` but we need `(Vec<[f64; 51]>, Vec<[f64; 51]>, Vec, Vec)`. Since the prefetcher is generic over `Send + 'static` closures, we return a tuple directly, bypassing the `PrefetchResult` type alias. The `EpochPrefetcher::spawn` accepts any `FnOnce() -> Result` — wait, it's typed to `PrefetchResult`. We need to check this. + +**Step 5: Check `EpochPrefetcher` type and adapt** + +`EpochPrefetcher::spawn` is hardcoded to `PrefetchResult = Vec<([f64; 51], Vec)>`. Our fold data is `(Vec<[f64; 51]>, Vec<[f64; 51]>, Vec, Vec)`. Two options: + +a) Make `EpochPrefetcher` generic over the result type +b) Use a plain `std::thread::spawn` + `mpsc::channel` in the binary + +Option (b) is simpler and doesn't modify library code. Replace the `EpochPrefetcher` usage with direct thread spawning in the binary: + +```rust +use std::sync::mpsc; + +type FoldData = (Vec<[f64; 51]>, Vec<[f64; 51]>, Vec, Vec); + +// In the loop, replace EpochPrefetcher::spawn with: +let prefetcher = if fold_idx + 1 < windows.len() { + let next_window = windows[fold_idx + 1].clone(); + let out_dir = args.output_dir.clone(); + let (tx, rx) = mpsc::channel::>(); + std::thread::Builder::new() + .name("fold-prefetch".into()) + .spawn(move || { + let result = prepare_fold_data(&next_window, &out_dir); + let _ = tx.send(result); + }) + .ok(); + Some(rx) +} else { + None +}; + +// Collecting: +if let Some(rx) = prefetcher { + match rx.recv() { + Ok(Some(data)) => { + info!(" Prefetched fold {} data ready ({} train features)", fold_idx + 2, data.0.len()); + prefetched_data = Some(data); + } + _ => { + warn!(" Prefetch for next fold failed (will load synchronously)"); + prefetched_data = None; + } + } +} +``` + +**Step 6: Run example check** + +Run: `SQLX_OFFLINE=true cargo check --example train_baseline_rl -p ml` +Expected: Compiles clean + +**Step 7: Commit** + +``` +git add crates/ml/examples/train_baseline_rl.rs +git commit -m "feat(ml): wire fold prefetching in walk-forward loop — overlap I/O with GPU training" +``` + +--- + +### Task 5: Verify WalkForwardWindow is Clone + +**Why:** Task 4 clones `windows[fold_idx + 1]` for the prefetcher's `Send + 'static` closure. `WalkForwardWindow` must derive or implement `Clone`. + +**Files:** +- Check: `crates/ml/examples/baseline_common/mod.rs` (or wherever `WalkForwardWindow` is defined) + +**Step 1: Find and check the struct** + +Run: `grep -rn "struct WalkForwardWindow" crates/ml/examples/` + +If it doesn't derive `Clone`, add `#[derive(Clone)]` to it. Also check `OHLCVBar` — it must be `Clone` too (likely already is since it's a data struct). + +**Step 2: Run example check** + +Run: `SQLX_OFFLINE=true cargo check --example train_baseline_rl -p ml` +Expected: Compiles clean + +**Step 3: Commit (if changes needed)** + +``` +git commit -m "fix(ml): derive Clone on WalkForwardWindow for prefetcher compatibility" +``` + +--- + +### Task 6: Full test suite + compile verification + +**Why:** Final validation that all changes integrate correctly. + +**Step 1: Full ml crate test suite** + +Run: `SQLX_OFFLINE=true cargo test -p ml --lib 2>&1 | tail -5` +Expected: 2444+ tests pass, 0 failures + +**Step 2: Training binary compile** + +Run: `SQLX_OFFLINE=true cargo check --example train_baseline_rl -p ml 2>&1 | tail -3` +Expected: Compiles clean + +**Step 3: Workspace compile** + +Run: `SQLX_OFFLINE=true cargo check --workspace 2>&1 | tail -3` +Expected: Compiles clean (may have unrelated warnings) + +--- + +### Task 7: Audit — verify all 3 I/O components are wired + +**Why:** Confirm the design goals are met. + +**Checklist:** +- [ ] `GpuBufferPool`: `self.buffer_pool.upload_dqn()` called in `train_with_data_full_loop()` (trainer.rs) +- [ ] `DoubleBufferedLoader`: field exists on DQNTrainer, initialized on CUDA, accessor methods available +- [ ] Fold prefetching: `train_baseline_rl.rs` walk-forward loop spawns background thread for fold N+1 +- [ ] No regressions: all ml tests pass, training binary compiles +- [ ] CPU fallback: all 3 paths degrade gracefully when not on CUDA (None checks, direct upload)