#![allow( clippy::assertions_on_constants, clippy::assertions_on_result_states, clippy::clone_on_copy, clippy::decimal_literal_representation, clippy::doc_markdown, clippy::empty_line_after_doc_comments, clippy::field_reassign_with_default, clippy::get_unwrap, clippy::identity_op, clippy::inconsistent_digit_grouping, clippy::indexing_slicing, clippy::integer_division, clippy::len_zero, clippy::let_underscore_must_use, clippy::manual_div_ceil, clippy::manual_let_else, clippy::manual_range_contains, clippy::modulo_arithmetic, clippy::needless_range_loop, clippy::non_ascii_literal, clippy::redundant_clone, clippy::shadow_reuse, clippy::shadow_same, clippy::shadow_unrelated, clippy::single_match_else, clippy::str_to_string, clippy::string_slice, clippy::tests_outside_test_module, clippy::too_many_lines, clippy::unnecessary_wraps, clippy::unseparated_literal_suffix, clippy::use_debug, clippy::useless_vec, clippy::wildcard_enum_match_arm, clippy::else_if_without_else, clippy::expect_used, clippy::missing_const_for_fn, clippy::similar_names, clippy::type_complexity, clippy::collapsible_else_if, clippy::doc_lazy_continuation, clippy::items_after_test_module, clippy::map_clone, clippy::multiple_unsafe_ops_per_block, clippy::unwrap_or_default, clippy::assign_op_pattern, clippy::needless_borrow, clippy::println_empty_string, clippy::unnecessary_cast, clippy::used_underscore_binding, clippy::create_dir, clippy::implicit_saturating_sub, clippy::exit, clippy::expect_fun_call, clippy::too_many_arguments, clippy::unnecessary_map_or, clippy::unwrap_used, dead_code, unused_imports, unused_variables, clippy::cloned_ref_to_slice_refs, clippy::neg_multiply, clippy::while_let_loop, clippy::bool_assert_comparison, clippy::excessive_precision, clippy::trivially_copy_pass_by_ref, clippy::op_ref, clippy::redundant_closure, clippy::unnecessary_lazy_evaluations, clippy::if_then_some_else_none, clippy::unnecessary_to_owned, clippy::single_component_path_imports, )] //! Precompute Features — DBN to .fxcache pipeline //! //! Reads OHLCV + MBP-10 + trades DBN data, runs the full feature extraction //! pipeline (standalone, no GPU required), and writes a `.fxcache` binary file for //! zero-overhead GPU loading during training. //! //! Usage: //! cargo run -p ml --example precompute_features --release -- \ //! --data-dir test_data/futures-baseline \ //! --mbp10-data-dir test_data/futures-baseline-mbp10 \ //! --trades-data-dir test_data/futures-baseline-trades \ //! --output-dir test_data/feature-cache \ //! --symbol ES.FUT \ //! --yes use std::path::{Path, PathBuf}; use std::time::Instant; use anyhow::{Context, Result}; use clap::Parser; use tracing::info; use ml::trainers::dqn::{ collect_dbn_files_filtered, collect_dbn_files_recursive, extract_features_from_bars, }; use ml::features::extraction::OHLCVBar; // --------------------------------------------------------------------------- // CLI // --------------------------------------------------------------------------- #[derive(Debug, Parser)] #[command( name = "precompute_features", about = "Precompute features from DBN data and write .fxcache binary" )] struct Opts { /// Directory containing OHLCV .dbn/.dbn.zst files (e.g. test_data/futures-baseline) #[arg(long)] data_dir: String, /// Directory containing MBP-10 .dbn/.dbn.zst files (optional) #[arg(long)] mbp10_data_dir: Option, /// Directory containing trades .dbn/.dbn.zst files (optional) #[arg(long)] trades_data_dir: Option, /// Output directory for .fxcache files (default: sibling of --data-dir) #[arg(long)] output_dir: Option, /// Symbol name (used in output filename) #[arg(long, default_value = "ES.FUT")] symbol: String, /// Data-source identifier mixed into the SHA256 cache key. /// /// MUST match the `data_source` field of the training profile that will /// consume this fxcache (e.g. `dqn-production.toml: data_source = "mbp10"`, /// `dqn-smoketest.toml: data_source = "mbp10"`). Mismatch produces a /// silent cache MISS at training time → DBN-direct fallback uploads /// un-normalised features → aux-head label_scale picks up raw-price /// magnitudes → cascade of broken metrics + impossible Sharpe (see /// 2026-04-27 incident: train-h5gxb epoch-0 Sharpe=141 from this exact /// path mismatch). /// /// Defaults to `"mbp10"` — production, smoke, and localdev profiles all /// use MBP-10 microstructure data. Pass `--data-source ohlcv` only when /// intentionally generating a 1-minute candle cache for a profile that /// overrides `data_source = "ohlcv"`. #[arg(long, default_value = "mbp10")] data_source: String, /// Imbalance bar formation threshold. ONLY used when `data_source == "mbp10"` /// and `mbp10_data_dir` is set; otherwise volume bars are produced. MUST /// match the trainer's `imbalance_bar_threshold` for fxcache HIT (the cache /// key includes this value as of 2026-05-09 architectural fix). Default /// matches `dqn-production.toml: imbalance_bar_threshold = 0.5`. #[arg(long, default_value_t = 0.5)] imbalance_bar_threshold: f64, /// Imbalance bar EWMA alpha (weight on OLD threshold in this codebase's /// reversed convention; α=1.0 disables adaptation). Default matches /// `dqn-production.toml: 0.1`. MUST match the trainer's value for cache HIT. #[arg(long, default_value_t = 0.1)] imbalance_bar_ewma_alpha: f64, /// Volume bar size: contracts of one-sided volume per bar. Used when /// `data_source != "mbp10"`. Default 100 matches `DEFAULT_VOLUME_BAR_SIZE`. /// MUST match the trainer's value for cache HIT. #[arg(long, default_value_t = 100)] volume_bar_size: u64, /// Skip confirmation prompt #[arg(long)] yes: bool, /// Row unit for the emitted fxcache. `bar` (default) emits one row per /// formed volume/imbalance bar with the 134-dim bar-level alpha stack. /// `snapshot` emits one row per MBP-10 snapshot event with the 81-dim /// snapshot-level stack (FoxhuntQ-Δ Phase 1c falsification test — /// canonical L2 book-update resolution, ~10× finer than bars). #[arg(long, default_value = "bar", value_parser = ["bar", "snapshot"])] row_unit: String, /// Snapshot warmup window (rows skipped at the head so the trailing /// trend scanner + stateful aggregators are warm before emission). /// Only used when `--row-unit snapshot`. #[arg(long, default_value_t = 100)] snapshot_warmup: usize, /// Purge the predecoded-sidecar cache before loading. /// /// Predecoded sidecars (`..predecoded.bin` under /// `/predecoded/`) skip zstd-DBN decompression on subsequent /// precompute runs. They self-invalidate when the source file's mtime or /// size changes, so this flag is normally unnecessary — pass it only after /// a format change to `Mbp10Snapshot` / `DbnTrade` (bumps fail to /// deserialize) or when explicitly testing the cold path. #[arg(long)] rebuild_predecoded: bool, } // --------------------------------------------------------------------------- // Cache key computation // --------------------------------------------------------------------------- /// Compute a SHA256 cache key from all DBN files across the given directories. /// // --------------------------------------------------------------------------- // Main // --------------------------------------------------------------------------- #[tokio::main] async fn main() -> Result<()> { // Initialize tracing tracing_subscriber::fmt() .with_env_filter( tracing_subscriber::EnvFilter::try_from_default_env() .unwrap_or_else(|_| tracing_subscriber::EnvFilter::new("info")), ) .init(); let opts = Opts::parse(); let data_dir = PathBuf::from(&opts.data_dir); // Default output_dir: sibling of data_dir (e.g. test_data/futures-baseline → test_data/feature-cache) let output_dir = match opts.output_dir { Some(ref d) => PathBuf::from(d), None => data_dir.parent().unwrap_or(&data_dir).join("feature-cache"), }; let mbp10_dir = opts.mbp10_data_dir.as_ref().map(PathBuf::from); let trades_dir = opts.trades_data_dir.as_ref().map(PathBuf::from); // Predecoded-sidecar cache lives under output_dir/predecoded/ so it sits // next to the .fxcache outputs (same writability assumptions). First call // for any source `.dbn.zst` pays zstd; subsequent calls deserialize the // sidecar and skip decompression entirely (zstd was 62% of wall time on // the 2026-05-16 1Q ES profile). let predecoded_dir = output_dir.join("predecoded"); std::fs::create_dir_all(&predecoded_dir) .with_context(|| format!("Failed to create predecoded dir {}", predecoded_dir.display()))?; if opts.rebuild_predecoded { let removed = ml::features::predecoded::purge_predecoded_dir(&predecoded_dir) .with_context(|| format!("Failed to purge {}", predecoded_dir.display()))?; info!("--rebuild-predecoded: purged {} sidecar(s) from {}", removed, predecoded_dir.display()); } // ── Validate inputs ────────────────────────────────────────────────────── if !data_dir.exists() { anyhow::bail!("Data directory not found: {}", data_dir.display()); } // ── Print config summary ───────────────────────────────────────────────── println!("================================================================================"); println!("Precompute Features — DBN to .fxcache Pipeline"); println!("================================================================================"); println!(); println!("Data dir: {}", data_dir.display()); if let Some(ref d) = mbp10_dir { println!("MBP-10 dir: {}", d.display()); } if let Some(ref d) = trades_dir { println!("Trades dir: {}", d.display()); } println!("Output dir: {}", output_dir.display()); println!("Symbol: {}", opts.symbol); println!("Format: f32 (v{}), OFI_DIM={}", ml::fxcache::FXCACHE_VERSION, ml::fxcache::OFI_DIM); println!(); // ── Confirmation ───────────────────────────────────────────────────────── if !opts.yes { print!("Proceed with feature extraction? (yes/no): "); std::io::Write::flush(&mut std::io::stdout())?; let mut input = String::new(); std::io::stdin().read_line(&mut input)?; let trimmed = input.trim(); if !trimmed.eq_ignore_ascii_case("yes") && !trimmed.eq_ignore_ascii_case("y") { println!("Cancelled."); return Ok(()); } println!(); } let t0 = Instant::now(); const WARMUP: usize = 50; // ── SP19 (2026-05-09) — producer-side multi-horizon reward blend ── // Path (B) of the multi-horizon reward augmentation spec. We blend // 1-bar / 5-bar / 30-bar log-returns into `tgt[1]` at fxcache write // time so kernel consumers see a single reward signal (no kernel // changes). `1/sqrt(N)` vol-scale correction inside the blend keeps // each horizon's log-return at unit-volatility-equivalent before // weighting under random-walk assumption. // // Why hardcoded weights here? `precompute_features` runs BEFORE // training so the ISV bus isn't initialised when the blend happens. // Slots [507..510) (`REWARD_HORIZON_WEIGHT_*BAR_INDEX`) are // reservations for a future Path (A) refactor that bumps TARGET_DIM // to carry per-horizon log-returns through the kernel pipeline; in // that world the consumer reads ISV at training-step time. For the // empirical hypothesis test ("does multi-horizon blend lift WR?") // equal-thirds is sufficient since varying weights requires fxcache // regen (~10-15 min on L40S) and per-batch adaptation is impossible // without the TARGET_DIM bump. // // Sentinel match: `SENTINEL_REWARD_HORIZON_WEIGHT_DEFAULT` in // `crates/ml/src/cuda_pipeline/sp14_isv_slots.rs` equals // `1.0_f32 / 3.0_f32 = 0.333_333_343`. The producer here uses f64 // `1.0 / 3.0` for numerical accuracy in the blend itself; the f32 // ISV slot value is what a future consumer reads, and matches // bit-identically when cast. const LOOKAHEAD_HORIZON_MAX: usize = 30; const SP19_HORIZON_BLEND_WEIGHTS: [f64; 3] = [ 1.0_f64 / 3.0_f64, 1.0_f64 / 3.0_f64, 1.0_f64 / 3.0_f64, ]; // ── Early exit if cache already exists AND matches current version ──────── let hex_key_early = ml::feature_cache::calculate_dbn_cache_key_full( &data_dir, mbp10_dir.as_deref(), trades_dir.as_deref(), &opts.symbol, &opts.data_source, opts.imbalance_bar_threshold, opts.imbalance_bar_ewma_alpha, opts.volume_bar_size, ).context("Failed to compute cache key")?; let early_check_path = output_dir.join(format!("{hex_key_early}.fxcache")); if early_check_path.exists() { match ml::fxcache::load_fxcache(&early_check_path) { Ok(cached) if cached.has_ofi || mbp10_dir.is_none() => { // has_ofi flag is necessary but not sufficient — verify actual OFI data is non-zero let ofi_nonzero = cached.ofi.iter().filter(|r| r.iter().any(|&v| v != 0.0)).count(); if mbp10_dir.is_some() && ofi_nonzero == 0 { println!("Cache has_ofi=true but OFI data is all zeros ({} bars) — deleting stale cache", cached.ofi.len()); std::fs::remove_file(&early_check_path).ok(); } else { let size_mb = std::fs::metadata(&early_check_path) .map(|m| m.len() as f64 / 1_048_576.0) .unwrap_or(0.0); println!("Cache already exists: {} ({:.1} MB, has_ofi={}, ofi_nonzero={}) — skipping extraction", early_check_path.display(), size_mb, cached.has_ofi, ofi_nonzero); return Ok(()); } } Ok(_) => { println!("Cache exists but has_ofi=false while MBP-10 data available — deleting stale cache"); std::fs::remove_file(&early_check_path).ok(); } Err(e) => { println!("Cache exists but failed validation: {e} — deleting stale cache"); std::fs::remove_file(&early_check_path).ok(); } } } // ── Step 1: Build volume bars from trades data ──────────────────────────── let t1 = Instant::now(); // Use trades_data_dir if provided, otherwise try data_dir for trades let trades_source = trades_dir.as_ref().unwrap_or(&data_dir); info!("Loading trades from {} (symbol={})...", trades_source.display(), opts.symbol); let mut trade_files = collect_dbn_files_filtered(trades_source, Some(&opts.symbol)); if trade_files.is_empty() { // Fallback: try data_dir if trades_dir was specified but had no files if trades_dir.is_some() { trade_files = collect_dbn_files_filtered(&data_dir, Some(&opts.symbol)); } if trade_files.is_empty() { anyhow::bail!( "No trades DBN files found for symbol '{}' in {} or {}", opts.symbol, trades_source.display(), data_dir.display() ); } } trade_files.sort(); info!("Found {} trades/OHLCV DBN files for symbol '{}'", trade_files.len(), opts.symbol); // Load trades PER FILE and filter front-month within each file. // Each quarterly DBN file has a different front-month contract (e.g. ESH24, ESM24). // Filtering globally would select only ONE quarter's contract. // // Parallel across quarters via rayon; each file goes through the // predecoded-sidecar cache so a re-run of `precompute_features` skips // zstd decompression entirely. let per_file: Vec<(usize, Vec)> = { use rayon::prelude::*; trade_files .par_iter() .filter_map(|file| { match ml::features::predecoded::load_or_predecode_trades(file, &predecoded_dir) { Ok(trades) => { let raw_count = trades.len(); let filtered = ml::features::filter_front_month(&trades); info!( " {} -> {} trades, {} front-month", file.file_name().unwrap_or_default().to_string_lossy(), raw_count, filtered.len() ); Some((raw_count, filtered)) } Err(e) => { tracing::warn!(" Failed {:?}: {e}", file.file_name()); None } } }) .collect() }; let total_raw: usize = per_file.iter().map(|(r, _)| r).sum(); let mut all_front_month_trades: Vec = Vec::with_capacity(per_file.iter().map(|(_, t)| t.len()).sum()); for (_, t) in per_file { all_front_month_trades.extend(t); } all_front_month_trades.sort_by_key(|t| t.timestamp); info!("Total trades: {} raw, {} front-month", total_raw, all_front_month_trades.len()); // Build bars — branch on data_source. Pre-2026-05-09 this was unconditional // build_volume_bars, silently ignoring data_source="mbp10" + the imbalance // bar threshold config. Audit confirmed 14 production runs all collided on // the same fxcache key regardless of TOML threshold value, so the imbalance // path was effectively dead. Wire it now: when "mbp10" + mbp10_dir is set, // construct imbalance bars from MBP-10 trades; otherwise fall back to volume // bars on the raw trade tape. let all_bars = if opts.data_source == "mbp10" && mbp10_dir.is_some() { let mbp10_path = mbp10_dir.as_deref().expect("mbp10_dir checked above"); info!( "Building imbalance bars from MBP-10: dir={}, threshold={}, ewma_alpha={}", mbp10_path.display(), opts.imbalance_bar_threshold, opts.imbalance_bar_ewma_alpha, ); let bars = ml::features::mbp10_loader::mbp10_to_imbalance_bars( mbp10_path, &opts.symbol, opts.imbalance_bar_threshold, opts.imbalance_bar_ewma_alpha, ).context("MBP-10 imbalance bar construction failed")?; info!( "Imbalance bars built: {} bars (threshold={}) in {:.1}s", bars.len(), opts.imbalance_bar_threshold, t1.elapsed().as_secs_f64(), ); bars } else { let bars = ml::features::build_volume_bars(&all_front_month_trades, opts.volume_bar_size); info!("Volume bars built: {} bars ({} contracts/bar) in {:.1}s", bars.len(), opts.volume_bar_size, t1.elapsed().as_secs_f64()); bars }; // SP19 Path (B) (2026-05-09): need at least `WARMUP + LOOKAHEAD_HORIZON_MAX + 1` // bars so the 30-bar log-return at `i = 0` (`all_bars[WARMUP + 30]`) is // valid AND we still produce at least one output row. if all_bars.len() < WARMUP + LOOKAHEAD_HORIZON_MAX + 1 { anyhow::bail!( "Insufficient data: {} bars, need at least {} (WARMUP={WARMUP} + LOOKAHEAD_HORIZON_MAX={LOOKAHEAD_HORIZON_MAX} + 1)", all_bars.len(), WARMUP + LOOKAHEAD_HORIZON_MAX + 1 ); } // ── Price continuity check (regression gate) ──────────────────────── let mut max_jump_pct = 0.0_f64; let mut jump_count = 0_usize; for w in all_bars.windows(2) { let pct = ((w[1].close - w[0].close) / w[0].close).abs() * 100.0; if pct > 0.5 { jump_count += 1; } max_jump_pct = max_jump_pct.max(pct); } info!("Price continuity: max_jump={:.3}%, jumps>0.5%={}/{}", max_jump_pct, jump_count, all_bars.len() - 1); if jump_count as f64 / all_bars.len() as f64 > 0.01 { anyhow::bail!( "PRICE CONTINUITY FAILED: {} of {} bars have >0.5% price jump — likely interleaved contracts", jump_count, all_bars.len() ); } // ── Step 2: Extract 42-dim feature vectors ────────────────────────────── let t1 = Instant::now(); info!("Extracting features..."); let feature_vectors = extract_features_from_bars(&all_bars) .context("Feature extraction failed")?; info!("Extracted {} feature vectors in {:.1}s", feature_vectors.len(), t1.elapsed().as_secs_f64()); // feature_vectors will be CONSUMED by the `to_vec()` slice copy below // and then dropped explicitly at the end of step 3. Without the drop // it'd live alongside the 6GB normalised features Vec through every // downstream step. // ── Step 3: Build targets [preproc_close, preproc_next, raw_close, raw_next, raw_open, mid_price_open] // [0:1] = log-return-normalized close/next (network input; matches contract // in experience_kernels.cu:1556 + cuda_pipeline/mod.rs:508) // [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 below) // // Prior implementation wrote raw_curr/raw_next to slots [0:1] in violation of // the documented contract — the network-input columns ended up containing raw // prices (~$5000+ for ES futures). Fixed here so the fxcache matches its own // documented schema. // // SP19 Path (B) (2026-05-09): `tgt[1]` (preproc_next) now holds the // multi-horizon-blended log-return rather than the 1-bar log-return. // The blend mixes 1-bar / 5-bar / 30-bar log-returns at equal-thirds // weights with `1/sqrt(N)` vol-scale correction. The valid range // shrinks by `LOOKAHEAD_HORIZON_MAX` bars (we need // `all_bars[i + WARMUP + 30]` for the 30-bar log-return). let n = feature_vectors.len().saturating_sub(LOOKAHEAD_HORIZON_MAX); let mut features: Vec<[f64; 42]> = feature_vectors[..n].to_vec(); // feature_vectors now duplicates `features` data; drop the original // so the ~6GB Vec doesn't shadow us through the OFI + alpha steps. drop(feature_vectors); 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; let raw_open = all_bars[i + WARMUP].open; // prev_close: 1-bar lag for log-return computation. WARMUP ≥ 50 so // all_bars[i + WARMUP - 1] is always valid. let prev_close = all_bars[i + WARMUP - 1].close; let preproc_close = if prev_close > 0.0 { (raw_curr / prev_close).ln() } else { 0.0 }; // ── SP19 Path (B) — multi-horizon blended log-return ── // 1-bar / 5-bar / 30-bar log-returns with `1/sqrt(N)` vol-scale // correction applied INSIDE the blend (before weighting). Tail // bars are excluded by the `n = len - LOOKAHEAD_HORIZON_MAX` // trim above so `all_bars[i + WARMUP + 30]` is always valid. let raw_5bar = all_bars[i + WARMUP + 5].close; let raw_30bar = all_bars[i + WARMUP + 30].close; let log_return_1bar = if raw_curr > 0.0 { (raw_next / raw_curr).ln() } else { 0.0 }; let log_return_5bar = if raw_curr > 0.0 { (raw_5bar / raw_curr).ln() / (5.0_f64).sqrt() } else { 0.0 }; let log_return_30bar = if raw_curr > 0.0 { (raw_30bar / raw_curr).ln() / (30.0_f64).sqrt() } else { 0.0 }; let preproc_next = SP19_HORIZON_BLEND_WEIGHTS[0] * log_return_1bar + SP19_HORIZON_BLEND_WEIGHTS[1] * log_return_5bar + SP19_HORIZON_BLEND_WEIGHTS[2] * log_return_30bar; // mid_price_open will be filled from MBP-10 data below if available; // default to raw_open (OHLCV-only fallback). targets.push([preproc_close, preproc_next, raw_curr, raw_next, raw_open, raw_open]); } // ── Step 3b: Build per-bar timestamps ───────────────────────────────── let timestamps: Vec = (0..n) .map(|i| all_bars[i + WARMUP].timestamp.timestamp_nanos_opt().unwrap_or(0)) .collect(); // Alpha feature inputs (hoisted to function scope so they survive past the // mbp10_dir if-let block for the alpha pipeline call below). Populated // INSIDE the MBP-10 branch (real data) or left empty for the volume-bar // path (bar-level alpha features only). let mut alpha_snapshots: Vec = Vec::new(); let mut alpha_trades: Vec = Vec::new(); // ── Step 4: Compute OFI from MBP-10 + trades ──────────────────────────── const OFI_DIM: usize = ml_core::state_layout::OFI_DIM; let t2 = Instant::now(); let ofi: Vec<[f64; OFI_DIM]> = if let Some(ref mbp10_path) = mbp10_dir { use ml::features::ofi_calculator::OFICalculator; info!("Loading MBP-10 snapshots from {}...", mbp10_path.display()); let mut mbp10_files = collect_dbn_files_recursive(mbp10_path); mbp10_files.sort(); if mbp10_files.is_empty() { info!("No MBP-10 files found, OFI will be zeros"); vec![[0.0; OFI_DIM]; n] } else { // Load MBP-10 snapshots (parallel, with predecoded sidecar) let per_file: Vec> = { use rayon::prelude::*; mbp10_files.par_iter().filter_map(|file| { // precompute_features is the multi-symbol feature-cache // builder; it intentionally consumes every instrument in // each file. No filter. match ml::features::predecoded::load_or_predecode_mbp10(file, &predecoded_dir, data::providers::databento::dbn_parser::InstrumentFilter::All) { Ok(snapshots) => { info!( " {} -> {} snapshots", file.file_name().unwrap_or_default().to_string_lossy(), snapshots.len() ); Some(snapshots) } Err(e) => { tracing::warn!(" Failed {:?}: {e}", file.file_name()); None } } }).collect() }; let mut all_snapshots = Vec::new(); for snaps in per_file { all_snapshots.extend(snaps); } all_snapshots.sort_by_key(|s| s.timestamp); info!("Loaded {} MBP-10 snapshots in {:.1}s", all_snapshots.len(), t2.elapsed().as_secs_f64()); // Load trades (parallel, with predecoded sidecar) and filter // front-month per-file. // // Pre-2026-05-09: this branch loaded trades unfiltered, leaking // off-contract trade activity into the per-bar OFI/VPIN/Kyle's // Lambda computations during contract rollover windows. Volume // bar formation already filters front-month at the trades-for-bars // load above; the OFI side did not. Audit confirmed the latent bug // affected every prior MBP-10+trades production run. // // Fix: mirror the per-file `filter_front_month` call from the // volume bar path so OFI windows only see in-month trades. // // Trade loads go through `load_or_predecode_trades` so the sidecar // cache is shared with the bar-formation path above (one decode of // each `.dbn.zst` per process — the OFI re-load is a sidecar HIT). let all_trades = if let Some(ref tdir) = trades_dir { let mut trade_files = collect_dbn_files_recursive(tdir); trade_files.sort(); let per_file_trades: Vec> = { use rayon::prelude::*; trade_files.par_iter().filter_map(|path| { match ml::features::predecoded::load_or_predecode_trades(path, &predecoded_dir) { Ok(t) => { let raw_count = t.len(); let filtered = ml::features::filter_front_month(&t); info!(" {} trades from {:?} ({} front-month)", raw_count, path.file_name(), filtered.len()); Some(filtered) } Err(e) => { tracing::warn!(" Failed {:?}: {e}", path.file_name()); None } } }).collect() }; let mut trades = Vec::new(); for t in per_file_trades { trades.extend(t); } trades.sort_by_key(|t| t.timestamp); if trades.is_empty() { None } else { Some(trades) } } else { None }; // ── Per-bar OFI extraction (SP20 parallel via time-bucket sharding) ── // // Replaces a sequential 132-line loop with a call into the // `ml-features` library function. Sharding strategy: K shards // process contiguous bar ranges in parallel; each shard owns a // freshly-cloned `OFICalculator` and replays a `WARMUP_BARS` prefix // of its range to populate rolling-window state (VPIN/Kyle/ // trade-imb/OFI-stats/prev-snap) to bit-identical sequential-walk // state at the warmup→core boundary. Post-loop passes (lag-1 // deltas, log_bar_duration) are pure functions of the full output // and run sequentially after concatenation. // // Bit-equivalence is verified in // `crates/ml-features/tests/ofi_parallel_bit_equiv_test.rs`. // // Shard count: from `OFI_SHARD_COUNT` env (override) or rayon's // current thread pool (`ci-compile-cpu` = 28 cores). Capped at // `n` shards. Set to 1 to disable parallelism (matches previous // sequential behaviour exactly). let shard_count: usize = std::env::var("OFI_SHARD_COUNT") .ok() .and_then(|s| s.parse::().ok()) .unwrap_or_else(rayon::current_num_threads); let inputs = ml::features::PerBarOfiInputs { bars: &all_bars, warmup: WARMUP, n, snapshots: &all_snapshots, trades: all_trades.as_deref(), }; let ofi_per_bar = ml::features::compute_ofi_per_bar_parallel( &inputs, shard_count, ml::features::DEFAULT_OFI_PARALLEL_WARMUP_BARS, ); info!( "OFI per-bar extraction: {} shards, warmup={} bars", shard_count.max(1).min(n), ml::features::DEFAULT_OFI_PARALLEL_WARMUP_BARS, ); let delta_nonzero = ofi_per_bar.iter().filter(|f| f[8..16].iter().any(|&v| v != 0.0)).count(); let book_nonzero = ofi_per_bar.iter().filter(|f| f[16] != 0.0).count(); let micro_nonzero = ofi_per_bar.iter().filter(|f| f[20..30].iter().any(|&v| v != 0.0)).count(); let tlob_nonzero = ofi_per_bar.iter().filter(|f| f[30] != 0.0 || f[31] != 0.0).count(); info!("OFI v5: deltas_nonzero={}/{}, book_aggression_nonzero={}/{}, microstructure_nonzero={}/{}, tlob_novel_nonzero={}/{}", delta_nonzero, n, book_nonzero, n, micro_nonzero, n, tlob_nonzero, n); // 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, if ofi_per_bar.is_empty() { 0.0 } else { non_zero as f64 / ofi_per_bar.len() as f64 * 100.0 }, t2.elapsed().as_secs_f64()); // Hand snapshots + trades off to the alpha pipeline (executes after // this if-let). DbnTrade (u64 ns timestamp, u64 volume) is converted // to Mbp10Trade (DateTime timestamp, f64 volume) — the alpha // pipeline uses the DateTime form for trade-flow ordering and the // f64 volume for size-weighted features. alpha_snapshots = all_snapshots; if let Some(t) = all_trades { // CONSUME the original DbnTrade Vec while building the // Mbp10Trade Vec. Using `.iter()` here previously kept // the original ~10GB Vec alive for the duration of the // transformation, so peak memory transiently doubled // to ~25GB (DbnTrade + Mbp10Trade Vecs simultaneously) // and OOM-killed the 9-quarter precompute at 56Gi // (2026-05-16). `.into_iter()` drops the source // allocation element-by-element as the destination // builds, capping peak at the larger of the two. alpha_trades = t .into_iter() .map(|dbn| ml::features::Mbp10Trade { price: dbn.price, volume: dbn.volume as f64, timestamp: chrono::DateTime::from_timestamp_nanos(dbn.timestamp as i64), is_buy: dbn.is_buy, instrument_id: dbn.instrument_id, }) .collect(); } ofi_per_bar } } else { info!("No MBP-10 directory, OFI will be zeros (log_bar_duration still computed)"); let mut ofi_per_bar = vec![[0.0_f64; OFI_DIM]; n]; // Even without MBP-10 data, compute log_bar_duration from bar timestamps for i in 0..n { let duration_secs = if i > 0 { let ts_cur = all_bars[i + WARMUP].timestamp; let ts_prev = all_bars[i + WARMUP - 1].timestamp; (ts_cur - ts_prev).num_seconds() as f64 } else { 60.0 // default 1 minute for first bar }; let log_duration = duration_secs.max(0.1).ln() / 10.0; ofi_per_bar[i][17] = log_duration; } ofi_per_bar }; let total_len = n; // ── Compute cache key ──────────────────────────────────────────────────── let hex_key = ml::feature_cache::calculate_dbn_cache_key_full( &data_dir, mbp10_dir.as_deref(), trades_dir.as_deref(), &opts.symbol, &opts.data_source, opts.imbalance_bar_threshold, opts.imbalance_bar_ewma_alpha, opts.volume_bar_size, ).context("Failed to compute cache key")?; let cache_key: [u8; 32] = hex::decode(&hex_key) .context("Invalid hex key")? .try_into() .map_err(|_| anyhow::anyhow!("Cache key wrong length"))?; // ── Normalize features (z-score) ────────────────────────────────────────── // Raw features contain OHLCV prices (up to 25,000). Normalize once at // precompute time so every consumer (training, hyperopt, inference) // gets consistent normalized features. Use the in-place variant so // we don't transiently hold pre- + post-normalised copies (~6GB // saving over `.normalize_batch(&features)` at this dataset size). let norm_stats = ml::walk_forward::NormStats::from_features(&features); norm_stats.normalize_batch_in_place(&mut features); info!("Features z-score normalized in place ({} bars × 42 dims)", features.len()); // Defence-in-depth gate: never write a poisoned fxcache to disk. If anything // upstream produced an out-of-bounds value, fail loudly here rather than // letting it silently propagate to every future training run that reads // this cache. Same gate the fxcache loader and DBN fallback enforce. ml::walk_forward::validate_normalized_features(&features, "precompute_features writer")?; // ── Write .fxcache ─────────────────────────────────────────────────────── std::fs::create_dir_all(&output_dir) .with_context(|| format!("Failed to create output directory: {}", output_dir.display()))?; let output_path = output_dir.join(format!("{hex_key}.fxcache")); // Save NormStats alongside cache for inference denormalization let norm_path = output_dir.join(format!("{hex_key}.norm_stats.json")); let norm_json = serde_json::to_string_pretty(&norm_stats) .context("Failed to serialize NormStats")?; std::fs::write(&norm_path, &norm_json) .with_context(|| format!("Failed to write NormStats: {}", norm_path.display()))?; info!("NormStats saved to {}", norm_path.display()); let t2 = Instant::now(); info!("Writing .fxcache to {}...", output_path.display()); let has_ofi = mbp10_dir.is_some(); // Alpha features (Stage 2): real 134-dim feature vector per bar produced // by `ml-features::alpha_pipeline::extract_alpha_features`. Combines Block // A-E factory extractors (Price/Volume/Statistical/ADX/RegimeADX/Time) // with the 80 authored Block F-W modules (multi-level OFI/GOFI, // Kyle's λ, frac-diff, microprice, Hasbrouck spread, Hawkes MLE, PIN, // LOB PCA, trend-scanning, etc.). Output dimensions are pinned by the // `ALPHA_FEATURE_DIM=134` invariant; runtime mismatch is caught by the // fxcache writer's per-row width check. // alpha_snapshots and alpha_trades populated inside the mbp10_dir branch // above (real MBP-10 data) or left empty for the volume-bar path // (bar-level alpha features only — F/G/H/T/N+V/O/CC/U/Y blocks degrade // gracefully to zeros when snapshots are empty). info!( "alpha pipeline inputs: {} snapshots, {} trades", alpha_snapshots.len(), alpha_trades.len() ); let (bytes_written, total_len_emitted, alpha_dim_emitted, has_ofi_emitted, row_unit_for_summary) = if opts.row_unit == "snapshot" { // ── Phase 1c falsification: per-MBP10-snapshot fxcache ── // Bypass the bar-level alpha pipeline; emit one row per snapshot // event using the 81-dim snapshot-native feature stack. let snap_warmup = opts.snapshot_warmup; if alpha_snapshots.len() <= snap_warmup { anyhow::bail!( "snapshot mode: only {} snapshots loaded, need > snapshot_warmup ({})", alpha_snapshots.len(), snap_warmup ); } let n_snap = alpha_snapshots.len() - snap_warmup; info!( "snapshot mode: emitting {} rows from {} snapshots (warmup={}, trades={})", n_snap, alpha_snapshots.len(), snap_warmup, alpha_trades.len() ); let snap_rows = ml::features::snapshot_pipeline::extract_snapshot_features( &alpha_snapshots, &alpha_trades, snap_warmup, n_snap, ); assert_eq!( snap_rows.len(), n_snap, "snapshot pipeline produced {} rows, expected {}", snap_rows.len(), n_snap ); let snap_dim = ml::features::SNAPSHOT_FEATURE_DIM; let snap_timestamps: Vec = snap_rows.iter().map(|r| r.timestamp_ns).collect(); let snap_features: Vec<[f64; 42]> = vec![[0.0_f64; 42]; n_snap]; let mut snap_targets: Vec<[f64; 6]> = vec![[0.0_f64; 6]; n_snap]; for (i, row) in snap_rows.iter().enumerate() { // COL_RAW_CLOSE = FEAT_DIM + 2 absolute → index 2 in the // 6-wide targets array. Snapshot mid-price becomes the // honest-direction label source. snap_targets[i][2] = row.mid_price; } let snap_ofi: Vec<[f64; OFI_DIM]> = vec![[0.0_f64; OFI_DIM]; n_snap]; let snap_alpha: Vec> = snap_rows.into_iter().map(|r| r.features).collect(); let bw = ml::fxcache::write_fxcache( &output_path, &snap_features, &snap_targets, &snap_ofi, &snap_timestamps, cache_key, false, // OFI column is zero-padded in snapshot mode Some(&snap_alpha), ) .context("Failed to write snapshot-level .fxcache file")?; (bw, n_snap, snap_dim, false, "snapshot") } else { let alpha_features = ml::features::alpha_pipeline::extract_alpha_features( &all_bars, &alpha_snapshots, &alpha_trades, WARMUP, features.len(), ); assert_eq!( alpha_features.len(), features.len(), "alpha pipeline produced {} rows, expected {}", alpha_features.len(), features.len() ); let bw = ml::fxcache::write_fxcache( &output_path, &features, &targets, &ofi, ×tamps, cache_key, has_ofi, Some(&alpha_features), ) .context("Failed to write .fxcache file")?; (bw, features.len(), ml::features::ALPHA_FEATURE_DIM, has_ofi, "bar") }; let write_secs = t2.elapsed().as_secs_f64(); let total_secs = t0.elapsed().as_secs_f64(); // ── Summary ────────────────────────────────────────────────────────────── println!(); println!("================================================================================"); println!("PRECOMPUTE SUMMARY"); println!("================================================================================"); println!(); println!("Row unit: {}", row_unit_for_summary); println!("Rows emitted: {}", total_len_emitted); println!("Alpha dim: {}", alpha_dim_emitted); println!("Targets: 6-dim"); println!("OFI: {}-dim ({})", OFI_DIM, if has_ofi_emitted { "from MBP-10" } else { "zero-padded" }); println!("Format: f32 (fxcache_v{}), OFI_DIM={}", ml::fxcache::FXCACHE_VERSION, ml::fxcache::OFI_DIM); println!("Cache key: {}", hex_key); println!("Output: {}", output_path.display()); println!("Size: {:.2} MB", bytes_written as f64 / 1_048_576.0); println!("Write time: {:.1}s", write_secs); println!("Total time: {:.1}s", total_secs); println!(); Ok(()) }