From 903e690e9a03947c688be9221f18c8cb725ef4ae Mon Sep 17 00:00:00 2001 From: jgrusewski Date: Tue, 3 Mar 2026 18:40:03 +0100 Subject: [PATCH] feat(hyperopt): wire TPE optimizer into pipeline with --optimizer flag Add optimize_with_tpe() function that uses Tree-Parzen Estimator for Bayesian hyperparameter optimization. Add --optimizer=tpe CLI flag to hyperopt_baseline_rl binary (default: pso for backward compat). Co-Authored-By: Claude Opus 4.6 --- crates/ml/examples/hyperopt_baseline_rl.rs | 85 +++++++----- crates/ml/src/hyperopt/mod.rs | 2 +- crates/ml/src/hyperopt/optimizer.rs | 143 ++++++++++++++++++++- 3 files changed, 196 insertions(+), 34 deletions(-) diff --git a/crates/ml/examples/hyperopt_baseline_rl.rs b/crates/ml/examples/hyperopt_baseline_rl.rs index 90cd6c05e..a337ff8c5 100644 --- a/crates/ml/examples/hyperopt_baseline_rl.rs +++ b/crates/ml/examples/hyperopt_baseline_rl.rs @@ -104,7 +104,11 @@ struct Args { #[arg(long, default_value = "1.0")] spread_ticks: f64, - /// Number of parallel trial evaluations. + /// Optimizer to use: "pso" (Particle Swarm) or "tpe" (Tree-Parzen Estimator) + #[arg(long, default_value = "pso")] + optimizer: String, + + /// Number of parallel trial evaluations (PSO only; TPE is always sequential). /// Each trial uses ~1 CPU core + shared GPU for forward/backward. /// 0 = auto-detect (CPUs - 1; GPU-bound trials need minimal CPU). /// 1 = sequential. N = N concurrent trials. @@ -160,24 +164,33 @@ fn run_dqn_hyperopt(args: &Args, parallel: usize, device: &candle_core::Device) info!("Training data preloaded and cached for all {} trials", args.trials); } - let optimizer = ArgminOptimizer::builder() - .max_trials(args.trials) - .n_initial(args.n_initial) - .seed(args.seed) - .build(); - training_metrics::set_hyperopt_trial("dqn", 0.0, args.trials as f64); let start = Instant::now(); - let result = if parallel > 1 { - info!("Using parallel optimization ({} threads)", parallel); - optimizer - .optimize_parallel(trainer) - .context("DQN parallel hyperopt optimization failed")? - } else { - optimizer - .optimize(trainer) - .context("DQN hyperopt optimization failed")? + let result = match args.optimizer.as_str() { + "tpe" => { + info!("Using TPE (Tree-Parzen Estimator) optimizer"); + ml::hyperopt::optimize_with_tpe(trainer, args.trials, args.n_initial, Some(args.seed)) + .context("DQN TPE hyperopt optimization failed")? + } + _ => { + let optimizer = ArgminOptimizer::builder() + .max_trials(args.trials) + .n_initial(args.n_initial) + .seed(args.seed) + .build(); + + if parallel > 1 { + info!("Using parallel PSO optimization ({} threads)", parallel); + optimizer + .optimize_parallel(trainer) + .context("DQN parallel hyperopt optimization failed")? + } else { + optimizer + .optimize(trainer) + .context("DQN hyperopt optimization failed")? + } + } }; let elapsed = start.elapsed().as_secs_f64(); @@ -230,24 +243,33 @@ fn run_ppo_hyperopt(args: &Args, parallel: usize, device: &candle_core::Device) info!("Training data preloaded and cached for all {} trials", args.trials); } - let optimizer = ArgminOptimizer::builder() - .max_trials(args.trials) - .n_initial(args.n_initial) - .seed(args.seed) - .build(); - training_metrics::set_hyperopt_trial("ppo", 0.0, args.trials as f64); let start = Instant::now(); - let result = if parallel > 1 { - info!("Using parallel optimization ({} threads)", parallel); - optimizer - .optimize_parallel(trainer) - .context("PPO parallel hyperopt optimization failed")? - } else { - optimizer - .optimize(trainer) - .context("PPO hyperopt optimization failed")? + let result = match args.optimizer.as_str() { + "tpe" => { + info!("Using TPE (Tree-Parzen Estimator) optimizer"); + ml::hyperopt::optimize_with_tpe(trainer, args.trials, args.n_initial, Some(args.seed)) + .context("PPO TPE hyperopt optimization failed")? + } + _ => { + let optimizer = ArgminOptimizer::builder() + .max_trials(args.trials) + .n_initial(args.n_initial) + .seed(args.seed) + .build(); + + if parallel > 1 { + info!("Using parallel PSO optimization ({} threads)", parallel); + optimizer + .optimize_parallel(trainer) + .context("PPO parallel hyperopt optimization failed")? + } else { + optimizer + .optimize(trainer) + .context("PPO hyperopt optimization failed")? + } + } }; let elapsed = start.elapsed().as_secs_f64(); @@ -301,6 +323,7 @@ fn main() -> Result<()> { info!(" Hyperopt Baseline Runner"); info!("========================================"); info!("Model: {}", args.model); + info!("Optimizer: {}", args.optimizer); info!("Symbol: {}", args.symbol); info!("Trials: {}", args.trials); info!("Initial LHS samples: {}", args.n_initial); diff --git a/crates/ml/src/hyperopt/mod.rs b/crates/ml/src/hyperopt/mod.rs index 578bcdf0a..a8aea3f24 100644 --- a/crates/ml/src/hyperopt/mod.rs +++ b/crates/ml/src/hyperopt/mod.rs @@ -58,7 +58,7 @@ mod tests_argmin; // New argmin tests // Re-exports for convenience pub use observer::TrialBudgetObserver; -pub use optimizer::{ArgminOptimizer, ArgminOptimizerBuilder, TwoPhaseObjective}; +pub use optimizer::{optimize_with_tpe, ArgminOptimizer, ArgminOptimizerBuilder, TwoPhaseObjective}; pub use optimizer::{EgoboxOptimizer, EgoboxOptimizerBuilder}; // Backward compatibility pub use traits::{ HardwareBudget, HyperoptStrategy, HyperparameterOptimizable, OptimizationResult, diff --git a/crates/ml/src/hyperopt/optimizer.rs b/crates/ml/src/hyperopt/optimizer.rs index 3e25a40ae..5bef2f62d 100644 --- a/crates/ml/src/hyperopt/optimizer.rs +++ b/crates/ml/src/hyperopt/optimizer.rs @@ -429,9 +429,12 @@ impl ArgminOptimizer { Ok(result) } - /// Evaluate a single point + /// Evaluate a single point in the parameter space. + /// + /// Converts `continuous_vec` to typed parameters, trains the model, + /// records the trial result, and returns the objective value. #[allow(clippy::unwrap_in_result)] - fn evaluate_point( + pub(crate) fn evaluate_point( continuous_vec: &[f64], model: &mut M, trial_results: &Arc>>>, @@ -823,6 +826,142 @@ impl ArgminOptimizer { } } +/// Run optimization using Tree-Parzen Estimator (TPE) instead of PSO. +/// +/// TPE builds separate density models for "good" and "bad" trials, +/// then suggests new points by maximizing Expected Improvement. +/// Better sample efficiency than PSO in 20-50D spaces. +/// +/// # Arguments +/// +/// * `model` - Model implementing `HyperparameterOptimizable` (consumed) +/// * `max_trials` - Total number of trials to evaluate +/// * `n_initial` - Number of initial LHS samples before TPE kicks in +/// * `seed` - Optional random seed for reproducibility +/// +/// # Returns +/// +/// `OptimizationResult` containing best parameters and full trial history +pub fn optimize_with_tpe( + mut model: M, + max_trials: usize, + n_initial: usize, + seed: Option, +) -> Result> +where + M: HyperparameterOptimizable + Send, + M::Params: ParameterSpace + Send, +{ + use crate::hyperopt::tpe::{TpeConfig, TpeOptimizer}; + + info!("╔═══════════════════════════════════════════════════════════╗"); + info!("║ Bayesian Hyperparameter Optimization (TPE) ║"); + info!("╚═══════════════════════════════════════════════════════════╝"); + + let budget = crate::hyperopt::traits::HardwareBudget::detect(); + let bounds = M::Params::continuous_bounds_for(&budget); + let n_params = bounds.len(); + + if n_params == 0 { + return Err(MLError::ConfigError { + reason: "Parameter space has zero dimensions".to_owned(), + } + .into()); + } + + // Clamp n_initial to valid range: at least 1, at most max_trials - 1 + let n_initial = n_initial.max(1).min(max_trials.saturating_sub(1).max(1)); + + info!("Configuration:"); + info!(" Optimizer: TPE (Tree-Parzen Estimator)"); + info!(" Max Trials: {}", max_trials); + info!(" Initial LHS Samples: {}", n_initial); + info!(" Parameters: {}", n_params); + info!(" Gamma (good quantile): 0.25"); + info!(" EI Candidates: 100"); + + let param_names = M::Params::param_names(); + for (i, name) in param_names.iter().enumerate() { + if let Some(&(lo, hi)) = bounds.get(i) { + info!(" {} - [{:.6}, {:.6}]", name, lo, hi); + } + } + + let tpe_config = TpeConfig { + n_dims: n_params, + max_trials, + n_initial, + gamma: 0.25, + n_candidates: 100, + seed, + }; + let mut tpe = TpeOptimizer::new(tpe_config); + + let trial_results = Arc::new(Mutex::new(Vec::new())); + let trial_counter = Arc::new(std::sync::atomic::AtomicUsize::new(0)); + + for trial_idx in 0..max_trials { + let continuous_vec = tpe.suggest(&bounds); + + info!( + "TPE Trial {}/{}: evaluating suggested point", + trial_idx + 1, + max_trials + ); + + ArgminOptimizer::evaluate_point( + &continuous_vec, + &mut model, + &trial_results, + &trial_counter, + ¶m_names, + )?; + + // Feed the objective back to TPE so it can update its density models + let results = trial_results.lock().map_err(|e| { + MLError::ConcurrencyError { + operation: format!("lock trial results: {}", e), + } + })?; + if let Some(last) = results.last() { + tpe.add_trial(continuous_vec, last.objective); + info!( + "TPE Trial {}/{}: objective = {:.6}", + trial_idx + 1, + max_trials, + last.objective + ); + } + } + + // Extract results + let trials = match Arc::try_unwrap(trial_results) { + Ok(mutex) => mutex.into_inner().map_err(|e| { + anyhow::anyhow!("Failed to unwrap trial results mutex: {}", e) + })?, + Err(arc) => arc.lock().map_err(|e| { + anyhow::anyhow!("Failed to lock trial results: {}", e) + })?.clone(), + }; + + if trials.is_empty() { + return Err( + MLError::ModelError("No valid trials completed".to_owned()).into(), + ); + } + + let result = OptimizationResult::from_trials(trials); + + info!("═══ TPE Optimization Complete ═══"); + info!( + " Best trial: objective {:.6}", + result.best_objective + ); + info!(" Total trials: {}", result.all_trials.len()); + + Ok(result) +} + /// Trait for models supporting two-phase objective switching. /// /// Used by [`ArgminOptimizer::optimize_two_phase()`] to switch between