From a958582eaa51728e25511c2419bbbed88d55b2bd Mon Sep 17 00:00:00 2001 From: jgrusewski Date: Sun, 1 Mar 2026 02:23:36 +0100 Subject: [PATCH] feat(common): add gRPC metrics Tower layer for automatic instrumentation Add a reusable GrpcMetricsLayer that can be applied to any tonic server via Server::builder().layer() to automatically emit three Prometheus metric families for every gRPC handler: started_total, handled_total (with status code), and handling_seconds histogram. Co-Authored-By: Claude Opus 4.6 --- Cargo.lock | 5 + crates/common/Cargo.toml | 5 +- crates/common/src/metrics/grpc_metrics.rs | 245 ++++++++++++++++++++++ crates/common/src/metrics/mod.rs | 3 + 4 files changed, 257 insertions(+), 1 deletion(-) create mode 100644 crates/common/src/metrics/grpc_metrics.rs diff --git a/Cargo.lock b/Cargo.lock index 05ae62893..a863a7d5e 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2291,6 +2291,7 @@ dependencies = [ "criterion", "fastrand", "futures", + "http 1.3.1", "jsonwebtoken", "lazy_static", "ml", @@ -2313,6 +2314,8 @@ dependencies = [ "tokio-test", "toml", "tonic 0.14.2", + "tower-layer", + "tower-service", "tracing", "tracing-appender", "tracing-opentelemetry", @@ -11543,6 +11546,8 @@ dependencies = [ "common", "futures-util", "jsonwebtoken", + "once_cell", + "prometheus", "prost 0.14.1", "serde", "serde_json", diff --git a/crates/common/Cargo.toml b/crates/common/Cargo.toml index 1c3f920c3..36bdce00e 100644 --- a/crates/common/Cargo.toml +++ b/crates/common/Cargo.toml @@ -47,8 +47,11 @@ opentelemetry.workspace = true opentelemetry-otlp.workspace = true opentelemetry_sdk.workspace = true -# gRPC (for correlation ID propagation) +# gRPC (for correlation ID propagation and metrics layer) tonic.workspace = true +http.workspace = true +tower-layer.workspace = true +tower-service.workspace = true # Configuration toml.workspace = true diff --git a/crates/common/src/metrics/grpc_metrics.rs b/crates/common/src/metrics/grpc_metrics.rs new file mode 100644 index 000000000..e2009ed5c --- /dev/null +++ b/crates/common/src/metrics/grpc_metrics.rs @@ -0,0 +1,245 @@ +//! gRPC metrics Tower layer for automatic request instrumentation. +//! +//! Provides `grpc_server_started_total`, `grpc_server_handled_total`, +//! and `grpc_server_handling_seconds` metrics for all gRPC handlers. +//! +//! # Usage +//! +//! ```rust,ignore +//! use common::metrics::grpc_metrics::GrpcMetricsLayer; +//! use tonic::transport::Server; +//! +//! Server::builder() +//! .layer(GrpcMetricsLayer::new("trading_service")) +//! .add_service(my_service) +//! .serve(addr) +//! .await?; +//! ``` + +use once_cell::sync::Lazy; +use prometheus::{register_counter_vec, register_histogram_vec, CounterVec, HistogramVec}; +use std::future::Future; +use std::pin::Pin; +use std::task::{Context, Poll}; +use tower_layer::Layer; +use tower_service::Service; + +/// Counter for gRPC requests started. +static GRPC_STARTED: Lazy = Lazy::new(|| { + register_counter_vec!( + "grpc_server_started_total", + "Total gRPC RPCs started on the server", + &["service", "grpc_method"] + ) + .unwrap_or_else(|e| { + // Metric registration failure at startup is unrecoverable. + // The process cannot operate without metrics instrumentation. + eprintln!("FATAL: failed to register grpc_server_started_total: {e}"); + std::process::abort() + }) +}); + +/// Counter for gRPC requests completed with status code. +static GRPC_HANDLED: Lazy = Lazy::new(|| { + register_counter_vec!( + "grpc_server_handled_total", + "Total gRPC RPCs completed on the server with status code", + &["service", "grpc_method", "grpc_code"] + ) + .unwrap_or_else(|e| { + eprintln!("FATAL: failed to register grpc_server_handled_total: {e}"); + std::process::abort() + }) +}); + +/// Histogram for gRPC request handling duration. +static GRPC_HANDLING_SECONDS: Lazy = Lazy::new(|| { + register_histogram_vec!( + "grpc_server_handling_seconds", + "Histogram of gRPC request handling duration in seconds", + &["service", "grpc_method"], + vec![0.001, 0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1.0, 5.0] + ) + .unwrap_or_else(|e| { + eprintln!("FATAL: failed to register grpc_server_handling_seconds: {e}"); + std::process::abort() + }) +}); + +/// Initialize gRPC metrics (force lazy registration). +/// +/// Call this early in service startup to ensure metrics are registered +/// before any requests arrive. It is safe to call multiple times. +pub fn init_grpc_metrics() { + let _ = &*GRPC_STARTED; + let _ = &*GRPC_HANDLED; + let _ = &*GRPC_HANDLING_SECONDS; +} + +/// Extract the method name from a gRPC URI path. +/// +/// gRPC paths follow the format `/package.Service/Method`. +/// This extracts the last segment (the method name). +fn extract_method(path: &str) -> &str { + match path.rsplit('/').next() { + Some(m) if !m.is_empty() => m, + _ => "unknown", + } +} + +/// Tower Layer that instruments gRPC handlers with Prometheus metrics. +/// +/// Records three metric families: +/// - `grpc_server_started_total` — incremented when a request begins +/// - `grpc_server_handled_total` — incremented when a request completes (with status code) +/// - `grpc_server_handling_seconds` — histogram of request duration +#[derive(Debug, Clone)] +pub struct GrpcMetricsLayer { + service_name: &'static str, +} + +impl GrpcMetricsLayer { + /// Create a new gRPC metrics layer for the given service. + /// + /// `service_name` is used as the `service` label on all emitted metrics. + pub fn new(service_name: &'static str) -> Self { + init_grpc_metrics(); + Self { service_name } + } +} + +impl Layer for GrpcMetricsLayer { + type Service = GrpcMetricsService; + + fn layer(&self, inner: S) -> Self::Service { + GrpcMetricsService { + inner, + service_name: self.service_name, + } + } +} + +/// Tower Service that records gRPC metrics around the inner service. +#[derive(Debug, Clone)] +pub struct GrpcMetricsService { + inner: S, + service_name: &'static str, +} + +impl Service> for GrpcMetricsService +where + S: Service, Response = http::Response> + + Clone + + Send + + 'static, + S::Future: Send + 'static, + S::Error: Send + 'static, + ReqBody: Send + 'static, +{ + type Response = S::Response; + type Error = S::Error; + type Future = Pin> + Send>>; + + fn poll_ready(&mut self, cx: &mut Context<'_>) -> Poll> { + self.inner.poll_ready(cx) + } + + fn call(&mut self, req: http::Request) -> Self::Future { + let service_name = self.service_name; + let method = extract_method(req.uri().path()).to_string(); + + GRPC_STARTED + .with_label_values(&[service_name, &method]) + .inc(); + + let start = std::time::Instant::now(); + // Clone inner before calling to satisfy tower's Service contract: + // poll_ready was called on `self`, so we must call on `self` + // (not the clone). We swap so the clone becomes `self` for + // the next call, and we consume the ready instance. + let mut inner = self.inner.clone(); + std::mem::swap(&mut self.inner, &mut inner); + + Box::pin(async move { + let result = inner.call(req).await; + let elapsed = start.elapsed().as_secs_f64(); + + let code = match &result { + Ok(resp) => resp + .headers() + .get("grpc-status") + .and_then(|v| v.to_str().ok()) + .map_or_else(|| "0".to_string(), |s| s.to_string()), + Err(_) => "13".to_string(), // INTERNAL + }; + + GRPC_HANDLED + .with_label_values(&[service_name, &method, &code]) + .inc(); + + GRPC_HANDLING_SECONDS + .with_label_values(&[service_name, &method]) + .observe(elapsed); + + result + }) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_extract_method_normal_path() { + assert_eq!(extract_method("/package.Service/MethodName"), "MethodName"); + } + + #[test] + fn test_extract_method_simple_path() { + assert_eq!(extract_method("/Method"), "Method"); + } + + #[test] + fn test_extract_method_empty_path() { + assert_eq!(extract_method(""), "unknown"); + } + + #[test] + fn test_extract_method_trailing_slash() { + // rsplit('/').next() on "foo/" yields "" (empty), so fallback + assert_eq!(extract_method("/Service/"), "unknown"); + } + + #[test] + fn test_extract_method_root_slash() { + assert_eq!(extract_method("/"), "unknown"); + } + + #[test] + fn test_metrics_layer_creates_service() { + let layer = GrpcMetricsLayer::new("test_svc"); + assert_eq!(layer.service_name, "test_svc"); + } + + #[test] + fn test_init_grpc_metrics_is_idempotent() { + // Calling init multiple times should not panic + init_grpc_metrics(); + init_grpc_metrics(); + init_grpc_metrics(); + } + + #[test] + fn test_metrics_registered() { + init_grpc_metrics(); + // Verify metrics are accessible via the global statics + GRPC_STARTED + .with_label_values(&["test_svc", "TestMethod"]) + .inc(); + let val = GRPC_STARTED + .with_label_values(&["test_svc", "TestMethod"]) + .get(); + assert!(val >= 1.0, "Counter should have been incremented"); + } +} diff --git a/crates/common/src/metrics/mod.rs b/crates/common/src/metrics/mod.rs index 7090d7f97..9297aceed 100644 --- a/crates/common/src/metrics/mod.rs +++ b/crates/common/src/metrics/mod.rs @@ -3,6 +3,7 @@ // Provides Prometheus-based metrics collection with a global registry // and helper functions for all services. +pub mod grpc_metrics; pub mod registry; // Re-export commonly used functions @@ -14,5 +15,7 @@ pub use registry::{ set_gauge, set_gauge_vec, REGISTRY, }; +pub use grpc_metrics::{init_grpc_metrics, GrpcMetricsLayer}; + #[cfg(test)] mod tests;