//! Query-Aware Metrics Tracking for M3.8 Optimization //! //! Tracks optimization progress and metrics per query_id, allowing clients //! to monitor compression ratios, latency, and progress in real-time. //! //! # Example //! //! ```ignore //! // Start tracking a query's optimization //! let metrics = QueryMetrics::new("query-123", "myproject"); //! //! // During optimization //! metrics.record_record_optimized("log", 1000, 300); //! metrics.record_record_optimized("text", 500, 250); //! //! // Query progress //! let progress = metrics.progress(); //! println!("{:.1}% complete, {:.1}% compression", //! progress.percent_complete, //! progress.compression_ratio()); //! //! // Get final metrics //! let final_metrics = metrics.to_summary(); //! ``` use std::sync::{Arc, Mutex}; use std::collections::HashMap; use serde::{Deserialize, Serialize}; /// Per-query optimization metrics and progress #[derive(Debug, Clone, Serialize, Deserialize)] pub struct QueryMetrics { /// Unique query identifier pub query_id: String, /// Project this query belongs to pub project: String, /// When optimization started pub started_at: String, /// Total records processed pub total_records: usize, /// Records completed (for progress tracking) pub records_completed: usize, /// Total bytes before optimization pub input_bytes_total: usize, /// Total bytes after optimization pub output_bytes_total: usize, /// Per-compressor breakdown pub per_compressor: HashMap, /// Per-content-type breakdown pub per_content_type: HashMap, /// Optimization status pub status: OptimizationStatus, /// Error message (if failed) pub error: Option, /// Estimated time remaining (seconds) pub eta_secs: Option, } /// Optimization status enum #[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)] pub enum OptimizationStatus { /// Not yet started Pending, /// Currently optimizing InProgress, /// Successfully completed Completed, /// Failed with error Failed, /// Paused (resumable) Paused, } impl QueryMetrics { /// Create new query metrics tracker pub fn new(query_id: impl Into, project: impl Into) -> Self { Self { query_id: query_id.into(), project: project.into(), started_at: format!("{}", time::OffsetDateTime::now_utc()), total_records: 0, records_completed: 0, input_bytes_total: 0, output_bytes_total: 0, per_compressor: HashMap::new(), per_content_type: HashMap::new(), status: OptimizationStatus::Pending, error: None, eta_secs: None, } } /// Record a successfully optimized record pub fn record_record_optimized( &mut self, compressor: impl Into, content_type: impl Into, input_bytes: usize, output_bytes: usize, ) { let compressor_name = compressor.into(); let content_type_name = content_type.into(); self.input_bytes_total += input_bytes; self.output_bytes_total += output_bytes; self.records_completed += 1; // Update per-compressor stats self.per_compressor .entry(compressor_name.clone()) .or_insert_with(|| CompressorMetrics { count: 0, input_bytes: 0, output_bytes: 0, }) .record(input_bytes, output_bytes); // Update per-content-type stats self.per_content_type .entry(content_type_name) .or_insert_with(|| ContentTypeMetrics { count: 0, input_bytes: 0, output_bytes: 0, }) .record(input_bytes, output_bytes); } /// Record a failed optimization (falls back to original) pub fn record_record_failed(&mut self, reason: impl Into) { self.error = Some(reason.into()); } /// Calculate compression ratio (%) pub fn compression_ratio(&self) -> f32 { if self.input_bytes_total == 0 { 0.0 } else { (self.output_bytes_total as f32 / self.input_bytes_total as f32) * 100.0 } } /// Calculate progress percentage (0-100) pub fn percent_complete(&self) -> f32 { if self.total_records == 0 { 0.0 } else { (self.records_completed as f32 / self.total_records as f32) * 100.0 } } /// Get progress snapshot pub fn progress(&self) -> ProgressSnapshot { ProgressSnapshot { query_id: self.query_id.clone(), project: self.project.clone(), status: self.status, percent_complete: self.percent_complete(), records_completed: self.records_completed, total_records: self.total_records, compression_ratio: self.compression_ratio(), input_bytes: self.input_bytes_total, output_bytes: self.output_bytes_total, eta_secs: self.eta_secs, } } /// Convert to summary (for storage/reporting) pub fn to_summary(&self) -> MetricsSummary { MetricsSummary { query_id: self.query_id.clone(), project: self.project.clone(), started_at: self.started_at.clone(), total_records: self.total_records, input_bytes_total: self.input_bytes_total, output_bytes_total: self.output_bytes_total, compression_ratio: self.compression_ratio(), per_compressor: self.per_compressor.clone(), per_content_type: self.per_content_type.clone(), status: self.status, error: self.error.clone(), } } } /// Progress snapshot for real-time monitoring #[derive(Debug, Clone, Serialize, Deserialize)] pub struct ProgressSnapshot { pub query_id: String, pub project: String, pub status: OptimizationStatus, /// Percentage complete (0-100) pub percent_complete: f32, /// Records processed so far pub records_completed: usize, /// Total records to process pub total_records: usize, /// Current compression ratio (%) pub compression_ratio: f32, /// Bytes before optimization pub input_bytes: usize, /// Bytes after optimization pub output_bytes: usize, /// Estimated seconds remaining pub eta_secs: Option, } /// Summary of optimization metrics (for storage) #[derive(Debug, Clone, Serialize, Deserialize)] pub struct MetricsSummary { pub query_id: String, pub project: String, pub started_at: String, pub total_records: usize, pub input_bytes_total: usize, pub output_bytes_total: usize, pub compression_ratio: f32, pub per_compressor: HashMap, pub per_content_type: HashMap, pub status: OptimizationStatus, pub error: Option, } /// Per-compressor statistics #[derive(Debug, Clone, Serialize, Deserialize)] pub struct CompressorMetrics { pub count: usize, pub input_bytes: usize, pub output_bytes: usize, } impl CompressorMetrics { fn record(&mut self, input_bytes: usize, output_bytes: usize) { self.count += 1; self.input_bytes += input_bytes; self.output_bytes += output_bytes; } pub fn compression_ratio(&self) -> f32 { if self.input_bytes == 0 { 0.0 } else { (self.output_bytes as f32 / self.input_bytes as f32) * 100.0 } } } /// Per-content-type statistics #[derive(Debug, Clone, Serialize, Deserialize)] pub struct ContentTypeMetrics { pub count: usize, pub input_bytes: usize, pub output_bytes: usize, } impl ContentTypeMetrics { fn record(&mut self, input_bytes: usize, output_bytes: usize) { self.count += 1; self.input_bytes += input_bytes; self.output_bytes += output_bytes; } pub fn compression_ratio(&self) -> f32 { if self.input_bytes == 0 { 0.0 } else { (self.output_bytes as f32 / self.input_bytes as f32) * 100.0 } } } /// Thread-safe metrics repository indexed by query_id #[derive(Debug, Clone)] pub struct QueryMetricsRepository { metrics: Arc>>, } impl QueryMetricsRepository { /// Create new metrics repository pub fn new() -> Self { Self { metrics: Arc::new(Mutex::new(HashMap::new())), } } /// Start tracking a new query pub fn create_query(&self, query_id: impl Into, project: impl Into) -> String { let query_id_str = query_id.into(); let metrics = QueryMetrics::new(query_id_str.clone(), project); let mut repo = self.metrics.lock().unwrap(); repo.insert(query_id_str.clone(), metrics); query_id_str } /// Get metrics for a specific query pub fn get_metrics(&self, query_id: &str) -> Option { let repo = self.metrics.lock().unwrap(); repo.get(query_id).cloned() } /// Update metrics for a query pub fn update_metrics(&self, query_id: &str, f: F) -> Result<(), String> where F: FnOnce(&mut QueryMetrics), { let mut repo = self.metrics.lock().unwrap(); repo.get_mut(query_id) .ok_or_else(|| format!("Query {} not found", query_id)) .map(|metrics| f(metrics)) } /// Get progress for a query pub fn get_progress(&self, query_id: &str) -> Option { self.get_metrics(query_id).map(|m| m.progress()) } /// List all active queries pub fn list_queries(&self) -> Vec { let repo = self.metrics.lock().unwrap(); repo.keys().cloned().collect() } /// Get metrics for all queries in a project pub fn get_project_metrics(&self, project: &str) -> Vec { let repo = self.metrics.lock().unwrap(); repo.values() .filter(|m| m.project == project) .cloned() .collect() } /// Clear completed query metrics (after storing to DB) pub fn remove_query(&self, query_id: &str) -> Option { let mut repo = self.metrics.lock().unwrap(); repo.remove(query_id) } } impl Default for QueryMetricsRepository { fn default() -> Self { Self::new() } } #[cfg(test)] mod tests { use super::*; #[test] fn test_query_metrics_creation() { let metrics = QueryMetrics::new("q1", "project1"); assert_eq!(metrics.query_id, "q1"); assert_eq!(metrics.project, "project1"); assert_eq!(metrics.status, OptimizationStatus::Pending); } #[test] fn test_record_optimized() { let mut metrics = QueryMetrics::new("q1", "project1"); metrics.total_records = 2; metrics.record_record_optimized("log", "text/plain", 1000, 300); metrics.record_record_optimized("text", "text/plain", 500, 250); assert_eq!(metrics.records_completed, 2); assert_eq!(metrics.input_bytes_total, 1500); assert_eq!(metrics.output_bytes_total, 550); assert!((metrics.compression_ratio() - 36.67).abs() < 0.1); } #[test] fn test_progress_calculation() { let mut metrics = QueryMetrics::new("q1", "project1"); metrics.total_records = 10; metrics.records_completed = 5; assert!((metrics.percent_complete() - 50.0).abs() < 0.1); } #[test] fn test_metrics_repository() { let repo = QueryMetricsRepository::new(); let query_id = repo.create_query("q1", "project1"); assert_eq!(query_id, "q1"); assert!(repo.get_metrics("q1").is_some()); assert!(repo.get_metrics("q2").is_none()); } #[test] fn test_repository_update() { let repo = QueryMetricsRepository::new(); repo.create_query("q1", "project1"); repo.update_metrics("q1", |m| { m.total_records = 100; m.record_record_optimized("log", "text/plain", 1000, 300); }).unwrap(); let metrics = repo.get_metrics("q1").unwrap(); assert_eq!(metrics.total_records, 100); assert_eq!(metrics.records_completed, 1); } #[test] fn test_progress_snapshot() { let mut metrics = QueryMetrics::new("q1", "project1"); metrics.total_records = 100; metrics.records_completed = 25; metrics.input_bytes_total = 10000; metrics.output_bytes_total = 3000; let progress = metrics.progress(); assert!((progress.percent_complete - 25.0).abs() < 0.1); assert!((progress.compression_ratio - 30.0).abs() < 0.1); } #[test] fn test_per_compressor_stats() { let mut metrics = QueryMetrics::new("q1", "project1"); metrics.record_record_optimized("log", "text/plain", 1000, 100); metrics.record_record_optimized("log", "text/plain", 500, 50); let log_stats = metrics.per_compressor.get("log").unwrap(); assert_eq!(log_stats.count, 2); assert_eq!(log_stats.input_bytes, 1500); assert_eq!(log_stats.output_bytes, 150); } }