Files
foxhunt/services/integration_tests/src/metrics_validation.rs
jgrusewski 83629f9ca8 feat(deployment): Complete Runpod GPU deployment infrastructure
Implement comprehensive Runpod deployment with S3 volume mount architecture for
FP32 ML model training on Tesla V100 GPUs.

## Infrastructure Components

### Deployment Scripts (scripts/)
- runpod_deploy.sh: Master deployment orchestrator (8-step workflow)
- runpod_upload.sh: S3 upload for binaries and test data
- upload_env_to_runpod.sh: Secure .env credentials upload
- runpod_deploy_test.sh: Prerequisites validation

### Docker Configuration
- Dockerfile.runpod: Multi-stage CUDA 12.1 runtime (~2GB, no binaries)
- entrypoint.sh: Volume verification and training execution
- Architecture: Volume mount (NO S3 downloads in pods)

### S3 Configuration
- Bucket: se3zdnb5o4 (Iceland region: eur-is-1)
- Endpoint: https://s3api-eur-is-1.runpod.io
- Structure: binaries/, test_data/, models/, .env

### OpenTofu Infrastructure (terraform/runpod/)
- main.tf: Pod and volume resources
- variables.tf: Configuration variables
- outputs.tf: Pod connection info
- Security: NO credentials in state (uses volume .env)

## Deployment Assets Uploaded

### Training Binaries (77MB)
- train_tft_parquet (23M) - TFT-225 features
- train_mamba2_parquet (22M) - MAMBA-2 state space
- train_dqn (22M) - Deep Q-Network
- train_ppo (13M) - Proximal Policy Optimization

### Test Data (13.8 MB)
- 9 Parquet files: ES.FUT, NQ.FUT, 6E.FUT, ZN.FUT (180-day datasets)

### Credentials
- .env file (1.5 KB, private access, chmod 600)

## Documentation

### Deployment Guides
- RUNPOD_DEPLOYMENT_READY_SUMMARY.md: Complete deployment status
- RUNPOD_VOLUME_DEPLOYMENT_GUIDE.md: Step-by-step guide (42KB)
- RUNPOD_DEPLOYMENT_QUICK_START.md: Quick reference
- RUNPOD_UPLOAD_GUIDE.md: S3 upload instructions
- RUNPOD_VOLUME_CONFIGURATION_COMPLETE.md: S3 setup report
- RUNPOD_S3_PARQUET_UPLOAD_REPORT.md: Data upload verification

### Architecture Documentation
- RUNPOD_VOLUME_MOUNT_ARCHITECTURE.md: Volume mount design
- RUNPOD_S3_ARCHITECTURE_DIAGRAM.txt: S3 API vs filesystem access
- DOCKERFILE_RUNPOD_FINAL_SUMMARY.md: Docker image specification

### Decision Documentation
- RUNPOD_DEPLOYMENT_CHECKLIST.md: Go/no-go decision matrix (27KB)
- RUNPOD_DEPLOYMENT_DECISION_TREE.md: Decision workflow
- FP32_RUNPOD_DEPLOYMENT_READY.md: FP32 deployment readiness

## QAT Enhancements

### Core QAT Infrastructure
- ml/src/memory_optimization/qat.rs: Enhanced QAT observer (+226 lines)
- ml/src/memory_optimization/auto_batch_size.rs: OOM recovery (+84 lines)
- ml/src/tft/qat_tft.rs: QAT TFT wrapper (+154 lines)
- ml/src/trainers/tft.rs: QAT training integration (+433 lines)
- ml/src/qat_metrics_exporter.rs: NEW - QAT metrics export

### QAT Testing
- ml/tests/qat_integration_tests.rs: NEW - Integration test suite
- ml/tests/qat_gradient_clipping_test.rs: NEW - Gradient clipping tests
- ml/tests/qat_device_consistency_test.rs: Device mismatch tests (+205 lines)
- ml/tests/qat_accuracy_validation_test.rs: Accuracy validation
- ml/tests/qat_tft_integration_test.rs: TFT QAT integration

### QAT Documentation
- ml/docs/QAT_GUIDE.md: Comprehensive QAT guide (+616 lines)
- ml/docs/QAT_GRADIENT_CHECKPOINTING_WORKAROUND.md: NEW - Workaround guide
- QAT_BLOCKERS_ROOT_CAUSE_ANALYSIS.md: P0 blocker analysis (44KB)
- QAT_ACCURACY_VALIDATION_REPORT.md: Accuracy comparison
- QAT_GRADIENT_CLIPPING_VALIDATION_REPORT.md: Clipping validation

