Files
foxhunt/crates/ml/examples/precompute_features.rs
jgrusewski a27cb40a9a fix(precompute): consume DbnTrade Vec during Mbp10Trade conversion
The 9-quarter precompute_features run OOM-killed at ~56Gi on the
ci-compile-cpu pool (POP2-HC-32C-64G) on 2026-05-16. Root cause:
lines 672-684's `t.iter().map(...).collect()` borrows the source
DbnTrade Vec while building the Mbp10Trade Vec — both alive
simultaneously, transient peak ~25GB just from this transformation
for the 199M-trade dataset.

`.into_iter()` consumes the source element-by-element so the
allocation drops as the destination grows, capping peak at the
larger of the two Vecs (~15GB) rather than their sum.

Should let the 9-quarter precompute fit comfortably on the existing
64GB ci-compile-cpu pool without provisioning a high-memory node.

Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
2026-05-16 10:05:25 +02:00

894 lines
42 KiB
Rust
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
#![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<String>,
/// Directory containing trades .dbn/.dbn.zst files (optional)
#[arg(long)]
trades_data_dir: Option<String>,
/// Output directory for .fxcache files (default: sibling of --data-dir)
#[arg(long)]
output_dir: Option<String>,
/// 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,
}
// ---------------------------------------------------------------------------
// 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);
// ── 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.
let mut all_front_month_trades: Vec<ml::features::DbnTrade> = Vec::new();
let mut total_raw = 0_usize;
for file in &trade_files {
match ml::features::load_trades_sync(file) {
Ok(trades) => {
let raw_count = trades.len();
total_raw += raw_count;
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());
all_front_month_trades.extend(filtered);
}
Err(e) => {
tracing::warn!(" Failed {:?}: {e}", file.file_name());
}
}
}
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());
// ── 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 features: Vec<[f64; 42]> = feature_vectors[..n].to_vec();
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<i64> = (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<data::providers::databento::mbp10::Mbp10Snapshot> = Vec::new();
let mut alpha_trades: Vec<ml::features::Mbp10Trade> = 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::mbp10_loader::load_ofi_features_parallel;
use ml::features::trades_loader::load_trades_sync;
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)
use data::providers::databento::dbn_parser::DbnParser;
let per_file: Vec<Vec<_>> = {
use rayon::prelude::*;
mbp10_files.par_iter().filter_map(|file| {
let parser = match DbnParser::new() {
Ok(p) => p,
Err(e) => { tracing::warn!("Parser init failed: {e}"); return None; }
};
let mut snapshots = Vec::new();
match parser.parse_mbp10_streaming(file, 100, |snap| {
snapshots.push(snap.clone());
}) {
Ok(_) => {
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) 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 line 354 of this
// file; 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.
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<Vec<_>> = {
use rayon::prelude::*;
trade_files.par_iter().filter_map(|path| {
match load_trades_sync(path) {
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::<usize>().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<Utc> 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.
let norm_stats = ml::walk_forward::NormStats::from_features(&features);
let features = norm_stats.normalize_batch(&features);
info!("Features z-score normalized ({} 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<i64> = 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<Vec<f32>> =
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,
&timestamps,
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(())
}