25+ files fixed: - hyperopt/adapters: duplicate Arc imports, Device→MlDevice - validation: regime_analysis rewritten for GpuTensor, ppo_adapter NativeDType - data_loaders: all 3 loaders migrated to GpuTensor::from_host - tft/training: MlDevice, GpuAdamW, StreamTensor, CPU loss tracking - training_pipeline: AdamW→GpuAdamW, NativeDevice fixes - features/multi_timeframe: removed to_candle_tensor - inference: ModelForward trait replaces candle_nn::Module - lib.rs: removed cuda_compat re-export - ppo/mod.rs: fixed re-export ~95 errors remain in: inference.rs, flash_attention, validation/ppo_adapter, hyperopt adapter method signatures (blocked on trainable adapter trait). Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
306 lines
11 KiB
Rust
306 lines
11 KiB
Rust
//! Shared data preparation utilities for hyperopt trainers.
|
|
//!
|
|
//! Provides common functions for loading DBN data, extracting features,
|
|
//! normalizing, and building (input, target) tensor pairs for supervised
|
|
//! model training during hyperparameter optimization.
|
|
|
|
use ml_core::device::MlDevice;
|
|
use ml_core::cuda_autograd::GpuTensor;
|
|
use serde::Serialize;
|
|
use std::fmt::Debug;
|
|
use std::path::Path;
|
|
use tracing::info;
|
|
|
|
use crate::features::extraction::OHLCVBar;
|
|
use crate::MLError;
|
|
|
|
/// Build flat (input, target) tensor pairs from extracted features and bars.
|
|
///
|
|
/// Each pair maps a single feature vector (51-dim) to a normalized close-price
|
|
/// return target. Features are z-score normalized using training-set statistics.
|
|
///
|
|
/// # Arguments
|
|
/// * `features` - Extracted feature vectors (51-dim each)
|
|
/// * `bars` - OHLCV bars aligned with features (offset by warmup period)
|
|
/// * `feature_dim` - Feature dimensionality (51)
|
|
/// * `device` - NativeDevice for tensor creation
|
|
pub fn build_flat_pairs(
|
|
features: &[[f64; 42]],
|
|
bars: &[OHLCVBar],
|
|
feature_dim: usize,
|
|
device: &MlDevice,
|
|
) -> Result<Vec<(GpuTensor, GpuTensor)>, MLError> {
|
|
if features.len() < 2 {
|
|
return Err(MLError::ModelError(
|
|
"Need at least 2 feature vectors to build pairs".to_owned(),
|
|
));
|
|
}
|
|
|
|
// The feature extractor has a warmup of 50 bars, so features[i] corresponds
|
|
// to bars[i + warmup]. We use the close-price return from the NEXT bar as target.
|
|
// Since we don't know the exact warmup offset here, we use price returns
|
|
// computed directly from consecutive features' corresponding bars.
|
|
//
|
|
// Simpler approach: target = normalized close-price return from consecutive features.
|
|
// features[i] -> target = (close[i+1] - close[i]) / close[i] where close comes
|
|
// from the bars that produced the features.
|
|
//
|
|
// Since we can't perfectly align bars<->features without the warmup constant,
|
|
// use the last (features.len()) bars for targets.
|
|
let n_bars = bars.len();
|
|
let n_features = features.len();
|
|
|
|
// features correspond to bars[offset..offset+n_features] where offset = n_bars - n_features
|
|
let bar_offset = n_bars.saturating_sub(n_features);
|
|
|
|
// Compute per-feature mean and std for z-score normalization
|
|
let mut means = vec![0.0_f64; feature_dim];
|
|
let mut vars = vec![0.0_f64; feature_dim];
|
|
let count = n_features as f64;
|
|
|
|
for feat in features.iter() {
|
|
for (j, val) in feat.iter().enumerate() {
|
|
if let Some(m) = means.get_mut(j) {
|
|
*m += val;
|
|
}
|
|
}
|
|
}
|
|
for m in means.iter_mut() {
|
|
*m /= count;
|
|
}
|
|
|
|
for feat in features.iter() {
|
|
for (j, val) in feat.iter().enumerate() {
|
|
if let (Some(v), Some(m)) = (vars.get_mut(j), means.get(j)) {
|
|
*v += (val - m) * (val - m);
|
|
}
|
|
}
|
|
}
|
|
let stds: Vec<f64> = vars
|
|
.iter()
|
|
.map(|v| (v / count).sqrt().max(1e-8))
|
|
.collect();
|
|
|
|
// Compute target return statistics for normalization
|
|
let mut returns = Vec::with_capacity(n_features.saturating_sub(1));
|
|
for i in 0..n_features.saturating_sub(1) {
|
|
let bar_idx = bar_offset + i;
|
|
let next_bar_idx = bar_offset + i + 1;
|
|
let close = bars.get(bar_idx).map(|b| b.close).unwrap_or(1.0);
|
|
let next_close = bars.get(next_bar_idx).map(|b| b.close).unwrap_or(1.0);
|
|
if close.abs() > 1e-12 {
|
|
returns.push((next_close - close) / close);
|
|
} else {
|
|
returns.push(0.0);
|
|
}
|
|
}
|
|
|
|
let ret_mean = returns.iter().sum::<f64>() / returns.len().max(1) as f64;
|
|
let ret_std = (returns
|
|
.iter()
|
|
.map(|r| (r - ret_mean) * (r - ret_mean))
|
|
.sum::<f64>()
|
|
/ returns.len().max(1) as f64)
|
|
.sqrt()
|
|
.max(1e-8);
|
|
|
|
let mut pairs = Vec::with_capacity(returns.len());
|
|
for (i, ret) in returns.iter().enumerate() {
|
|
let feat = features.get(i).ok_or_else(|| {
|
|
MLError::ModelError(format!("Feature index {} out of bounds", i))
|
|
})?;
|
|
|
|
// Z-score normalize features
|
|
let normalized: Vec<f32> = feat
|
|
.iter()
|
|
.enumerate()
|
|
.map(|(j, val)| {
|
|
let m = means.get(j).copied().unwrap_or(0.0);
|
|
let s = stds.get(j).copied().unwrap_or(1.0);
|
|
((val - m) / s) as f32
|
|
})
|
|
.collect();
|
|
|
|
let target_normalized = ((ret - ret_mean) / ret_std) as f32;
|
|
|
|
let stream = device.cuda_stream()
|
|
.map_err(|e| MLError::ModelError(format!("CUDA stream: {}", e)))?;
|
|
let input = GpuTensor::from_host(normalized.as_slice(), vec![normalized.len()], stream)
|
|
.map_err(|e| MLError::ModelError(format!("Input tensor: {}", e)))?
|
|
.reshape(vec![1, feature_dim])
|
|
.map_err(|e| MLError::ModelError(format!("Input reshape: {}", e)))?;
|
|
|
|
let target = GpuTensor::from_host(&[target_normalized], vec![1], stream)
|
|
.map_err(|e| MLError::ModelError(format!("Target tensor: {}", e)))?
|
|
.reshape(vec![1, 1])
|
|
.map_err(|e| MLError::ModelError(format!("Target reshape: {}", e)))?;
|
|
|
|
pairs.push((input, target));
|
|
}
|
|
|
|
info!(
|
|
"Built {} (input, target) pairs, feature_dim={}, ret_mean={:.6}, ret_std={:.6}",
|
|
pairs.len(),
|
|
feature_dim,
|
|
ret_mean,
|
|
ret_std
|
|
);
|
|
|
|
Ok(pairs)
|
|
}
|
|
|
|
/// Build sequenced (input, target) tensor pairs for models requiring sequence input.
|
|
///
|
|
/// Each input has shape `(1, seq_len, feature_dim)` and target is `(1, 1)`.
|
|
///
|
|
/// # Arguments
|
|
/// * `features` - Extracted feature vectors
|
|
/// * `bars` - OHLCV bars
|
|
/// * `feature_dim` - Feature dimensionality (51)
|
|
/// * `seq_len` - Sequence length
|
|
/// * `device` - MlDevice for tensor creation
|
|
pub fn build_sequence_pairs(
|
|
features: &[[f64; 42]],
|
|
bars: &[OHLCVBar],
|
|
feature_dim: usize,
|
|
seq_len: usize,
|
|
device: &MlDevice,
|
|
) -> Result<Vec<(GpuTensor, GpuTensor)>, MLError> {
|
|
if features.len() < seq_len + 1 {
|
|
return Err(MLError::ModelError(format!(
|
|
"Need at least {} features for seq_len={}, got {}",
|
|
seq_len + 1,
|
|
seq_len,
|
|
features.len()
|
|
)));
|
|
}
|
|
|
|
let n_bars = bars.len();
|
|
let n_features = features.len();
|
|
let bar_offset = n_bars.saturating_sub(n_features);
|
|
|
|
// Compute per-feature stats for normalization
|
|
let mut means = vec![0.0_f64; feature_dim];
|
|
let count = n_features as f64;
|
|
for feat in features.iter() {
|
|
for (j, val) in feat.iter().enumerate() {
|
|
if let Some(m) = means.get_mut(j) {
|
|
*m += val;
|
|
}
|
|
}
|
|
}
|
|
for m in means.iter_mut() {
|
|
*m /= count;
|
|
}
|
|
|
|
let mut vars = vec![0.0_f64; feature_dim];
|
|
for feat in features.iter() {
|
|
for (j, val) in feat.iter().enumerate() {
|
|
if let (Some(v), Some(m)) = (vars.get_mut(j), means.get(j)) {
|
|
*v += (val - m) * (val - m);
|
|
}
|
|
}
|
|
}
|
|
let stds: Vec<f64> = vars
|
|
.iter()
|
|
.map(|v| (v / count).sqrt().max(1e-8))
|
|
.collect();
|
|
|
|
// Target: close-price returns
|
|
let mut returns = Vec::with_capacity(n_features.saturating_sub(1));
|
|
for i in 0..n_features.saturating_sub(1) {
|
|
let bar_idx = bar_offset + i;
|
|
let next_bar_idx = bar_offset + i + 1;
|
|
let close = bars.get(bar_idx).map(|b| b.close).unwrap_or(1.0);
|
|
let next_close = bars.get(next_bar_idx).map(|b| b.close).unwrap_or(1.0);
|
|
if close.abs() > 1e-12 {
|
|
returns.push((next_close - close) / close);
|
|
} else {
|
|
returns.push(0.0);
|
|
}
|
|
}
|
|
|
|
let ret_mean = returns.iter().sum::<f64>() / returns.len().max(1) as f64;
|
|
let ret_std = (returns
|
|
.iter()
|
|
.map(|r| (r - ret_mean) * (r - ret_mean))
|
|
.sum::<f64>()
|
|
/ returns.len().max(1) as f64)
|
|
.sqrt()
|
|
.max(1e-8);
|
|
|
|
let num_samples = n_features.saturating_sub(seq_len);
|
|
let mut pairs = Vec::with_capacity(num_samples);
|
|
|
|
for i in 0..num_samples {
|
|
// Build sequence of normalized features
|
|
let mut seq_data = Vec::with_capacity(seq_len * feature_dim);
|
|
for t in 0..seq_len {
|
|
let feat_idx = i + t;
|
|
let feat = features.get(feat_idx).ok_or_else(|| {
|
|
MLError::ModelError(format!("Feature index {} out of bounds", feat_idx))
|
|
})?;
|
|
for (j, val) in feat.iter().enumerate() {
|
|
let m = means.get(j).copied().unwrap_or(0.0);
|
|
let s = stds.get(j).copied().unwrap_or(1.0);
|
|
seq_data.push(((val - m) / s) as f32);
|
|
}
|
|
}
|
|
|
|
// Target: return at the step after the sequence
|
|
let ret_idx = i + seq_len - 1;
|
|
let target_val = returns
|
|
.get(ret_idx)
|
|
.map(|r| ((r - ret_mean) / ret_std) as f32)
|
|
.unwrap_or(0.0);
|
|
|
|
let stream = device.cuda_stream()
|
|
.map_err(|e| MLError::ModelError(format!("CUDA stream: {}", e)))?;
|
|
let input = GpuTensor::from_host(seq_data.as_slice(), vec![seq_data.len()], stream)
|
|
.map_err(|e| MLError::ModelError(format!("Seq input tensor: {}", e)))?
|
|
.reshape(vec![1, seq_len, feature_dim])
|
|
.map_err(|e| MLError::ModelError(format!("Seq input reshape: {}", e)))?;
|
|
|
|
let target = GpuTensor::from_host(&[target_val], vec![1], stream)
|
|
.map_err(|e| MLError::ModelError(format!("Seq target tensor: {}", e)))?
|
|
.reshape(vec![1, 1])
|
|
.map_err(|e| MLError::ModelError(format!("Seq target reshape: {}", e)))?;
|
|
|
|
pairs.push((input, target));
|
|
}
|
|
|
|
info!(
|
|
"Built {} sequence pairs (seq_len={}, feature_dim={})",
|
|
pairs.len(),
|
|
seq_len,
|
|
feature_dim
|
|
);
|
|
|
|
Ok(pairs)
|
|
}
|
|
|
|
/// Write a trial result to a JSON file (appending to existing trials).
|
|
pub fn write_trial_result_json<P: Serialize + Debug>(
|
|
hyperopt_dir: &Path,
|
|
trial_result: &crate::hyperopt::traits::TrialResult<P>,
|
|
) -> Result<(), std::io::Error> {
|
|
let trials_file = hyperopt_dir.join("trials.json");
|
|
|
|
let mut all_trials = if trials_file.exists() {
|
|
let content = std::fs::read_to_string(&trials_file)?;
|
|
serde_json::from_str::<Vec<serde_json::Value>>(&content).unwrap_or_default()
|
|
} else {
|
|
Vec::new()
|
|
};
|
|
|
|
let trial_json = serde_json::to_value(trial_result)
|
|
.map_err(|e| std::io::Error::new(std::io::ErrorKind::Other, e))?;
|
|
all_trials.push(trial_json);
|
|
|
|
let content = serde_json::to_string_pretty(&all_trials)
|
|
.map_err(|e| std::io::Error::new(std::io::ErrorKind::Other, e))?;
|
|
std::fs::write(&trials_file, content)?;
|
|
|
|
Ok(())
|
|
}
|