### QAT Monitoring
- config/grafana/dashboards/qat-training-metrics.json: NEW - Grafana dashboard

## AWS CLI Configuration

### Credentials Setup
- ~/.aws/credentials: Runpod profile configured
  - Access Key: user_2xxA3XcIFj16yfL3aBon9niiSpr
  - Secret Key: (from RUNPOD_S3_SECRET)
- ~/.aws/config: Iceland region (eur-is-1)

## Production Readiness

### FP32 Models:  READY FOR DEPLOYMENT
- DQN: 15-20s training, ~6MB GPU memory
- PPO: 7-10s training, ~145MB GPU memory
- MAMBA-2: 2-3 min training, ~164MB GPU memory
- TFT-225: 3-5 min training, ~500MB GPU memory
- Total GPU Budget: 815MB (fits on 4GB+ Tesla V100)

### QAT Models: 🔴 BLOCKED
- 24 tests implemented but DO NOT COMPILE (11 errors)
- 3 P0 blockers: device mismatch, gradient checkpointing, OOM recovery
- Timeline: 1-2 weeks to fix (13h P0 fixes + validation)

### Wave D Features:  OPERATIONAL
- 225 features fully integrated
- Feature extraction: 5.10μs/bar (196x faster than target)
- Wave D backtest: Sharpe 2.00, Win Rate 60%, Drawdown 15%
- Database migration 045: Applied cleanly, zero conflicts

## Cost Analysis

### One-Time Setup
- Network Volume: $4/month (50GB SSD)
- Upload costs: FREE (S3 API included)

### Per Training Run (TFT-225)
- GPU: Tesla V100-PCIE-16GB @ $0.29/hr
- Training Time: ~4 hours
- Cost per run: $1.16

### Monthly (20 Training Runs)
- Storage: $4.00/month
- Training: $23.20/month (20 runs × $1.16)
- Total: $27.20/month

## Security

### Credentials Management
-  NO credentials in Docker image
-  NO credentials in Terraform state
-  .env gitignored and not committed
-  .env file private on S3 (HTTP 401 on public access)
-  Docker Hub repository PRIVATE (jgrusewski/foxhunt)

### Access Control
- S3 API: Local client uploads only
- Volume mount: Pod filesystem access only
- Authentication: AWS CLI with Runpod profile required

## Next Steps

1.  COMPLETE: Build Docker image
2.  PENDING: Push to Docker Hub
3.  PENDING: Deploy pod via Runpod console
4.  PENDING: Validate training on Tesla V100

## Performance Targets

- Build time: 5-10 min
- Upload time: ~20 sec (90MB total)
- Pod startup: ~30 sec
- Training time: 3-5 min (TFT-225)
- Total deployment: ~40 min from start to first training run

## Test Status

- FP32 tests: 597/608 passing (98.2%)
- QAT tests: 0/24 passing (compilation errors)
- Overall: 2,062/2,086 passing (98.8% excluding QAT)

