Files
foxhunt/crates/ml-ensemble/src/weight_optimizer.rs
jgrusewski db6462ba7a fix(clippy): resolve all clippy warnings across entire workspace (--all-targets)
Systematic fix of 360+ clippy errors across 37+ crates covering lib,
test, bench, and example targets. Key changes:

- Add targeted #[allow(...)] on #[cfg(test)] modules for test-only lints
  (assertions_on_result_states, float_cmp, str_to_string, indexing, etc.)
- Feature-gate broken integration tests behind __<crate>_integration flags
  where public APIs changed (trading-service, backtesting-service, etc.)
- Remove dead [[test]] entries from Cargo.toml files pointing to deleted files
- Fix production code: field_reassign_with_default, manual_range_contains,
  assert!(false) → panic!(), format!("{}") simplification, len() > 0 → !is_empty()
- Delete truly unused code (Order struct, unused methods/fields/variants)
- Convert sqlx::query!() to sqlx::query() for SQLX_OFFLINE compatibility

Result: cargo clippy --workspace --all-targets -- -D warnings = 0 errors, 0 warnings

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-13 10:18:35 +01:00

432 lines
14 KiB
Rust

//! Autonomous weight optimizer for ensemble models
//!
//! Adjusts model weights based on rolling Sharpe ratios using
//! exponentially-weighted moving average. Bounded by safety rails.
use std::collections::HashMap;
use std::time::{Duration, Instant};
use tracing::{info, warn};
/// Minimum weight for any model (prevents zeroing out)
const MIN_WEIGHT: f64 = 0.05;
/// Maximum weight for any model (prevents dominance)
const MAX_WEIGHT: f64 = 0.40;
/// Maximum weight change per optimization cycle
const MAX_STEP: f64 = 0.03;
/// Minimum number of trades before first adjustment
const MIN_OBSERVATIONS: u64 = 100;
/// Cooldown between adjustments
const COOLDOWN: Duration = Duration::from_secs(24 * 3600); // 24 hours
/// Grace period for newly deployed models
const GRACE_PERIOD: Duration = Duration::from_secs(7 * 24 * 3600); // 7 days
/// Per-model rolling performance metrics (fed from QuestDB queries)
#[derive(Debug, Clone)]
pub struct ModelRollingMetrics {
pub model_id: String,
pub sharpe_30d: f64,
pub win_rate_30d: f64,
pub prediction_accuracy: f64,
pub trade_count: u64,
pub deployed_at: Instant,
}
/// Weight optimizer configuration
#[derive(Debug, Clone)]
pub struct WeightOptimizerConfig {
pub ema_alpha: f64,
pub min_weight: f64,
pub max_weight: f64,
pub max_step: f64,
pub min_observations: u64,
pub cooldown: Duration,
pub grace_period: Duration,
}
impl Default for WeightOptimizerConfig {
fn default() -> Self {
Self {
ema_alpha: 0.1,
min_weight: MIN_WEIGHT,
max_weight: MAX_WEIGHT,
max_step: MAX_STEP,
min_observations: MIN_OBSERVATIONS,
cooldown: COOLDOWN,
grace_period: GRACE_PERIOD,
}
}
}
/// Proposed weight changes from the optimizer
#[derive(Debug, Clone)]
pub struct WeightAdjustment {
pub model_id: String,
pub old_weight: f64,
pub new_weight: f64,
pub reason: String,
}
/// Result of an optimization cycle
#[derive(Debug, Clone)]
pub enum OptimizationResult {
/// Weights adjusted successfully
Adjusted(Vec<WeightAdjustment>),
/// Not enough observations yet
InsufficientData { total_trades: u64, required: u64 },
/// Still in cooldown from last adjustment
Cooldown { remaining: Duration },
/// Kill switch is active, no adjustments allowed
KillSwitchActive,
}
/// Autonomous weight optimizer
#[derive(Debug)]
pub struct WeightOptimizer {
config: WeightOptimizerConfig,
last_adjustment: Option<Instant>,
ema_sharpe: HashMap<String, f64>,
frozen: bool,
}
impl WeightOptimizer {
pub fn new(config: WeightOptimizerConfig) -> Self {
Self {
config,
last_adjustment: None,
ema_sharpe: HashMap::new(),
frozen: false,
}
}
/// Freeze all adjustments (kill switch)
pub fn freeze(&mut self) {
self.frozen = true;
warn!("Weight optimizer frozen by kill switch");
}
/// Unfreeze adjustments (human acknowledgment)
pub fn unfreeze(&mut self) {
self.frozen = false;
info!("Weight optimizer unfrozen");
}
pub fn is_frozen(&self) -> bool {
self.frozen
}
/// Run one optimization cycle
///
/// Takes current weights and rolling metrics, returns proposed adjustments.
pub fn optimize(
&mut self,
current_weights: &HashMap<String, f64>,
metrics: &[ModelRollingMetrics],
) -> OptimizationResult {
if self.frozen {
return OptimizationResult::KillSwitchActive;
}
// Check cooldown
if let Some(last) = self.last_adjustment {
let elapsed = last.elapsed();
if elapsed < self.config.cooldown {
return OptimizationResult::Cooldown {
remaining: self.config.cooldown - elapsed,
};
}
}
// Check minimum observations
let total_trades: u64 = metrics.iter().map(|m| m.trade_count).sum();
if total_trades < self.config.min_observations {
return OptimizationResult::InsufficientData {
total_trades,
required: self.config.min_observations,
};
}
// Filter out models in grace period
let eligible: Vec<&ModelRollingMetrics> = metrics
.iter()
.filter(|m| m.deployed_at.elapsed() >= self.config.grace_period)
.collect();
if eligible.is_empty() {
return OptimizationResult::InsufficientData {
total_trades,
required: self.config.min_observations,
};
}
// Update EMA of Sharpe ratios
for m in &eligible {
let prev = self
.ema_sharpe
.get(&m.model_id)
.copied()
.unwrap_or(m.sharpe_30d);
let new_ema =
self.config.ema_alpha * m.sharpe_30d + (1.0 - self.config.ema_alpha) * prev;
self.ema_sharpe.insert(m.model_id.clone(), new_ema);
}
// Calculate raw weights from EMA Sharpe (shift to positive range)
let min_ema = self
.ema_sharpe
.values()
.cloned()
.fold(f64::INFINITY, f64::min);
let shift = if min_ema < 0.0 {
min_ema.abs() + 0.1
} else {
0.0
};
let raw_weights: HashMap<String, f64> = self
.ema_sharpe
.iter()
.map(|(id, &ema)| (id.clone(), (ema + shift).max(0.01)))
.collect();
let total_raw: f64 = raw_weights.values().sum();
if total_raw <= 0.0 {
return OptimizationResult::InsufficientData {
total_trades,
required: self.config.min_observations,
};
}
// Normalize and clamp
let mut target_weights: HashMap<String, f64> = raw_weights
.iter()
.map(|(id, &raw)| {
let normalized = raw / total_raw;
let clamped = normalized.clamp(self.config.min_weight, self.config.max_weight);
(id.clone(), clamped)
})
.collect();
// Re-normalize after clamping
let total_clamped: f64 = target_weights.values().sum();
if total_clamped > 0.0 {
for w in target_weights.values_mut() {
*w /= total_clamped;
}
}
// Apply max step limit
let mut adjustments = Vec::new();
for (model_id, &target) in &target_weights {
let current = current_weights.get(model_id).copied().unwrap_or(target);
let delta = (target - current).clamp(-self.config.max_step, self.config.max_step);
let new_weight =
(current + delta).clamp(self.config.min_weight, self.config.max_weight);
if (new_weight - current).abs() > 1e-6 {
adjustments.push(WeightAdjustment {
model_id: model_id.clone(),
old_weight: current,
new_weight,
reason: format!(
"EMA Sharpe: {:.3}, target: {:.3}",
self.ema_sharpe.get(model_id).unwrap_or(&0.0),
target
),
});
}
}
if !adjustments.is_empty() {
self.last_adjustment = Some(Instant::now());
info!("Weight optimizer adjusted {} models", adjustments.len());
}
OptimizationResult::Adjusted(adjustments)
}
}
#[cfg(test)]
#[allow(
clippy::field_reassign_with_default,
clippy::float_cmp_const,
clippy::unchecked_duration_subtraction
)]
mod tests {
use super::*;
fn make_metrics(models: &[(&str, f64, u64)]) -> Vec<ModelRollingMetrics> {
models
.iter()
.map(|(id, sharpe, trades)| ModelRollingMetrics {
model_id: id.to_string(),
sharpe_30d: *sharpe,
win_rate_30d: 0.55,
prediction_accuracy: 0.60,
trade_count: *trades,
deployed_at: Instant::now() - Duration::from_secs(30 * 24 * 3600),
})
.collect()
}
fn equal_weights(models: &[&str]) -> HashMap<String, f64> {
let w = 1.0 / models.len() as f64;
models.iter().map(|id| (id.to_string(), w)).collect()
}
#[test]
fn test_insufficient_data() {
let mut opt = WeightOptimizer::new(WeightOptimizerConfig::default());
let weights = equal_weights(&["dqn", "ppo"]);
let metrics = make_metrics(&[("dqn", 1.0, 10), ("ppo", 0.5, 10)]);
let result = opt.optimize(&weights, &metrics);
assert!(matches!(
result,
OptimizationResult::InsufficientData { .. }
));
}
#[test]
fn test_cooldown_enforced() {
let mut opt = WeightOptimizer::new(WeightOptimizerConfig::default());
let weights = equal_weights(&["dqn", "ppo"]);
let metrics = make_metrics(&[("dqn", 1.5, 200), ("ppo", 0.5, 200)]);
// First optimization succeeds
let result = opt.optimize(&weights, &metrics);
assert!(matches!(result, OptimizationResult::Adjusted(_)));
// Second immediately after should be cooldown
let result2 = opt.optimize(&weights, &metrics);
assert!(matches!(result2, OptimizationResult::Cooldown { .. }));
}
#[test]
fn test_kill_switch_freezes() {
let mut opt = WeightOptimizer::new(WeightOptimizerConfig::default());
opt.freeze();
let weights = equal_weights(&["dqn"]);
let metrics = make_metrics(&[("dqn", 1.0, 200)]);
let result = opt.optimize(&weights, &metrics);
assert!(matches!(result, OptimizationResult::KillSwitchActive));
}
#[test]
fn test_unfreeze_allows_optimization() {
let mut opt = WeightOptimizer::new(WeightOptimizerConfig::default());
opt.freeze();
assert!(opt.is_frozen());
opt.unfreeze();
assert!(!opt.is_frozen());
}
#[test]
fn test_weights_bounded() {
let mut config = WeightOptimizerConfig::default();
config.cooldown = Duration::ZERO;
let mut opt = WeightOptimizer::new(config);
let weights = equal_weights(&["dqn", "ppo", "tft"]);
let metrics = make_metrics(&[("dqn", 5.0, 200), ("ppo", -2.0, 200), ("tft", 0.5, 200)]);
let result = opt.optimize(&weights, &metrics);
if let OptimizationResult::Adjusted(adjustments) = result {
for adj in &adjustments {
assert!(
adj.new_weight >= MIN_WEIGHT,
"Weight below minimum: {}",
adj.new_weight
);
assert!(
adj.new_weight <= MAX_WEIGHT,
"Weight above maximum: {}",
adj.new_weight
);
}
}
}
#[test]
fn test_max_step_enforced() {
let mut config = WeightOptimizerConfig::default();
config.cooldown = Duration::ZERO;
let mut opt = WeightOptimizer::new(config);
// Both weights within [MIN_WEIGHT, MAX_WEIGHT] range
let mut weights = HashMap::new();
weights.insert("dqn".into(), 0.30);
weights.insert("ppo".into(), 0.30);
weights.insert("tft".into(), 0.40);
let metrics = make_metrics(&[("dqn", 5.0, 200), ("ppo", -1.0, 200), ("tft", 0.5, 200)]);
let result = opt.optimize(&weights, &metrics);
if let OptimizationResult::Adjusted(adjustments) = result {
for adj in &adjustments {
let delta = (adj.new_weight - adj.old_weight).abs();
// Step should be bounded by MAX_STEP, unless weight bounds force a correction
assert!(
delta <= MAX_STEP + 1e-6 || adj.new_weight == MIN_WEIGHT || adj.new_weight == MAX_WEIGHT,
"Step too large: {} for {} (old={}, new={})",
delta,
adj.model_id,
adj.old_weight,
adj.new_weight
);
}
}
}
#[test]
fn test_grace_period_excludes_new_models() {
let mut config = WeightOptimizerConfig::default();
config.cooldown = Duration::ZERO;
let mut opt = WeightOptimizer::new(config);
let weights = equal_weights(&["dqn", "new_model"]);
let mut metrics = make_metrics(&[("dqn", 1.0, 200)]);
// New model deployed 1 day ago (within 7-day grace period)
metrics.push(ModelRollingMetrics {
model_id: "new_model".into(),
sharpe_30d: 2.0,
win_rate_30d: 0.70,
prediction_accuracy: 0.80,
trade_count: 50,
deployed_at: Instant::now() - Duration::from_secs(24 * 3600),
});
let result = opt.optimize(&weights, &metrics);
if let OptimizationResult::Adjusted(adjustments) = result {
// new_model should not be in adjustments (grace period)
assert!(adjustments.iter().all(|a| a.model_id != "new_model"));
}
}
#[test]
fn test_weights_sum_approximately_one() {
let mut config = WeightOptimizerConfig::default();
config.cooldown = Duration::ZERO;
let mut opt = WeightOptimizer::new(config);
let weights = equal_weights(&["a", "b", "c", "d"]);
let metrics = make_metrics(&[
("a", 1.0, 200),
("b", 1.5, 200),
("c", 0.5, 200),
("d", 0.8, 200),
]);
let result = opt.optimize(&weights, &metrics);
if let OptimizationResult::Adjusted(adjustments) = result {
let mut new_weights = weights.clone();
for adj in &adjustments {
new_weights.insert(adj.model_id.clone(), adj.new_weight);
}
let sum: f64 = new_weights.values().sum();
// Due to step limits, sum may not be exactly 1.0 but should be close
assert!(
(sum - 1.0).abs() < 0.2,
"Weights sum to {} (expected ~1.0)",
sum
);
}
}
}