feat: M3.8.3 complete — metrics & monitoring (7 tests)
MetricsCollector implementation: - Per-project aggregation of OptimizationMetrics - Structured logging via tracing (log_all_projects) - Prometheus export format (prometheus_export) - Per-compressor stat tracking 7 new tests (all passing): - test_collector_merge_single_project - test_collector_merge_multiple_projects - test_collector_merge_aggregates - test_collector_nonexistent_project - test_collector_per_compressor_stats - test_prometheus_export_format - test_prometheus_compression_ratio Ready to integrate into rebuild.rs: let collector = MetricsCollector::new(); ... collector.merge_project(project_id, metrics); collector.log_all_projects(); Total M3.8 progress: - M3.8.1: ✅ 62 tests (core compressors) - M3.8.2: ✅ 5 tests (ingest helpers) - M3.8.3: ✅ 7 tests (metrics & monitoring) - M3.8.4: ✅ IMPLICIT (no query compression needed) - M3.8.5: ⏳ Benchmarks - M3.8.6: ⏳ Gate 79 tests passing total (62+5+7+5 from optimizer_sink)
This commit is contained in:
@@ -3,9 +3,11 @@ pub mod claude_transcript;
|
||||
pub mod doc_corpus;
|
||||
pub mod derived_filter;
|
||||
pub mod optimizer_sink;
|
||||
pub mod optimizer_metrics;
|
||||
|
||||
pub use pi_session::PiSessionSource;
|
||||
pub use claude_transcript::ClaudeTranscriptSource;
|
||||
pub use doc_corpus::{DocCorpusSource, DocSection, DryRunReport};
|
||||
pub use derived_filter::{ArtifactRecord, DerivedFilter, DerivedMatch};
|
||||
pub use optimizer_sink::{OptimizationMetrics, CompressorStats, optimize_record_with_metrics};
|
||||
pub use optimizer_metrics::MetricsCollector;
|
||||
|
||||
@@ -0,0 +1,306 @@
|
||||
//! M3.8.3 Metrics & Monitoring
|
||||
//!
|
||||
//! Tracks and exports optimization metrics for observability.
|
||||
//! Emits structured logs via tracing and provides Prometheus-ready counters.
|
||||
|
||||
use crate::OptimizationMetrics;
|
||||
use std::collections::HashMap;
|
||||
use std::sync::{Arc, Mutex};
|
||||
|
||||
/// Per-project optimization metrics with aggregation.
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct MetricsCollector {
|
||||
/// Metrics keyed by project_id
|
||||
by_project: Arc<Mutex<HashMap<String, OptimizationMetrics>>>,
|
||||
}
|
||||
|
||||
impl MetricsCollector {
|
||||
/// Create a new collector.
|
||||
pub fn new() -> Self {
|
||||
MetricsCollector {
|
||||
by_project: Arc::new(Mutex::new(HashMap::new())),
|
||||
}
|
||||
}
|
||||
|
||||
/// Merge metrics for a project.
|
||||
pub fn merge_project(&self, project: &str, metrics: OptimizationMetrics) {
|
||||
let mut by_proj = self.by_project.lock().unwrap();
|
||||
let entry = by_proj.entry(project.to_string()).or_insert_with(|| {
|
||||
OptimizationMetrics {
|
||||
total_records: 0,
|
||||
input_bytes_total: 0,
|
||||
output_bytes_total: 0,
|
||||
per_compressor: HashMap::new(),
|
||||
}
|
||||
});
|
||||
|
||||
entry.total_records += metrics.total_records;
|
||||
entry.input_bytes_total += metrics.input_bytes_total;
|
||||
entry.output_bytes_total += metrics.output_bytes_total;
|
||||
|
||||
// Merge per-compressor stats
|
||||
for (name, stats) in metrics.per_compressor {
|
||||
let entry_stats = entry
|
||||
.per_compressor
|
||||
.entry(name)
|
||||
.or_insert_with(|| crate::CompressorStats {
|
||||
count: 0,
|
||||
input_bytes: 0,
|
||||
output_bytes: 0,
|
||||
});
|
||||
|
||||
entry_stats.count += stats.count;
|
||||
entry_stats.input_bytes += stats.input_bytes;
|
||||
entry_stats.output_bytes += stats.output_bytes;
|
||||
}
|
||||
}
|
||||
|
||||
/// Get metrics for a specific project.
|
||||
pub fn get_project(&self, project: &str) -> Option<OptimizationMetrics> {
|
||||
self.by_project
|
||||
.lock()
|
||||
.unwrap()
|
||||
.get(project)
|
||||
.map(|m| m.clone())
|
||||
}
|
||||
|
||||
/// Get all project metrics.
|
||||
pub fn all_projects(&self) -> HashMap<String, OptimizationMetrics> {
|
||||
self.by_project.lock().unwrap().clone()
|
||||
}
|
||||
|
||||
/// Log all metrics with structured output.
|
||||
pub fn log_all_projects(&self) {
|
||||
let projects = self.all_projects();
|
||||
|
||||
tracing::info!(
|
||||
project_count = projects.len(),
|
||||
"M3.8 optimization metrics (all projects)"
|
||||
);
|
||||
|
||||
for (project, metrics) in projects {
|
||||
metrics.log_summary(&project);
|
||||
}
|
||||
}
|
||||
|
||||
/// Emit Prometheus-style metrics.
|
||||
pub fn prometheus_export(&self) -> String {
|
||||
let mut output = String::new();
|
||||
let projects = self.all_projects();
|
||||
|
||||
// Counter: total records processed
|
||||
output.push_str("# HELP m3_8_optimization_records_total Total records optimized\n");
|
||||
output.push_str("# TYPE m3_8_optimization_records_total counter\n");
|
||||
for (project, metrics) in &projects {
|
||||
output.push_str(&format!(
|
||||
"m3_8_optimization_records_total{{project=\"{}\"}} {}\n",
|
||||
project, metrics.total_records
|
||||
));
|
||||
}
|
||||
output.push('\n');
|
||||
|
||||
// Gauge: input bytes
|
||||
output.push_str("# HELP m3_8_optimization_input_bytes_total Input bytes before compression\n");
|
||||
output.push_str("# TYPE m3_8_optimization_input_bytes_total gauge\n");
|
||||
for (project, metrics) in &projects {
|
||||
output.push_str(&format!(
|
||||
"m3_8_optimization_input_bytes_total{{project=\"{}\"}} {}\n",
|
||||
project, metrics.input_bytes_total
|
||||
));
|
||||
}
|
||||
output.push('\n');
|
||||
|
||||
// Gauge: output bytes
|
||||
output.push_str("# HELP m3_8_optimization_output_bytes_total Output bytes after compression\n");
|
||||
output.push_str("# TYPE m3_8_optimization_output_bytes_total gauge\n");
|
||||
for (project, metrics) in &projects {
|
||||
output.push_str(&format!(
|
||||
"m3_8_optimization_output_bytes_total{{project=\"{}\"}} {}\n",
|
||||
project, metrics.output_bytes_total
|
||||
));
|
||||
}
|
||||
output.push('\n');
|
||||
|
||||
// Gauge: compression ratio
|
||||
output.push_str("# HELP m3_8_optimization_compression_ratio Compression ratio (output/input * 100)\n");
|
||||
output.push_str("# TYPE m3_8_optimization_compression_ratio gauge\n");
|
||||
for (project, metrics) in &projects {
|
||||
let ratio = metrics.compression_ratio();
|
||||
output.push_str(&format!(
|
||||
"m3_8_optimization_compression_ratio{{project=\"{}\"}} {:.2}\n",
|
||||
project, ratio
|
||||
));
|
||||
}
|
||||
output.push('\n');
|
||||
|
||||
// Per-compressor stats
|
||||
output.push_str("# HELP m3_8_optimization_compressor_records Records processed by compressor\n");
|
||||
output.push_str("# TYPE m3_8_optimization_compressor_records counter\n");
|
||||
for (project, metrics) in &projects {
|
||||
for (compressor, stats) in &metrics.per_compressor {
|
||||
output.push_str(&format!(
|
||||
"m3_8_optimization_compressor_records{{project=\"{}\",compressor=\"{}\"}} {}\n",
|
||||
project, compressor, stats.count
|
||||
));
|
||||
}
|
||||
}
|
||||
output.push('\n');
|
||||
|
||||
output.push_str("# HELP m3_8_optimization_compressor_ratio Compression ratio by compressor\n");
|
||||
output.push_str("# TYPE m3_8_optimization_compressor_ratio gauge\n");
|
||||
for (project, metrics) in &projects {
|
||||
for (compressor, stats) in &metrics.per_compressor {
|
||||
if stats.input_bytes > 0 {
|
||||
let ratio = (stats.output_bytes as f32 / stats.input_bytes as f32) * 100.0;
|
||||
output.push_str(&format!(
|
||||
"m3_8_optimization_compressor_ratio{{project=\"{}\",compressor=\"{}\"}} {:.2}\n",
|
||||
project, compressor, ratio
|
||||
));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
output
|
||||
}
|
||||
}
|
||||
|
||||
impl Default for MetricsCollector {
|
||||
fn default() -> Self {
|
||||
Self::new()
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::CompressorStats;
|
||||
|
||||
fn make_test_metrics(project: &str, records: usize, input: usize, output: usize) -> OptimizationMetrics {
|
||||
OptimizationMetrics {
|
||||
total_records: records,
|
||||
input_bytes_total: input,
|
||||
output_bytes_total: output,
|
||||
per_compressor: {
|
||||
let mut map = HashMap::new();
|
||||
map.insert(
|
||||
"LogCompressor".to_string(),
|
||||
CompressorStats {
|
||||
count: records,
|
||||
input_bytes: input,
|
||||
output_bytes: output,
|
||||
},
|
||||
);
|
||||
map
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_collector_merge_single_project() {
|
||||
let collector = MetricsCollector::new();
|
||||
let metrics = make_test_metrics("proj1", 10, 1000, 500);
|
||||
|
||||
collector.merge_project("proj1", metrics);
|
||||
|
||||
let result = collector.get_project("proj1").unwrap();
|
||||
assert_eq!(result.total_records, 10);
|
||||
assert_eq!(result.input_bytes_total, 1000);
|
||||
assert_eq!(result.output_bytes_total, 500);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_collector_merge_multiple_projects() {
|
||||
let collector = MetricsCollector::new();
|
||||
|
||||
collector.merge_project("proj1", make_test_metrics("proj1", 10, 1000, 500));
|
||||
collector.merge_project("proj2", make_test_metrics("proj2", 5, 500, 300));
|
||||
|
||||
let all = collector.all_projects();
|
||||
assert_eq!(all.len(), 2);
|
||||
assert_eq!(all.get("proj1").unwrap().total_records, 10);
|
||||
assert_eq!(all.get("proj2").unwrap().total_records, 5);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_collector_merge_aggregates() {
|
||||
let collector = MetricsCollector::new();
|
||||
|
||||
// Merge same project twice
|
||||
collector.merge_project("proj1", make_test_metrics("proj1", 5, 500, 250));
|
||||
collector.merge_project("proj1", make_test_metrics("proj1", 5, 500, 250));
|
||||
|
||||
let result = collector.get_project("proj1").unwrap();
|
||||
assert_eq!(result.total_records, 10, "should aggregate records");
|
||||
assert_eq!(result.input_bytes_total, 1000, "should sum input bytes");
|
||||
assert_eq!(result.output_bytes_total, 500, "should sum output bytes");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_prometheus_export_format() {
|
||||
let collector = MetricsCollector::new();
|
||||
collector.merge_project("test", make_test_metrics("test", 5, 100, 50));
|
||||
|
||||
let output = collector.prometheus_export();
|
||||
|
||||
// Should contain expected metric lines
|
||||
assert!(output.contains("m3_8_optimization_records_total"));
|
||||
assert!(output.contains("m3_8_optimization_input_bytes_total"));
|
||||
assert!(output.contains("m3_8_optimization_output_bytes_total"));
|
||||
assert!(output.contains("m3_8_optimization_compression_ratio"));
|
||||
assert!(output.contains("project=\"test\""));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_prometheus_compression_ratio() {
|
||||
let collector = MetricsCollector::new();
|
||||
collector.merge_project("test", make_test_metrics("test", 5, 1000, 500));
|
||||
|
||||
let output = collector.prometheus_export();
|
||||
|
||||
// Should have 50.00% ratio
|
||||
assert!(output.contains("m3_8_optimization_compression_ratio{project=\"test\"} 50.00"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_collector_per_compressor_stats() {
|
||||
let collector = MetricsCollector::new();
|
||||
let mut metrics = OptimizationMetrics {
|
||||
total_records: 5,
|
||||
input_bytes_total: 1000,
|
||||
output_bytes_total: 750,
|
||||
per_compressor: HashMap::new(),
|
||||
};
|
||||
|
||||
metrics.per_compressor.insert(
|
||||
"LogCompressor".to_string(),
|
||||
CompressorStats {
|
||||
count: 3,
|
||||
input_bytes: 600,
|
||||
output_bytes: 400,
|
||||
},
|
||||
);
|
||||
metrics.per_compressor.insert(
|
||||
"TextCompressor".to_string(),
|
||||
CompressorStats {
|
||||
count: 2,
|
||||
input_bytes: 400,
|
||||
output_bytes: 350,
|
||||
},
|
||||
);
|
||||
|
||||
collector.merge_project("test", metrics);
|
||||
let result = collector.get_project("test").unwrap();
|
||||
|
||||
assert_eq!(result.per_compressor.len(), 2);
|
||||
assert_eq!(
|
||||
result.per_compressor.get("LogCompressor").unwrap().count,
|
||||
3
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_collector_nonexistent_project() {
|
||||
let collector = MetricsCollector::new();
|
||||
assert!(collector.get_project("nonexistent").is_none());
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user