diff --git a/crates/ml/examples/precompute_features.rs b/crates/ml/examples/precompute_features.rs index 664530205..27000b151 100644 --- a/crates/ml/examples/precompute_features.rs +++ b/crates/ml/examples/precompute_features.rs @@ -309,14 +309,17 @@ async fn main() -> Result<()> { .context("Feature extraction failed")?; info!("Extracted {} feature vectors in {:.1}s", feature_vectors.len(), t1.elapsed().as_secs_f64()); - // ── Step 3: Build targets [preproc_close, preproc_next, raw_close, raw_next] + // ── Step 3: Build targets [preproc_close, preproc_next, raw_close, raw_next, raw_open, mid_price_open] let n = feature_vectors.len().saturating_sub(1); // need next bar for target let features: Vec<[f64; 42]> = feature_vectors[..n].to_vec(); - let mut targets: Vec<[f64; 4]> = Vec::with_capacity(n); + let mut targets: Vec<[f64; 6]> = Vec::with_capacity(n); for i in 0..n { let raw_curr = all_bars[i + WARMUP].close; let raw_next = all_bars[i + WARMUP + 1].close; - targets.push([raw_curr, raw_next, raw_curr, raw_next]); + let raw_open = all_bars[i + WARMUP].open; + // mid_price_open will be filled from MBP-10 data below if available; + // default to raw_open (OHLCV-only fallback). + targets.push([raw_curr, raw_next, raw_curr, raw_next, raw_open, raw_open]); } // ── Step 3b: Build per-bar timestamps ───────────────────────────────── @@ -452,6 +455,24 @@ async fn main() -> Result<()> { } } + // Fill targets[i][5] with MBP-10 midpoint at bar open (mid_price_open) + let mut mid_fill_count = 0_usize; + for i in 0..n { + let bar = &all_bars[i + WARMUP]; + let bar_ts = bar.timestamp.timestamp_nanos_opt().unwrap_or(0) as u64; + + // Find first MBP-10 snapshot at or after bar open + let snap_idx = all_snapshots.partition_point(|s| s.timestamp < bar_ts); + if let Some(snap) = all_snapshots.get(snap_idx) { + let mid = snap.mid_price(); + if mid > 0.0 { + targets[i][5] = mid; + mid_fill_count += 1; + } + } + } + info!("mid_price_open filled from MBP-10: {}/{} bars", mid_fill_count, n); + let non_zero = ofi_per_bar.iter().filter(|f| f.iter().any(|&v| v != 0.0)).count(); info!("OFI computed: {} bars, {} non-zero ({:.1}%) in {:.1}s", ofi_per_bar.len(), non_zero, @@ -527,7 +548,7 @@ async fn main() -> Result<()> { println!(); println!("Bars: {}", total_len); println!("Features: 42-dim"); - println!("Targets: 4-dim"); + println!("Targets: 6-dim"); println!("OFI: 20-dim ({})", if has_ofi { "from MBP-10" } else { "zero-padded" }); println!("Format: f32 (v{}), OFI_DIM={}", ml::fxcache::FXCACHE_VERSION, ml::fxcache::OFI_DIM); println!("Cache key: {}", hex_key); diff --git a/crates/ml/examples/train_baseline_rl.rs b/crates/ml/examples/train_baseline_rl.rs index f75517aa5..cdd5ea513 100644 --- a/crates/ml/examples/train_baseline_rl.rs +++ b/crates/ml/examples/train_baseline_rl.rs @@ -447,8 +447,8 @@ fn train_dqn_fold( fold: usize, train_features: &[[f64; 42]], val_features: &[[f64; 42]], - train_targets: &[[f64; 4]], - val_targets: &[[f64; 4]], + train_targets: &[[f64; 6]], + val_targets: &[[f64; 6]], range: &ml::walk_forward::FoldRange, output_dir: &Path, checkpoint_prefix: &str, @@ -580,8 +580,8 @@ fn run_training(args: &Args) -> Result> { let timestamps: Vec = aligned_bars.iter() .map(|b| b.timestamp.timestamp_nanos_opt().unwrap_or(0)) .collect(); - let targets: Vec<[f64; 4]> = aligned_bars.iter() - .map(|b| [b.close, b.close, b.close, b.close]) + let targets: Vec<[f64; 6]> = aligned_bars.iter() + .map(|b| [b.close, b.close, b.close, b.close, b.open, b.open]) .collect(); let ofi = vec![[0.0_f64; 20]; n]; diff --git a/crates/ml/src/cuda_pipeline/experience_kernels.cu b/crates/ml/src/cuda_pipeline/experience_kernels.cu index d80da28d0..defd262a6 100644 --- a/crates/ml/src/cuda_pipeline/experience_kernels.cu +++ b/crates/ml/src/cuda_pipeline/experience_kernels.cu @@ -1170,11 +1170,13 @@ extern "C" __global__ void experience_env_step( } /* ---- Read raw prices from targets ---- */ - /* Target layout: [preproc_close, preproc_next, raw_close, raw_next] + /* Target layout: [preproc_close, preproc_next, raw_close, raw_next, raw_open, mid_price_open] * [0:1] = preprocessed (log-return normalized) — for Q-network * [2] = raw close — for portfolio simulation P&L + tx costs - * [3] = raw_next — NOT used in reward/P&L/equity (future info) */ - const float* tgt = targets + (long long)bar_idx * 4; + * [3] = raw_next — NOT used in reward/P&L/equity (future info) + * [4] = raw open price + * [5] = mid-price at bar open (MBP-10 midpoint; fallback to raw_open) */ + const float* tgt = targets + (long long)bar_idx * 6; /* Trade physics (execute_trade, compute_tx_cost, etc.) use float internally * because they involve multi-step accumulation where bf16 precision is * insufficient (e.g. cash -= delta * price; cash -= cost). We convert diff --git a/crates/ml/src/cuda_pipeline/gpu_walk_forward.rs b/crates/ml/src/cuda_pipeline/gpu_walk_forward.rs index 02f936f04..4fa49193c 100644 --- a/crates/ml/src/cuda_pipeline/gpu_walk_forward.rs +++ b/crates/ml/src/cuda_pipeline/gpu_walk_forward.rs @@ -525,7 +525,7 @@ impl std::fmt::Debug for GpuWalkForwardData { pub struct GpuWalkForwardData { /// Market features [total_bars * FEATURE_DIM] (FEATURE_DIM = 42). pub features: CudaSlice, - /// Target prices [total_bars * TARGET_DIM] (TARGET_DIM = 4). + /// Target prices [total_bars * TARGET_DIM] (TARGET_DIM = 6). pub targets: CudaSlice, /// OFI features [total_bars * 8] (optional, from MBP-10). pub ofi: Option>, @@ -533,7 +533,7 @@ pub struct GpuWalkForwardData { pub total_bars: usize, /// Feature dimension (42). pub feature_dim: usize, - /// Target dimension (4). + /// Target dimension (6). pub target_dim: usize, /// Generated fold ranges. pub folds: Vec, @@ -561,7 +561,7 @@ impl GpuWalkForwardData { } let feature_dim = 42; - let target_dim = 4; + let target_dim = 6; // Flatten features: [total_bars * 42] let mut flat_features = Vec::with_capacity(total_bars * feature_dim); @@ -571,7 +571,7 @@ impl GpuWalkForwardData { } } - // Flatten targets: [total_bars * 4] + // Flatten targets: [total_bars * 6] let mut flat_targets = Vec::with_capacity(total_bars * target_dim); for (_, targets) in training_data { for i in 0..target_dim { diff --git a/crates/ml/src/cuda_pipeline/mod.rs b/crates/ml/src/cuda_pipeline/mod.rs index dc9f29daa..635a88125 100644 --- a/crates/ml/src/cuda_pipeline/mod.rs +++ b/crates/ml/src/cuda_pipeline/mod.rs @@ -207,14 +207,14 @@ pub fn clone_cuda_slice_u32( /// Pre-uploaded GPU training data for DQN trainer. /// -/// Holds market features [N * 42] and target prices [N * 4] as GPU-resident +/// Holds market features [N * 42] and target prices [N * 6] as GPU-resident /// `CudaSlice` buffers. Portfolio features (3 dims) are computed per-bar /// by the trainer and concatenated on-device. OFI features (8 dims from MBP-10 /// order book) are optionally uploaded and appended to the state vector. pub struct DqnGpuData { /// Market features [num_bars * 42] on GPU (f32, row-major) pub features: CudaSlice, - /// Target prices [num_bars * 4] on GPU (f32, row-major) + /// Target prices [num_bars * 6] on GPU (f32, row-major) pub targets: CudaSlice, /// OFI features [num_bars * 8] on GPU (f32), from MBP-10 data pub ofi_features: Option>, @@ -252,7 +252,7 @@ impl DqnGpuData { } let feature_dim = 42; - let target_dim = 4; + let target_dim = 6; let estimated_bytes = estimate_vram_bytes(num_bars * (feature_dim + target_dim)); if estimated_bytes > MAX_UPLOAD_BYTES { @@ -297,7 +297,7 @@ impl DqnGpuData { /// OFI is uploaded only when at least one value is non-zero. pub fn upload_slices( features: &[[f64; 42]], - targets: &[[f64; 4]], + targets: &[[f64; 6]], ofi: &[[f64; 20]], stream: &Arc, ) -> Result { @@ -307,7 +307,7 @@ impl DqnGpuData { } let feature_dim = 42; - let target_dim = 4; + let target_dim = 6; let estimated_bytes = estimate_vram_bytes(num_bars * (feature_dim + target_dim)); if estimated_bytes > MAX_UPLOAD_BYTES { @@ -400,7 +400,7 @@ impl DqnGpuData { /// Estimated VRAM usage in bytes. pub fn vram_bytes(&self) -> usize { let ofi_dim: usize = 20; - estimate_vram_bytes(self.num_bars * (self.feature_dim + 4 + ofi_dim)) + estimate_vram_bytes(self.num_bars * (self.feature_dim + 6 + ofi_dim)) } /// D2D subrange copy: extract `count` f32 elements at `offset` into a new `CudaSlice`. @@ -442,23 +442,23 @@ impl DqnGpuData { Self::d2d_subrange(&self.features, start * self.feature_dim, count * self.feature_dim, stream) } - /// Get target prices for a single bar as a [4] `CudaSlice`. + /// Get target prices for a single bar as a [6] `CudaSlice`. pub fn bar_targets(&self, bar_idx: usize, stream: &Arc) -> Result, MLError> { if bar_idx >= self.num_bars { return Err(MLError::ModelError(format!("Bar {bar_idx} target out of range (num_bars={})", self.num_bars))); } - Self::d2d_subrange(&self.targets, bar_idx * 4, 4, stream) + Self::d2d_subrange(&self.targets, bar_idx * 6, 6, stream) } - /// Get target prices for a batch as [batch_size * 4] `CudaSlice`. + /// Get target prices for a batch as [batch_size * 6] `CudaSlice`. pub fn batch_targets(&self, start: usize, count: usize, stream: &Arc) -> Result, MLError> { let count = count.min(self.num_bars.saturating_sub(start)); - Self::d2d_subrange(&self.targets, start * 4, count * 4, stream) + Self::d2d_subrange(&self.targets, start * 6, count * 6, stream) } - /// Get target values for a bar as a GPU-resident [4] `CudaSlice`. + /// Get target values for a bar as a GPU-resident [6] `CudaSlice`. /// - /// Tensor order: [current_close_preproc, next_close_preproc, current_close_raw, next_close_raw]. + /// Tensor order: [preproc_close, preproc_next, raw_close, raw_next, raw_open, mid_price_open]. /// Stays on GPU. Callers that need CPU scalars should extract at their own boundary. pub fn bar_target_values(&self, bar_idx: usize, stream: &Arc) -> Result, MLError> { self.bar_targets(bar_idx, stream) @@ -629,7 +629,7 @@ impl DqnGpuData { /// /// # Usage /// ```ignore -/// let mut pool = GpuBufferPool::new(300_000, 51, 4); +/// let mut pool = GpuBufferPool::new(300_000, 51, 6); /// for fold in folds { /// let gpu_data = pool.upload_dqn(&fold_data, &device)?; /// // ... train on gpu_data ... @@ -1030,12 +1030,12 @@ mod tests { #[test] fn test_gpu_buffer_pool_basic() { let stream = cuda_stream(); - let mut pool = GpuBufferPool::new(1000, 42, 4); + let mut pool = GpuBufferPool::new(1000, 42, 6); assert_eq!(pool.max_bars, 1000); - assert_eq!(pool.staging_bytes(), (1000 * 42 + 1000 * 4) * 4); + assert_eq!(pool.staging_bytes(), (1000 * 42 + 1000 * 6) * 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])) + .map(|i| ([i as f64 * 0.01; 42], vec![100.0, 101.0, 100.5, 101.5, 99.5, 100.0])) .collect(); let gpu_data = pool.upload_dqn(&data, &stream).expect("upload"); assert_eq!(gpu_data.num_bars, 50); @@ -1045,23 +1045,23 @@ mod tests { #[test] fn test_gpu_buffer_pool_reuse() { let stream = cuda_stream(); - let mut pool = GpuBufferPool::new(500, 42, 4); + let mut pool = GpuBufferPool::new(500, 42, 6); 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])) + .map(|i| ([i as f64 * 0.01; 42], vec![100.0, 101.0, 100.5, 101.5, 99.5, 100.0])) .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])) + .map(|i| ([i as f64 * 0.02; 42], vec![200.0, 201.0, 200.5, 201.5, 199.5, 200.0])) .collect(); let g2 = pool.upload_dqn(&data2, &stream).expect("upload2"); assert_eq!(g2.num_bars, 200); // Download target values and verify let t = g2.bar_target_values(0, &stream).expect("targets"); - let mut host = vec![0.0_f32; 4]; + let mut host = vec![0.0_f32; 6]; dtoh_f32(&stream, &t, &mut host).expect("DtoH"); assert!((host[0] - 200.0).abs() < 1.0); } @@ -1069,10 +1069,10 @@ mod tests { #[test] fn test_gpu_buffer_pool_grow() { let stream = cuda_stream(); - let mut pool = GpuBufferPool::new(10, 42, 4); + let mut pool = GpuBufferPool::new(10, 42, 6); let data: Vec<([f64; 42], Vec)> = (0..50) - .map(|i| ([i as f64; 42], vec![1.0, 2.0, 3.0, 4.0])) + .map(|i| ([i as f64; 42], vec![1.0, 2.0, 3.0, 4.0, 0.5, 1.0])) .collect(); let gpu_data = pool.upload_dqn(&data, &stream).expect("upload"); assert_eq!(gpu_data.num_bars, 50); @@ -1082,7 +1082,7 @@ mod tests { #[test] fn test_gpu_buffer_pool_empty() { let stream = cuda_stream(); - let mut pool = GpuBufferPool::new(100, 42, 4); + let mut pool = GpuBufferPool::new(100, 42, 6); let data: Vec<([f64; 42], Vec)> = vec![]; assert!(pool.upload_dqn(&data, &stream).is_err()); } diff --git a/crates/ml/src/fxcache.rs b/crates/ml/src/fxcache.rs index da5951f43..a4553e13d 100644 --- a/crates/ml/src/fxcache.rs +++ b/crates/ml/src/fxcache.rs @@ -10,9 +10,9 @@ //! ┌───────────────────────────────────────────────────┐ //! │ FxCacheHeader (64 bytes) │ //! │ magic [u8; 8] = b"FXCACHE\0" │ -//! │ version u16 = 1 (f32) │ +//! │ version u16 = 3 (f32) │ //! │ feat_dim u16 = 42 │ -//! │ target_dim u16 = 4 │ +//! │ target_dim u16 = 6 │ //! │ ofi_dim u16 = 20 │ //! │ bar_count u64 │ //! │ cache_key [u8; 32] (SHA256 raw bytes) │ @@ -20,7 +20,7 @@ //! └───────────────────────────────────────────────────┘ //! │ Body (bar_count records) │ //! │ Each record starts with an i64 timestamp (ns). │ -//! │ Version 1: [i64 ts][66 × f32] = 272 bytes/bar │ +//! │ Version 3: [i64 ts][68 × f32] = 280 bytes/bar │ //! └───────────────────────────────────────────────────┘ //! ``` @@ -40,18 +40,21 @@ const HEADER_SIZE: usize = 64; /// Cache format version. Bump on ANY format change (OFI_DIM, FEAT_DIM, etc.). /// Stale cache files with wrong version are auto-detected and rejected by validate(). /// The ensure-fxcache Argo step catches the error and regenerates. -pub const FXCACHE_VERSION: u16 = 2; // v2: OFI_DIM 8→20 (20 microstructure features) +pub const FXCACHE_VERSION: u16 = 3; // v3: TARGET_DIM 4→6 (added raw_open, mid_price_open) /// Feature vector dimensionality. const FEAT_DIM: usize = 42; -/// Target vector dimensionality (close, next_close, raw_close, raw_next). -const TARGET_DIM: usize = 4; +/// Target vector dimensionality (close, next_close, raw_close, raw_next, raw_open, mid_price_open). +const TARGET_DIM: usize = 6; + +/// Legacy v2 target dimension (4 slots). +const LEGACY_TARGET_DIM: usize = 4; /// OFI vector dimensionality (8 base + 12 extended microstructure features). pub const OFI_DIM: usize = 20; -/// Total f64 values per record: features + targets + OFI = 42 + 4 + 20 = 66. +/// Total f64 values per record: features + targets + OFI = 42 + 6 + 20 = 68. const RECORD_F64_COUNT: usize = FEAT_DIM + TARGET_DIM + OFI_DIM; /// Total f32 values per record (same count as f64, no padding needed). @@ -68,7 +71,7 @@ pub struct FxCacheHeader { pub version: u16, /// Feature dimension (42). pub feat_dim: u16, - /// Target dimension (4). + /// Target dimension (6; legacy v2 files have 4). pub target_dim: u16, /// OFI dimension (20). pub ofi_dim: u16, @@ -98,6 +101,9 @@ impl FxCacheHeader { } /// Validate header integrity. + /// + /// Accepts both v2 (target_dim=4) and v3 (target_dim=6) files. + /// v2 files are read with legacy padding (tgt[4]=raw_close, tgt[5]=raw_close). pub fn validate(&self) -> Result<()> { if self.magic != FXCACHE_MAGIC { bail!( @@ -106,9 +112,10 @@ impl FxCacheHeader { self.magic ); } - if self.version != FXCACHE_VERSION { + // Accept v2 (legacy 4-target) and v3 (current 6-target) + if self.version != FXCACHE_VERSION && self.version != 2 { bail!( - "Stale FxCache version: {} (expected {}). Delete and regenerate.", + "Stale FxCache version: {} (expected {} or 2). Delete and regenerate.", self.version, FXCACHE_VERSION ); } @@ -119,11 +126,14 @@ impl FxCacheHeader { self.feat_dim ); } - if self.target_dim as usize != TARGET_DIM { + // v2 files have target_dim=4, v3 files have target_dim=6 + let td = self.target_dim as usize; + if td != TARGET_DIM && td != LEGACY_TARGET_DIM { bail!( - "Target dimension mismatch: expected {}, got {}", + "Target dimension mismatch: expected {} or {}, got {}", TARGET_DIM, - self.target_dim + LEGACY_TARGET_DIM, + td ); } if self.ofi_dim as usize != OFI_DIM { @@ -194,7 +204,7 @@ pub struct FxCacheData { pub timestamps: Vec, /// Feature vectors, one per bar (42 elements each). pub features: Vec<[f64; FEAT_DIM]>, - /// Target vectors, one per bar (4 elements each). + /// Target vectors, one per bar (6 elements each). pub targets: Vec<[f64; TARGET_DIM]>, /// OFI vectors, one per bar (20 elements each). pub ofi: Vec<[f64; OFI_DIM]>, @@ -215,7 +225,7 @@ pub struct FxCacheData { /// /// * `path` — Output file path (parent directories are created automatically) /// * `features` — Slice of 42-element feature vectors -/// * `targets` — Slice of 4-element target vectors +/// * `targets` — Slice of 6-element target vectors /// * `ofi` — Slice of 20-element OFI vectors /// * `timestamps` — Per-bar timestamps (nanoseconds since Unix epoch) /// * `cache_key` — SHA256 key (raw 32 bytes) @@ -282,7 +292,7 @@ pub fn write_fxcache( Ok(total_bytes) } -/// Write body in f32 format (version 1): [i64 ts] + 66 × f32 = 272 bytes per bar. +/// Write body in f32 format: [i64 ts] + (42+6+20) × f32 = 280 bytes per bar. fn write_body_f32( writer: &mut BufWriter, features: &[[f64; FEAT_DIM]], @@ -341,9 +351,12 @@ pub fn load_fxcache(path: &Path) -> Result { header.validate()?; let bar_count = header.bar_count as usize; + let is_legacy_v2 = header.version == 2 && header.target_dim as usize == LEGACY_TARGET_DIM; + let on_disk_target_dim = header.target_dim as usize; - // Sanity-check file size: [i64 ts] + 66 × f32 = 272 bytes/bar - let expected_body = bar_count as u64 * (8 + RECORD_F32_COUNT as u64 * 4); + // Sanity-check file size + let on_disk_record_f32 = FEAT_DIM + on_disk_target_dim + OFI_DIM; + let expected_body = bar_count as u64 * (8 + on_disk_record_f32 as u64 * 4); let expected_total = HEADER_SIZE as u64 + expected_body; if file_len < expected_total { bail!( @@ -353,12 +366,18 @@ pub fn load_fxcache(path: &Path) -> Result { ); } - // Read body (f32 format) - let (timestamps, features, targets, ofi) = read_body_f32(&mut reader, bar_count)?; + // Read body (f32 format) — handles v2 (4-target) and v3 (6-target) + let (timestamps, features, targets, ofi) = if is_legacy_v2 { + read_body_f32_legacy_v2(&mut reader, bar_count)? + } else { + read_body_f32(&mut reader, bar_count)? + }; info!( - "FxCache loaded: {} bars, v1 (f32) from {:?}", + "FxCache loaded: {} bars, v{} (f32{}) from {:?}", bar_count, + header.version, + if is_legacy_v2 { ", legacy 4→6 padded" } else { "" }, path ); @@ -375,7 +394,7 @@ pub fn load_fxcache(path: &Path) -> Result { }) } -/// Read body in f32 format (version 1), converting f32 back to f64 on load. +/// Read body in f32 format (v3 / current), converting f32 back to f64 on load. fn read_body_f32( reader: &mut BufReader, bar_count: usize, @@ -416,6 +435,62 @@ fn read_body_f32( Ok((timestamps, features, targets, ofi)) } +/// Read body from legacy v2 files (target_dim=4), padding to 6-slot targets. +/// +/// v2 layout: [preproc_close, preproc_next, raw_close, raw_next] +/// Padded to: [preproc_close, preproc_next, raw_close, raw_next, raw_close, raw_close] +/// tgt[4] = raw_close (fallback for raw_open — unavailable in v2) +/// tgt[5] = raw_close (fallback for mid_price_open — unavailable in v2) +fn read_body_f32_legacy_v2( + reader: &mut BufReader, + bar_count: usize, +) -> Result<(Vec, Vec<[f64; FEAT_DIM]>, Vec<[f64; TARGET_DIM]>, Vec<[f64; OFI_DIM]>)> { + let mut timestamps = Vec::with_capacity(bar_count); + let mut features = Vec::with_capacity(bar_count); + let mut targets = Vec::with_capacity(bar_count); + let mut ofi = Vec::with_capacity(bar_count); + let mut i64_buf = [0u8; 8]; + let mut f32_buf = [0u8; 4]; + + for _ in 0..bar_count { + reader.read_exact(&mut i64_buf)?; + timestamps.push(i64::from_le_bytes(i64_buf)); + + let mut feat = [0.0_f64; FEAT_DIM]; + for slot in &mut feat { + reader.read_exact(&mut f32_buf)?; + *slot = f32::from_le_bytes(f32_buf) as f64; + } + features.push(feat); + + // Read 4 legacy target slots + let mut legacy_tgt = [0.0_f64; LEGACY_TARGET_DIM]; + for slot in &mut legacy_tgt { + reader.read_exact(&mut f32_buf)?; + *slot = f32::from_le_bytes(f32_buf) as f64; + } + // Pad to 6: tgt[4]=raw_close, tgt[5]=raw_close + let raw_close = legacy_tgt[2]; + let mut tgt = [0.0_f64; TARGET_DIM]; + tgt[0] = legacy_tgt[0]; + tgt[1] = legacy_tgt[1]; + tgt[2] = legacy_tgt[2]; + tgt[3] = legacy_tgt[3]; + tgt[4] = raw_close; // raw_open fallback + tgt[5] = raw_close; // mid_price_open fallback + targets.push(tgt); + + let mut ofi_row = [0.0_f64; OFI_DIM]; + for slot in &mut ofi_row { + reader.read_exact(&mut f32_buf)?; + *slot = f32::from_le_bytes(f32_buf) as f64; + } + ofi.push(ofi_row); + } + + Ok((timestamps, features, targets, ofi)) +} + // ── Finder ─────────────────────────────────────────────────────────────────── #[cfg(test)] @@ -427,7 +502,7 @@ mod has_ofi_tests { let dir = tempfile::tempdir().unwrap(); let path = dir.path().join("test_ofi_true.fxcache"); let features = vec![[1.0_f64; 42]; 10]; - let targets = vec![[0.0_f64; 4]; 10]; + let targets = vec![[0.0_f64; 6]; 10]; let ofi = vec![[0.5_f64; OFI_DIM]; 10]; let timestamps = vec![1_i64; 10]; let key = [0u8; 32]; @@ -443,7 +518,7 @@ mod has_ofi_tests { let dir = tempfile::tempdir().unwrap(); let path = dir.path().join("test_ofi_false.fxcache"); let features = vec![[1.0_f64; 42]; 10]; - let targets = vec![[0.0_f64; 4]; 10]; + let targets = vec![[0.0_f64; 6]; 10]; let ofi = vec![[0.0_f64; OFI_DIM]; 10]; let timestamps = vec![1_i64; 10]; let key = [0u8; 32]; @@ -594,7 +669,7 @@ mod discover_tests { "0000000000000000000000000000000000000000000000000000000000000000.fxcache", ); let features = vec![[1.0_f64; 42]; 5]; - let targets = vec![[0.0_f64; 4]; 5]; + let targets = vec![[0.0_f64; 6]; 5]; let ofi = vec![[0.0_f64; OFI_DIM]; 5]; let timestamps = vec![1_i64; 5]; write_fxcache(&path, &features, &targets, &ofi, ×tamps, [0u8; 32], false) diff --git a/crates/ml/src/hyperopt/adapters/dqn.rs b/crates/ml/src/hyperopt/adapters/dqn.rs index a63763f0b..6e0b59319 100644 --- a/crates/ml/src/hyperopt/adapters/dqn.rs +++ b/crates/ml/src/hyperopt/adapters/dqn.rs @@ -910,8 +910,8 @@ impl DQNTrainer { let mut targets = Vec::with_capacity(total); for (feat, tgt) in train_data.iter().chain(val_data.iter()) { features.push(*feat); - let mut t = [0.0_f64; 4]; - for (i, v) in tgt.iter().enumerate().take(4) { + let mut t = [0.0_f64; 6]; + for (i, v) in tgt.iter().enumerate().take(6) { t[i] = *v; } targets.push(t); diff --git a/crates/ml/src/trainers/dqn/data_loading.rs b/crates/ml/src/trainers/dqn/data_loading.rs index cf47afbd6..c498b994f 100644 --- a/crates/ml/src/trainers/dqn/data_loading.rs +++ b/crates/ml/src/trainers/dqn/data_loading.rs @@ -423,25 +423,29 @@ impl DQNTrainer { } // Create training data pairs (features, target) - // Target layout: [preproc_close, preproc_next, raw_close, raw_next] + // Target layout: [preproc_close, preproc_next, raw_close, raw_next, raw_open, mid_price_open] // [0:1] = network input (raw prices here; preprocessing applied in feature vector) // [2:3] = raw dollar prices for portfolio simulation P&L + tx costs + // [4] = raw open price + // [5] = mid-price at bar open (MBP-10 midpoint; fallback to raw_open) let mut training_data = Vec::new(); for i in 0..feature_vectors.len().saturating_sub(1) { let raw_current = all_ohlcv_bars[i + 50].close; let raw_next = all_ohlcv_bars[i + 1 + 50].close; + let raw_open = all_ohlcv_bars[i + 50].open; training_data.push(( feature_vectors[i], - vec![raw_current, raw_next, raw_current, raw_next], + vec![raw_current, raw_next, raw_current, raw_next, raw_open, raw_open], )); } // Last sample targets itself if !feature_vectors.is_empty() { let idx = all_ohlcv_bars.len() - 1; let raw_close = all_ohlcv_bars[idx].close; + let raw_open = all_ohlcv_bars[idx].open; training_data.push(( feature_vectors[feature_vectors.len() - 1], - vec![raw_close, raw_close, raw_close, raw_close], + vec![raw_close, raw_close, raw_close, raw_close, raw_open, raw_open], )); } diff --git a/crates/ml/src/trainers/dqn/trainer/constructor.rs b/crates/ml/src/trainers/dqn/trainer/constructor.rs index 1d75b1337..9c13f6cd9 100644 --- a/crates/ml/src/trainers/dqn/trainer/constructor.rs +++ b/crates/ml/src/trainers/dqn/trainer/constructor.rs @@ -525,8 +525,8 @@ impl DQNTrainer { // GPU pipeline: pre-allocate staging buffers on CUDA devices let buffer_pool = matches!(device, MlDevice::Cuda { .. }).then(|| { - info!("GpuBufferPool: pre-allocated staging buffers (100k bars, 42 features, 4 targets)"); - crate::cuda_pipeline::GpuBufferPool::new(100_000, 42, 4) + info!("GpuBufferPool: pre-allocated staging buffers (100k bars, 42 features, 6 targets)"); + crate::cuda_pipeline::GpuBufferPool::new(100_000, 42, 6) }); // GPU pipeline: double-buffered loader for zero-downtime fold transitions diff --git a/crates/ml/src/trainers/dqn/trainer/mod.rs b/crates/ml/src/trainers/dqn/trainer/mod.rs index 980c0e9e6..b7656d7eb 100644 --- a/crates/ml/src/trainers/dqn/trainer/mod.rs +++ b/crates/ml/src/trainers/dqn/trainer/mod.rs @@ -639,9 +639,9 @@ impl DQNTrainer { f.copy_from_slice(&fv[..42]); f }).collect(); - let targets: Vec<[f64; 4]> = training_data.iter().map(|(_, tgt)| { - let mut t = [0.0_f64; 4]; - for (i, &v) in tgt.iter().enumerate().take(4) { t[i] = v; } + let targets: Vec<[f64; 6]> = training_data.iter().map(|(_, tgt)| { + let mut t = [0.0_f64; 6]; + for (i, &v) in tgt.iter().enumerate().take(6) { t[i] = v; } t }).collect(); let val_features: Vec<[f64; 42]> = val_data.iter().map(|(fv, _)| { @@ -649,9 +649,9 @@ impl DQNTrainer { f.copy_from_slice(&fv[..42]); f }).collect(); - let val_targets: Vec<[f64; 4]> = val_data.iter().map(|(_, tgt)| { - let mut t = [0.0_f64; 4]; - for (i, &v) in tgt.iter().enumerate().take(4) { t[i] = v; } + let val_targets: Vec<[f64; 6]> = val_data.iter().map(|(_, tgt)| { + let mut t = [0.0_f64; 6]; + for (i, &v) in tgt.iter().enumerate().take(6) { t[i] = v; } t }).collect(); @@ -664,14 +664,14 @@ impl DQNTrainer { } /// Train one fold from contiguous feature/target slices. - /// Zero Vec allocation — targets read as fixed-size &[[f64; 4]]. + /// Zero Vec allocation — targets read as fixed-size &[[f64; 6]]. /// /// Caller must call `set_val_data_from_slices` and `set_training_range` before /// this method. Val data and OFI offset are already set on `self`. pub async fn train_fold_from_slices( &mut self, features: &[[f64; 42]], - targets: &[[f64; 4]], + targets: &[[f64; 6]], checkpoint_callback: F, ) -> Result where @@ -685,8 +685,8 @@ impl DQNTrainer { info!("train_fold_from_slices: {} bars", features.len()); // Build lightweight training_data with fixed-size arrays — zero heap alloc per element. - // Both [f64; 42] and [f64; 4] are Copy types that live entirely on the stack. - let training_data: Vec<([f64; 42], [f64; 4])> = features + // Both [f64; 42] and [f64; 6] are Copy types that live entirely on the stack. + let training_data: Vec<([f64; 42], [f64; 6])> = features .iter() .zip(targets.iter()) .map(|(f, t)| (*f, *t)) @@ -722,9 +722,9 @@ impl DQNTrainer { f.copy_from_slice(&fv[..42]); f }).collect(); - let targets: Vec<[f64; 4]> = training_data.iter().map(|(_, tgt)| { - let mut t = [0.0_f64; 4]; - for (i, &v) in tgt.iter().enumerate().take(4) { t[i] = v; } + let targets: Vec<[f64; 6]> = training_data.iter().map(|(_, tgt)| { + let mut t = [0.0_f64; 6]; + for (i, &v) in tgt.iter().enumerate().take(6) { t[i] = v; } t }).collect(); @@ -1211,7 +1211,7 @@ impl DQNTrainer { pub async fn init_from_fxcache( &mut self, features: &[[f64; 42]], - targets: &[[f64; 4]], + targets: &[[f64; 6]], ofi: &[[f64; 20]], ) -> Result<()> { let stream = self.cuda_stream.as_ref() @@ -1229,7 +1229,7 @@ impl DQNTrainer { tracing::info!( "init_from_fxcache: {} bars uploaded to GPU ({:.1} MB, OFI=true)", features.len(), - (features.len() * (42 + 4 + 20) * 4) as f64 / 1_048_576.0, + (features.len() * (42 + 6 + 20) * 4) as f64 / 1_048_576.0, ); self.gpu_data = Some(gpu_data); @@ -1251,7 +1251,7 @@ impl DQNTrainer { pub fn set_val_data_from_slices( &mut self, features: &[[f64; 42]], - targets: &[[f64; 4]], + targets: &[[f64; 6]], ofi_val_offset: usize, ) { self.val_data = features.iter().zip(targets.iter()) diff --git a/crates/ml/src/trainers/dqn/trainer/tests.rs b/crates/ml/src/trainers/dqn/trainer/tests.rs index 85e540854..34f4575be 100644 --- a/crates/ml/src/trainers/dqn/trainer/tests.rs +++ b/crates/ml/src/trainers/dqn/trainer/tests.rs @@ -439,7 +439,7 @@ 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<([f64; 42], [f64; 6])> = vec![]; let checkpoint_callback = |_, _, _| Ok(String::new()); let result = trainer diff --git a/crates/ml/src/trainers/dqn/trainer/training_loop.rs b/crates/ml/src/trainers/dqn/trainer/training_loop.rs index e8d1d0e78..3957e47b7 100644 --- a/crates/ml/src/trainers/dqn/trainer/training_loop.rs +++ b/crates/ml/src/trainers/dqn/trainer/training_loop.rs @@ -55,13 +55,13 @@ impl DQNTrainer { // Main training loop — fixed-size arrays, zero heap alloc per bar // ═══════════════════════════════════════════════════════════════════════ - /// Main training loop accepting `&[([f64; 42], [f64; 4])]`. + /// Main training loop accepting `&[([f64; 42], [f64; 6])]`. /// /// Requires `gpu_data` and `targets_raw_cuda` to already be populated /// (via `init_from_fxcache`). The data is already resident on GPU. pub(crate) async fn train_with_data_full_loop_slices( &mut self, - training_data: &[([f64; 42], [f64; 4])], + training_data: &[([f64; 42], [f64; 6])], mut checkpoint_callback: F, ) -> Result where @@ -849,7 +849,7 @@ impl DQNTrainer { pub(crate) async fn init_gpu_raw_buffers_from_slices( &mut self, features: &[[f64; 42]], - targets: &[[f64; 4]], + targets: &[[f64; 6]], ) -> Result<()> { if self.targets_raw_cuda.is_some() { return Ok(()); @@ -876,7 +876,7 @@ impl DQNTrainer { .map_err(|e| anyhow::anyhow!("features_raw upload: {e}"))? ); - tracing::info!("CUDA raw buffers uploaded from slices: {} bars x 42+4", num_bars); + tracing::info!("CUDA raw buffers uploaded from slices: {} bars x 42+6", num_bars); Ok(()) } @@ -1026,12 +1026,12 @@ impl DQNTrainer { // Helper: Phase 3 — GPU experience collection // ═══════════════════════════════════════════════════════════════════════ - /// Collect GPU experiences from fixed-size `&[([f64; 42], [f64; 4])]` slices. - /// Only `.len()` and `.0` (features) are accessed — the `[f64; 4]` targets + /// Collect GPU experiences from fixed-size `&[([f64; 42], [f64; 6])]` slices. + /// Only `.len()` and `.0` (features) are accessed — the `[f64; 6]` targets /// are never read here (GPU raw buffers are already uploaded). pub(crate) async fn collect_gpu_experiences_slices( &mut self, - training_data: &[([f64; 42], [f64; 4])], + training_data: &[([f64; 42], [f64; 6])], ) -> Result { // Feature normalization stats calculation epoch (kept for API compat) let _stats_collection_epochs = { diff --git a/crates/ml/tests/fxcache_roundtrip_test.rs b/crates/ml/tests/fxcache_roundtrip_test.rs index 98b6b2710..a0ed3f369 100644 --- a/crates/ml/tests/fxcache_roundtrip_test.rs +++ b/crates/ml/tests/fxcache_roundtrip_test.rs @@ -22,12 +22,12 @@ fn make_features(n: usize) -> Vec<[f64; 42]> { } /// Deterministic target vector for bar `i`. -fn make_targets(n: usize) -> Vec<[f64; 4]> { +fn make_targets(n: usize) -> Vec<[f64; 6]> { (0..n) .map(|i| { - let mut row = [0.0_f64; 4]; - for j in 0..4 { - row[j] = 5000.0 + (i * 4 + j) as f64 * 0.25; + let mut row = [0.0_f64; 6]; + for j in 0..6 { + row[j] = 5000.0 + (i * 6 + j) as f64 * 0.25; } row }) @@ -109,7 +109,7 @@ fn test_fxcache_f32_roundtrip() { data.features[i][j] ); } - for j in 0..4 { + for j in 0..6 { let diff = (data.targets[i][j] - targets[i][j]).abs(); // f32 at ~5000 has resolution of ~0.0005, allow generous margin. assert!( @@ -186,7 +186,7 @@ fn test_fxcache_empty() { // Writer rejects 0 bars. let features: Vec<[f64; 42]> = vec![]; - let targets: Vec<[f64; 4]> = vec![]; + let targets: Vec<[f64; 6]> = vec![]; let ofi: Vec<[f64; 20]> = vec![]; let timestamps: Vec = vec![]; let key = test_cache_key(); @@ -202,9 +202,9 @@ fn test_fxcache_empty() { // Hand-craft a valid-magic header with bar_count=0. let mut header = [0u8; 64]; header[0..8].copy_from_slice(b"FXCACHE\0"); - header[8..10].copy_from_slice(&1u16.to_le_bytes()); // version=1 + header[8..10].copy_from_slice(&3u16.to_le_bytes()); // version=3 header[10..12].copy_from_slice(&42u16.to_le_bytes()); // feat_dim - header[12..14].copy_from_slice(&4u16.to_le_bytes()); // target_dim + header[12..14].copy_from_slice(&6u16.to_le_bytes()); // target_dim header[14..16].copy_from_slice(&20u16.to_le_bytes()); // ofi_dim header[16..24].copy_from_slice(&0u64.to_le_bytes()); // bar_count=0 // cache_key + reserved stay zeroed. diff --git a/docs/superpowers/plans/2026-04-19-kernel-optimization-phase3.md b/docs/superpowers/plans/2026-04-19-kernel-optimization-phase3.md new file mode 100644 index 000000000..ba45cfe1f --- /dev/null +++ b/docs/superpowers/plans/2026-04-19-kernel-optimization-phase3.md @@ -0,0 +1,419 @@ +# Phase 3 Kernel Optimization Implementation Plan + +> **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking. + +**Goal:** Eliminate the 3 remaining serial-batch bottleneck kernels found by nsys Phase 2 profiling — `reward_rank_normalize` (18.3%), `bias_grad_reduce_f32_kernel` (11.4%), and `curiosity_bias_grad_reduce` (0.9%) — to reach <40ms/step on H100. + +**Architecture:** (1) Replace O(N²) rank-normalize with radix-sort-based O(N log N) rank. (2) Convert the main DQN backward bias-grad-reduce from serial batch loop to 2-phase shared-memory reduction (same proven pattern used in IQN/IQL/attention). (3) Convert curiosity bias-grad-reduce to the same 2-phase pattern. + +**Tech Stack:** CUDA 12.4, CUB radix sort (device-wide), cudarc vendored bindings, existing 2-phase reduce pattern from `iqn_dual_head_kernel.cu`. + +**Dimensions:** B=8192, SH2=256, out_dim varies (3–256), CUR_OUTPUT=42, CUR_HIDDEN=128. + +--- + +## File Structure + +| File | Responsibility | +|------|---------------| +| `crates/ml/src/cuda_pipeline/reward_shaping_kernel.cu` | Reward rank-normalize kernel (rewrite to sort-based) | +| `crates/ml/src/cuda_pipeline/gpu_experience_collector.rs` | Rust wiring for reward rank-normalize launch | +| `crates/ml/src/cuda_pipeline/backward_kernels.cu` | Main DQN backward bias-grad-reduce (rewrite to 2-phase) | +| `crates/ml/src/cuda_pipeline/batched_backward.rs` | Rust wiring for bias-grad-reduce launch + partials buffer | +| `crates/ml/src/cuda_pipeline/curiosity_training_kernel.cu` | Curiosity bias-grad-reduce (rewrite to 2-phase) | +| `crates/ml/src/cuda_pipeline/gpu_curiosity_trainer.rs` | Rust wiring for curiosity bias-grad-reduce | + +--- + +### Task 1: Rewrite reward_rank_normalize — sort-based O(N log N) rank (18.3%, 10.6s) + +The current kernel is O(N²): each of N threads scans all N elements to count how many are ≤ its value. With N=~4M (4096 episodes × ~1000 steps), this takes 10.6 seconds for a single call. + +**Fix:** Sort the absolute Sharpe values using CUB `DeviceRadixSort::SortPairs`, then compute rank from sorted position. O(N log N) via radix sort — should complete in <50ms for 4M elements on H100. + +**Files:** +- Modify: `crates/ml/src/cuda_pipeline/reward_shaping_kernel.cu` +- Modify: `crates/ml/src/cuda_pipeline/gpu_experience_collector.rs` + +- [ ] **Step 1: Write 2 new helper kernels in `reward_shaping_kernel.cu`** + +Keep the old `reward_rank_normalize` kernel. Add these BEFORE it: + +```c +/* reward_compute_abs_sharpe — Compute |sharpe[i]| and original index for sort. + * Grid: ceil(N/256), Block: 256. + */ +extern "C" __global__ void reward_compute_abs_sharpe( + const float* __restrict__ rewards_in, + float* __restrict__ abs_sharpe_out, /* [N] sort keys */ + int* __restrict__ indices_out, /* [N] original indices */ + int N, float std_ema, float return_mean_ema +) { + int i = blockIdx.x * blockDim.x + threadIdx.x; + if (i >= N) return; + float raw = rewards_in[i]; + float sharpe = (std_ema > 1e-6f) ? (raw - return_mean_ema) / std_ema : raw; + abs_sharpe_out[i] = fabsf(sharpe); + indices_out[i] = i; +} + +/* reward_scatter_rank — After sort, write rank = position/N to output, preserving sign. + * sorted_indices[pos] = original_index. rank = pos/N. + * Grid: ceil(N/256), Block: 256. + */ +extern "C" __global__ void reward_scatter_rank( + const float* __restrict__ rewards_in, + const int* __restrict__ sorted_indices, + float* __restrict__ rewards_out, + int N, float std_ema, float return_mean_ema +) { + int pos = blockIdx.x * blockDim.x + threadIdx.x; + if (pos >= N) return; + int orig_idx = sorted_indices[pos]; + float raw = rewards_in[orig_idx]; + float sharpe = (std_ema > 1e-6f) ? (raw - return_mean_ema) / std_ema : raw; + int is_trade = (fabsf(sharpe) > 0.001f) ? 1 : 0; + if (!is_trade) { + rewards_out[orig_idx] = 0.0f; + return; + } + float rank = (float)(pos + 1) / (float)N; /* 1-based rank / N */ + float sign = (sharpe > 0.0f) ? 1.0f : -1.0f; + rewards_out[orig_idx] = sign * rank; +} +``` + +- [ ] **Step 2: Add CUB sort buffers and new kernel handles in `gpu_experience_collector.rs`** + +In the struct, add: +```rust +reward_abs_sharpe_buf: CudaSlice, // [max_N] sort keys +reward_indices_buf: CudaSlice, // [max_N] original indices +reward_sorted_keys_buf: CudaSlice, // [max_N] sorted output +reward_sorted_indices_buf: CudaSlice,// [max_N] sorted indices +reward_sort_temp_buf: CudaSlice, // CUB sort workspace +reward_compute_abs_sharpe_kernel: CudaFunction, +reward_scatter_rank_kernel: CudaFunction, +``` + +max_N = gpu_n_episodes × max_bars (the maximum number of rewards to rank). Allocate CUB temp storage using `cub::DeviceRadixSort::SortPairs` temp size query (via cudarc raw API or pre-computed upper bound of `2 * N * sizeof(f32)` which is safe for radix sort). + +In the constructor, load `reward_compute_abs_sharpe` and `reward_scatter_rank` from the reward_shaping cubin. Allocate the scratch buffers. + +- [ ] **Step 3: Rewrite `shape_rewards()` to use sort-based pipeline** + +Replace the single `reward_rank_normalize` launch with: +1. `reward_compute_abs_sharpe` kernel — compute keys + indices +2. CUB `DeviceRadixSort::SortPairs` — sort keys, permute indices +3. `reward_scatter_rank` kernel — scatter rank to original positions + +For CUB sort, use `cudarc::driver::sys` to call `cub::DeviceRadixSort::SortPairs` via raw CUDA API. If CUB is not directly available through cudarc, use the alternative: `thrust::sort_by_key` via a small custom kernel that calls `__device__` sort, OR implement a simple GPU merge sort. + +Simplest practical approach: use `cudarc`'s built-in `CudaSlice` sort if available, otherwise implement bitonic sort as a fallback (O(N log²N), still far better than O(N²)). + +**Fallback if CUB not available:** Write a 3-kernel bitonic sort: +- `bitonic_compare_swap` kernel, launched log²(N) times with varying step/stage parameters +- Total kernel launches: ~20 for N=4M (log2(4M) ≈ 22 stages × log2(stage) sub-passes) +- Each launch: grid=N/2, block=256, 0.05ms → total ~1ms vs 10.6s + +- [ ] **Step 4: Run smoke tests** + +```bash +SQLX_OFFLINE=true FOXHUNT_TEST_DATA=test_data/futures-baseline cargo test -p ml --lib -- smoke_tests --ignored --nocapture +``` +Expected: 20/20 pass. + +- [ ] **Step 5: Commit** + +```bash +git add crates/ml/src/cuda_pipeline/reward_shaping_kernel.cu crates/ml/src/cuda_pipeline/gpu_experience_collector.rs +git commit -m "perf: reward_rank_normalize O(N²) → sort-based O(N log N) — 10.6s → <50ms" +``` + +--- + +### Task 2: Rewrite bias_grad_reduce_f32_kernel — 2-phase shared-memory reduce (11.4%, 6.6s) + +The main DQN backward bias gradient kernel in `backward_kernels.cu` uses the serial anti-pattern: 1 thread per `out_dim` element, loops over `batch_size=8192`. Called 28,321 times during 993 training steps (~28 per step for all FC layers). Total: 6.6 seconds. + +**Fix:** Same 2-phase pattern proven in IQN/IQL/attention bias grad reduces. + +**Files:** +- Modify: `crates/ml/src/cuda_pipeline/backward_kernels.cu` +- Modify: `crates/ml/src/cuda_pipeline/batched_backward.rs` + +- [ ] **Step 1: Write 2-phase kernels in `backward_kernels.cu`** + +Replace the old `bias_grad_reduce_f32_kernel` with: + +```c +/* Phase 1: block-level shared-memory reduce. + * Grid: (num_blocks, out_dim), Block: 256, shared: 256*sizeof(float). + */ +extern "C" __global__ void bias_grad_reduce_f32_p1( + const float* __restrict__ dy, + float* __restrict__ partials, + int out_dim, int batch_size +) { + extern __shared__ float sdata[]; + int j = blockIdx.y; + int tid = threadIdx.x; + int gid = blockIdx.x * blockDim.x + tid; + float val = (gid < batch_size) ? dy[gid * out_dim + j] : 0.0f; + sdata[tid] = val; + __syncthreads(); + for (int s = blockDim.x / 2; s > 0; s >>= 1) { + if (tid < s) sdata[tid] += sdata[tid + s]; + __syncthreads(); + } + if (tid == 0) partials[blockIdx.x * out_dim + j] = sdata[0]; +} + +/* Phase 2: final reduce across block partials + accumulate into db. + * Grid: ceil(out_dim/256), Block: 256. + * NOTE: db[j] += sum (accumulate, not overwrite) to match existing behavior. + */ +extern "C" __global__ void bias_grad_reduce_f32_p2( + const float* __restrict__ partials, + float* __restrict__ db, + int out_dim, int num_blocks +) { + int j = blockIdx.x * blockDim.x + threadIdx.x; + if (j >= out_dim) return; + float sum = 0.0f; + for (int i = 0; i < num_blocks; i++) sum += partials[i * out_dim + j]; + db[j] += sum; +} +``` + +Note: The old kernel uses `db[j] += sum` (accumulate), so the p2 kernel must also accumulate. + +- [ ] **Step 2: Update `batched_backward.rs` struct + constructor** + +In the `BatchedBackward` struct (or equivalent name), add: +```rust +bias_grad_p1_kernel: CudaFunction, +bias_grad_p2_kernel: CudaFunction, +bias_grad_partials_buf: CudaSlice, // [num_blocks * max_out_dim] +bias_grad_num_blocks: u32, +``` + +`max_out_dim` = maximum `out_dim` across all FC layers = SH2 = 256. +`num_blocks` = `(batch_size + 255) / 256` = 32. +Partials buffer: 32 × 256 × 4 = 32KB — trivial. + +In the constructor, load `bias_grad_reduce_f32_p1` and `bias_grad_reduce_f32_p2` from the backward_kernels cubin. + +- [ ] **Step 3: Rewrite `launch_bias_grad()` to 2-phase pattern** + +```rust +fn launch_bias_grad( + &self, stream: &Arc, + dy: u64, db: u64, out_dim: usize, batch: usize, +) -> Result<(), MLError> { + let out_dim_i32 = out_dim as i32; + let batch_i32 = batch as i32; + let num_blocks = self.bias_grad_num_blocks; + let partials_ptr = self.bias_grad_partials_buf.raw_ptr(); + // Phase 1 + unsafe { + stream.launch_builder(&self.bias_grad_p1_kernel) + .arg(&dy).arg(&partials_ptr) + .arg(&out_dim_i32).arg(&batch_i32) + .launch(LaunchConfig { + grid_dim: (num_blocks, out_dim as u32, 1), + block_dim: (256, 1, 1), + shared_mem_bytes: 256 * 4, + }) + .map_err(|e| MLError::ModelError(format!("bias_grad_reduce_f32_p1: {e}")))?; + } + // Phase 2 + let num_blocks_i32 = num_blocks as i32; + let p2_blocks = ((out_dim + 255) / 256) as u32; + unsafe { + stream.launch_builder(&self.bias_grad_p2_kernel) + .arg(&partials_ptr).arg(&db) + .arg(&out_dim_i32).arg(&num_blocks_i32) + .launch(LaunchConfig { + grid_dim: (p2_blocks, 1, 1), + block_dim: (256, 1, 1), + shared_mem_bytes: 0, + }) + .map_err(|e| MLError::ModelError(format!("bias_grad_reduce_f32_p2: {e}")))?; + } + Ok(()) +} +``` + +- [ ] **Step 4: Run smoke tests** + +```bash +SQLX_OFFLINE=true FOXHUNT_TEST_DATA=test_data/futures-baseline cargo test -p ml --lib -- smoke_tests --ignored --nocapture +``` +Expected: 20/20 pass. + +- [ ] **Step 5: Commit** + +```bash +git add crates/ml/src/cuda_pipeline/backward_kernels.cu crates/ml/src/cuda_pipeline/batched_backward.rs +git commit -m "perf: bias_grad_reduce_f32 — 2-phase shared-memory reduce (0.24ms×28K → <0.01ms×28K)" +``` + +--- + +### Task 3: Rewrite curiosity_bias_grad_reduce — 2-phase pattern (0.9%, 134ms/call) + +Same serial anti-pattern: `out_dim` threads loop over N=8192. Called 4 times (2 bias layers × 2 calls). 134ms per call is extreme for a simple reduction. + +**Fix:** Same 2-phase shared-memory reduce pattern. + +**Files:** +- Modify: `crates/ml/src/cuda_pipeline/curiosity_training_kernel.cu` +- Modify: `crates/ml/src/cuda_pipeline/gpu_curiosity_trainer.rs` + +- [ ] **Step 1: Write 2-phase kernels in `curiosity_training_kernel.cu`** + +Replace the old `curiosity_bias_grad_reduce` with: + +```c +/* Phase 1: block-level shared-memory reduce. + * Grid: (num_blocks, out_dim), Block: 256, shared: 256*sizeof(float). + */ +extern "C" __global__ void curiosity_bias_grad_reduce_p1( + const float* __restrict__ dy, + float* __restrict__ partials, + int N, int out_dim +) { + extern __shared__ float sdata[]; + int d = blockIdx.y; + int tid = threadIdx.x; + int gid = blockIdx.x * blockDim.x + tid; + float val = (gid < N) ? dy[gid * out_dim + d] : 0.0f; + sdata[tid] = val; + __syncthreads(); + for (int s = blockDim.x / 2; s > 0; s >>= 1) { + if (tid < s) sdata[tid] += sdata[tid + s]; + __syncthreads(); + } + if (tid == 0) partials[blockIdx.x * out_dim + d] = sdata[0]; +} + +/* Phase 2: final reduce across block partials. + * Grid: ceil(out_dim/256), Block: 256. + */ +extern "C" __global__ void curiosity_bias_grad_reduce_p2( + const float* __restrict__ partials, + float* __restrict__ grad_b, + int out_dim, int num_blocks +) { + int d = blockIdx.x * blockDim.x + threadIdx.x; + if (d >= out_dim) return; + float sum = 0.0f; + for (int i = 0; i < num_blocks; i++) sum += partials[i * out_dim + d]; + grad_b[d] = sum; +} +``` + +- [ ] **Step 2: Update `gpu_curiosity_trainer.rs` struct + constructor** + +In the struct, add: +```rust +bias_grad_p1_func: CudaFunction, +bias_grad_p2_func: CudaFunction, +bias_grad_partials: CudaSlice, // [num_blocks * max(CUR_OUTPUT, CUR_HIDDEN)] +bias_grad_num_blocks: u32, +``` + +In the constructor, load `curiosity_bias_grad_reduce_p1` and `curiosity_bias_grad_reduce_p2`. Allocate partials buffer: `num_blocks * max(CUR_OUTPUT=42, CUR_HIDDEN=128) = 32 * 128 = 4096 floats = 16KB`. + +- [ ] **Step 3: Replace both `curiosity_bias_grad_reduce` call sites** + +At line ~810 (b2) and ~850 (b1), replace the single kernel launch with 2-phase pattern: + +```rust +// Phase 1: block reduce +let partials_ptr = self.bias_grad_partials.raw_ptr(); +let num_blocks = self.bias_grad_num_blocks; +unsafe { + stream.launch_builder(&self.bias_grad_p1_func) + .arg(&pred_ptr).arg(&partials_ptr) + .arg(&n_i32).arg(&cur_output_i32) + .launch(LaunchConfig { + grid_dim: (num_blocks, CUR_OUTPUT as u32, 1), + block_dim: (256, 1, 1), + shared_mem_bytes: 256 * 4, + }) + .map_err(|e| MLError::ModelError(format!("curiosity_bias_grad_reduce_p1 b2: {e}")))?; +} +// Phase 2: final reduce +let num_blocks_i32 = num_blocks as i32; +unsafe { + stream.launch_builder(&self.bias_grad_p2_func) + .arg(&partials_ptr).arg(&grad_b2_ptr) + .arg(&cur_output_i32).arg(&num_blocks_i32) + .launch(LaunchConfig { + grid_dim: (((CUR_OUTPUT + 255) / 256) as u32, 1, 1), + block_dim: (256, 1, 1), + shared_mem_bytes: 0, + }) + .map_err(|e| MLError::ModelError(format!("curiosity_bias_grad_reduce_p2 b2: {e}")))?; +} +``` + +Repeat the same pattern for the b1 call site with `d_hidden_ptr`, `grad_b1_ptr`, `CUR_HIDDEN`. + +- [ ] **Step 4: Run smoke tests** + +```bash +SQLX_OFFLINE=true FOXHUNT_TEST_DATA=test_data/futures-baseline cargo test -p ml --lib -- smoke_tests --ignored --nocapture +``` +Expected: 20/20 pass. + +- [ ] **Step 5: Commit** + +```bash +git add crates/ml/src/cuda_pipeline/curiosity_training_kernel.cu crates/ml/src/cuda_pipeline/gpu_curiosity_trainer.rs +git commit -m "perf: curiosity_bias_grad_reduce — 2-phase shared-memory reduce (134ms → <0.5ms)" +``` + +--- + +### Task 4: Remove dead old kernels + deploy nsys validation + +**Files:** +- Modify: `crates/ml/src/cuda_pipeline/reward_shaping_kernel.cu` (remove old kernel) +- Modify: `crates/ml/src/cuda_pipeline/backward_kernels.cu` (remove old kernel) +- Modify: `crates/ml/src/cuda_pipeline/curiosity_training_kernel.cu` (remove old kernel) + +- [ ] **Step 1: Remove old `reward_rank_normalize` kernel body** (keep comment noting replacement) +- [ ] **Step 2: Remove old `bias_grad_reduce_f32_kernel`** (replaced by p1+p2) +- [ ] **Step 3: Remove old `curiosity_bias_grad_reduce`** (replaced by p1+p2) + +- [ ] **Step 4: Run smoke tests** + +```bash +SQLX_OFFLINE=true FOXHUNT_TEST_DATA=test_data/futures-baseline cargo test -p ml --lib -- smoke_tests --ignored --nocapture +``` + +- [ ] **Step 5: Commit, push, deploy with nsys** + +```bash +git add -A && git commit -m "cleanup: remove old serial kernels replaced by batch-parallel versions" +git push origin main +./scripts/argo-train.sh dqn --sanitizer nsys --epochs 1 --baseline +``` + +Target: `bias_grad_reduce_f32_kernel` gone from top 10. `reward_rank_normalize` from 10.6s → <50ms. Step time <40ms. + +--- + +## Expected Impact + +| Kernel | Before | After | Savings | +|--------|--------|-------|---------| +| reward_rank_normalize | 10,652ms (1 call) | <50ms | 10.6s one-time | +| bias_grad_reduce_f32_kernel | 6,659ms (28K calls) | <200ms | 6.5s total | +| curiosity_bias_grad_reduce | 539ms (4 calls) | <2ms | 537ms total | +| **Total savings** | **17,850ms** | **<252ms** | **~17.6s** | + +Per-step training time (993 steps in 60s capture): current ~60ms/step → target ~40ms/step after bias_grad_reduce fix. The reward_rank_normalize is one-time per epoch, not per-step. diff --git a/docs/superpowers/specs/2026-04-19-ofi-momentum-dense-reward-design.md b/docs/superpowers/specs/2026-04-19-ofi-momentum-dense-reward-design.md index 82846dabd..59e8e1cc5 100644 --- a/docs/superpowers/specs/2026-04-19-ofi-momentum-dense-reward-design.md +++ b/docs/superpowers/specs/2026-04-19-ofi-momentum-dense-reward-design.md @@ -100,20 +100,46 @@ if (!segment_complete && fabsf(position) > 0.001f && micro_reward_scale > 0.0f) /* Informed flow intensity */ float vol_informed = ofi_cur[6] * ofi_cur[7]; /* VPIN × Kyle's_lambda */ - /* Composite quality — performance-adaptive cost */ + /* ── Pearl 2: Order book center of mass (aggression signal) ── + * Uses all 10 MBP-10 depth levels to detect aggressive vs passive flow. + * buy_com near 1 = aggressive buying (pressure at top of book). + * Precomputed in fxcache as a single feature (book_aggression). */ + float book_aggression = (state_dim >= 83) ? batch_states[(long long)i * state_dim + 82] : 0.0f; + + /* Composite quality — performance-adaptive cost + book aggression */ float quality = flow_momentum * (1.0f + vol_informed) + 0.5f * price_confirm + + 0.3f * sign_pos * book_aggression /* Pearl 2: aggressive flow in our direction */ - spread_penalty; /* adaptive: costly when losing, cheap when winning */ + + /* ── Pearl 3: Retrospective hold quality bonus ── + * Evaluate whether the PREVIOUS bar's Hold decision was correct. + * If model held and mid-price moved favorably > spread, bonus. */ + float prev_mid = ps[PREV_MID_SLOT]; + float mid_now = tgt[5]; /* MBP-10 mid-price at this bar's formation */ + if (prev_mid > 0.0f && mid_now > 0.0f) { + float hold_return = (mid_now - prev_mid) / fmaxf(prev_mid, 1.0f); + float hold_quality = sign_pos * hold_return / fmaxf(atr_pct, 1e-6f); + quality += 0.2f * hold_quality; /* retrospective: was holding correct? */ + } + reward = micro_reward_scale * tanhf(quality / 3.0f); } -/* Update prev-OFI in portfolio state for next bar's delta */ -if (market_dim >= 50) { +/* Update prev-OFI and prev-mid in portfolio state for next bar's delta */ +if (batch_states != NULL && state_dim >= 74) { const float* ofi_cur = batch_states + (long long)i * state_dim + 66; for (int k = 0; k < 8; k++) ps[30 + k] = ofi_cur[k]; } +ps[PREV_MID_SLOT] = tgt[5]; /* MBP-10 mid-price at this decision point */ +ps[PREV_CLOSE_SLOT] = raw_close; /* for DSR actual per-bar returns */ ``` +**PREV_MID_SLOT:** Use slot 22 (currently reserved). Define `#define PREV_MID_SLOT 22`. +**PREV_CLOSE_SLOT:** Use slot 6 (currently reserved). Define `#define PREV_CLOSE_SLOT 6`. + +**MFT mark-to-market:** The micro-reward uses MBP-10 mid-prices at bar BOUNDARIES (imbalance threshold crossings) — the model's actual decision points. This is the TRUE mark-to-market for MFT-style trading, not bar OHLCV. Each imbalance bar represents a burst of directional order flow; the mid-price at formation is where the model would actually execute. + **Note:** `batch_states` must be passed to `env_step_batch` as a new kernel parameter. Currently the kernel receives `features` (42-dim market features) but NOT the full `batch_states` (96-dim with OFI). Add `const float* __restrict__ batch_states` and `int state_dim` parameters. ### Phase D: Rank Normalization Fix @@ -130,13 +156,57 @@ Or simpler: lower threshold to `1e-6f` (only filter true zeros). The original 0. **Recommendation:** Lower to `1e-5f`. This preserves the zero-filtering for truly flat bars while passing all micro-rewards through. +### Phase D2: Novel Features — Pearls + +**Pearl 1: Bar Duration Encoding in Mamba2** + +Imbalance bars are variable-length. 8 bars spanning 30 seconds = volatile regime shift. 8 bars over 2 hours = calm accumulation. Encode `bar_duration_seconds` as a feature alongside the OFI embed in the Mamba2 history input. + +- Precompute bar duration during fxcache creation: `duration = timestamp[bar] - timestamp[bar-1]` +- Store as a normalized feature in the state vector (log-scale, clipped): `log_duration = log(max(duration_secs, 0.1)) / 10.0` +- Add to state layout at index 83 (within existing 96-dim padding) +- Mamba2 input becomes `[h_s2(SH2); ofi_embed(8); log_duration(1)]` = SH2+9 + +The SSM's forget gate learns duration-dependent decay — short bars maintain state (fast regime), long bars decay state (calm, previous patterns less relevant). This is a continuous-time SSM for free. + +**Pearl 2: Order Book Center of Mass (Book Aggression)** + +Precompute from all 10 MBP-10 depth levels during fxcache creation: +```python +buy_com = sum(level_i * bid_depth_i for i in 1..10) / sum(bid_depth_i) +sell_com = sum(level_i * ask_depth_i for i in 1..10) / sum(ask_depth_i) +book_aggression = sell_com - buy_com # positive = aggressive buying (pressure near L1) +``` + +- Store as feature at state index 82 +- Used directly in micro-reward formula (see Phase C above) +- Also feeds into OFI embed MLP: extend input from 16-dim to 18-dim (add book_aggression + log_duration) + +**Pearl 3: Retrospective Hold Quality Bonus** + +At bar t, evaluate whether the PREVIOUS bar's Hold was correct using mid-price change between decision points. Implemented directly in the micro-reward formula (see Phase C above). Teaches Mamba2 which temporal patterns precede profitable holds. + +**State vector layout (final):** +``` +[0..42) market features (OHLCV + technicals) +[42..50) portfolio features (position, equity, etc.) +[50..66) multi-timeframe features +[66..74) OFI raw (8 features from MBP-10) +[74..82) OFI delta (8 pre-computed temporal derivatives) +[82] book_aggression (Pearl 2: order book center of mass) +[83] log_bar_duration (Pearl 1: imbalance bar formation time) +[84..96) padding (12 slots free for future) +``` + +This fits within the existing 96-dim padded state. No state_dim change needed. Mamba2 input widens from SH2+10 to SH2+10 (ofi_embed(8) + book_aggression(1) + log_duration(1)). + ### Phase E: OFI Momentum Embedding MLP -**Purpose:** Compress [raw_ofi(8) ; delta_ofi(8)] = 16-dim into 8-dim learned representation. +**Purpose:** Compress [raw_ofi(8) ; delta_ofi(8) ; book_aggression(1) ; log_duration(1)] = 18-dim into 10-dim learned representation (8 for OFI embed + 2 for book/duration pass-through, or learned 10-dim compression). -**Implementation:** cuBLAS SGEMM `[8, 16] × [16, B] → [8, B]` + elementwise bias+ReLU. +**Implementation:** cuBLAS SGEMM `[10, 18] × [18, B] → [10, B]` + elementwise bias+ReLU. - Same pattern as `trade_plan_forward` (already proven at 45ms/step) -- 136 trainable params (16×8 + 8 bias). Xavier init. +- 190 trainable params (18×10 + 10 bias). Xavier init. - Runs during training forward pass, NOT during experience collection - Must be in graph-captured CUmodule (same isolation as other forward ops) @@ -161,13 +231,13 @@ This fits within the existing 96-dim padded state. No state_dim change needed. ### Phase F: Mamba2 History Enrichment **Current:** `h_history[B, K, SH2]` where SH2=256. -**New:** `h_history[B, K, SH2+8]` — OFI embed (8-dim) appended per timestep. +**New:** `h_history[B, K, SH2+10]` — OFI embed (8-dim) appended per timestep. **Changes required:** 1. `mamba2_update_history` kernel: add `ofi_embed` pointer param, copy `[h_s2(SH2); ofi_embed(8)]` per sample -2. `mamba2_h_history` buffer: reallocate as `B × K × (SH2+8)` -3. W_A, W_B parameter buffers: grow from `[SH2, STATE_D]` to `[SH2+8, STATE_D]` -4. Mamba2 projection GEMM descriptors: destroy and recreate with k=SH2+8 +2. `mamba2_h_history` buffer: reallocate as `B × K × (SH2+10)` +3. W_A, W_B parameter buffers: grow from `[SH2, STATE_D]` to `[SH2+10, STATE_D]` +4. Mamba2 projection GEMM descriptors: destroy and recreate with k=SH2+10 5. `mamba2_scan_projected_fwd/bwd`: these operate on pre-projected values (A_proj, B_proj) — unchanged, STATE_D stays 16 ### Phase G: Mamba2 Backward — Add d_h_history Gradient @@ -176,35 +246,35 @@ This fits within the existing 96-dim padded state. No state_dim change needed. **Fix:** After the existing dW GEMMs, add: ``` -d_h_history[SH2+8, B*K] = W_A^T @ d_gate + W_B^T @ d_x +d_h_history[SH2+10, B*K] = W_A^T @ d_gate + W_B^T @ d_x ``` This is 2 additional cuBLAS GEMMs: -- `d_h_from_A = W_A[SH2+8, STATE_D] @ d_gate[STATE_D, B*K]` → TRANSA=N, TRANSB=N -- `d_h_from_B = W_B[SH2+8, STATE_D] @ d_x[STATE_D, B*K]` → TRANSA=N, TRANSB=N +- `d_h_from_A = W_A[SH2+10, STATE_D] @ d_gate[STATE_D, B*K]` → TRANSA=N, TRANSB=N +- `d_h_from_B = W_B[SH2+10, STATE_D] @ d_x[STATE_D, B*K]` → TRANSA=N, TRANSB=N - `d_h_history = d_h_from_A + d_h_from_B` (elementwise add kernel) Then extract the last 8 dims of d_h_history as `d_ofi_embed_from_mamba2`, and accumulate across K timesteps (sum reduce). **New buffers:** -- `d_h_history[SH2+8, B*K]` — scratch for d_h gradient +- `d_h_history[SH2+10, B*K]` — scratch for d_h gradient - `d_ofi_embed_mamba2[8, B]` — accumulated OFI embed gradient from Mamba2 ### Phase H: Attention Input Enrichment + Backward Gradient Exposure **Current:** Attention input is `saved_input = h_s2[B, D=256]`. -**New:** Attention input is `[h_s2; ofi_embed] = [B, D+8=264]`. +**New:** Attention input is `[h_s2; ofi_embed] = [B, D+10=264]`. **Forward changes:** -- W_Q, W_K, W_V: input dimension grows from D to D+8. Matrices become `[D, D+8]`. -- `saved_input` buffer grows from `[D, B]` to `[D+8, B]` -- All 4 forward GEMM descriptors (Q/K/V/O): W_Q/K/V get k=D+8, W_O stays k=D +- W_Q, W_K, W_V: input dimension grows from D to D+10. Matrices become `[D, D+10]`. +- `saved_input` buffer grows from `[D, B]` to `[D+10, B]` +- All 4 forward GEMM descriptors (Q/K/V/O): W_Q/K/V get k=D+10, W_O stays k=D - Layer norm input: stays D (operates on projected+residual, not raw input) **Backward changes:** -- dW_Q/K/V GEMMs: k dimension of right-hand side changes from D to D+8 -- `d_input_scratch` grows from `[D, B]` to `[D+8, B]` +- dW_Q/K/V GEMMs: k dimension of right-hand side changes from D to D+10 +- `d_input_scratch` grows from `[D, B]` to `[D+10, B]` - **NEW:** After `backward()` completes, expose `d_input_scratch` via a public accessor -- The caller (fused_training.rs) extracts `d_input_scratch[D..D+8, :]` as `d_ofi_embed_from_attn` +- The caller (fused_training.rs) extracts `d_input_scratch[D..D+10, :]` as `d_ofi_embed_from_attn` ### Phase I: OFI Embed MLP Backward @@ -248,7 +318,7 @@ No code change needed — PopArt handles this correctly by design. |--------|------|------| | ofi_embed_buf [8, B] | 0.5MB | Forward scratch | | ofi_input_buf [16, B] | 1MB | Forward scratch | -| d_h_history [SH2+8, B*K] | 17MB | Backward scratch | +| d_h_history [SH2+10, B*K] | 17MB | Backward scratch | | d_ofi_embed [8, B] × 2 | 1MB | Gradient accumulators | | h_history growth (+8 per K step) | +2MB | Existing buffer wider | | W_A, W_B growth (+8 rows) | +1KB | Params | @@ -264,8 +334,8 @@ No code change needed — PopArt handles this correctly by design. | `experience_kernels.cu` | PORTFOLIO_STRIDE=38, micro-reward body with tgt[4]/tgt[5], OFI delta storage, batch_states param | | `gpu_experience_collector.rs` | Wire ofi_gpu into state_gather, pass batch_states to env_step | | `gpu_dqn_trainer.rs` | Fix concat_ofi indices 42→66, OFI embed MLP, wider Mamba2 params, d_h_history GEMMs | -| `mamba2_temporal_kernel.cu` | Update history kernel to copy SH2+8 dims with ofi_embed param | -| `gpu_attention.rs` | Wider W_Q/K/V (k=D+8), expose d_input_scratch, wider saved_input | +| `mamba2_temporal_kernel.cu` | Update history kernel to copy SH2+10 dims with ofi_embed param | +| `gpu_attention.rs` | Wider W_Q/K/V (k=D+10), expose d_input_scratch, wider saved_input | | `batched_forward.rs` | Concatenate ofi_embed to attention input | | `batched_backward.rs` | Wider backward dimensions, extract d_ofi_embed | | `fused_training.rs` | OFI embed forward/backward calls, gradient accumulation | @@ -303,9 +373,70 @@ micro_reward_scale = 0.1 # was 0.0 (disabled) 7. Gradient through OFI embed MLP is non-zero (grad_norm check) 8. Mamba2 d_h_history gradient flows (not all-zero) +### Phase K: 4th Direction — Hold Action (b0_size 3→4) + +**Problem:** Direction branch has 3 options: Short(-1), Flat(0), Long(+1). Flat means "close position to zero" — there's no "keep current position unchanged" option. The model is forced to actively re-decide every bar, biasing toward action. Every action except Flat has transaction cost. + +**Fix:** Add Hold as a 4th direction. Action space: `direction(4) × magnitude(3) × order(3) × urgency(3) = 108` (was 81). + +| Direction | idx | Position Effect | Cost | +|-----------|-----|-----------------|------| +| Short | 0 | Open/maintain short | Spread + commission | +| Hold | 1 | Keep current position unchanged | **Zero** | +| Long | 2 | Open/maintain long | Spread + commission | +| Flat | 3 | Close position to zero | Spread + commission | + +**Changes required:** +1. `config`: `branch_0_size` from 3 to 4 in production config +2. `trade_physics.cuh`: `compute_target_position_4branch` — Hold (dir_idx=1) returns current position unchanged +3. `experience_kernels.cu`: `env_step_batch` — Hold action skips trade execution, zero tx cost +4. `gpu_action_selector.rs`: epsilon-greedy handles 4 directions. Action encoding: `dir*27 + mag*9 + ord*3 + urg` (was `dir*27`, now dir range 0-3) +5. `gpu_dqn_trainer.rs`: direction branch head output changes from 3 to 4 (affects w_b0 weight tensor) +6. Counterfactual: direction mirror for 4 actions — mirror of Short(0) is Long(2), mirror of Hold(1) is Hold(1), mirror of Flat(3) is Flat(3) +7. `evaluate_baseline.rs`: backtest action decoding for 4 directions +8. IQL: per-branch advantage now has 4 direction values +9. `ExposureLevel`: needs Hold variant + +**Synergy with dense micro-reward:** Hold + micro-reward creates a natural feedback loop: +- Micro-reward says "this position has positive OFI quality" → Hold (zero cost) +- Micro-reward says "negative quality" → Flat (close, pay spread once) +- Without Hold, the model must re-choose Long/Short every bar, paying spread implicitly through magnitude re-evaluation + +### Phase L: Fix DSR Sharpe EMA (always zero) + +**Problem:** `price_change_dsr = 0.0f` at line 1733 of experience_kernels.cu. The running Sharpe estimate `ps[DSR_A_SLOT]` is permanently zero, making the adaptive cost tolerance always 0.1 (floor). + +**Fix:** Compute actual per-bar return using stored prev_close: +```c +float prev_close_val = ps[PREV_CLOSE_SLOT]; /* slot 6 */ +float price_change_dsr = (prev_close_val > 0.0f) + ? (raw_close - prev_close_val) / prev_close_val + : 0.0f; +``` +Store `ps[PREV_CLOSE_SLOT] = raw_close` after DSR update (already in the spec for micro-reward). + +### Phase M: Fix Counterfactual Magnitude/Order Sign with do_flip + +**Problem:** We fixed `cf_cycle==0` (direction) for do_flip interaction. But `cf_cycle==1` (magnitude) does `cf_reward = reward * (alt_mag / actual_mag)` and `cf_cycle==2` (order) does `cf_reward = reward + cost_delta`. When `do_flip=true`, `reward` is already negated. Both magnitude scaling and order cost addition preserve the wrong sign. + +**Fix:** +```c +// cf_cycle==1 (magnitude): use original reward sign +float base_reward = do_flip ? -reward : reward; /* undo flip for CF base */ +cf_reward = base_reward * (alt_mag_frac / actual_mag_frac); +if (do_flip) cf_reward = -cf_reward; /* re-apply flip */ + +// cf_cycle==2 (order): same treatment +float base_reward_ord = do_flip ? -reward : reward; +cf_reward = base_reward_ord + cost_delta_scaled; +if (do_flip) cf_reward = -cf_reward; +``` + +Or more simply: compute counterfactual reward from the PRE-FLIP reward, then apply flip at the end. + ### Non-Goals -- No IQN/IQL architecture changes -- No new exploration mechanisms +- No new exploration mechanisms beyond Hold - No hyperopt integration (fixed micro_reward_scale=0.1 initially) - No changes to C51 atom support (PopArt handles scaling) +- No walk-forward window changes (separate investigation)