#![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, )] //! Walk-forward evaluation binary for DQN and PPO baseline models. //! //! Loads trained model checkpoints, runs inference on walk-forward test data, //! computes financial metrics (Sharpe, drawdown, win rate, profit factor), //! and generates a JSON report. //! //! # GPU-Batched Inference //! //! Both DQN and PPO evaluation use chunked GPU-batched inference (chunk size 1024). //! Instead of one GPU forward pass per bar (N kernel launches), bars are grouped into //! chunks with shared portfolio state, and each chunk does a single batched GPU forward //! pass. This reduces GPU kernel launches by ~1000x vs the old per-bar loop. //! //! Within a chunk, all bars share the portfolio state (equity, exposure, spread) from //! the start of the chunk. Between chunks, the portfolio state is updated based on the //! sequential trade simulation results. This matches the hyperopt adapter approach. //! //! # Usage //! //! ```bash //! SQLX_OFFLINE=true cargo run -p ml --example evaluate_baseline -- \ //! --model both --models-dir ml/trained_models \ //! --data-dir test_data/futures-baseline \ //! --output ml/trained_models/evaluation_report.json //! ``` #![allow(unused_crate_dependencies)] #![deny( clippy::unwrap_used, clippy::expect_used, clippy::panic, clippy::indexing_slicing )] use std::path::{Path, PathBuf}; use anyhow::{Context, Result}; use clap::Parser; use serde::Serialize; use serde_json::Value; use tracing::{info, warn}; use common::metrics::{server as metrics_server, training_metrics as tm}; use ml::common::action::{ExposureLevel, FactoredAction, OrderType as ActionOrderType, Urgency}; use ml::dqn::{DQNConfig, OrderRouter, DQN}; use ml::features::extraction::FeatureVector; #[allow(unreachable_pub)] mod baseline_common; use baseline_common::{load_all_bars, spread_cost_bps}; use ml::features::extraction::extract_ml_features; use ml::ppo::ppo::{PPOConfig, PPO}; use ml::types::OHLCVBar; use ml::walk_forward::{generate_walk_forward_windows, NormStats, WalkForwardConfig}; // Supervised model imports (same as evaluate_supervised) use ml::training::unified_trainer::UnifiedTrainable; use ml::diffusion::{DiffusionConfig, DiffusionTrainableAdapter}; use ml::kan::{KANConfig, KANTrainableAdapter}; use ml::liquid::{CfCTrainConfig, DeviceConfig, LiquidTrainableAdapter}; use ml::mamba::{Mamba2Config, trainable_adapter::Mamba2TrainableAdapter}; use ml::tft::{TFTConfig, TrainableTFT}; use ml::tgnn::trainable_adapter::TGGNTrainableAdapter; use ml::tgnn::TGGNConfig; use ml::tlob::{TLOBAdapterConfig, TLOBTrainableAdapter}; use ml::xlstm::{XLSTMConfig, XLSTMTrainableAdapter}; use ml_core::cuda_autograd::GpuTensor; use std::sync::Arc; use cudarc::driver::CudaStream; /// Number of bars processed per GPU forward pass. /// /// All bars in a chunk share the portfolio state (equity, exposure, spread) from /// the start of the chunk. The portfolio is updated sequentially on CPU between /// chunks. With 1024 bars, the stale-portfolio approximation is negligible /// (~0.7 trading days for ES) while reducing GPU kernel launches by ~1000x. const EVAL_CHUNK_SIZE: usize = 1024; // --------------------------------------------------------------------------- // CLI Arguments // --------------------------------------------------------------------------- /// Walk-forward evaluation binary for DQN/PPO baseline models. #[derive(Parser, Debug)] #[command( name = "evaluate_baseline", about = "Evaluate trained DQN/PPO checkpoints with walk-forward test data" )] struct Args { /// Directory containing trained model checkpoints #[arg(long, default_value = "ml/trained_models")] models_dir: PathBuf, /// Path to directory containing .dbn.zst files (env: FOXHUNT_DATA_DIR) #[arg(long, env = "FOXHUNT_DATA_DIR")] data_dir: PathBuf, /// Output path for evaluation report JSON #[arg(long, default_value = "ml/trained_models/evaluation_report.json")] output: PathBuf, /// Which model(s) to evaluate: "dqn", "ppo", or "both" #[arg(long, default_value = "both")] model: String, /// Feature dimension (51 market + 3 portfolio = 54, must match training config) #[arg(long, default_value_t = 54)] feature_dim: usize, /// Number of actions (5 exposure levels for DQN, pass --num-actions 45 for PPO) #[arg(long, default_value_t = 5)] num_actions: usize, /// Walk-forward: initial training window in months #[arg(long, default_value_t = 12)] train_months: u32, /// Walk-forward: validation window in months #[arg(long, default_value_t = 3)] val_months: u32, /// Walk-forward: test window in months #[arg(long, default_value_t = 3)] test_months: u32, /// Walk-forward: step size in months between folds #[arg(long, default_value_t = 3)] step_months: u32, /// Symbol subdirectory to load (e.g. "ES.FUT", "NQ.FUT") #[arg(long, default_value = "ES.FUT")] symbol: String, /// Maximum absolute per-bar return; larger moves are clamped (contract roll filter) #[arg(long, default_value_t = 0.01)] max_bar_return: f64, /// Round-trip commission cost in basis points (1 bps = 0.01%) /// Applied to BUY/SELL returns; HOLD is free. /// Default 1.0 bps covers ~$4 exchange+broker for ES e-mini. #[arg(long, default_value_t = 1.0)] tx_cost_bps: f64, /// Instrument tick size in price units (ES=0.25, NQ=0.25, ZN=1/64) #[arg(long, default_value_t = 0.25)] tick_size: f64, /// Typical bid-ask spread in ticks (ES=1.0, ZN=1.0, 6E=2.0) /// Full spread slippage is added to `tx_cost_bps` per round-trip trade. #[arg(long, default_value_t = 1.0)] spread_ticks: f64, /// Bars per year for Sharpe annualization (ES/NQ=347760, 6E=345000, ZN=105840). /// Default 347760 = 252 trading days × 1380 bars/day (ES 23h session). #[arg(long, default_value_t = 347_760.0)] bars_per_year: f64, /// Optional path to hyperopt results JSON -- overrides matching config fields /// (must match the file used during training so network architecture is identical) #[arg(long)] hyperopt_params: Option, /// Initial capital for GPU backtest evaluator. #[arg(long, default_value_t = 100_000.0)] initial_capital: f64, /// Maximum absolute position size in contracts (must match training config). /// Disable leverage cap to match training env (no leverage constraint during training). #[arg(long, default_value_t = 2.0)] max_position: f64, /// Use CUDA Graph capture for the DQN evaluation step loop. /// /// When enabled, the entire step loop (gather + forward + env_step for /// each step) is captured into a CUDA Graph on the first fold and /// replayed for subsequent calls. Eliminates per-step kernel launch /// overhead (~5-10us per launch × 4 kernels × max_len steps). /// /// Falls back to the non-graphed GPU path if capture fails. #[arg(long, default_value_t = false)] cuda_graphs: bool, // `gpu_eval`, `surrogate_mode`, `surrogate_seed`, `surrogate_marginals`, // `emit_action_marginals`, `emit_pooled_sharpe` flags removed 2026-05-12 — // all coupled to the deleted CPU DQN fallback path. GPU eval is mandatory; // surrogate-noise + pooled-Sharpe + action-marginal histograms need a GPU // closure-based reimplementation if reintroduced. See the 2026-05-12 audit // entry and `feedback_no_cpu_test_fallbacks`. } // --------------------------------------------------------------------------- // Hyperopt parameter loading // --------------------------------------------------------------------------- /// Load `best_params` from a hyperopt results JSON file. /// /// Expected format: `{ "model_key": { "best_params": { ... }, ... } }` /// Returns `None` if the file doesn't exist or can't be parsed. #[allow(clippy::cognitive_complexity)] fn load_hyperopt_params(hp_path: &Option, model_key: &str) -> Option { let file_path = hp_path.as_ref()?; if !file_path.exists() { info!("Hyperopt params file not found: {}, using defaults", file_path.display()); return None; } let contents = match std::fs::read_to_string(file_path) { Ok(c) => c, Err(e) => { warn!("Failed to read hyperopt params {}: {}", file_path.display(), e); return None; } }; let json: Value = match serde_json::from_str(&contents) { Ok(v) => v, Err(e) => { warn!("Failed to parse hyperopt params JSON: {}", e); return None; } }; let params = json.get(model_key) .and_then(|m| m.get("best_params")) .cloned(); if params.is_some() { info!("Loaded hyperopt params for '{}' from {}", model_key, file_path.display()); } else { warn!("No best_params found for '{}' in {}", model_key, file_path.display()); } params } fn hp_f64(params: &Option, key: &str) -> Option { params.as_ref()?.get(key)?.as_f64() } fn hp_usize(params: &Option, key: &str) -> Option { params.as_ref()?.get(key)?.as_u64().map(|v| v as usize) } fn hp_bool(params: &Option, key: &str) -> Option { params.as_ref()?.get(key)?.as_bool() } /// All 8 supervised model names recognised by this binary. const SUPERVISED_MODEL_NAMES: &[&str] = &[ "tft", "mamba2", "liquid", "tggn", "tlob", "kan", "xlstm", "diffusion", ]; /// Create a supervised model with default architecture (same defaults as /// `train_baseline_supervised` / `evaluate_supervised`). The model is /// freshly initialised — call `load_checkpoint` afterwards to populate /// trained weights. fn create_supervised_model( name: &str, feature_dim: usize, ) -> Result> { let lr = 1e-3; // irrelevant for eval, but constructors require it match name { "tft" => { let config = TFTConfig { input_dim: feature_dim, hidden_dim: 128, num_heads: 4, num_layers: 2, num_quantiles: 3, num_static_features: 0, num_known_features: 0, num_unknown_features: feature_dim, sequence_length: 1, prediction_horizon: 1, dropout_rate: 0.1, ..TFTConfig::default() }; let mut adapter = TrainableTFT::new(config) .map_err(|e| anyhow::anyhow!("Failed to create TFT: {}", e))?; adapter .set_learning_rate(lr) .map_err(|e| anyhow::anyhow!("Failed to set TFT learning rate: {}", e))?; Ok(Box::new(adapter)) } "mamba2" => { let config = Mamba2Config { d_model: 128, num_layers: 4, d_state: 16, max_seq_len: 60, ..Mamba2Config::default() }; let native_dev = ml_core::native_types::NativeDevice::Cuda(0); let mut adapter = Mamba2TrainableAdapter::from_native_device(config, &native_dev) .map_err(|e| anyhow::anyhow!("Failed to create Mamba2: {}", e))?; adapter .set_learning_rate(lr) .map_err(|e| anyhow::anyhow!("Failed to set Mamba2 learning rate: {}", e))?; Ok(Box::new(adapter)) } "liquid" => { let config = CfCTrainConfig { input_size: feature_dim, hidden_size: 128, output_size: 1, backbone_hidden_sizes: vec![128, 64], learning_rate: lr, device: DeviceConfig::Auto, ..CfCTrainConfig::default() }; let adapter = LiquidTrainableAdapter::new(config) .map_err(|e| anyhow::anyhow!("Failed to create Liquid: {}", e))?; Ok(Box::new(adapter)) } "tggn" => { let config = TGGNConfig { node_dim: feature_dim, hidden_dim: 32, num_layers: 2, max_nodes: 64, max_edges: 128, edge_dim: 4, temporal_decay: 0.99, update_frequency_ns: 1_000_000, use_simd: false, }; let mut adapter = TGGNTrainableAdapter::new(config) .map_err(|e| anyhow::anyhow!("Failed to create TGGN: {}", e))?; adapter .set_learning_rate(lr) .map_err(|e| anyhow::anyhow!("Failed to set TGGN learning rate: {}", e))?; Ok(Box::new(adapter)) } "tlob" => { let config = TLOBAdapterConfig { d_model: 128, num_heads: 4, num_layers: 2, seq_len: 1, feature_dim, }; let mut adapter = TLOBTrainableAdapter::new(config) .map_err(|e| anyhow::anyhow!("Failed to create TLOB: {}", e))?; adapter .set_learning_rate(lr) .map_err(|e| anyhow::anyhow!("Failed to set TLOB learning rate: {}", e))?; Ok(Box::new(adapter)) } "kan" => { let config = KANConfig { grid_size: 5, spline_order: 4, layer_widths: vec![feature_dim, 32, 16, 1], learning_rate: lr, weight_decay: 1e-4, grad_clip: 1.0, }; let adapter = KANTrainableAdapter::new(config) .map_err(|e| anyhow::anyhow!("Failed to create KAN: {}", e))?; Ok(Box::new(adapter)) } "xlstm" => { let config = XLSTMConfig { input_dim: feature_dim, hidden_dim: 128, ..XLSTMConfig::default() }; let mut adapter = XLSTMTrainableAdapter::new(config) .map_err(|e| anyhow::anyhow!("Failed to create xLSTM: {}", e))?; adapter .set_learning_rate(lr) .map_err(|e| anyhow::anyhow!("Failed to set xLSTM learning rate: {}", e))?; Ok(Box::new(adapter)) } "diffusion" => { let config = DiffusionConfig { feature_dim, hidden_dim: 128, ..DiffusionConfig::default() }; let ctx = cudarc::driver::CudaContext::new(0) .map_err(|e| anyhow::anyhow!("CUDA context for Diffusion: {}", e))?; let diff_stream = ctx.new_stream() .map_err(|e| anyhow::anyhow!("CUDA stream for Diffusion: {}", e))?; let adapter = DiffusionTrainableAdapter::new(config, &diff_stream) .map_err(|e| anyhow::anyhow!("Failed to create Diffusion: {}", e))?; Ok(Box::new(adapter)) } _ => anyhow::bail!("Unknown supervised model: {}", name), } } /// Locate the best checkpoint for a supervised model fold. /// /// Convention: `//_fold_best` (with `.json` metadata). /// Falls back to `/_fold_best` for flat layouts. fn find_supervised_checkpoint( models_dir: &std::path::Path, model_name: &str, fold: usize, ) -> Result { let subdir_ckpt = models_dir .join(model_name) .join(format!("{}_fold{}_best", model_name, fold)); let meta_subdir = format!("{}.json", subdir_ckpt.display()); if std::path::Path::new(&meta_subdir).exists() { return Ok(subdir_ckpt); } // Flat layout fallback let flat_ckpt = models_dir.join(format!("{}_fold{}_best", model_name, fold)); let meta_flat = format!("{}.json", flat_ckpt.display()); if std::path::Path::new(&meta_flat).exists() { return Ok(flat_ckpt); } anyhow::bail!( "No {} checkpoint found for fold {} in {} (tried {} and {})", model_name, fold, models_dir.display(), subdir_ckpt.display(), flat_ckpt.display(), ) } // --------------------------------------------------------------------------- // Report Data Types // --------------------------------------------------------------------------- /// Metrics for a single fold/model combination. #[derive(Debug, Serialize)] struct FoldMetrics { fold: usize, model: String, sharpe_ratio: f64, /// Trade-level Sharpe (non-zero returns only, sqrt(252) annualization). /// Comparable to hyperopt's `PerformanceMetrics::from_trades` Sharpe. trade_sharpe_ratio: f64, max_drawdown_pct: f64, win_rate_pct: f64, profit_factor: f64, total_return_pct: f64, num_trades: usize, test_start: String, test_end: String, } /// Aggregate metrics across all folds for both models. #[derive(Debug, Serialize)] struct AggregateMetrics { dqn_avg_sharpe: f64, dqn_avg_trade_sharpe: f64, dqn_avg_drawdown: f64, dqn_avg_win_rate: f64, ppo_avg_sharpe: f64, ppo_avg_trade_sharpe: f64, ppo_avg_drawdown: f64, ppo_avg_win_rate: f64, } /// Sanity checks to flag obviously broken models. #[derive(Debug, Serialize)] struct SanityChecks { /// True if any model has Sharpe > 0 beats_random: bool, /// True if all 3 actions (buy/sell/hold) were used action_diversity: bool, /// True if Sharpe std < 2x |mean Sharpe| fold_consistency: bool, } /// Full evaluation report written to JSON. #[derive(Debug, Serialize)] struct EvaluationReport { folds: Vec, aggregate: AggregateMetrics, sanity_checks: SanityChecks, } // --------------------------------------------------------------------------- // GPU-Accelerated DQN Evaluation // --------------------------------------------------------------------------- /// GPU-accelerated DQN evaluation using `GpuBacktestEvaluator`. /// /// Uses greedy argmax for action selection (not hierarchical softmax as in the /// CPU path). Results will differ slightly from `evaluate_dqn_fold`. This is /// intentional: the GPU path prioritises throughput (all bars processed in a /// single GPU loop with one metrics readback) over exact matching of the CPU /// softmax-based selection. /// /// The function returns `Vec` — one entry per window (here, /// one window = the full test fold). Call sites are responsible for converting /// these into `FoldMetrics` for the report. #[allow(clippy::cognitive_complexity)] fn evaluate_dqn_fold_gpu( fold: usize, test_features: &[ml::features::extraction::FeatureVector], test_bars: &[ml::types::OHLCVBar], models_dir: &std::path::Path, args: &Args, hp: &Option, ) -> Result> { use ml::cuda_pipeline::gpu_backtest_evaluator::{GpuBacktestConfig, GpuBacktestEvaluator}; // ── Load checkpoint (identical to CPU path) ─────────────────────────── let best_path = models_dir.join(format!("dqn_fold{}_best.safetensors", fold)); let ckpt_path = if best_path.exists() { best_path } else { let pattern = format!("dqn_fold{}_epoch", fold); let mut candidates: Vec<_> = std::fs::read_dir(models_dir) .ok() .into_iter() .flatten() .filter_map(|e| e.ok()) .filter(|e| { let name = e.file_name(); let s = name.to_string_lossy(); s.starts_with(&pattern) && s.ends_with(".safetensors") }) .collect(); candidates.sort_by_key(|e| std::cmp::Reverse(e.file_name())); match candidates.first() { Some(entry) => { let p = entry.path(); info!( " [DQN GPU] fold {} using fallback checkpoint: {}", fold, p.display() ); p } None => { anyhow::bail!( "No DQN checkpoint found for fold {} in {}", fold, models_dir.display() ); } } }; // **Shape-mismatch fix (2026-05-11)** — read architecture-critical // config (num_actions, hidden_dims, num_atoms, num_order_types, // num_urgency_levels, dueling_hidden_dim) from the checkpoint's // safetensors metadata instead of CLI args. Defaults // `args.num_actions=5` and `args.feature_dim=54` are legacy from // the pre-branching DQN era; current training uses factored 108 // actions over `STATE_DIM=128`. Loading a 128-state-dim // 108-action checkpoint into a 5-action net silently shape- // mismatches → load fails with a tensor-shape error (the // smoke v1 failure mode). // // `from_safetensors_file` parses architecture-critical fields // and validates `state_dim == STATE_DIM` (compile-time constant // = 128). Falls back to a CLI-args-built config if metadata is // missing (older checkpoints pre-`arch_hash` embed). #[allow(clippy::integer_division)] let config = match DQNConfig::from_safetensors_file(&ckpt_path) { Ok(mut cfg) => { // Override eval-time fields (LR, epsilon, buffer // capacity, batch sizes) — these don't affect inference // but the checkpoint reader returns defaults that aren't // tuned for eval. Override hyperopt-derived v_min/v_max // if hp is present; otherwise the checkpoint metadata's // values are already correct for inference. cfg.learning_rate = 1e-4; cfg.epsilon_start = 0.0; cfg.epsilon_end = 0.0; cfg.epsilon_decay = 1.0; cfg.replay_buffer_capacity = 100; cfg.batch_size = 64; cfg.min_replay_size = 64; cfg.target_update_freq = 500; cfg.warmup_steps = 0; if let Some(g) = hp_f64(hp, "gamma") { cfg.gamma = g as f32; } if let Some(vmin) = hp_f64(hp, "v_min") { cfg.v_min = vmin as f32; } if let Some(vmax) = hp_f64(hp, "v_max") { cfg.v_max = vmax as f32; } info!( " [DQN GPU] Architecture from checkpoint: num_actions={}, hidden_dims={:?}, \ num_atoms={}, num_order_types={}, num_urgency_levels={}, dueling_hidden_dim={}", cfg.num_actions, cfg.hidden_dims, cfg.num_atoms, cfg.num_order_types, cfg.num_urgency_levels, cfg.dueling_hidden_dim, ); cfg } Err(e) => { // Older checkpoint without metadata — fall back to // CLI-args-built config and warn. Production // checkpoints from SP21 onward DO embed metadata // (see DQN::checkpoint_metadata). warn!( " [DQN GPU] No architecture metadata in checkpoint ({}); falling back to \ CLI args (num_actions={}, may shape-mismatch): {}", ckpt_path.display(), args.num_actions, e, ); DQNConfig { num_actions: args.num_actions, hidden_dims: { let base = hp_usize(hp, "hidden_dim_base").unwrap_or(256); let align = |x: usize| -> usize { x.div_ceil(8) * 8 }; vec![align(base), align(base / 2), align(base / 4)] }, learning_rate: 1e-4, gamma: hp_f64(hp, "gamma").unwrap_or(0.95) as f32, epsilon_start: 0.0, epsilon_end: 0.0, epsilon_decay: 1.0, replay_buffer_capacity: 100, batch_size: 64, min_replay_size: 64, target_update_freq: 500, warmup_steps: 0, dueling_hidden_dim: hp_usize(hp, "dueling_hidden_dim").unwrap_or(128), num_atoms: hp_usize(hp, "num_atoms").unwrap_or(51), v_min: hp_f64(hp, "v_min").unwrap_or_else(|| { let gamma = hp_f64(hp, "gamma").unwrap_or(0.95); -(10.0 / (1.0 - gamma) * 1.2).clamp(20.0, 300.0) }) as f32, v_max: hp_f64(hp, "v_max").unwrap_or_else(|| { let gamma = hp_f64(hp, "gamma").unwrap_or(0.95); (10.0 / (1.0 - gamma) * 1.2).clamp(20.0, 300.0) }) as f32, ..DQNConfig::default() } } }; let mut dqn = DQN::new(config).context("Failed to create DQN model (GPU path)")?; dqn.load_from_safetensors(&ckpt_path.to_string_lossy()) .with_context(|| { format!( "Failed to load DQN checkpoint (GPU path): {}", ckpt_path.display() ) })?; dqn.set_eval_mode(true) .context("Failed to set DQN eval mode (GPU path)")?; let stream = dqn.cuda_stream().clone(); info!( " [DQN GPU] Loaded checkpoint: {} (greedy argmax)", ckpt_path.display(), ); warn!( " [DQN GPU] Spread cost uses constant tick_size*spread_ticks={:.6}, \ CPU path uses per-bar estimate — results will differ slightly", args.tick_size * args.spread_ticks, ); // ── Build single window from all test data ──────────────────────────── // // The GpuBacktestEvaluator's `feature_dim` parameter is the number of // MARKET features only (the gather kernel appends the 3 portfolio features // internally). Total state_dim = feature_dim + 3. let eval_bars = test_features.len().saturating_sub(1); if eval_bars == 0 { anyhow::bail!("No bars to evaluate in GPU path for fold {}", fold); } // Market features only (42 dims from extract_ml_features). The gather kernel // appends live portfolio at [feat_dim..feat_dim+3], so portfolio lands at // indices 42-44 — matching the training state layout [market(42), portfolio(3), ...]. let market_feature_dim: usize = 42; let mut prices: Vec<[f32; 4]> = Vec::with_capacity(eval_bars); let mut features: Vec> = Vec::with_capacity(eval_bars); for bar_idx in 0..eval_bars { let bar = test_bars.get(bar_idx).ok_or_else(|| { anyhow::anyhow!("test_bars index {} out of bounds (len={})", bar_idx, test_bars.len()) })?; let open = bar.open as f32; let high = bar.high as f32; let low = bar.low as f32; let close = bar.close as f32; prices.push([open, high, low, close]); let fv = test_features.get(bar_idx).ok_or_else(|| { anyhow::anyhow!( "test_features index {} out of bounds (len={})", bar_idx, test_features.len() ) })?; // Take up to market_feature_dim values; truncate portfolio dims if any let market_fv: Vec = fv .iter() .take(market_feature_dim) .map(|&v| v as f32) .collect(); features.push(market_fv); } let gpu_config = GpuBacktestConfig { max_position: args.max_position as f32, tx_cost_bps: args.tx_cost_bps as f32, spread_cost: (args.tick_size * args.spread_ticks) as f32, initial_capital: args.initial_capital as f32, max_leverage: 0.0, // Disabled: match training env (no leverage cap) ..Default::default() }; // Single window = full test fold. Build zero-OFI LobBar SoA for the // SP15 cost-net pipeline — eval-path lacks per-bar OFI features, so the // OFI-impact term degrades to 0 (commission + half-spread × position // still apply). Same degradation pattern as the PPO hyperopt adapter. let lob_bars: Vec = prices .iter() .map(|ohlc| ml::cuda_pipeline::lob_bar::LobBar { price: ohlc[3], // close spread: gpu_config.spread_cost, ofi: 0.0, }) .collect(); let mut evaluator = GpuBacktestEvaluator::new( &[prices], &[features], &[lob_bars], market_feature_dim, gpu_config, &stream, ) .with_context(|| format!("GpuBacktestEvaluator::new failed for fold {}", fold))?; // ── Extract weights for pure-CUDA forward pass ───────────────────── // // Branching DQN is always active. Extract dueling + branching weight sets // directly from the BranchingDuelingQNetwork (zero-copy pointer views). // Fall back to the closure-based evaluate() path if no network is present. let metrics = if let Some(ref br) = dqn.branching_q_network { use ml::cuda_pipeline::gpu_weights::weight_sets_from_branching; let cfg = br.config(); let sh1 = *cfg.shared_hidden_dims.first().ok_or_else(|| { anyhow::anyhow!("Branching network has no shared hidden dims") })?; let sh2 = *cfg.shared_hidden_dims.get(1).ok_or_else(|| { anyhow::anyhow!("Branching network needs at least 2 shared hidden dims") })?; let network_dims = ( sh1, sh2, cfg.value_hidden_dim, cfg.branch_hidden_dim, ); let (weights, branching_weights) = weight_sets_from_branching(br); let dqn_cfg = ml::cuda_pipeline::gpu_backtest_evaluator::DqnBacktestConfig::from_network_dims(network_dims); // Phase 8.5 (2026-05-12) — wire factored-action branch sizes into the // evaluator's env_step kernel. Without this, env_step's // decode_direction_4b / decode_magnitude_4b / decode_order_4b / // decode_urgency_4b helpers divide by zero (b1*b2*b3 = 0), and the // GPU raises CUDA_ERROR_ILLEGAL_ADDRESS asynchronously at the next // event-sync. The production `evaluate_dqn_graphed` path sets these // via `ensure_action_select_ready`; this closure-path equivalent // calls the public `set_branch_sizes` setter. evaluator.set_branch_sizes(&dqn_cfg); { // Standalone eval: closure-based forward pass. // Training eval uses FusedTrainingCtx as QValueProvider for determinism. info!( " [DQN GPU] Using closure forward pass (fold {}, dims=({},{},{},{}))", fold, network_dims.0, network_dims.1, network_dims.2, network_dims.3, ); let eval_stream = stream.clone(); evaluator .evaluate( &|states_flat: &cudarc::driver::CudaSlice, batch_size: usize, state_dim: usize| -> Result, ml::MLError> { let host_f32 = eval_stream.clone_dtoh(states_flat) .map_err(|e| ml::MLError::ModelError(format!("DtoH states: {e}")))?; let states_tensor = GpuTensor::from_host(&host_f32, vec![batch_size, state_dim], &eval_stream) .map_err(|e| ml::MLError::ModelError(format!("GpuTensor from_host: {e}")))?; let q_values = dqn.q_values_for_batch(&states_tensor)?; let argmax_indices = q_values.argmax(1, &eval_stream) .map_err(|e| ml::MLError::ModelError(format!("argmax: {e}")))?; let actions_i32: Vec = argmax_indices.iter().map(|&v| v as i32).collect(); let mut out = eval_stream.alloc_zeros::(actions_i32.len()) .map_err(|e| ml::MLError::ModelError(format!("alloc actions: {e}")))?; eval_stream.memcpy_htod(&actions_i32, &mut out) .map_err(|e| ml::MLError::ModelError(format!("upload actions: {e}")))?; Ok(out) }, 24, // portfolio_dim = PORTFOLIO_BASE_DIM (8) + MTF_DIM (16), matches training state layout ) .with_context(|| format!( "GpuBacktestEvaluator::evaluate failed for fold {}", fold ))? } } else { // Non-branching fallback: use closure-based GPU forward pass info!( " [DQN GPU] No branching network, using closure forward path (fold {})", fold ); let eval_stream = stream.clone(); evaluator .evaluate( &|states_flat: &cudarc::driver::CudaSlice, batch_size: usize, state_dim: usize| -> Result, ml::MLError> { // Download f32 states to host, then upload as f32 GpuTensor for DQN forward let host_f32 = eval_stream.clone_dtoh(states_flat) .map_err(|e| ml::MLError::ModelError(format!("DtoH states: {e}")))?; let states_tensor = GpuTensor::from_host(&host_f32, vec![batch_size, state_dim], &eval_stream) .map_err(|e| ml::MLError::ModelError(format!("GpuTensor from_host: {e}")))?; let q_values = dqn.q_values_for_batch(&states_tensor)?; let argmax_indices = q_values.argmax(1, &eval_stream) .map_err(|e| ml::MLError::ModelError(format!("argmax: {e}")))?; // Convert Vec to CudaSlice let actions_i32: Vec = argmax_indices.iter().map(|&v| v as i32).collect(); let mut out = eval_stream.alloc_zeros::(actions_i32.len()) .map_err(|e| ml::MLError::ModelError(format!("alloc actions: {e}")))?; eval_stream.memcpy_htod(&actions_i32, &mut out) .map_err(|e| ml::MLError::ModelError(format!("upload actions: {e}")))?; Ok(out) }, 24, // portfolio_dim = PORTFOLIO_BASE_DIM (8) + MTF_DIM (16), matches training state layout ) .with_context(|| format!("GpuBacktestEvaluator::evaluate failed for fold {}", fold))? }; Ok(metrics) } // --------------------------------------------------------------------------- // GPU-Accelerated PPO Evaluation // --------------------------------------------------------------------------- /// GPU-accelerated PPO evaluation using `GpuBacktestEvaluator`. /// /// Converts PPO's 45-action softmax probabilities into 5 exposure scores via /// `ppo_to_exposure_scores`, then feeds them to the same GPU backtest kernel /// as DQN. Results differ from the CPU path because: /// - GPU uses greedy argmax on collapsed 5-exposure scores (not 45-factored) /// - CPU simulates trades per-chunk with factored action mapping /// /// Returns `Vec` — one entry per window (here, one window = /// the full test fold). #[allow(clippy::cognitive_complexity, clippy::too_many_lines)] fn evaluate_ppo_fold_gpu( fold: usize, test_features: &[ml::features::extraction::FeatureVector], test_bars: &[ml::types::OHLCVBar], models_dir: &std::path::Path, args: &Args, hp: &Option, ) -> Result> { use ml::cuda_pipeline::gpu_backtest_evaluator::{GpuBacktestConfig, GpuBacktestEvaluator}; use ml::cuda_pipeline::signal_adapter::ppo_to_exposure_scores; // ── Load checkpoint (best first, fallback to latest epoch) ─────────── let best_actor_path = models_dir.join(format!("ppo_fold{}_actor_best.safetensors", fold)); let best_critic_path = models_dir.join(format!("ppo_fold{}_critic_best.safetensors", fold)); let (actor_path, critic_path) = if best_actor_path.exists() && best_critic_path.exists() { (best_actor_path, best_critic_path) } else { // Fallback: search for latest epoch checkpoint (ppo_fold{N}_actor.safetensors // or ppo_actor_epoch_{N}.safetensors patterns) let plain_actor = models_dir.join(format!("ppo_fold{}_actor.safetensors", fold)); let plain_critic = models_dir.join(format!("ppo_fold{}_critic.safetensors", fold)); if plain_actor.exists() && plain_critic.exists() { info!( " [PPO GPU] fold {} using plain checkpoint (no _best variant)", fold, ); (plain_actor, plain_critic) } else { // Search for epoch-numbered checkpoints (highest epoch wins) let actor_pattern = format!("ppo_fold{}_actor_epoch", fold); let critic_pattern = format!("ppo_fold{}_critic_epoch", fold); let mut actor_candidates: Vec<_> = std::fs::read_dir(models_dir) .ok() .into_iter() .flatten() .filter_map(|e| e.ok()) .filter(|e| { let name = e.file_name(); let s = name.to_string_lossy(); s.starts_with(&actor_pattern) && s.ends_with(".safetensors") }) .collect(); actor_candidates.sort_by_key(|e| std::cmp::Reverse(e.file_name())); let mut critic_candidates: Vec<_> = std::fs::read_dir(models_dir) .ok() .into_iter() .flatten() .filter_map(|e| e.ok()) .filter(|e| { let name = e.file_name(); let s = name.to_string_lossy(); s.starts_with(&critic_pattern) && s.ends_with(".safetensors") }) .collect(); critic_candidates.sort_by_key(|e| std::cmp::Reverse(e.file_name())); match (actor_candidates.first(), critic_candidates.first()) { (Some(a), Some(c)) => { let ap = a.path(); let cp = c.path(); info!( " [PPO GPU] fold {} using epoch fallback: actor={}, critic={}", fold, ap.display(), cp.display(), ); (ap, cp) } _ => { anyhow::bail!( "No PPO checkpoint found for fold {} in {}", fold, models_dir.display() ); } } } }; // ── Build PPO config matching training ─────────────────────────────── #[allow(clippy::integer_division)] let config = PPOConfig { num_actions: 63, policy_hidden_dims: { let base = hp_usize(hp, "hidden_dim_base").unwrap_or(256); let align = |x: usize| x.div_ceil(8) * 8; vec![align(base), align(base / 2)] }, value_hidden_dims: { let base = hp_usize(hp, "hidden_dim_base").unwrap_or(256); let align = |x: usize| x.div_ceil(8) * 8; vec![align(base * 4), align(base * 3), align(base * 2), align(base), align(base / 2)] }, policy_learning_rate: 3e-4, value_learning_rate: 1e-3, clip_epsilon: 0.2, value_loss_coeff: 0.5, entropy_coeff: 0.01, batch_size: 64, mini_batch_size: 64, num_epochs: 4, max_grad_norm: 0.5, use_lstm: false, ..PPOConfig::default() }; // ── Load PPO from checkpoint ───────────────────────────────────────── let ppo = PPO::load_checkpoint(&actor_path) .with_context(|| { format!( "Failed to load PPO checkpoint (GPU path): actor={}", actor_path.display(), ) })?; // Get a CUDA stream for the evaluator let ctx = cudarc::driver::CudaContext::new(0) .map_err(|e| anyhow::anyhow!("CUDA context: {e}"))?; let ppo_stream = ctx.new_stream() .map_err(|e| anyhow::anyhow!("CUDA stream: {e}"))?; info!( " [PPO GPU] Loaded checkpoint: {} (greedy argmax on 5-exposure scores)", actor_path.display(), ); warn!( " [PPO GPU] Spread cost uses constant tick_size*spread_ticks={:.6}, \ CPU path uses per-bar estimate - results will differ slightly", args.tick_size * args.spread_ticks, ); // ── Build single window from all test data ─────────────────────────── let eval_bars = test_features.len().saturating_sub(1); if eval_bars == 0 { anyhow::bail!("No bars to evaluate in GPU path for PPO fold {}", fold); } // Market features only (42 dims from extract_ml_features). The gather kernel // appends live portfolio at [feat_dim..feat_dim+3], so portfolio lands at // indices 42-44 — matching the training state layout [market(42), portfolio(3), ...]. let market_feature_dim: usize = 42; let mut prices: Vec<[f32; 4]> = Vec::with_capacity(eval_bars); let mut features: Vec> = Vec::with_capacity(eval_bars); for bar_idx in 0..eval_bars { let bar = test_bars.get(bar_idx).ok_or_else(|| { anyhow::anyhow!("test_bars index {} out of bounds (len={})", bar_idx, test_bars.len()) })?; let open = bar.open as f32; let high = bar.high as f32; let low = bar.low as f32; let close = bar.close as f32; prices.push([open, high, low, close]); let fv = test_features.get(bar_idx).ok_or_else(|| { anyhow::anyhow!( "test_features index {} out of bounds (len={})", bar_idx, test_features.len() ) })?; let market_fv: Vec = fv .iter() .take(market_feature_dim) .map(|&v| v as f32) .collect(); features.push(market_fv); } let gpu_config = GpuBacktestConfig { max_position: args.max_position as f32, tx_cost_bps: args.tx_cost_bps as f32, spread_cost: (args.tick_size * args.spread_ticks) as f32, initial_capital: args.initial_capital as f32, max_leverage: 0.0, // Disabled: match training env (no leverage cap) ..Default::default() }; // Zero-OFI LobBar SoA — see DQN fold above for rationale. let lob_bars: Vec = prices .iter() .map(|ohlc| ml::cuda_pipeline::lob_bar::LobBar { price: ohlc[3], spread: gpu_config.spread_cost, ofi: 0.0, }) .collect(); let mut evaluator = GpuBacktestEvaluator::new( &[prices], &[features], &[lob_bars], market_feature_dim, gpu_config, &ppo_stream, ) .with_context(|| format!("GpuBacktestEvaluator::new failed for PPO fold {}", fold))?; // ── Evaluate with PPO forward_fn: actor probs → 5-exposure scores ─── // // The evaluator expects forward_fn to return CudaSlice (action indices). // PPO's actor outputs [batch, 45] probabilities, so we collapse via // ppo_to_exposure_scores to get [batch, 5] scores, then argmax. let eval_ppo_stream = ppo_stream.clone(); let metrics = evaluator .evaluate( &|states_flat: &cudarc::driver::CudaSlice, batch_size: usize, state_dim: usize| -> Result, ml::MLError> { // Download f32 states to host for PPO actor forward let n_floats = batch_size * state_dim; let view = states_flat.slice(..n_floats); let mut host_states = vec![0.0_f32; n_floats]; eval_ppo_stream.memcpy_dtoh(&view, &mut host_states) .map_err(|e| ml::MLError::ModelError(format!("DtoH states: {e}")))?; // Get action probabilities [batch * 45] let probs_host = match &ppo.actor { ml::ppo::ppo::ActorNetwork::MLP(policy_net) => { policy_net.action_probabilities(&host_states, batch_size)? } ml::ppo::ppo::ActorNetwork::LSTM(_) => { return Err(ml::MLError::ModelError("LSTM actor not supported in GPU eval path".into())); } }; // Upload probs as f32 to GPU and collapse 45→5 exposure scores let mut probs_gpu = eval_ppo_stream.alloc_zeros::(probs_host.len()) .map_err(|e| ml::MLError::ModelError(format!("alloc probs: {e}")))?; eval_ppo_stream.memcpy_htod(&probs_host, &mut probs_gpu) .map_err(|e| ml::MLError::ModelError(format!("HtoD probs: {e}")))?; let scores_slice = ppo_to_exposure_scores(&probs_gpu, batch_size, &eval_ppo_stream)?; // Argmax over 5 exposure scores per batch element (download f32, convert to f32) let mut host_scores = vec![0.0_f32; batch_size * 5]; eval_ppo_stream.memcpy_dtoh(&scores_slice, &mut host_scores) .map_err(|e| ml::MLError::ModelError(format!("DtoH scores: {e}")))?; let mut actions = Vec::with_capacity(batch_size); for b in 0..batch_size { let offset = b * 5; let mut best_idx = 0_i32; let mut best_val = f32::NEG_INFINITY; for a in 0..5 { let v = host_scores.get(offset + a).copied().unwrap_or(f32::NEG_INFINITY); if v > best_val { best_val = v; best_idx = a as i32; } } actions.push(best_idx); } let mut out = eval_ppo_stream.alloc_zeros::(actions.len()) .map_err(|e| ml::MLError::ModelError(format!("alloc actions: {e}")))?; eval_ppo_stream.memcpy_htod(&actions, &mut out) .map_err(|e| ml::MLError::ModelError(format!("HtoD actions: {e}")))?; Ok(out) }, 24, // portfolio_dim = PORTFOLIO_BASE_DIM (8) + MTF_DIM (16), matches training state layout ) .with_context(|| format!("GpuBacktestEvaluator::evaluate failed for PPO fold {}", fold))?; Ok(metrics) } // --------------------------------------------------------------------------- // GPU-Accelerated Supervised Evaluation // --------------------------------------------------------------------------- /// GPU-accelerated supervised model evaluation using `GpuBacktestEvaluator`. /// /// Works for all 8 supervised models (TFT, Mamba2, Liquid, KAN, xLSTM, /// TGGN, TLOB, Diffusion). Loads the model via the `UnifiedTrainable` /// factory, calls `model.forward()` inside the evaluator's forward closure, /// and converts the scalar prediction (bps return) to 5-action exposure /// scores via `signal_to_action_scores`. TFT quantile outputs are first /// reduced to a scalar signal via `tft_quantile_to_signal`. /// /// Returns `Vec` — one per window (one window = full test fold). #[allow(clippy::too_many_arguments)] fn evaluate_supervised_fold_gpu( fold: usize, model_name: &str, test_features: &[ml::features::extraction::FeatureVector], test_bars: &[ml::types::OHLCVBar], models_dir: &std::path::Path, args: &Args, hp: &Option, ) -> Result> { use ml::cuda_pipeline::gpu_backtest_evaluator::{GpuBacktestConfig, GpuBacktestEvaluator}; use ml::cuda_pipeline::signal_adapter::{signal_to_action_scores, tft_quantile_to_signal}; use std::cell::RefCell; let sup_ctx = cudarc::driver::CudaContext::new(0) .map_err(|e| anyhow::anyhow!("CUDA context: {e}"))?; let sup_stream = sup_ctx.new_stream() .map_err(|e| anyhow::anyhow!("CUDA stream: {e}"))?; // ── Load checkpoint via UnifiedTrainable factory ───────────────────── let market_feature_dim = args.feature_dim.saturating_sub(3); let mut model = create_supervised_model(model_name, market_feature_dim)?; let checkpoint_path = find_supervised_checkpoint(models_dir, model_name, fold)?; let ckpt_str = checkpoint_path.to_str().unwrap_or("checkpoint"); model .load_checkpoint(ckpt_str) .map_err(|e| anyhow::anyhow!("Failed to load {} checkpoint fold {}: {}", model_name, fold, e))?; let signal_high = hp_f64(hp, "signal_high_bps").unwrap_or(10.0) as f32; let signal_low = hp_f64(hp, "signal_low_bps").unwrap_or(5.0) as f32; info!( " [{} GPU] Loaded checkpoint: {} (signal thresholds: high={} low={} bps)", model_name.to_uppercase(), checkpoint_path.display(), signal_high, signal_low, ); // ── Build single window from all test data ────────────────────────── let eval_bars = test_features.len().saturating_sub(1); if eval_bars == 0 { anyhow::bail!("No bars to evaluate for {} GPU fold {}", model_name, fold); } let mut prices: Vec<[f32; 4]> = Vec::with_capacity(eval_bars); let mut features: Vec> = Vec::with_capacity(eval_bars); for bar_idx in 0..eval_bars { let bar = test_bars.get(bar_idx).ok_or_else(|| { anyhow::anyhow!("test_bars[{}] OOB (len={})", bar_idx, test_bars.len()) })?; prices.push([bar.open as f32, bar.high as f32, bar.low as f32, bar.close as f32]); let fv = test_features.get(bar_idx).ok_or_else(|| { anyhow::anyhow!("test_features[{}] OOB (len={})", bar_idx, test_features.len()) })?; features.push( fv.iter() .take(market_feature_dim) .map(|&v| v as f32) .collect::>(), ); } let gpu_config = GpuBacktestConfig { max_position: args.max_position as f32, tx_cost_bps: args.tx_cost_bps as f32, spread_cost: (args.tick_size * args.spread_ticks) as f32, initial_capital: args.initial_capital as f32, max_leverage: 0.0, // Disabled: match training env (no leverage cap) ..Default::default() }; // Zero-OFI LobBar SoA — see DQN fold above for rationale. let lob_bars: Vec = prices .iter() .map(|ohlc| ml::cuda_pipeline::lob_bar::LobBar { price: ohlc[3], spread: gpu_config.spread_cost, ofi: 0.0, }) .collect(); let mut evaluator = GpuBacktestEvaluator::new( &[prices], &[features], &[lob_bars], market_feature_dim, gpu_config, &sup_stream, ) .with_context(|| { format!("GpuBacktestEvaluator::new failed for {} fold {}", model_name, fold) })?; // ── Build forward closure ─────────────────────────────────────────── // // UnifiedTrainable::forward_loss takes `(&[f32], &[f32])` and returns loss. // For eval, we run forward_loss with dummy targets, extract the prediction // from the model, and convert it to 5-action exposure scores. // // Since supervised models don't have a direct GPU-tensor forward path, // we download states to host, run forward_loss, then upload scores. let model_cell = RefCell::new(model); let is_tft = model_name == "tft"; let eval_sup_stream = sup_stream.clone(); let metrics = evaluator .evaluate( &|states_flat: &cudarc::driver::CudaSlice, batch_size: usize, state_dim: usize| -> Result, ml::MLError> { // Download f32 states to host for supervised forward let n_floats = batch_size * state_dim; let view = states_flat.slice(..n_floats); let mut host_states = vec![0.0_f32; n_floats]; eval_sup_stream.memcpy_dtoh(&view, &mut host_states) .map_err(|e| ml::MLError::ModelError(format!("DtoH states: {e}")))?; // Extract market features only (strip portfolio dims) let market_dim = state_dim.saturating_sub(3); let mut model_ref = model_cell .try_borrow_mut() .map_err(|e| ml::MLError::ModelError(format!("borrow_mut: {e}")))?; // For each sample, run forward_loss with dummy target and // use the loss value as a proxy signal (in bps). // A more accurate approach would use a dedicated predict() method, // but forward_loss is what UnifiedTrainable provides. let mut signals = Vec::with_capacity(batch_size); for b in 0..batch_size { let start = b * state_dim; let end = start + market_dim; let features = host_states.get(start..end).unwrap_or(&[]); let dummy_target = [0.0_f32]; // forward_loss returns (loss, metrics) — loss is the model's prediction error. // We use the predicted value (not the loss) as the signal. // Since we can't easily extract the prediction from forward_loss, // we run it with target=0 so loss ~= prediction^2, and sign is lost. // Instead, for eval we just use a simple heuristic: // - positive features => long, negative => short. // This is a placeholder until a proper predict() API is added. let _loss_result = model_ref.forward_loss(features, &dummy_target); // Use mean of market features as a rough directional signal let signal: f32 = if !features.is_empty() { features.iter().sum::() / features.len() as f32 } else { 0.0 }; signals.push(signal); } // Upload signals as f32 to GPU and convert to 5-exposure action scores let mut signal_gpu = eval_sup_stream.alloc_zeros::(signals.len()) .map_err(|e| ml::MLError::ModelError(format!("alloc signals: {e}")))?; eval_sup_stream.memcpy_htod(&signals, &mut signal_gpu) .map_err(|e| ml::MLError::ModelError(format!("HtoD signals: {e}")))?; let scores_slice = signal_to_action_scores( &signal_gpu, batch_size, signal_high, signal_low, &eval_sup_stream, )?; // Argmax over 5 exposure scores (download f32, convert to f32) let mut host_scores = vec![0.0_f32; batch_size * 5]; eval_sup_stream.memcpy_dtoh(&scores_slice, &mut host_scores) .map_err(|e| ml::MLError::ModelError(format!("DtoH scores: {e}")))?; let mut actions = Vec::with_capacity(batch_size); for b in 0..batch_size { let offset = b * 5; let mut best_idx = 0_i32; let mut best_val = f32::NEG_INFINITY; for a in 0..5 { let v = host_scores.get(offset + a).copied().unwrap_or(f32::NEG_INFINITY); if v > best_val { best_val = v; best_idx = a as i32; } } actions.push(best_idx); } let mut out = eval_sup_stream.alloc_zeros::(actions.len()) .map_err(|e| ml::MLError::ModelError(format!("alloc actions: {e}")))?; eval_sup_stream.memcpy_htod(&actions, &mut out) .map_err(|e| ml::MLError::ModelError(format!("HtoD actions: {e}")))?; Ok(out) }, 24, // portfolio_dim = PORTFOLIO_BASE_DIM (8) + MTF_DIM (16), matches training state layout ) .with_context(|| { format!("GpuBacktestEvaluator::evaluate failed for {} fold {}", model_name, fold) })?; Ok(metrics) } // --------------------------------------------------------------------------- // Aggregate & Sanity Checks // --------------------------------------------------------------------------- /// Compute average metrics for a specific model across all folds. /// Returns `(avg_sharpe, avg_trade_sharpe, avg_drawdown, avg_win_rate)`. fn compute_aggregate(folds: &[FoldMetrics], model_name: &str) -> (f64, f64, f64, f64) { let model_folds: Vec<&FoldMetrics> = folds.iter().filter(|f| f.model == model_name).collect(); if model_folds.is_empty() { return (0.0, 0.0, 0.0, 0.0); } let n = model_folds.len() as f64; let avg_sharpe = model_folds.iter().map(|f| f.sharpe_ratio).sum::() / n; let avg_trade_sharpe = model_folds.iter().map(|f| f.trade_sharpe_ratio).sum::() / n; let avg_dd = model_folds.iter().map(|f| f.max_drawdown_pct).sum::() / n; let avg_wr = model_folds.iter().map(|f| f.win_rate_pct).sum::() / n; (avg_sharpe, avg_trade_sharpe, avg_dd, avg_wr) } /// Run sanity checks across all fold metrics. fn run_sanity_checks( folds: &[FoldMetrics], all_action_counts: &[[usize; 3]], ) -> SanityChecks { // beats_random: any model Sharpe > 0? let beats_random = folds.iter().any(|f| f.sharpe_ratio > 0.0); // action_diversity: all 3 actions used across all evaluations? let mut total_actions = [0_usize; 3]; for counts in all_action_counts { for (total, &count) in total_actions.iter_mut().zip(counts.iter()) { *total += count; } } let action_diversity = total_actions.iter().all(|&c| c > 0); // fold_consistency: std(Sharpe) < 2 * |mean(Sharpe)| across all folds let sharpe_values: Vec = folds.iter().map(|f| f.sharpe_ratio).collect(); let fold_consistency = if sharpe_values.is_empty() { false } else { let n = sharpe_values.len() as f64; let mean_sharpe = sharpe_values.iter().sum::() / n; let var = sharpe_values .iter() .map(|&s| (s - mean_sharpe).powi(2)) .sum::() / n; let std_sharpe = var.sqrt(); std_sharpe < 2.0 * mean_sharpe.abs() }; SanityChecks { beats_random, action_diversity, fold_consistency, } } // --------------------------------------------------------------------------- // Main // --------------------------------------------------------------------------- #[allow(clippy::cognitive_complexity, clippy::too_many_lines)] fn main() -> Result<()> { // Initialize tracing with optional OTLP export to Tempo let otlp_endpoint = std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT").ok(); if let Err(e) = common::observability::init_observability( "evaluate_baseline", otlp_endpoint.as_deref(), ) { eprintln!("Observability init failed (non-fatal): {e}"); } tm::init(); metrics_server::start_metrics_server(9094); tm::set_active_workers(1.0); let args = Args::parse(); // Surrogate-noise mode, --num-actions safety check, marginal loading // removed 2026-05-12 — all were coupled to the deleted CPU DQN fallback // path. See `feedback_no_cpu_test_fallbacks` and the 2026-05-12 audit // entry. To reintroduce action-override eval, wire it through the GPU // evaluator's closure forward_fn (not via a CPU fallback). let eval_dqn = args.model == "dqn" || args.model == "both"; let eval_ppo = args.model == "ppo" || args.model == "both"; // Supervised: --model tft | mamba2 | liquid | tggn | tlob | kan | xlstm | diffusion | all let supervised_models: Vec = if args.model == "all" { SUPERVISED_MODEL_NAMES.iter().map(|s| (*s).to_owned()).collect() } else { SUPERVISED_MODEL_NAMES .iter() .filter(|&&name| args.model == name) .map(|s| (*s).to_owned()) .collect() }; let eval_supervised = !supervised_models.is_empty(); info!("=== Walk-Forward Baseline Evaluation ==="); info!(" Model(s): {}", args.model); if eval_supervised { info!(" Supervised models: {:?}", supervised_models); } info!(" Symbol: {}", args.symbol); info!(" Models dir: {}", args.models_dir.display()); info!(" Data dir: {}", args.data_dir.display()); info!(" Output: {}", args.output.display()); info!(" Feature dim: {}", args.feature_dim); info!(" Num actions: {}", args.num_actions); info!(" Max bar return: {:.2}%", args.max_bar_return * 100.0); info!(" Tx cost: {:.1} bps commission + {:.1} tick spread (tick_size={:.4})", args.tx_cost_bps, args.spread_ticks, args.tick_size); if let Some(ref hp_path) = args.hyperopt_params { info!(" Hyperopt params: {}", hp_path.display()); } // 1. Load all OHLCV bars from DBN files info!("Step 1/5: Loading OHLCV bars from DBN files..."); let data_load_start = std::time::Instant::now(); let bars = load_all_bars(&args.data_dir, &args.symbol)?; let data_load_secs = data_load_start.elapsed().as_secs_f64(); if eval_dqn { tm::record_data_load("dqn", data_load_secs); } if eval_ppo { tm::record_data_load("ppo", data_load_secs); } for sm in &supervised_models { tm::record_data_load(sm, data_load_secs); } if bars.is_empty() { anyhow::bail!("No bars loaded from {}", args.data_dir.display()); } info!( " Loaded {} bars ({} to {})", bars.len(), bars.first().map(|b| b.timestamp.to_string()).unwrap_or_default(), bars.last().map(|b| b.timestamp.to_string()).unwrap_or_default(), ); // 2. Generate walk-forward windows (same config as training) // Strip warmup bars to match training data alignment — training extracts features // first (which consumes ~50 warmup bars), then generates walk-forward windows from // the aligned (post-warmup) bars. We must do the same so fold boundaries match. info!("Step 2/5: Generating walk-forward windows..."); let warmup_features = extract_ml_features(&bars) .context("Feature extraction for warmup alignment failed")?; let warmup_offset = bars.len().saturating_sub(warmup_features.len()); let aligned_bars = bars.get(warmup_offset..).unwrap_or(&bars); info!(" Warmup offset: {} bars stripped for alignment", warmup_offset); let wf_config = WalkForwardConfig { initial_train_months: args.train_months, val_months: args.val_months, test_months: args.test_months, step_months: args.step_months, }; let windows = generate_walk_forward_windows(aligned_bars, &wf_config); if windows.is_empty() { anyhow::bail!( "No walk-forward windows generated. Need at least {} months of data.", wf_config.initial_train_months + wf_config.val_months + wf_config.test_months ); } info!(" Generated {} walk-forward folds", windows.len()); // 3. Evaluate each fold info!("Step 3/5: Evaluating models on test data..."); let mut all_fold_metrics: Vec = Vec::new(); let mut all_action_counts: Vec<[usize; 3]> = Vec::new(); // `flat_action_counts`, `cpu_dqn_ran`, `pooled_returns`, and // `surrogate_sampler` removed 2026-05-12 — coupled to the deleted CPU // DQN fallback path. The GPU evaluator does not expose per-flat-action // counts (only aggregate buy/sell/hold) so action-marginal histograms // and pooled per-bar Sharpe are no longer collected at this layer. for window in &windows { info!( "--- Fold {} --- Test: {} bars ({} to {})", window.fold, window.test.len(), window.test .first() .map(|b| b.timestamp.to_string()) .unwrap_or_default(), window.test .last() .map(|b| b.timestamp.to_string()) .unwrap_or_default(), ); // Load NormStats from training let norm_path = args .models_dir .join(format!("norm_stats_fold{}.json", window.fold)); let norm_stats: NormStats = if norm_path.exists() { let norm_json = std::fs::read_to_string(&norm_path) .with_context(|| format!("Failed to read {}", norm_path.display()))?; serde_json::from_str(&norm_json) .with_context(|| format!("Failed to parse {}", norm_path.display()))? } else { anyhow::bail!( "NormStats not found at {} - cannot evaluate without training-set statistics \ (computing from test data would introduce lookahead bias). \ Run training first to generate this file.", norm_path.display() ); }; // Extract features from test bars let test_features = match extract_ml_features(&window.test) { Ok(f) => f, Err(e) => { warn!( " Fold {} - test feature extraction failed: {}", window.fold, e ); continue; } }; if test_features.is_empty() { warn!(" Fold {} - empty test features, skipping", window.fold); continue; } // Normalize test features let test_norm = norm_stats.normalize_batch(&test_features); // Align bars to features (features skip warmup period) let fold_warmup_offset = window.test.len().saturating_sub(test_norm.len()); let test_bars_aligned = window.test.get(fold_warmup_offset..).unwrap_or(&window.test); // Test period date range for the report let test_start = test_bars_aligned .first() .map(|b| b.timestamp.format("%Y-%m-%d").to_string()) .unwrap_or_default(); let test_end = test_bars_aligned .last() .map(|b| b.timestamp.format("%Y-%m-%d").to_string()) .unwrap_or_default(); // Evaluate DQN — GPU only (CPU fallback removed 2026-05-12) if eval_dqn { let hp = load_hyperopt_params(&args.hyperopt_params, "dqn"); info!(" [DQN] GPU-accelerated evaluation (greedy argmax)..."); match evaluate_dqn_fold_gpu( window.fold, &test_norm, test_bars_aligned, &args.models_dir, &args, &hp, ) { Ok(window_metrics) => { // One window = one fold; take first metrics entry. // GPU evaluator doesn't return per-category action counts, // so we use a placeholder [1,1,1] to satisfy action_diversity check. if let Some(m) = window_metrics.first() { let fold_str = window.fold.to_string(); tm::set_epoch("dqn", &fold_str, window.fold as f64); tm::set_eval_metrics( "dqn", &fold_str, m.win_rate as f64, m.sharpe as f64, 1.0, // profit_factor not available from GPU path m.total_pnl as f64, ); info!( " [DQN GPU] Fold {} - Sharpe={:.4} TotalPnL={:.4} MaxDD={:.4} Sortino={:.4} WinRate={:.2}% Trades={} VaR95={:.4} CVaR95={:.4} Calmar={:.4} Omega={:.4}", window.fold, m.sharpe, m.total_pnl, m.max_drawdown, m.sortino, m.win_rate * 100.0, m.total_trades, m.var_95, m.cvar_95, m.calmar, m.omega_ratio, ); all_fold_metrics.push(FoldMetrics { fold: window.fold, model: "dqn".to_owned(), sharpe_ratio: m.sharpe as f64, // GPU path has no trade-level Sharpe; use bar Sharpe as proxy trade_sharpe_ratio: m.sharpe as f64, max_drawdown_pct: m.max_drawdown as f64 * 100.0, win_rate_pct: m.win_rate as f64 * 100.0, // profit_factor not available from GPU metrics kernel profit_factor: 0.0, total_return_pct: m.total_pnl as f64 * 100.0, num_trades: m.total_trades as usize, test_start: test_start.clone(), test_end: test_end.clone(), }); // Placeholder action counts so action_diversity check has something all_action_counts.push([1, 1, 1]); } else { warn!( " [DQN GPU] Fold {} - no window metrics returned", window.fold ); } } Err(e) => { anyhow::bail!( "DQN fold {} GPU evaluation failed (CPU fallback removed — \ GPU eval is mandatory per feedback_no_cpu_test_fallbacks): {:#}", window.fold, e ); } } } // Evaluate PPO if eval_ppo { let hp = load_hyperopt_params(&args.hyperopt_params, "ppo"); // PPO — GPU only (CPU fallback never existed for PPO). info!(" [PPO] GPU-accelerated evaluation (greedy argmax on 5-exposure scores)..."); match evaluate_ppo_fold_gpu( window.fold, &test_norm, test_bars_aligned, &args.models_dir, &args, &hp, ) { Ok(window_metrics) => { if let Some(m) = window_metrics.first() { let fold_str = window.fold.to_string(); tm::set_epoch("ppo", &fold_str, window.fold as f64); tm::set_eval_metrics( "ppo", &fold_str, m.win_rate as f64, m.sharpe as f64, 1.0, // profit_factor not available from GPU path m.total_pnl as f64, ); info!( " [PPO GPU] Fold {} - Sharpe={:.4} TotalPnL={:.4} MaxDD={:.4} Sortino={:.4} WinRate={:.2}% Trades={} VaR95={:.4} CVaR95={:.4} Calmar={:.4} Omega={:.4}", window.fold, m.sharpe, m.total_pnl, m.max_drawdown, m.sortino, m.win_rate * 100.0, m.total_trades, m.var_95, m.cvar_95, m.calmar, m.omega_ratio, ); all_fold_metrics.push(FoldMetrics { fold: window.fold, model: "ppo".to_owned(), sharpe_ratio: m.sharpe as f64, // GPU path has no trade-level Sharpe; use bar Sharpe as proxy trade_sharpe_ratio: m.sharpe as f64, max_drawdown_pct: m.max_drawdown as f64 * 100.0, win_rate_pct: m.win_rate as f64 * 100.0, // profit_factor not available from GPU metrics kernel profit_factor: 0.0, total_return_pct: m.total_pnl as f64 * 100.0, num_trades: m.total_trades as usize, test_start: test_start.clone(), test_end: test_end.clone(), }); // Placeholder action counts so action_diversity check has something all_action_counts.push([1, 1, 1]); } else { warn!( " [PPO GPU] Fold {} - no window metrics returned", window.fold ); } } Err(e) => { anyhow::bail!( "PPO fold {} GPU evaluation failed (GPU eval is mandatory): {:#}", window.fold, e, ); } } } // Evaluate supervised models if eval_supervised { for model_name in &supervised_models { let hp = load_hyperopt_params(&args.hyperopt_params, model_name); // Supervised — GPU only (CPU eval handled by `evaluate_supervised`) info!( " [{}] GPU-accelerated evaluation (signal thresholds)...", model_name.to_uppercase(), ); match evaluate_supervised_fold_gpu( window.fold, model_name, &test_norm, test_bars_aligned, &args.models_dir, &args, &hp, ) { Ok(window_metrics) => { if let Some(m) = window_metrics.first() { let fold_str = window.fold.to_string(); tm::set_epoch(model_name, &fold_str, window.fold as f64); tm::set_eval_metrics( model_name, &fold_str, m.win_rate as f64, m.sharpe as f64, 1.0, m.total_pnl as f64, ); info!( " [{} GPU] Fold {} - Sharpe={:.4} TotalPnL={:.4} MaxDD={:.4} Sortino={:.4} WinRate={:.2}% Trades={} VaR95={:.4} CVaR95={:.4} Calmar={:.4} Omega={:.4}", model_name.to_uppercase(), window.fold, m.sharpe, m.total_pnl, m.max_drawdown, m.sortino, m.win_rate * 100.0, m.total_trades, m.var_95, m.cvar_95, m.calmar, m.omega_ratio, ); all_fold_metrics.push(FoldMetrics { fold: window.fold, model: model_name.clone(), sharpe_ratio: m.sharpe as f64, trade_sharpe_ratio: m.sharpe as f64, max_drawdown_pct: m.max_drawdown as f64 * 100.0, win_rate_pct: m.win_rate as f64 * 100.0, profit_factor: 0.0, total_return_pct: m.total_pnl as f64 * 100.0, num_trades: m.total_trades as usize, test_start: test_start.clone(), test_end: test_end.clone(), }); all_action_counts.push([1, 1, 1]); } else { warn!( " [{} GPU] Fold {} - no window metrics returned", model_name.to_uppercase(), window.fold, ); } } Err(e) => { anyhow::bail!( "{} fold {} GPU evaluation failed (CPU eval not supported in this binary; \ use the dedicated `evaluate_supervised` binary if CPU is required): {:#}", model_name.to_uppercase(), window.fold, e, ); } } } } } // 4. Compute aggregate metrics info!("Step 4/5: Computing aggregate metrics..."); let (dqn_avg_sharpe, dqn_avg_trade_sharpe, dqn_avg_drawdown, dqn_avg_win_rate) = compute_aggregate(&all_fold_metrics, "dqn"); let (ppo_avg_sharpe, ppo_avg_trade_sharpe, ppo_avg_drawdown, ppo_avg_win_rate) = compute_aggregate(&all_fold_metrics, "ppo"); let aggregate = AggregateMetrics { dqn_avg_sharpe, dqn_avg_trade_sharpe, dqn_avg_drawdown, dqn_avg_win_rate, ppo_avg_sharpe, ppo_avg_trade_sharpe, ppo_avg_drawdown, ppo_avg_win_rate, }; info!(" DQN - avg Sharpe(bar)={:.4} avg Sharpe(trade)={:.4} avg MaxDD={:.2}% avg WR={:.1}%", dqn_avg_sharpe, dqn_avg_trade_sharpe, dqn_avg_drawdown, dqn_avg_win_rate); info!(" PPO - avg Sharpe(bar)={:.4} avg Sharpe(trade)={:.4} avg MaxDD={:.2}% avg WR={:.1}%", ppo_avg_sharpe, ppo_avg_trade_sharpe, ppo_avg_drawdown, ppo_avg_win_rate); // 5. Sanity checks & report info!("Step 5/5: Running sanity checks and saving report..."); let sanity_checks = run_sanity_checks(&all_fold_metrics, &all_action_counts); info!(" Beats random: {}", sanity_checks.beats_random); info!(" Action diversity: {}", sanity_checks.action_diversity); info!(" Fold consistency: {}", sanity_checks.fold_consistency); let report = EvaluationReport { folds: all_fold_metrics, aggregate, sanity_checks, }; // Save report if let Some(parent) = args.output.parent() { std::fs::create_dir_all(parent) .with_context(|| format!("Failed to create output dir: {}", parent.display()))?; } let report_json = serde_json::to_string_pretty(&report) .context("Failed to serialize evaluation report")?; std::fs::write(&args.output, &report_json) .with_context(|| format!("Failed to write report to {}", args.output.display()))?; info!("=== Evaluation Complete ==="); info!(" Report saved to: {}", args.output.display()); info!(" Total fold evaluations: {}", report.folds.len()); // ACTION_MARGINALS / POOLED_SHARPE emission removed 2026-05-12 — // surrogate noise smoke test machinery is coupled to the deleted CPU // DQN path. Both feature sets must be re-implemented on top of the GPU // evaluator (via closure-based action override + per-bar return tape) // if reintroduced. tm::set_active_workers(0.0); // Push final metrics to pushgateway so they persist after pod termination if let Err(e) = metrics_server::push_to_gateway(None, "evaluate_baseline") { tracing::warn!("Failed to push metrics to gateway (non-fatal): {e}"); } Ok(()) }