diff --git a/docs/plans/2026-02-23-production-pipeline-wiring-implementation.md b/docs/plans/2026-02-23-production-pipeline-wiring-implementation.md new file mode 100644 index 000000000..249a87c48 --- /dev/null +++ b/docs/plans/2026-02-23-production-pipeline-wiring-implementation.md @@ -0,0 +1,1807 @@ +# Production Pipeline Wiring — Implementation Plan + +> **For Claude:** REQUIRED SUB-SKILL: Use superpowers:executing-plans to implement this plan task-by-task. + +**Goal:** Wire all 10 ML models into the ensemble coordinator, prove the critical path with an E2E integration test, and validate Docker/CI infrastructure. + +**Architecture:** 5 new inference adapters follow two patterns: single-step (KAN, TGGN) and sequence-buffered (xLSTM, TLOB, Diffusion). Each implements `ModelInferenceAdapter` trait, maps model output to direction[-1,1] + confidence[0,1] via sigmoid. All 10 registered in main.rs with equal 0.10 weights. + +**Tech Stack:** Rust, candle 0.9.1 (VarMap/VarBuilder), tokio, safetensors, cargo test + +**Constraints:** +- Build: `SQLX_OFFLINE=true cargo check --workspace` +- Test: `SQLX_OFFLINE=true cargo test -p ml --lib` +- Clippy: `#![deny(clippy::unwrap_used, clippy::expect_used, clippy::panic, clippy::indexing_slicing)]` +- GPU: RTX 3050 Ti 4GB — hidden dims 32-64 max +- No overlap with `feat/trading-universe-data-org` branch + +**Reference adapters:** +- Single-step pattern: `ml/src/ensemble/adapters/dqn.rs` (DQN), `ml/src/ensemble/adapters/liquid.rs` (Liquid) +- Sequence-buffered pattern: `ml/src/ensemble/adapters/mamba2.rs` (Mamba2), `ml/src/ensemble/adapters/tft.rs` (TFT) + +**Common sigmoid→direction+confidence mapping** (used by all 5 new adapters): +```rust +// For models with scalar output (1-dim): +let prob = 1.0 / (1.0 + (-raw_val).exp()); +let direction = (2.0 * prob - 1.0).clamp(-1.0, 1.0); +let confidence = ((prob - 0.5).abs() * 2.0).clamp(0.0, 1.0); +``` + +--- + +### Task 1: KAN Inference Adapter + +**Files:** +- Create: `ml/src/ensemble/adapters/kan.rs` +- Modify: `ml/src/ensemble/adapters/mod.rs` + +KAN is the simplest new adapter: single-step (no buffer), VarMap-based, `KANNetwork::forward([1, 51]) → [1, 1]`. + +**Step 1: Write the adapter and tests** + +Create `ml/src/ensemble/adapters/kan.rs`: + +```rust +//! KAN (Kolmogorov-Arnold Network) inference adapter for ensemble prediction +//! +//! Wraps a `KANNetwork` and normalizes its scalar output into a directional +//! signal + confidence via sigmoid for ensemble aggregation. + +use std::sync::Mutex; + +use candle_core::{DType, Device, Tensor}; +use candle_nn::{VarBuilder, VarMap}; + +use crate::ensemble::inference_adapter::{ + EnsemblePrediction, FeatureVector, ModelInferenceAdapter, PredictionMeta, +}; +use crate::gpu::DeviceConfig; +use crate::kan::network::KANNetwork; +use crate::kan::config::KANConfig; +use crate::{MLError, MLResult}; + +/// Inference adapter that wraps a KAN (B-spline activation) network. +/// +/// Converts the scalar output through sigmoid into direction [-1, 1] +/// and confidence [0, 1] for ensemble aggregation. +#[allow(missing_debug_implementations)] +pub struct KanInferenceAdapter { + model: Mutex, + #[allow(dead_code)] + varmap: VarMap, + device: Device, + input_dim: usize, +} + +// SAFETY: KANNetwork uses candle tensors which are Send+Sync. +// The Mutex provides exclusive access for inference calls. +#[allow(unsafe_code)] +unsafe impl Send for KanInferenceAdapter {} +#[allow(unsafe_code)] +unsafe impl Sync for KanInferenceAdapter {} + +impl KanInferenceAdapter { + /// Create a new KAN inference adapter from configuration. + pub fn new(config: KANConfig) -> MLResult { + let device = DeviceConfig::Auto.resolve().unwrap_or(Device::Cpu); + let input_dim = config.layer_widths.first().copied().unwrap_or(51); + let varmap = VarMap::new(); + let vb = VarBuilder::from_varmap(&varmap, DType::F32, &device); + let network = KANNetwork::new(&config, vb)?; + + Ok(Self { + model: Mutex::new(network), + varmap, + device, + input_dim, + }) + } + + /// Create a KAN inference adapter and load weights from a safetensors checkpoint. + pub fn from_checkpoint(config: KANConfig, path: &str) -> MLResult { + let device = DeviceConfig::Auto.resolve().unwrap_or(Device::Cpu); + let input_dim = config.layer_widths.first().copied().unwrap_or(51); + let mut varmap = VarMap::new(); + let vb = VarBuilder::from_varmap(&varmap, DType::F32, &device); + let network = KANNetwork::new(&config, vb)?; + + varmap + .load(path) + .map_err(|e| MLError::ModelError(format!("Failed to load KAN checkpoint: {e}")))?; + + Ok(Self { + model: Mutex::new(network), + varmap, + device, + input_dim, + }) + } + + /// Pad or truncate feature values to match input_dim. + fn pad_features(&self, values: &[f64]) -> Vec { + let mut padded = vec![0.0f32; self.input_dim]; + let copy_len = values.len().min(self.input_dim); + for i in 0..copy_len { + if let Some(v) = values.get(i) { + if let Some(slot) = padded.get_mut(i) { + *slot = *v as f32; + } + } + } + padded + } +} + +impl ModelInferenceAdapter for KanInferenceAdapter { + fn model_name(&self) -> &str { + "KAN" + } + + fn predict(&self, features: &FeatureVector) -> MLResult { + let start = std::time::Instant::now(); + + let padded = self.pad_features(&features.values); + let input = Tensor::from_vec(padded, (1, self.input_dim), &self.device) + .map_err(|e| MLError::ModelError(format!("Failed to create KAN input tensor: {e}")))?; + + let model = self + .model + .lock() + .map_err(|e| MLError::LockError(format!("KAN model lock poisoned: {e}")))?; + let output = model.forward(&input)?; + + let squeezed = output + .squeeze(0) + .map_err(|e| MLError::ModelError(format!("Failed to squeeze KAN output: {e}")))?; + let raw: Vec = squeezed + .to_vec1() + .map_err(|e| MLError::ModelError(format!("Failed to extract KAN output: {e}")))?; + + let raw_val = raw.first().copied().unwrap_or(0.0) as f64; + + // Sigmoid → direction + confidence + let prob = 1.0 / (1.0 + (-raw_val).exp()); + let direction = (2.0 * prob - 1.0).clamp(-1.0, 1.0); + let confidence = ((prob - 0.5).abs() * 2.0).clamp(0.0, 1.0); + + let latency_us = start.elapsed().as_micros() as u64; + + Ok(EnsemblePrediction { + model_name: "KAN".to_string(), + direction, + confidence, + metadata: PredictionMeta { + latency_us, + quantiles: None, + attention_weights: None, + q_values: None, + }, + }) + } + + fn is_ready(&self) -> bool { + self.model.lock().is_ok() + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn test_config() -> KANConfig { + KANConfig { + layer_widths: vec![51, 32, 16, 1], + grid_size: 3, + spline_order: 3, + ..Default::default() + } + } + + #[test] + fn test_kan_adapter_creation() { + let adapter = KanInferenceAdapter::new(test_config()); + assert!(adapter.is_ok(), "KAN adapter creation failed: {:?}", adapter.err()); + let adapter = adapter.unwrap(); + assert_eq!(adapter.model_name(), "KAN"); + assert!(adapter.is_ready()); + } + + #[test] + fn test_kan_adapter_predict_direction_range() { + let adapter = KanInferenceAdapter::new(test_config()).unwrap(); + let fv = FeatureVector { + values: vec![0.1; 51], + timestamp: 1700000000_000_000, + }; + let pred = adapter.predict(&fv).unwrap(); + assert!( + pred.direction >= -1.0 && pred.direction <= 1.0, + "direction {} out of [-1,1]", + pred.direction + ); + assert!( + pred.confidence >= 0.0 && pred.confidence <= 1.0, + "confidence {} out of [0,1]", + pred.confidence + ); + } + + #[test] + fn test_kan_adapter_deterministic() { + let adapter = KanInferenceAdapter::new(test_config()).unwrap(); + let fv = FeatureVector { + values: vec![0.1; 51], + timestamp: 1700000000_000_000, + }; + let pred1 = adapter.predict(&fv).unwrap(); + let pred2 = adapter.predict(&fv).unwrap(); + assert_eq!( + pred1.direction, pred2.direction, + "Deterministic predictions should have same direction" + ); + } +} +``` + +**Step 2: Register the module** + +Add to `ml/src/ensemble/adapters/mod.rs`: +```rust +pub mod kan; +pub use kan::KanInferenceAdapter; +``` + +**Step 3: Build and test** + +Run: `SQLX_OFFLINE=true cargo test -p ml --lib ensemble::adapters::kan -- --nocapture` +Expected: 3 tests PASS + +**Note:** If `KANNetwork::new` or `forward` signatures differ (e.g. `&config` vs `config`, `&mut self` vs `&self`), fix the adapter to match. Check `ml/src/kan/network.rs` and `ml/src/kan/config.rs` for exact signatures. + +**Step 4: Commit** + +```bash +git add ml/src/ensemble/adapters/kan.rs ml/src/ensemble/adapters/mod.rs +git commit -m "feat(ml): add KAN inference adapter for ensemble" +``` + +--- + +### Task 2: TGGN Inference Adapter + +**Files:** +- Create: `ml/src/ensemble/adapters/tggn.rs` +- Modify: `ml/src/ensemble/adapters/mod.rs` + +TGGN's base model uses ndarray (not candle). The `TGGNTrainableAdapter` adds a 2-layer candle projection (`input_linear: node_dim→hidden_dim`, `output_linear: hidden_dim→1`). The inference adapter replicates this projection with VarMap for checkpoint compatibility. + +**Step 1: Write the adapter and tests** + +Create `ml/src/ensemble/adapters/tggn.rs`: + +```rust +//! TGGN inference adapter for ensemble prediction +//! +//! Builds a 2-layer candle projection matching the TGGNTrainableAdapter +//! architecture (input→hidden→1) for checkpoint-compatible inference. + +use std::sync::Mutex; + +use candle_core::{DType, Device, Tensor}; +use candle_nn::{linear, Linear, VarBuilder, VarMap}; + +use crate::ensemble::inference_adapter::{ + EnsemblePrediction, FeatureVector, ModelInferenceAdapter, PredictionMeta, +}; +use crate::gpu::DeviceConfig; +use crate::{MLError, MLResult}; + +/// 2-layer projection matching TGGNTrainableAdapter's candle architecture. +struct TggnProjection { + linear1: Linear, + linear2: Linear, +} + +impl TggnProjection { + fn new(input_dim: usize, hidden_dim: usize, vb: VarBuilder<'_>) -> MLResult { + let linear1 = linear(input_dim, hidden_dim, vb.pp("linear1")) + .map_err(|e| MLError::ModelError(format!("TGGN linear1 init: {e}")))?; + let linear2 = linear(hidden_dim, 1, vb.pp("linear2")) + .map_err(|e| MLError::ModelError(format!("TGGN linear2 init: {e}")))?; + Ok(Self { linear1, linear2 }) + } + + fn forward(&self, input: &Tensor) -> MLResult { + let h = self + .linear1 + .forward(input) + .map_err(|e| MLError::ModelError(format!("TGGN forward linear1: {e}")))?; + let h = h + .relu() + .map_err(|e| MLError::ModelError(format!("TGGN forward relu: {e}")))?; + self.linear2 + .forward(&h) + .map_err(|e| MLError::ModelError(format!("TGGN forward linear2: {e}"))) + } +} + +/// Inference adapter for the TGGN (Temporal Graph Gated Network). +#[allow(missing_debug_implementations)] +pub struct TggnInferenceAdapter { + model: Mutex, + #[allow(dead_code)] + varmap: VarMap, + device: Device, + input_dim: usize, +} + +#[allow(unsafe_code)] +unsafe impl Send for TggnInferenceAdapter {} +#[allow(unsafe_code)] +unsafe impl Sync for TggnInferenceAdapter {} + +impl TggnInferenceAdapter { + /// Create a new TGGN inference adapter. + /// + /// # Arguments + /// * `input_dim` - Feature vector dimension (typically 51) + /// * `hidden_dim` - Hidden layer width (32-64 for RTX 3050 Ti) + pub fn new(input_dim: usize, hidden_dim: usize) -> MLResult { + let device = DeviceConfig::Auto.resolve().unwrap_or(Device::Cpu); + let varmap = VarMap::new(); + let vb = VarBuilder::from_varmap(&varmap, DType::F32, &device); + let projection = TggnProjection::new(input_dim, hidden_dim, vb)?; + + Ok(Self { + model: Mutex::new(projection), + varmap, + device, + input_dim, + }) + } + + /// Load from a safetensors checkpoint. + pub fn from_checkpoint( + input_dim: usize, + hidden_dim: usize, + path: &str, + ) -> MLResult { + let device = DeviceConfig::Auto.resolve().unwrap_or(Device::Cpu); + let mut varmap = VarMap::new(); + let vb = VarBuilder::from_varmap(&varmap, DType::F32, &device); + let projection = TggnProjection::new(input_dim, hidden_dim, vb)?; + + varmap + .load(path) + .map_err(|e| MLError::ModelError(format!("Failed to load TGGN checkpoint: {e}")))?; + + Ok(Self { + model: Mutex::new(projection), + varmap, + device, + input_dim, + }) + } + + fn pad_features(&self, values: &[f64]) -> Vec { + let mut padded = vec![0.0f32; self.input_dim]; + let copy_len = values.len().min(self.input_dim); + for i in 0..copy_len { + if let Some(v) = values.get(i) { + if let Some(slot) = padded.get_mut(i) { + *slot = *v as f32; + } + } + } + padded + } +} + +impl ModelInferenceAdapter for TggnInferenceAdapter { + fn model_name(&self) -> &str { + "TGGN" + } + + fn predict(&self, features: &FeatureVector) -> MLResult { + let start = std::time::Instant::now(); + + let padded = self.pad_features(&features.values); + let input = Tensor::from_vec(padded, (1, self.input_dim), &self.device) + .map_err(|e| MLError::ModelError(format!("TGGN input tensor: {e}")))?; + + let model = self + .model + .lock() + .map_err(|e| MLError::LockError(format!("TGGN lock poisoned: {e}")))?; + let output = model.forward(&input)?; + + let squeezed = output + .squeeze(0) + .map_err(|e| MLError::ModelError(format!("TGGN squeeze: {e}")))?; + let raw: Vec = squeezed + .to_vec1() + .map_err(|e| MLError::ModelError(format!("TGGN extract: {e}")))?; + + let raw_val = raw.first().copied().unwrap_or(0.0) as f64; + let prob = 1.0 / (1.0 + (-raw_val).exp()); + let direction = (2.0 * prob - 1.0).clamp(-1.0, 1.0); + let confidence = ((prob - 0.5).abs() * 2.0).clamp(0.0, 1.0); + + let latency_us = start.elapsed().as_micros() as u64; + + Ok(EnsemblePrediction { + model_name: "TGGN".to_string(), + direction, + confidence, + metadata: PredictionMeta { + latency_us, + quantiles: None, + attention_weights: None, + q_values: None, + }, + }) + } + + fn is_ready(&self) -> bool { + self.model.lock().is_ok() + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_tggn_adapter_creation() { + let adapter = TggnInferenceAdapter::new(51, 64); + assert!(adapter.is_ok(), "TGGN adapter creation failed: {:?}", adapter.err()); + let adapter = adapter.unwrap(); + assert_eq!(adapter.model_name(), "TGGN"); + assert!(adapter.is_ready()); + } + + #[test] + fn test_tggn_adapter_predict_direction_range() { + let adapter = TggnInferenceAdapter::new(51, 64).unwrap(); + let fv = FeatureVector { + values: vec![0.1; 51], + timestamp: 1700000000_000_000, + }; + let pred = adapter.predict(&fv).unwrap(); + assert!( + pred.direction >= -1.0 && pred.direction <= 1.0, + "direction {} out of [-1,1]", pred.direction + ); + assert!( + pred.confidence >= 0.0 && pred.confidence <= 1.0, + "confidence {} out of [0,1]", pred.confidence + ); + } + + #[test] + fn test_tggn_adapter_deterministic() { + let adapter = TggnInferenceAdapter::new(51, 64).unwrap(); + let fv = FeatureVector { + values: vec![0.1; 51], + timestamp: 1700000000_000_000, + }; + let pred1 = adapter.predict(&fv).unwrap(); + let pred2 = adapter.predict(&fv).unwrap(); + assert_eq!(pred1.direction, pred2.direction, + "Deterministic predictions should have same direction"); + } +} +``` + +**Step 2: Register the module** + +Add to `ml/src/ensemble/adapters/mod.rs`: +```rust +pub mod tggn; +pub use tggn::TggnInferenceAdapter; +``` + +**Step 3: Build and test** + +Run: `SQLX_OFFLINE=true cargo test -p ml --lib ensemble::adapters::tggn -- --nocapture` +Expected: 3 tests PASS + +**Step 4: Commit** + +```bash +git add ml/src/ensemble/adapters/tggn.rs ml/src/ensemble/adapters/mod.rs +git commit -m "feat(ml): add TGGN inference adapter for ensemble" +``` + +--- + +### Task 3: xLSTM Inference Adapter + +**Files:** +- Create: `ml/src/ensemble/adapters/xlstm.rs` +- Modify: `ml/src/ensemble/adapters/mod.rs` + +xLSTM is a sequence model (sLSTM + mLSTM blocks). Follows the Mamba2 adapter pattern: ring buffer of feature vectors, neutral prediction until buffer full, then forward pass on full sequence. + +**Step 1: Write the adapter and tests** + +Create `ml/src/ensemble/adapters/xlstm.rs`: + +```rust +//! xLSTM inference adapter for ensemble prediction +//! +//! Wraps an `XLSTMNetwork` with a sequence buffer. Feature vectors accumulate +//! until `sequence_length` frames are collected, then the sLSTM+mLSTM blocks +//! run a forward pass and the scalar output is mapped via sigmoid to +//! direction + confidence. + +use std::collections::VecDeque; +use std::sync::Mutex; + +use candle_core::{DType, Device, Tensor}; +use candle_nn::{VarBuilder, VarMap}; + +use crate::ensemble::inference_adapter::{ + EnsemblePrediction, FeatureVector, ModelInferenceAdapter, PredictionMeta, +}; +use crate::gpu::DeviceConfig; +use crate::xlstm::config::XLSTMConfig; +use crate::xlstm::network::XLSTMNetwork; +use crate::{MLError, MLResult}; + +/// Inference adapter for xLSTM (sLSTM + mLSTM blocks) with sequence buffering. +#[allow(missing_debug_implementations)] +pub struct XlstmInferenceAdapter { + model: Mutex, + buffer: Mutex>>, + #[allow(dead_code)] + varmap: VarMap, + sequence_length: usize, + input_dim: usize, + device: Device, +} + +#[allow(unsafe_code)] +unsafe impl Send for XlstmInferenceAdapter {} +#[allow(unsafe_code)] +unsafe impl Sync for XlstmInferenceAdapter {} + +impl XlstmInferenceAdapter { + /// Create a new xLSTM inference adapter. + /// + /// # Arguments + /// * `config` - xLSTM model configuration + /// * `sequence_length` - Number of feature vectors to buffer before inference + pub fn new(config: XLSTMConfig, sequence_length: usize) -> MLResult { + let device = DeviceConfig::Auto.resolve().unwrap_or(Device::Cpu); + let input_dim = config.input_dim; + let varmap = VarMap::new(); + let vb = VarBuilder::from_varmap(&varmap, DType::F32, &device); + let network = XLSTMNetwork::new(&config, vb)?; + + Ok(Self { + model: Mutex::new(network), + buffer: Mutex::new(VecDeque::with_capacity(sequence_length)), + varmap, + sequence_length, + input_dim, + device, + }) + } + + /// Pad or truncate to input_dim, converting f64→f32. + fn pad_features(&self, values: &[f64]) -> Vec { + let mut padded = vec![0.0f32; self.input_dim]; + let copy_len = values.len().min(self.input_dim); + for i in 0..copy_len { + if let Some(v) = values.get(i) { + if let Some(slot) = padded.get_mut(i) { + *slot = *v as f32; + } + } + } + padded + } +} + +impl ModelInferenceAdapter for XlstmInferenceAdapter { + fn model_name(&self) -> &str { + "xLSTM" + } + + fn predict(&self, features: &FeatureVector) -> MLResult { + let start = std::time::Instant::now(); + + let padded = self.pad_features(&features.values); + + // Buffer phase (short-lived lock) + let buffer_ready = { + let mut buf = self + .buffer + .lock() + .map_err(|e| MLError::LockError(format!("xLSTM buffer lock poisoned: {e}")))?; + buf.push_back(padded.clone()); + while buf.len() > self.sequence_length { + buf.pop_front(); + } + buf.len() >= self.sequence_length + }; + + if !buffer_ready { + let latency_us = start.elapsed().as_micros() as u64; + return Ok(EnsemblePrediction { + model_name: "xLSTM".to_string(), + direction: 0.0, + confidence: 0.0, + metadata: PredictionMeta { + latency_us, + quantiles: None, + attention_weights: None, + q_values: None, + }, + }); + } + + // Build [1, seq_len, input_dim] tensor from buffer + let flat_data = { + let buf = self + .buffer + .lock() + .map_err(|e| MLError::LockError(format!("xLSTM buffer lock poisoned: {e}")))?; + let mut data = Vec::with_capacity(self.sequence_length * self.input_dim); + for frame in buf.iter() { + data.extend_from_slice(frame); + } + data + }; + + let input = Tensor::from_vec( + flat_data, + (1, self.sequence_length, self.input_dim), + &self.device, + ) + .map_err(|e| MLError::ModelError(format!("xLSTM input tensor: {e}")))?; + + let mut model = self + .model + .lock() + .map_err(|e| MLError::LockError(format!("xLSTM model lock poisoned: {e}")))?; + let output = model.forward(&input)?; + + // Output: [1, output_dim] → squeeze → [output_dim] + let squeezed = output + .squeeze(0) + .map_err(|e| MLError::ModelError(format!("xLSTM squeeze: {e}")))?; + let raw: Vec = squeezed + .to_vec1() + .map_err(|e| MLError::ModelError(format!("xLSTM extract: {e}")))?; + + let raw_val = raw.first().copied().unwrap_or(0.0) as f64; + let prob = 1.0 / (1.0 + (-raw_val).exp()); + let direction = (2.0 * prob - 1.0).clamp(-1.0, 1.0); + let confidence = ((prob - 0.5).abs() * 2.0).clamp(0.0, 1.0); + + let latency_us = start.elapsed().as_micros() as u64; + + Ok(EnsemblePrediction { + model_name: "xLSTM".to_string(), + direction, + confidence, + metadata: PredictionMeta { + latency_us, + quantiles: None, + attention_weights: None, + q_values: None, + }, + }) + } + + fn is_ready(&self) -> bool { + let buf_ok = self + .buffer + .lock() + .map(|buf| buf.len() >= self.sequence_length) + .unwrap_or(false); + let model_ok = self.model.lock().is_ok(); + buf_ok && model_ok + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn test_config() -> XLSTMConfig { + XLSTMConfig { + input_dim: 51, + hidden_dim: 32, + num_blocks: 2, + num_heads: 2, + output_dim: 1, + dropout: 0.0, + ..Default::default() + } + } + + const TEST_SEQ_LEN: usize = 4; + + #[test] + fn test_xlstm_adapter_creation() { + let adapter = XlstmInferenceAdapter::new(test_config(), TEST_SEQ_LEN); + assert!(adapter.is_ok(), "xLSTM adapter creation failed: {:?}", adapter.err()); + let adapter = adapter.unwrap(); + assert_eq!(adapter.model_name(), "xLSTM"); + assert!(!adapter.is_ready(), "Should not be ready with empty buffer"); + } + + #[test] + fn test_xlstm_adapter_buffers_and_predicts() { + let adapter = XlstmInferenceAdapter::new(test_config(), TEST_SEQ_LEN).unwrap(); + + for i in 0..TEST_SEQ_LEN { + let fv = FeatureVector { + values: vec![0.1 * (i as f64 + 1.0); 51], + timestamp: 1700000000_000_000 + i as i64, + }; + let pred = adapter.predict(&fv).unwrap(); + if i < TEST_SEQ_LEN - 1 { + assert_eq!(pred.direction, 0.0, "Neutral direction while buffering"); + } else { + assert!(pred.direction >= -1.0 && pred.direction <= 1.0, + "direction {} out of [-1,1]", pred.direction); + assert!(pred.confidence >= 0.0 && pred.confidence <= 1.0, + "confidence {} out of [0,1]", pred.confidence); + } + } + assert!(adapter.is_ready()); + } + + #[test] + fn test_xlstm_adapter_deterministic() { + let adapter = XlstmInferenceAdapter::new(test_config(), TEST_SEQ_LEN).unwrap(); + // Fill buffer + for _ in 0..TEST_SEQ_LEN { + let fv = FeatureVector { values: vec![0.1; 51], timestamp: 1700000000 }; + let _ = adapter.predict(&fv); + } + let fv = FeatureVector { values: vec![0.1; 51], timestamp: 1700000000 }; + let pred1 = adapter.predict(&fv).unwrap(); + let pred2 = adapter.predict(&fv).unwrap(); + assert!( + (pred1.direction - pred2.direction).abs() < 1e-6, + "direction mismatch: {} vs {}", pred1.direction, pred2.direction + ); + } +} +``` + +**Step 2: Register and test** + +Add to `ml/src/ensemble/adapters/mod.rs`: +```rust +pub mod xlstm; +pub use xlstm::XlstmInferenceAdapter; +``` + +Run: `SQLX_OFFLINE=true cargo test -p ml --lib ensemble::adapters::xlstm -- --nocapture` +Expected: 3 tests PASS + +**Note:** If `XLSTMNetwork::new` takes `config` by value instead of `&config`, or `forward` takes `&self` instead of `&mut self`, adjust accordingly. Check `ml/src/xlstm/network.rs`. + +**Step 3: Commit** + +```bash +git add ml/src/ensemble/adapters/xlstm.rs ml/src/ensemble/adapters/mod.rs +git commit -m "feat(ml): add xLSTM inference adapter for ensemble" +``` + +--- + +### Task 4: TLOB Inference Adapter + +**Files:** +- Create: `ml/src/ensemble/adapters/tlob.rs` +- Modify: `ml/src/ensemble/adapters/mod.rs` + +TLOB is a sequence transformer. Like xLSTM, it uses a ring buffer. Uses a candle projection matching `TLOBTrainableAdapter`'s architecture (flatten sequence → linear layers → scalar). + +**Step 1: Write the adapter and tests** + +Create `ml/src/ensemble/adapters/tlob.rs`. Same structure as xLSTM adapter but with: +- Model name: `"TLOB"` +- Config: `input_dim: usize, hidden_dim: usize, seq_len: usize` (no separate config struct — uses simple candle projection) +- Internal `TlobProjection` struct (3-layer MLP: `seq_len*input_dim → hidden_dim → hidden_dim/2 → 1`) +- Buffer stores `seq_len` feature vectors, flattens to `[1, seq_len * input_dim]` for forward pass + +```rust +//! TLOB inference adapter for ensemble prediction +//! +//! Wraps a candle projection matching TLOBTrainableAdapter's architecture. +//! Buffers seq_len feature vectors, flattens, and projects to scalar. + +use std::collections::VecDeque; +use std::sync::Mutex; + +use candle_core::{DType, Device, Tensor}; +use candle_nn::{linear, Linear, VarBuilder, VarMap}; + +use crate::ensemble::inference_adapter::{ + EnsemblePrediction, FeatureVector, ModelInferenceAdapter, PredictionMeta, +}; +use crate::gpu::DeviceConfig; +use crate::{MLError, MLResult}; + +struct TlobProjection { + linear1: Linear, + linear2: Linear, + linear3: Linear, +} + +impl TlobProjection { + fn new(flat_dim: usize, hidden_dim: usize, vb: VarBuilder<'_>) -> MLResult { + let linear1 = linear(flat_dim, hidden_dim, vb.pp("linear1")) + .map_err(|e| MLError::ModelError(format!("TLOB linear1: {e}")))?; + let linear2 = linear(hidden_dim, hidden_dim / 2, vb.pp("linear2")) + .map_err(|e| MLError::ModelError(format!("TLOB linear2: {e}")))?; + let linear3 = linear(hidden_dim / 2, 1, vb.pp("linear3")) + .map_err(|e| MLError::ModelError(format!("TLOB linear3: {e}")))?; + Ok(Self { linear1, linear2, linear3 }) + } + + fn forward(&self, input: &Tensor) -> MLResult { + let h = self.linear1.forward(input) + .map_err(|e| MLError::ModelError(format!("TLOB fwd l1: {e}")))?; + let h = h.relu() + .map_err(|e| MLError::ModelError(format!("TLOB fwd relu1: {e}")))?; + let h = self.linear2.forward(&h) + .map_err(|e| MLError::ModelError(format!("TLOB fwd l2: {e}")))?; + let h = h.relu() + .map_err(|e| MLError::ModelError(format!("TLOB fwd relu2: {e}")))?; + self.linear3.forward(&h) + .map_err(|e| MLError::ModelError(format!("TLOB fwd l3: {e}"))) + } +} + +/// Inference adapter for the TLOB (Transformer Limit Order Book) model. +#[allow(missing_debug_implementations)] +pub struct TlobInferenceAdapter { + model: Mutex, + buffer: Mutex>>, + #[allow(dead_code)] + varmap: VarMap, + sequence_length: usize, + feature_dim: usize, + device: Device, +} + +#[allow(unsafe_code)] +unsafe impl Send for TlobInferenceAdapter {} +#[allow(unsafe_code)] +unsafe impl Sync for TlobInferenceAdapter {} + +impl TlobInferenceAdapter { + pub fn new(feature_dim: usize, hidden_dim: usize, sequence_length: usize) -> MLResult { + let device = DeviceConfig::Auto.resolve().unwrap_or(Device::Cpu); + let flat_dim = sequence_length * feature_dim; + let varmap = VarMap::new(); + let vb = VarBuilder::from_varmap(&varmap, DType::F32, &device); + let projection = TlobProjection::new(flat_dim, hidden_dim, vb)?; + + Ok(Self { + model: Mutex::new(projection), + buffer: Mutex::new(VecDeque::with_capacity(sequence_length)), + varmap, + sequence_length, + feature_dim, + device, + }) + } + + fn pad_features(&self, values: &[f64]) -> Vec { + let mut padded = vec![0.0f32; self.feature_dim]; + let copy_len = values.len().min(self.feature_dim); + for i in 0..copy_len { + if let Some(v) = values.get(i) { + if let Some(slot) = padded.get_mut(i) { + *slot = *v as f32; + } + } + } + padded + } +} + +impl ModelInferenceAdapter for TlobInferenceAdapter { + fn model_name(&self) -> &str { + "TLOB" + } + + fn predict(&self, features: &FeatureVector) -> MLResult { + let start = std::time::Instant::now(); + let padded = self.pad_features(&features.values); + + let buffer_ready = { + let mut buf = self.buffer.lock() + .map_err(|e| MLError::LockError(format!("TLOB buffer lock: {e}")))?; + buf.push_back(padded); + while buf.len() > self.sequence_length { buf.pop_front(); } + buf.len() >= self.sequence_length + }; + + if !buffer_ready { + let latency_us = start.elapsed().as_micros() as u64; + return Ok(EnsemblePrediction { + model_name: "TLOB".to_string(), + direction: 0.0, confidence: 0.0, + metadata: PredictionMeta { latency_us, ..Default::default() }, + }); + } + + // Flatten buffer to [1, seq_len * feature_dim] + let flat_data = { + let buf = self.buffer.lock() + .map_err(|e| MLError::LockError(format!("TLOB buffer lock: {e}")))?; + let mut data = Vec::with_capacity(self.sequence_length * self.feature_dim); + for frame in buf.iter() { data.extend_from_slice(frame); } + data + }; + + let flat_dim = self.sequence_length * self.feature_dim; + let input = Tensor::from_vec(flat_data, (1, flat_dim), &self.device) + .map_err(|e| MLError::ModelError(format!("TLOB input tensor: {e}")))?; + + let model = self.model.lock() + .map_err(|e| MLError::LockError(format!("TLOB model lock: {e}")))?; + let output = model.forward(&input)?; + + let squeezed = output.squeeze(0) + .map_err(|e| MLError::ModelError(format!("TLOB squeeze: {e}")))?; + let raw: Vec = squeezed.to_vec1() + .map_err(|e| MLError::ModelError(format!("TLOB extract: {e}")))?; + + let raw_val = raw.first().copied().unwrap_or(0.0) as f64; + let prob = 1.0 / (1.0 + (-raw_val).exp()); + let direction = (2.0 * prob - 1.0).clamp(-1.0, 1.0); + let confidence = ((prob - 0.5).abs() * 2.0).clamp(0.0, 1.0); + let latency_us = start.elapsed().as_micros() as u64; + + Ok(EnsemblePrediction { + model_name: "TLOB".to_string(), + direction, confidence, + metadata: PredictionMeta { latency_us, ..Default::default() }, + }) + } + + fn is_ready(&self) -> bool { + let buf_ok = self.buffer.lock() + .map(|buf| buf.len() >= self.sequence_length).unwrap_or(false); + buf_ok && self.model.lock().is_ok() + } +} + +#[cfg(test)] +mod tests { + use super::*; + + const FEAT_DIM: usize = 51; + const HIDDEN: usize = 64; + const SEQ_LEN: usize = 4; + + #[test] + fn test_tlob_adapter_creation() { + let adapter = TlobInferenceAdapter::new(FEAT_DIM, HIDDEN, SEQ_LEN); + assert!(adapter.is_ok(), "TLOB creation failed: {:?}", adapter.err()); + assert_eq!(adapter.unwrap().model_name(), "TLOB"); + } + + #[test] + fn test_tlob_adapter_buffers_and_predicts() { + let adapter = TlobInferenceAdapter::new(FEAT_DIM, HIDDEN, SEQ_LEN).unwrap(); + for i in 0..SEQ_LEN { + let fv = FeatureVector { values: vec![0.1; 51], timestamp: i as i64 }; + let pred = adapter.predict(&fv).unwrap(); + if i < SEQ_LEN - 1 { + assert_eq!(pred.direction, 0.0); + } else { + assert!(pred.direction >= -1.0 && pred.direction <= 1.0); + assert!(pred.confidence >= 0.0 && pred.confidence <= 1.0); + } + } + assert!(adapter.is_ready()); + } + + #[test] + fn test_tlob_adapter_deterministic() { + let adapter = TlobInferenceAdapter::new(FEAT_DIM, HIDDEN, SEQ_LEN).unwrap(); + for _ in 0..SEQ_LEN { + let fv = FeatureVector { values: vec![0.1; 51], timestamp: 0 }; + let _ = adapter.predict(&fv); + } + let fv = FeatureVector { values: vec![0.1; 51], timestamp: 0 }; + let p1 = adapter.predict(&fv).unwrap(); + let p2 = adapter.predict(&fv).unwrap(); + assert!((p1.direction - p2.direction).abs() < 1e-6); + } +} +``` + +**Step 2: Register, test, commit** (same pattern as Tasks 1-3) + +Add to `mod.rs`: `pub mod tlob; pub use tlob::TlobInferenceAdapter;` + +Run: `SQLX_OFFLINE=true cargo test -p ml --lib ensemble::adapters::tlob -- --nocapture` + +```bash +git add ml/src/ensemble/adapters/tlob.rs ml/src/ensemble/adapters/mod.rs +git commit -m "feat(ml): add TLOB inference adapter for ensemble" +``` + +--- + +### Task 5: Diffusion Inference Adapter + +**Files:** +- Create: `ml/src/ensemble/adapters/diffusion.rs` +- Modify: `ml/src/ensemble/adapters/mod.rs` + +Diffusion uses the `Denoiser` network with a minimal timestep (t=1) to process features. The mean of the denoiser output provides the directional signal. + +**Step 1: Write the adapter and tests** + +Create `ml/src/ensemble/adapters/diffusion.rs`: + +```rust +//! Diffusion model inference adapter for ensemble prediction +//! +//! Wraps a `Denoiser` and runs a single forward pass at t=1 +//! (minimal noise level). The mean of the denoiser output is +//! mapped via sigmoid to direction + confidence. + +use std::sync::Mutex; + +use candle_core::{DType, Device, Tensor}; +use candle_nn::{VarBuilder, VarMap}; + +use crate::diffusion::denoiser::Denoiser; +use crate::diffusion::config::DiffusionConfig; +use crate::ensemble::inference_adapter::{ + EnsemblePrediction, FeatureVector, ModelInferenceAdapter, PredictionMeta, +}; +use crate::gpu::DeviceConfig; +use crate::{MLError, MLResult}; + +/// Inference adapter for the Diffusion (DDPM/DDIM) denoising model. +#[allow(missing_debug_implementations)] +pub struct DiffusionInferenceAdapter { + model: Mutex, + #[allow(dead_code)] + varmap: VarMap, + device: Device, + data_dim: usize, +} + +#[allow(unsafe_code)] +unsafe impl Send for DiffusionInferenceAdapter {} +#[allow(unsafe_code)] +unsafe impl Sync for DiffusionInferenceAdapter {} + +impl DiffusionInferenceAdapter { + /// Create a new Diffusion inference adapter. + /// + /// Uses `config.seq_len * config.feature_dim` as the data dimension. + /// For 51-dim feature input: set `feature_dim=51, seq_len=1`. + pub fn new(config: DiffusionConfig) -> MLResult { + let device = DeviceConfig::Auto.resolve().unwrap_or(Device::Cpu); + let data_dim = config.seq_len * config.feature_dim; + let varmap = VarMap::new(); + let vb = VarBuilder::from_varmap(&varmap, DType::F32, &device); + let denoiser = Denoiser::new( + data_dim, + config.hidden_dim, + config.num_layers, + config.time_embed_dim, + vb.pp("denoiser"), + &device, + )?; + + Ok(Self { + model: Mutex::new(denoiser), + varmap, + device, + data_dim, + }) + } + + /// Load from a safetensors checkpoint. + pub fn from_checkpoint(config: DiffusionConfig, path: &str) -> MLResult { + let device = DeviceConfig::Auto.resolve().unwrap_or(Device::Cpu); + let data_dim = config.seq_len * config.feature_dim; + let mut varmap = VarMap::new(); + let vb = VarBuilder::from_varmap(&varmap, DType::F32, &device); + let denoiser = Denoiser::new( + data_dim, + config.hidden_dim, + config.num_layers, + config.time_embed_dim, + vb.pp("denoiser"), + &device, + )?; + + varmap + .load(path) + .map_err(|e| MLError::ModelError(format!("Failed to load Diffusion checkpoint: {e}")))?; + + Ok(Self { + model: Mutex::new(denoiser), + varmap, + device, + data_dim, + }) + } + + fn pad_features(&self, values: &[f64]) -> Vec { + let mut padded = vec![0.0f32; self.data_dim]; + let copy_len = values.len().min(self.data_dim); + for i in 0..copy_len { + if let Some(v) = values.get(i) { + if let Some(slot) = padded.get_mut(i) { + *slot = *v as f32; + } + } + } + padded + } +} + +impl ModelInferenceAdapter for DiffusionInferenceAdapter { + fn model_name(&self) -> &str { + "Diffusion" + } + + fn predict(&self, features: &FeatureVector) -> MLResult { + let start = std::time::Instant::now(); + + let padded = self.pad_features(&features.values); + let input = Tensor::from_vec(padded, (1, self.data_dim), &self.device) + .map_err(|e| MLError::ModelError(format!("Diffusion input tensor: {e}")))?; + + // Timestep t=1 (minimal noise level for feature processing) + let t = Tensor::from_vec(vec![1u32], (1,), &self.device) + .map_err(|e| MLError::ModelError(format!("Diffusion timestep tensor: {e}")))?; + + let model = self + .model + .lock() + .map_err(|e| MLError::LockError(format!("Diffusion lock poisoned: {e}")))?; + let output = model.forward(&input, &t)?; + + // Mean of denoiser output → sigmoid → direction + confidence + let mean_val: f32 = output + .mean_all() + .map_err(|e| MLError::ModelError(format!("Diffusion mean: {e}")))? + .to_scalar() + .map_err(|e| MLError::ModelError(format!("Diffusion scalar: {e}")))?; + + let raw_val = mean_val as f64; + let prob = 1.0 / (1.0 + (-raw_val).exp()); + let direction = (2.0 * prob - 1.0).clamp(-1.0, 1.0); + let confidence = ((prob - 0.5).abs() * 2.0).clamp(0.0, 1.0); + + let latency_us = start.elapsed().as_micros() as u64; + + Ok(EnsemblePrediction { + model_name: "Diffusion".to_string(), + direction, + confidence, + metadata: PredictionMeta { + latency_us, + quantiles: None, + attention_weights: None, + q_values: None, + }, + }) + } + + fn is_ready(&self) -> bool { + self.model.lock().is_ok() + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::diffusion::config::NoiseSchedule; + + fn test_config() -> DiffusionConfig { + DiffusionConfig { + feature_dim: 51, + seq_len: 1, + hidden_dim: 64, + num_layers: 2, + time_embed_dim: 32, + num_timesteps: 100, + sampling_steps: 5, + schedule: NoiseSchedule::Cosine, + ..Default::default() + } + } + + #[test] + fn test_diffusion_adapter_creation() { + let adapter = DiffusionInferenceAdapter::new(test_config()); + assert!(adapter.is_ok(), "Diffusion creation failed: {:?}", adapter.err()); + let adapter = adapter.unwrap(); + assert_eq!(adapter.model_name(), "Diffusion"); + assert!(adapter.is_ready()); + } + + #[test] + fn test_diffusion_adapter_predict_direction_range() { + let adapter = DiffusionInferenceAdapter::new(test_config()).unwrap(); + let fv = FeatureVector { + values: vec![0.1; 51], + timestamp: 1700000000_000_000, + }; + let pred = adapter.predict(&fv).unwrap(); + assert!(pred.direction >= -1.0 && pred.direction <= 1.0, + "direction {} out of [-1,1]", pred.direction); + assert!(pred.confidence >= 0.0 && pred.confidence <= 1.0, + "confidence {} out of [0,1]", pred.confidence); + } + + #[test] + fn test_diffusion_adapter_deterministic() { + let adapter = DiffusionInferenceAdapter::new(test_config()).unwrap(); + let fv = FeatureVector { values: vec![0.1; 51], timestamp: 1700000000 }; + let pred1 = adapter.predict(&fv).unwrap(); + let pred2 = adapter.predict(&fv).unwrap(); + assert_eq!(pred1.direction, pred2.direction, + "Deterministic predictions should have same direction"); + } +} +``` + +**Step 2: Register, test, commit** + +Add to `mod.rs`: `pub mod diffusion; pub use diffusion::DiffusionInferenceAdapter;` + +Run: `SQLX_OFFLINE=true cargo test -p ml --lib ensemble::adapters::diffusion -- --nocapture` + +**Note:** If `Denoiser::new` or `Denoiser::forward` signatures differ from expected, check `ml/src/diffusion/denoiser.rs`. The timestep tensor type (u32 vs i64) may need adjustment. + +```bash +git add ml/src/ensemble/adapters/diffusion.rs ml/src/ensemble/adapters/mod.rs +git commit -m "feat(ml): add Diffusion inference adapter for ensemble" +``` + +--- + +### Task 6: Register All 10 Models in main.rs + +**Files:** +- Modify: `services/trading_service/src/main.rs:240-407` + +**Step 1: Add 5 new adapter registrations and rebalance weights** + +After the existing Liquid-CfC registration block (~line 397), add 5 new blocks following the same pattern. Change all weights from current (0.25/0.25/0.20/0.15/0.15) to equal (0.10 each). + +Weight changes for existing 5: +```rust +// DQN: 0.25 → 0.10 +// PPO: 0.25 → 0.10 +// TFT: 0.20 → 0.10 +// MAMBA2: 0.15 → 0.10 +// Liquid: 0.15 → 0.10 +``` + +New 5 registration blocks (add after Liquid-CfC block): + +```rust +// TGGN adapter (51-dim input, 2-layer projection) +match ml::ensemble::adapters::TggnInferenceAdapter::new(51, 64) { + Ok(adapter) => { + let bridge = Arc::new(InferenceAdapterBridge::new( + Box::new(adapter), + "TGGN".to_string(), + ml::ModelType::TGGN, + )); + if let Err(e) = coordinator + .register_loaded_model("TGGN".to_string(), bridge, 0.10) + .await + { + warn!("Failed to register TGGN adapter: {}", e); + } else { + info!("Registered TGGN inference adapter (51-dim, candle projection)"); + } + } + Err(e) => warn!("Failed to create TGGN inference adapter: {}", e), +} + +// TLOB adapter (51-dim input, sequence buffer len=10) +match ml::ensemble::adapters::TlobInferenceAdapter::new(51, 64, 10) { + Ok(adapter) => { + let bridge = Arc::new(InferenceAdapterBridge::new( + Box::new(adapter), + "TLOB".to_string(), + ml::ModelType::TLOB, + )); + if let Err(e) = coordinator + .register_loaded_model("TLOB".to_string(), bridge, 0.10) + .await + { + warn!("Failed to register TLOB adapter: {}", e); + } else { + info!("Registered TLOB inference adapter (51-dim, candle sequence)"); + } + } + Err(e) => warn!("Failed to create TLOB inference adapter: {}", e), +} + +// KAN adapter (51→32→16→1 B-spline network) +match ml::ensemble::adapters::KanInferenceAdapter::new(ml::kan::config::KANConfig { + layer_widths: vec![51, 32, 16, 1], + grid_size: 5, + spline_order: 4, + ..Default::default() +}) { + Ok(adapter) => { + let bridge = Arc::new(InferenceAdapterBridge::new( + Box::new(adapter), + "KAN".to_string(), + ml::ModelType::KAN, + )); + if let Err(e) = coordinator + .register_loaded_model("KAN".to_string(), bridge, 0.10) + .await + { + warn!("Failed to register KAN adapter: {}", e); + } else { + info!("Registered KAN inference adapter (51-dim, B-spline)"); + } + } + Err(e) => warn!("Failed to create KAN inference adapter: {}", e), +} + +// xLSTM adapter (51-dim, seq_len=10, 2 sLSTM+mLSTM blocks) +match ml::ensemble::adapters::XlstmInferenceAdapter::new( + ml::xlstm::config::XLSTMConfig { + input_dim: 51, + hidden_dim: 32, + num_blocks: 2, + num_heads: 2, + output_dim: 1, + dropout: 0.0, + ..Default::default() + }, + 10, +) { + Ok(adapter) => { + let bridge = Arc::new(InferenceAdapterBridge::new( + Box::new(adapter), + "xLSTM".to_string(), + ml::ModelType::XLSTM, + )); + if let Err(e) = coordinator + .register_loaded_model("xLSTM".to_string(), bridge, 0.10) + .await + { + warn!("Failed to register xLSTM adapter: {}", e); + } else { + info!("Registered xLSTM inference adapter (51-dim, sLSTM+mLSTM)"); + } + } + Err(e) => warn!("Failed to create xLSTM inference adapter: {}", e), +} + +// Diffusion adapter (51-dim, denoiser with t=1) +match ml::ensemble::adapters::DiffusionInferenceAdapter::new( + ml::diffusion::config::DiffusionConfig { + feature_dim: 51, + seq_len: 1, + hidden_dim: 64, + num_layers: 2, + time_embed_dim: 32, + num_timesteps: 100, + ..Default::default() + }, +) { + Ok(adapter) => { + let bridge = Arc::new(InferenceAdapterBridge::new( + Box::new(adapter), + "Diffusion".to_string(), + ml::ModelType::Diffusion, + )); + if let Err(e) = coordinator + .register_loaded_model("Diffusion".to_string(), bridge, 0.10) + .await + { + warn!("Failed to register Diffusion adapter: {}", e); + } else { + info!("Registered Diffusion inference adapter (51-dim, denoiser)"); + } + } + Err(e) => warn!("Failed to create Diffusion inference adapter: {}", e), +} +``` + +Update the model_count log line (~line 402): +```rust +info!( + "Ensemble coordinator initialized with {} real inference adapters \ + (10 models, equal weight 0.10)", + model_count +); +``` + +**Step 2: Build check** + +Run: `SQLX_OFFLINE=true cargo check -p trading_service` +Expected: compiles with 0 errors + +**Note:** Import paths for config structs may differ (e.g. `ml::kan::KANConfig` vs `ml::kan::config::KANConfig`). Fix to match actual module exports. + +**Step 3: Commit** + +```bash +git add services/trading_service/src/main.rs +git commit -m "feat(trading_service): register all 10 ML models with equal weights" +``` + +--- + +### Task 7: Wire Market Data → EnsembleCoordinator + +**Files:** +- Modify: `services/trading_service/src/main.rs` (after ensemble init) + +**Step 1: Add market data feed background task** + +After `let ensemble_coordinator = { ... };` and `let service_state = ...`, add a background tokio task that feeds the test market data generator into the ensemble coordinator: + +```rust +// Wire market data feed to ensemble coordinator for feature extraction +{ + let coordinator = Arc::clone(&ensemble_coordinator); + tokio::spawn(async move { + use trading_service::test_market_data_generator::TestMarketDataGenerator; + + let mut generator = TestMarketDataGenerator::new(); + let symbols = ["ES.FUT", "NQ.FUT", "ZN.FUT", "6E.FUT"]; + let mut bars_fed = 0u64; + + info!("Starting market data feed to ensemble coordinator (warmup ~51 bars)"); + + loop { + for symbol in &symbols { + let bar = generator.next_bar(symbol); + coordinator + .update_market_data( + symbol.to_string(), + bar.close, + bar.volume, + bar.timestamp, + ) + .await; + } + bars_fed += 1; + + if bars_fed == 51 { + info!("Ensemble coordinator warmup complete ({} bars per symbol)", bars_fed); + } + + // 1-second interval for test generator; production uses real bar arrival + tokio::time::sleep(std::time::Duration::from_secs(1)).await; + } + }); +} +``` + +**Step 2: Verify the generator module exists** + +Run: `SQLX_OFFLINE=true cargo check -p trading_service` + +**Note:** The exact module path and bar struct fields may differ. Check `services/trading_service/src/test_market_data_generator.rs` for the actual API. If the generator doesn't exist as a module, check for `dbn_market_data_generator` instead. Adapt the code to match whatever generator is available. + +**Step 3: Commit** + +```bash +git add services/trading_service/src/main.rs +git commit -m "feat(trading_service): wire market data feed to ensemble coordinator" +``` + +--- + +### Task 8: End-to-End Integration Test + +**Files:** +- Create: `services/trading_service/tests/e2e_pipeline_test.rs` + +**Step 1: Write the E2E test** + +```rust +//! End-to-end pipeline integration test +//! +//! Proves the critical path: OHLCV → features → ensemble → gates → paper trade +//! using all 10 real candle inference adapters with random (untrained) weights. +//! No external dependencies (no DB, no Docker). + +use std::sync::Arc; + +#[tokio::test(flavor = "multi_thread")] +async fn test_e2e_pipeline_ohlcv_to_paper_trade() { + use trading_service::adapter_bridge::InferenceAdapterBridge; + use trading_service::ensemble_coordinator::EnsembleCoordinator; + + let coordinator = EnsembleCoordinator::new(); + + // Register all 10 models with equal weights + let models: Vec<(&str, Box)> = vec![ + ("DQN", Box::new( + ml::ensemble::adapters::DqnInferenceAdapter::new(ml::dqn::dqn::DQNConfig { + state_dim: 51, num_actions: 45, hidden_dims: vec![64, 64], + ..Default::default() + }).expect("DQN adapter") + )), + ("PPO", Box::new( + ml::ensemble::adapters::PpoInferenceAdapter::new(ml::ppo::ppo::PPOConfig { + state_dim: 51, num_actions: 3, + policy_hidden_dims: vec![64, 64], value_hidden_dims: vec![64, 64], + ..Default::default() + }).expect("PPO adapter") + )), + ("TFT", Box::new( + ml::ensemble::adapters::TftInferenceAdapter::new( + ml::tft::TFTConfig { + input_dim: 51, hidden_dim: 32, num_heads: 4, num_layers: 1, + prediction_horizon: 1, sequence_length: 10, num_quantiles: 3, + num_static_features: 3, num_known_features: 8, num_unknown_features: 40, + ..Default::default() + }, 10, + ).expect("TFT adapter") + )), + ("MAMBA2", Box::new( + ml::ensemble::adapters::Mamba2InferenceAdapter::new( + ml::mamba::Mamba2Config { + d_model: 64, d_state: 16, d_head: 16, num_heads: 2, + expand: 1, num_layers: 1, ..Default::default() + }, 10, + ).expect("MAMBA2 adapter") + )), + ("Liquid-CfC", Box::new( + ml::ensemble::adapters::LiquidInferenceAdapter::new( + ml::liquid::candle_cfc::CfCTrainConfig { + input_size: 51, hidden_size: 32, output_size: 3, + backbone_hidden_sizes: vec![32], seq_len: 1, + device: ml::liquid::candle_cfc::DeviceConfig::Cpu, + ..ml::liquid::candle_cfc::CfCTrainConfig::default() + }, + ).expect("Liquid adapter") + )), + // 5 new adapters + ("TGGN", Box::new( + ml::ensemble::adapters::TggnInferenceAdapter::new(51, 64) + .expect("TGGN adapter") + )), + ("TLOB", Box::new( + ml::ensemble::adapters::TlobInferenceAdapter::new(51, 64, 10) + .expect("TLOB adapter") + )), + ("KAN", Box::new( + ml::ensemble::adapters::KanInferenceAdapter::new(ml::kan::config::KANConfig { + layer_widths: vec![51, 32, 16, 1], grid_size: 3, spline_order: 3, + ..Default::default() + }).expect("KAN adapter") + )), + ("xLSTM", Box::new( + ml::ensemble::adapters::XlstmInferenceAdapter::new( + ml::xlstm::config::XLSTMConfig { + input_dim: 51, hidden_dim: 32, num_blocks: 2, num_heads: 2, + output_dim: 1, dropout: 0.0, ..Default::default() + }, 10, + ).expect("xLSTM adapter") + )), + ("Diffusion", Box::new( + ml::ensemble::adapters::DiffusionInferenceAdapter::new( + ml::diffusion::config::DiffusionConfig { + feature_dim: 51, seq_len: 1, hidden_dim: 64, num_layers: 2, + time_embed_dim: 32, num_timesteps: 100, ..Default::default() + }, + ).expect("Diffusion adapter") + )), + ]; + + for (name, adapter) in models { + let model_type = match name { + "DQN" => ml::ModelType::DQN, + "PPO" => ml::ModelType::PPO, + "TFT" => ml::ModelType::TFT, + "MAMBA2" => ml::ModelType::MAMBA, + "Liquid-CfC" => ml::ModelType::LNN, + "TGGN" => ml::ModelType::TGGN, + "TLOB" => ml::ModelType::TLOB, + "KAN" => ml::ModelType::KAN, + "xLSTM" => ml::ModelType::XLSTM, + "Diffusion" => ml::ModelType::Diffusion, + _ => ml::ModelType::DQN, + }; + let bridge = Arc::new(InferenceAdapterBridge::new( + adapter, name.to_string(), model_type, + )); + coordinator + .register_loaded_model(name.to_string(), bridge, 0.10) + .await + .unwrap_or_else(|e| panic!("Failed to register {name}: {e}")); + } + + let model_count = coordinator.model_count().await; + assert_eq!(model_count, 10, "Should have 10 registered models"); + + // Feed 60 synthetic OHLCV bars (warmup period for sequence models) + let base_price = 5000.0; + for i in 0..60 { + let price = base_price + (i as f64 * 0.5).sin() * 50.0; + let volume = 1000.0 + (i as f64 * 100.0); + let timestamp = 1700000000_000_000i64 + (i * 60_000_000); + + coordinator + .update_market_data("ES.FUT".to_string(), price, volume, timestamp) + .await; + } + + // Verify feature extraction works + let features = coordinator.fetch_features_for_symbol("ES.FUT").await; + assert!(features.is_some(), "Should have features after 60 bars"); + + if let Some(fv) = &features { + assert_eq!(fv.values.len(), 51, "Feature vector should be 51-dim"); + // Verify no NaN or Inf + for (idx, val) in fv.values.iter().enumerate() { + assert!(val.is_finite(), "Feature[{}] is not finite: {}", idx, val); + } + } + + // Run ensemble prediction + let decision = coordinator.predict("ES.FUT").await; + assert!(decision.is_ok(), "Prediction should succeed: {:?}", decision.err()); + + if let Ok(ref d) = decision { + // Verify decision structure + assert!( + d.confidence >= 0.0 && d.confidence <= 1.0, + "Ensemble confidence {} out of [0,1]", d.confidence + ); + // Action should be one of Buy, Sell, Hold + // Model votes should have entries (at least from single-step models) + println!( + "E2E result: action={:?}, confidence={:.4}, model_votes={}", + d.action, d.confidence, + d.model_votes.as_ref().map(|v| v.len()).unwrap_or(0) + ); + } +} +``` + +**Step 2: Run the test** + +Run: `SQLX_OFFLINE=true cargo test -p trading_service --test e2e_pipeline_test -- --nocapture` + +**Note:** This test may need adjustments based on the actual `EnsembleCoordinator` API: +- `predict()` might be `predict_for_symbol()` or `generate_prediction()` +- The return type might differ from `EnsembleDecision` +- Feature extraction method might return `Result` instead of `Option` +- Check `services/trading_service/src/ensemble_coordinator.rs` for exact signatures + +**Step 3: Commit** + +```bash +git add services/trading_service/tests/e2e_pipeline_test.rs +git commit -m "test(trading_service): add E2E pipeline integration test (10 models)" +``` + +--- + +### Task 9: Docker Build Validation + +**Files:** +- Modify: Service Dockerfiles as needed (fix-only) + +**Step 1: Validate primary service Docker build** + +```bash +docker build -t foxhunt-trading-service:test -f services/trading_service/Dockerfile . +``` + +If it fails, fix the Dockerfile (missing deps, stale paths, wrong features). + +**Step 2: Validate remaining services** + +```bash +docker build -t foxhunt-api-gateway:test -f services/api_gateway/Dockerfile . +docker build -t foxhunt-ml-training:test -f services/ml_training_service/Dockerfile . +docker build -t foxhunt-backtesting:test -f services/backtesting_service/Dockerfile . +docker build -t foxhunt-trading-agent:test -f services/trading_agent_service/Dockerfile . +docker build -t foxhunt-broker-gateway:test -f services/broker_gateway_service/Dockerfile . +``` + +**Step 3: Verify binaries start** + +```bash +docker run --rm foxhunt-trading-service:test --help || true +docker run --rm foxhunt-api-gateway:test --help || true +``` + +**Step 4: Commit any fixes** + +```bash +git add services/*/Dockerfile Dockerfile.foxhunt-build +git commit -m "fix(docker): update Dockerfiles for 37-crate workspace" +``` + +--- + +### Task 10: CI Pipeline Triage + +**Files:** +- Modify: `.github/workflows/` (fix-only, no new files) + +**Step 1: Identify essential workflows** + +Check `.github/workflows/` for the 3 core workflows: +1. Build/compile (`ci.yml` or `build.yml`) +2. Test (`test.yml`) +3. Lint/clippy (`compilation-guard.yml`) + +**Step 2: Validate each uses correct settings** + +For each essential workflow, verify: +- `SQLX_OFFLINE=true` is set +- Rust toolchain is correct (stable or nightly as needed) +- Workspace paths are correct +- Cache keys include `Cargo.lock` hash +- Actions versions are current (e.g. `actions/checkout@v4`) + +**Step 3: Fix broken steps and delete duplicates** + +- Fix missing env vars +- Fix wrong cache keys +- Delete obvious duplicates (e.g. `comprehensive-testing.yml` vs `comprehensive_testing.yml`) + +**Step 4: Commit fixes** + +```bash +git add .github/workflows/ +git commit -m "fix(ci): update essential workflows for current workspace" +``` + +--- + +## Execution Order + +Tasks 1-5 are independent and can be parallelized (separate adapter files, no conflicts). +Task 6 depends on Tasks 1-5 (needs all adapters to register). +Task 7 depends on Task 6 (needs registered models). +Task 8 depends on Tasks 6-7 (E2E test uses all models). +Tasks 9-10 are independent of Tasks 1-8 (infrastructure-only). + +``` +Tasks 1-5 (parallel) ──→ Task 6 ──→ Task 7 ──→ Task 8 +Tasks 9-10 (parallel, independent) +``` + +## Verification + +After all tasks complete: + +```bash +# Full workspace build +SQLX_OFFLINE=true cargo check --workspace + +# ML adapter tests (15 new tests) +SQLX_OFFLINE=true cargo test -p ml --lib ensemble::adapters -- --nocapture + +# E2E pipeline test +SQLX_OFFLINE=true cargo test -p trading_service --test e2e_pipeline_test -- --nocapture + +# Clippy clean +SQLX_OFFLINE=true cargo clippy -p ml -p trading_service -- -D warnings +```