From e9b98e56694bb5f6bdd22f82bc8ce74fb6f282c4 Mon Sep 17 00:00:00 2001 From: Story Crater Bot <19826264+Riotpiaole@users.noreply.github.com> Date: Fri, 28 Aug 2026 11:45:12 -0700 Subject: [PATCH] =?UTF-8?q?feat:=20M3.8.3=20complete=20=E2=80=94=20metrics?= =?UTF-8?q?=20&=20monitoring=20(7=20tests)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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) --- crates/mem-ingest/src/lib.rs | 2 + crates/mem-ingest/src/optimizer_metrics.rs | 306 +++++++++++++++++++++ 2 files changed, 308 insertions(+) create mode 100644 crates/mem-ingest/src/optimizer_metrics.rs diff --git a/crates/mem-ingest/src/lib.rs b/crates/mem-ingest/src/lib.rs index ec93471..385259c 100644 --- a/crates/mem-ingest/src/lib.rs +++ b/crates/mem-ingest/src/lib.rs @@ -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; diff --git a/crates/mem-ingest/src/optimizer_metrics.rs b/crates/mem-ingest/src/optimizer_metrics.rs new file mode 100644 index 0000000..9c703de --- /dev/null +++ b/crates/mem-ingest/src/optimizer_metrics.rs @@ -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>>, +} + +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 { + self.by_project + .lock() + .unwrap() + .get(project) + .map(|m| m.clone()) + } + + /// Get all project metrics. + pub fn all_projects(&self) -> HashMap { + 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()); + } +}