🤖 Generated with Claude Code (https://claude.com/claude-code)

Co-Authored-By: Claude <noreply@anthropic.com>
2025-10-24 01:11:43 +02:00

517 lines
17 KiB
Rust

//! Comprehensive Metrics Pipeline Validation Tests
//!
//! This module validates the end-to-end metrics collection pipeline including:
//! - Prometheus metrics export correctness
//! - Metric schema validation
//! - Metric cardinality checks
//! - Performance overhead measurement
//! - `InfluxDB` client functionality (when enabled)
use reqwest::Client;
use std::collections::HashMap;
use std::time::Duration;
use tokio::time::timeout;
/// Result of metrics validation
#[derive(Debug, Clone)]
#[allow(clippy::module_name_repetitions)]
pub struct MetricsValidationResult {
pub service_name: String,
pub metrics_endpoint: String,
pub total_metrics: usize,
pub missing_required_metrics: Vec<String>,
pub invalid_metrics: Vec<String>,
pub cardinality_warnings: Vec<String>,
pub scrape_duration_ms: u128,
pub success: bool,
}
/// Required metrics per service
pub struct RequiredMetrics {
pub api_gateway: Vec<&'static str>,
pub trading_service: Vec<&'static str>,
pub backtesting_service: Vec<&'static str>,
pub ml_training_service: Vec<&'static str>,
}
impl Default for RequiredMetrics {
fn default() -> Self {
Self::new()
}
}
impl RequiredMetrics {
pub fn new() -> Self {
Self {
api_gateway: vec![
"api_gateway_auth_attempts_total",
"api_gateway_auth_duration_milliseconds",
"api_gateway_backend_requests_total",
"api_gateway_backend_request_duration_milliseconds",
"api_gateway_circuit_breaker_state",
"api_gateway_connection_pool_active",
"api_gateway_health_status",
"api_gateway_rate_limit_hits",
],
trading_service: vec![
"trading_orders_submitted_total",
"trading_order_submission_latency_us",
"trading_buffer_capacity",
"trading_buffer_used",
"trading_service_uptime_seconds",
],
backtesting_service: vec![
"backtesting_runs_total",
"backtesting_duration_seconds",
"backtesting_market_events_processed",
],
ml_training_service: vec![
"ml_inference_latency_microseconds",
"ml_model_health_status",
"ml_model_accuracy_percent",
"ml_predictions_total",
],
}
}
pub fn get_required(&self, service: &str) -> Vec<&'static str> {
match service {
"api_gateway" => self.api_gateway.clone(),
"trading_service" => self.trading_service.clone(),
"backtesting_service" => self.backtesting_service.clone(),
"ml_training_service" => self.ml_training_service.clone(),
_ => Vec::new(),
}
}
}
/// Parse Prometheus metrics exposition format
#[derive(Debug, Clone)]
pub struct PrometheusMetric {
pub name: String,
pub labels: HashMap<String, String>,
pub value: f64,
pub metric_type: Option<String>,
pub help: Option<String>,
}
pub struct MetricsParser;
impl MetricsParser {
/// Parse Prometheus exposition format
pub fn parse(content: &str) -> Vec<PrometheusMetric> {
let mut metrics = Vec::new();
let mut current_type: Option<String> = None;
let mut current_help: Option<String> = None;
let mut current_name: Option<String> = None;
for line in content.lines() {
let line = line.trim();
// Skip empty lines
if line.is_empty() {
continue;
}
// Parse TYPE declaration
if line.starts_with("# TYPE ") {
let parts: Vec<&str> = line.splitn(4, ' ').collect();
if parts.len() >= 4 {
current_name = parts.get(2).copied().map(str::to_owned);
current_type = parts.get(3).copied().map(str::to_owned);
}
continue;
}
// Parse HELP declaration
if line.starts_with("# HELP ") {
let parts: Vec<&str> = line.splitn(3, ' ').collect();
if parts.len() >= 3 {
current_help = parts.get(2).copied().map(str::to_owned);
}
continue;
}
// Skip other comments
if line.starts_with('#') {
continue;
}
// Parse metric line
if let Some((name_labels, value_str)) = line.split_once(' ') {
// Extract name and labels
let (name, labels) = if let Some(open_brace) = name_labels.find('{') {
let name = &name_labels[..open_brace];
let start = open_brace.saturating_add(1);
let end = name_labels.len().saturating_sub(1);
let labels_str = &name_labels[start..end];
let labels = Self::parse_labels(labels_str);
(name.to_owned(), labels)
} else {
(name_labels.to_owned(), HashMap::new())
};
// Parse value
if let Ok(value) = value_str.trim().parse::<f64>() {
metrics.push(PrometheusMetric {
name: name.clone(),
labels,
value,
metric_type: if Some(&name) == current_name.as_ref() {
current_type.clone()
} else {
None
},
help: if Some(&name) == current_name.as_ref() {
current_help.clone()
} else {
None
},
});
}
}
}
metrics
}
/// Parse label string like: key1="value1",key2="value2"
fn parse_labels(labels_str: &str) -> HashMap<String, String> {
let mut labels = HashMap::new();
for pair in labels_str.split(',') {
if let Some((key, value)) = pair.split_once('=') {
let key = key.trim();
let value = value.trim().trim_matches('"');
labels.insert(key.to_owned(), value.to_owned());
}
}
labels
}
}
/// Validate metrics from a service endpoint
///
/// # Errors
/// Returns error if the operation fails
pub async fn validate_service_metrics(
service_name: &str,
endpoint: &str,
required_metrics: &[&str],
) -> Result<MetricsValidationResult, Box<dyn std::error::Error>> {
let client = Client::builder().timeout(Duration::from_secs(5)).build()?;
let start = std::time::Instant::now();
// Fetch metrics with timeout
let response = timeout(Duration::from_secs(10), client.get(endpoint).send()).await??;
let scrape_duration_ms = start.elapsed().as_millis();
if !response.status().is_success() {
return Ok(MetricsValidationResult {
service_name: service_name.to_owned(),
metrics_endpoint: endpoint.to_owned(),
total_metrics: 0,
missing_required_metrics: required_metrics.iter().map(|s| (*s).to_owned()).collect(),
invalid_metrics: Vec::new(),
cardinality_warnings: Vec::new(),
scrape_duration_ms,
success: false,
});
}
let content = response.text().await?;
let metrics = MetricsParser::parse(&content);
// Check for required metrics
let metric_names: Vec<String> = metrics.iter().map(|m| m.name.clone()).collect();
let unique_metrics: std::collections::HashSet<String> = metric_names.iter().cloned().collect();
let mut missing_required = Vec::new();
for required in required_metrics {
if !unique_metrics.contains(*required) {
missing_required.push((*required).to_owned());
}
}
// Check for invalid metrics (basic validation)
let mut invalid_metrics = Vec::new();
for metric in &metrics {
// Check for NaN or Inf values
if metric.value.is_nan() || metric.value.is_infinite() {
invalid_metrics.push(format!(
"{} has invalid value: {}",
metric.name, metric.value
));
}
// Check for negative values in counters (by convention)
if metric.metric_type.as_ref().is_some_and(|t| t == "counter") && metric.value < 0.0_f64 {
invalid_metrics.push(format!("{} is a counter with negative value", metric.name));
}
}
// Check cardinality (warn if too many unique label combinations)
let mut cardinality_warnings = Vec::new();
let mut cardinality_map: HashMap<String, usize> = HashMap::new();
for metric in &metrics {
cardinality_map
.entry(metric.name.clone())
.and_modify(|c| *c = c.saturating_add(1))
.or_insert(1);
}
for (name, count) in &cardinality_map {
if *count > 1000 {
cardinality_warnings.push(format!(
"{} has {} unique series (high cardinality)",
name, count
));
}
}
let success = missing_required.is_empty() && invalid_metrics.is_empty();
Ok(MetricsValidationResult {
service_name: service_name.to_owned(),
metrics_endpoint: endpoint.to_owned(),
total_metrics: metrics.len(),
missing_required_metrics: missing_required,
invalid_metrics,
cardinality_warnings,
scrape_duration_ms,
success,
})
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_metrics_parser() {
let content = r#"
# HELP api_gateway_requests_total Total HTTP requests
# TYPE api_gateway_requests_total counter
api_gateway_requests_total{method="GET",status="200"} 1234
api_gateway_requests_total{method="POST",status="201"} 567
# HELP api_gateway_latency_seconds Request latency
# TYPE api_gateway_latency_seconds histogram
api_gateway_latency_seconds{quantile="0.5"} 0.05
api_gateway_latency_seconds{quantile="0.99"} 0.5
"#;
let metrics = MetricsParser::parse(content);
assert_eq!(metrics.len(), 4);
// Check first metric
assert_eq!(metrics[0].name, "api_gateway_requests_total");
assert_eq!(metrics[0].value, 1234.0);
assert_eq!(metrics[0].metric_type, Some("counter".to_owned()));
assert_eq!(metrics[0].labels.get("method"), Some(&"GET".to_owned()));
assert_eq!(metrics[0].labels.get("status"), Some(&"200".to_owned()));
// Check histogram metric
assert_eq!(metrics[2].name, "api_gateway_latency_seconds");
assert_eq!(metrics[2].labels.get("quantile"), Some(&"0.5".to_owned()));
assert_eq!(metrics[2].value, 0.05);
}
#[test]
fn test_required_metrics() {
let required = RequiredMetrics::new();
let api_metrics = required.get_required("api_gateway");
assert!(api_metrics.contains(&"api_gateway_auth_attempts_total"));
assert!(api_metrics.contains(&"api_gateway_circuit_breaker_state"));
let trading_metrics = required.get_required("trading_service");
assert!(trading_metrics.contains(&"trading_orders_submitted_total"));
assert!(trading_metrics.contains(&"trading_buffer_capacity"));
}
#[tokio::test]
async fn test_metrics_parser_edge_cases() {
// Test empty content
let metrics = MetricsParser::parse("");
assert_eq!(metrics.len(), 0);
// Test metric without labels
let content = "simple_counter 42.0\n";
let metrics = MetricsParser::parse(content);
assert_eq!(metrics.len(), 1);
assert_eq!(metrics[0].name, "simple_counter");
assert_eq!(metrics[0].value, 42.0);
assert!(metrics[0].labels.is_empty());
// Test invalid value (NaN is actually a valid f64 in Rust)
let content = "invalid_metric NaN\n";
let metrics = MetricsParser::parse(content);
// NaN parses successfully as f64::NAN
assert_eq!(metrics.len(), 1);
assert!(metrics[0].value.is_nan());
}
// Integration tests (require services running)
#[tokio::test]
#[ignore = "Only run with `cargo test -- --ignored`"]
async fn test_api_gateway_metrics() {
let required = RequiredMetrics::new();
let result = validate_service_metrics(
"api_gateway",
"http://localhost:9091/metrics",
&required.api_gateway,
)
.await;
if let Ok(validation) = result {
println!("API Gateway Metrics Validation:");
println!(" Total Metrics: {}", validation.total_metrics);
println!(" Scrape Duration: {}ms", validation.scrape_duration_ms);
println!(" Success: {}", validation.success);
if !validation.missing_required_metrics.is_empty() {
println!(" Missing Required Metrics:");
for metric in &validation.missing_required_metrics {
println!(" - {}", metric);
}
}
if !validation.invalid_metrics.is_empty() {
println!(" Invalid Metrics:");
for metric in &validation.invalid_metrics {
println!(" - {}", metric);
}
}
assert!(validation.success, "API Gateway metrics validation failed");
}
}
#[tokio::test]
#[ignore = "Integration test - run with --ignored"]
async fn test_trading_service_metrics() {
let required = RequiredMetrics::new();
let result = validate_service_metrics(
"trading_service",
"http://localhost:9092/metrics",
&required.trading_service,
)
.await;
if let Ok(validation) = result {
println!("Trading Service Metrics Validation:");
println!(" Total Metrics: {}", validation.total_metrics);
println!(" Scrape Duration: {}ms", validation.scrape_duration_ms);
println!(" Success: {}", validation.success);
assert!(
validation.success,
"Trading Service metrics validation failed"
);
}
}
#[tokio::test]
#[ignore = "Integration test - run with --ignored"]
async fn test_all_services_metrics() {
let required = RequiredMetrics::new();
let services = vec![
(
"api_gateway",
"http://localhost:9091/metrics",
required.api_gateway,
),
(
"trading_service",
"http://localhost:9092/metrics",
required.trading_service,
),
(
"backtesting_service",
"http://localhost:9093/metrics",
required.backtesting_service,
),
(
"ml_training_service",
"http://localhost:9094/metrics",
required.ml_training_service,
),
];
let mut all_success = true;
for (name, endpoint, required_metrics) in services {
match validate_service_metrics(name, endpoint, &required_metrics).await {
Ok(validation) => {
println!("\n{} Metrics Validation:", name);
println!(
" Status: {}",
if validation.success {
"✓ PASS"
} else {
"✗ FAIL"
}
);
println!(" Total Metrics: {}", validation.total_metrics);
println!(" Scrape Duration: {}ms", validation.scrape_duration_ms);
if !validation.missing_required_metrics.is_empty() {
println!(
" Missing Metrics: {:?}",
validation.missing_required_metrics
);
}
if !validation.cardinality_warnings.is_empty() {
println!(" Cardinality Warnings:");
for warning in &validation.cardinality_warnings {
println!(" - {}", warning);
}
}
all_success &= validation.success;
},
Err(e) => {
println!("\n{} Metrics Validation: ✗ ERROR", name);
println!(" Error: {}", e);
all_success = false;
},
}
}
assert!(all_success, "Not all services passed metrics validation");
}
#[tokio::test]
#[ignore = "Integration test - run with --ignored"]
async fn test_metrics_scrape_performance() {
// Test that scraping is fast enough for Prometheus defaults (5s interval)
let endpoint = "http://localhost:9091/metrics";
let start = std::time::Instant::now();
let client = Client::builder()
.timeout(Duration::from_secs(2))
.build()
.unwrap();
let response = client.get(endpoint).send().await;
let scrape_duration = start.elapsed();
assert!(response.is_ok(), "Failed to scrape metrics");
assert!(
scrape_duration < Duration::from_millis(500),
"Metrics scrape took too long: {:?}",
scrape_duration
);
println!("Metrics scrape performance: {:?}", scrape_duration);
}
}