exercises

Unnamed repository; edit this file 'description' to name the repository.
Log | Files | Refs | README

commit 6a4c130576362f984076231453ccf75325ed2bc6
parent b7827798bf35fc0f34dd7fd8ef67546c34427a1c
Author: ling0x <ling0x@users.noreply.github.com>
Date:   Sat, 29 Aug 2026 02:31:55 +0100

add metrics

Diffstat:
Mdistributed_observability/Cargo.lock | 24+++++++++++++++++++++---
Mdistributed_observability/Cargo.toml | 1+
Mdistributed_observability/inventory/Cargo.toml | 1+
Mdistributed_observability/inventory/src/main.rs | 16+++++++++++++---
Mdistributed_observability/inventory/src/telemetry.rs | 184+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--------------
Mdistributed_observability/orders/Cargo.toml | 1+
Mdistributed_observability/orders/src/main.rs | 16+++++++++++++---
Mdistributed_observability/orders/src/telemetry.rs | 184+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--------------
Mdistributed_observability/otelmart/Cargo.toml | 1+
Mdistributed_observability/otelmart/src/main.rs | 16+++++++++++++---
Mdistributed_observability/otelmart/src/telemetry.rs | 186++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---------------
Mdistributed_observability/products/Cargo.toml | 1+
Mdistributed_observability/products/src/main.rs | 16+++++++++++++---
Mdistributed_observability/products/src/telemetry.rs | 196++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---------
14 files changed, 706 insertions(+), 137 deletions(-)

diff --git a/distributed_observability/Cargo.lock b/distributed_observability/Cargo.lock @@ -170,6 +170,20 @@ dependencies = [ ] [[package]] +name = "axum-otel-metrics" +version = "0.14.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d0d74d40b1ff9bd3e8f782aa6798d22ff6a7d1b53049e0e239f5332dab2dd83d" +dependencies = [ + "axum", + "http-body", + "opentelemetry", + "opentelemetry-semantic-conventions", + "pin-project-lite", + "tower", +] + +[[package]] name = "axum-tracing-opentelemetry" version = "0.38.0" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -1181,6 +1195,7 @@ version = "0.1.0" dependencies = [ "anyhow", "axum", + "axum-otel-metrics", "axum-tracing-opentelemetry", "bigdecimal", "chrono", @@ -1643,6 +1658,7 @@ version = "0.1.0" dependencies = [ "anyhow", "axum", + "axum-otel-metrics", "axum-tracing-opentelemetry", "bigdecimal", "chrono", @@ -1673,6 +1689,7 @@ version = "0.1.0" dependencies = [ "anyhow", "axum", + "axum-otel-metrics", "axum-tracing-opentelemetry", "bcrypt", "chrono", @@ -1906,6 +1923,7 @@ version = "0.1.0" dependencies = [ "anyhow", "axum", + "axum-otel-metrics", "axum-tracing-opentelemetry", "bigdecimal", "chrono", @@ -2015,7 +2033,7 @@ dependencies = [ "once_cell", "socket2", "tracing", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -2351,7 +2369,7 @@ dependencies = [ "security-framework", "security-framework-sys", "webpki-root-certs", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -3602,7 +3620,7 @@ version = "0.1.11" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22" dependencies = [ - "windows-sys 0.48.0", + "windows-sys 0.61.2", ] [[package]] diff --git a/distributed_observability/Cargo.toml b/distributed_observability/Cargo.toml @@ -98,6 +98,7 @@ tracing-opentelemetry = "0.33.0" axum-tracing-opentelemetry = "0.38.0" reqwest-middleware = { version = "0.5.2", features = ["json"] } reqwest-tracing = { version = "0.7.1", features = ["opentelemetry_0_32"] } +axum-otel-metrics = "0.14.1" [profile.release] opt-level = 3 diff --git a/distributed_observability/inventory/Cargo.toml b/distributed_observability/inventory/Cargo.toml @@ -54,3 +54,4 @@ tracing-opentelemetry = { workspace = true } axum-tracing-opentelemetry = { workspace = true } reqwest-middleware = { workspace = true } reqwest-tracing = { workspace = true } +axum-otel-metrics = { workspace = true } diff --git a/distributed_observability/inventory/src/main.rs b/distributed_observability/inventory/src/main.rs @@ -45,6 +45,7 @@ use axum::{ routing::{get, post, put}, Router, }; +use axum_otel_metrics::HttpMetricsLayerBuilder; use axum_tracing_opentelemetry::middleware::{OtelAxumLayer, OtelInResponseLayer}; use std::net::SocketAddr; use tower_http::cors::CorsLayer; @@ -59,7 +60,7 @@ async fn main() -> Result<()> { dotenvy::dotenv().ok(); // Initialize telemetry (tracing + OpenTelemetry) - let tracer_provider = telemetry::init_telemetry("inventory"); + let _telemetry_guard = telemetry::init_telemetry("inventory"); // Load configuration from config.toml let config = Config::load()?; @@ -74,6 +75,13 @@ async fn main() -> Result<()> { // Initialize database connection pool let db = Database::new(&config.database.url, config.database.max_connections).await?; + // Register observable gauges for connection pool health metrics + let meter = opentelemetry::global::meter("inventory-service"); + telemetry::register_pool_metrics(&meter, db.pool().clone()); + + // Build automatic HTTP RED metrics layer + let metrics = HttpMetricsLayerBuilder::new().build(); + // Build the application router with all endpoints let app = Router::new() // Inventory endpoints @@ -93,6 +101,8 @@ async fn main() -> Result<()> { get(handlers::get_pricing_by_product), ) .route("/pricing/{product_uuid}", put(handlers::upsert_pricing)) + // Automatic RED metrics (request rate, error rate, duration) + .layer(metrics) // Include trace context as header into the response .layer(OtelInResponseLayer::default()) // Start OpenTelemetry trace on incoming request @@ -124,8 +134,8 @@ async fn main() -> Result<()> { let listener = tokio::net::TcpListener::bind(addr).await?; axum::serve(listener, app).await?; - // Shutdown telemetry on exit (flush pending spans) - telemetry::shutdown_telemetry(tracer_provider); + // Telemetry is flushed and shut down by the `Drop` impl on `_telemetry` + // when this function returns. Ok(()) } diff --git a/distributed_observability/inventory/src/telemetry.rs b/distributed_observability/inventory/src/telemetry.rs @@ -1,54 +1,106 @@ //! Telemetry initialization for the Inventory service. //! -//! This module configures OpenTelemetry tracing with OTLP export to Jaeger. -//! It sets up a layered tracing subscriber that outputs both to stdout and -//! sends spans to the configured OTLP endpoint. +//! This module configures OpenTelemetry tracing and metrics pipelines. +//! Traces are exported via gRPC to Jaeger, and metrics are pushed via +//! OTLP HTTP to Prometheus. It sets up a layered tracing subscriber +//! that outputs both to stdout and sends spans to the configured endpoint. use opentelemetry::global; use opentelemetry::trace::TracerProvider as _; +use opentelemetry::KeyValue; use opentelemetry_otlp::WithExportConfig; -use opentelemetry_sdk::{propagation::TraceContextPropagator, trace::SdkTracerProvider, Resource}; +use opentelemetry_sdk::{ + metrics::SdkMeterProvider, propagation::TraceContextPropagator, trace::SdkTracerProvider, + Resource, +}; use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt, EnvFilter}; -/// Initializes the telemetry pipeline for the service. -/// -/// This function sets up: -/// - An OTLP exporter that sends traces to Jaeger (via OTEL_EXPORTER_OTLP_ENDPOINT) -/// - A tracing subscriber with environment-based filtering (RUST_LOG) -/// - Console output for local debugging +/// Holds both the tracer and meter providers, ensuring they are +/// shut down gracefully when the application exits. Dropping the guard +/// flushes pending telemetry and shuts down each provider in turn, so +/// callers only need to keep the value alive for the lifetime of the +/// application (e.g. `let _telemetry = telemetry::init_telemetry("inventory");`). +pub struct TelemetryGuard { + tracer_provider: SdkTracerProvider, + meter_provider: SdkMeterProvider, +} + +impl Drop for TelemetryGuard { + fn drop(&mut self) { + tracing::info!("Shutting down telemetry..."); + if let Err(e) = self.meter_provider.shutdown() { + eprintln!("Error shutting down meter provider: {:?}", e); + } + if let Err(e) = self.tracer_provider.shutdown() { + eprintln!("Error shutting down tracer provider: {:?}", e); + } + tracing::info!("Telemetry shutdown complete"); + } +} + +/// Initializes the full telemetry pipeline (tracing + metrics). /// /// # Arguments -/// * `service_name` - The name of this service, used to identify traces in Jaeger +/// * `service_name` - The name of this service, used to identify traces and metrics /// /// # Returns -/// * `SdkTracerProvider` - The tracer provider, which should be kept alive and -/// shut down gracefully when the application exits +/// * `TelemetryGuard` - Holds both providers; keep alive for the lifetime of the app /// /// # Panics -/// Panics if the OTLP exporter or tracing subscriber cannot be initialized. -pub fn init_telemetry(service_name: &str) -> SdkTracerProvider { +/// Panics if the OTLP exporters or tracing subscriber cannot be initialized. +pub fn init_telemetry(service_name: &str) -> TelemetryGuard { // Set up W3C Trace Context propagator for cross-service trace correlation global::set_text_map_propagator(TraceContextPropagator::new()); - // Get the OTLP endpoint from environment, defaulting to localhost for local dev - let otlp_endpoint = std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT") - .unwrap_or_else(|_| "http://localhost:4317".to_string()); + // Build a shared resource describing this service + let resource = Resource::builder() + .with_service_name(service_name.to_string()) + .with_attributes([ + KeyValue::new("service.version", env!("CARGO_PKG_VERSION")), + KeyValue::new( + "deployment.environment", + std::env::var("ENVIRONMENT").unwrap_or_else(|_| "development".into()), + ), + ]) + .build(); + + // Resolve separate endpoints for traces (gRPC) and metrics (HTTP) + let traces_endpoint = std::env::var("OTEL_EXPORTER_OTLP_TRACES_ENDPOINT") + .or_else(|_| std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT")) + .unwrap_or_else(|_| "http://localhost:4317".into()); + + let metrics_endpoint = std::env::var("OTEL_EXPORTER_OTLP_METRICS_ENDPOINT") + .or_else(|_| std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT")) + .unwrap_or_else(|_| "http://localhost:9090/api/v1/otlp/v1/metrics".into()); + + // Initialize both providers + let tracer_provider = init_tracer_provider(service_name, &traces_endpoint, &resource); + let meter_provider = init_meter_provider(&metrics_endpoint, &resource); + + TelemetryGuard { + tracer_provider, + meter_provider, + } +} +/// Creates the OTLP gRPC trace exporter and tracer provider, +/// then wires it into the global tracing subscriber. +fn init_tracer_provider( + service_name: &str, + endpoint: &str, + resource: &Resource, +) -> SdkTracerProvider { // Create the OTLP exporter configured to send traces via gRPC let exporter = opentelemetry_otlp::SpanExporter::builder() .with_tonic() - .with_endpoint(&otlp_endpoint) + .with_endpoint(endpoint) .build() - .expect("Failed to create OTLP exporter"); + .expect("Failed to create OTLP span exporter"); // Build the tracer provider with the exporter and service resource let provider = SdkTracerProvider::builder() .with_batch_exporter(exporter) - .with_resource( - Resource::builder_empty() - .with_service_name(service_name.to_string()) - .build(), - ) + .with_resource(resource.clone()) .build(); // Create a tracer from the provider for the OpenTelemetry layer @@ -78,15 +130,81 @@ pub fn init_telemetry(service_name: &str) -> SdkTracerProvider { provider } -/// Shuts down the telemetry pipeline gracefully. +/// Creates the OTLP HTTP metric exporter and meter provider, +/// then registers it as the global meter provider. +fn init_meter_provider(endpoint: &str, resource: &Resource) -> SdkMeterProvider { + // Create the OTLP HTTP exporter targeting Prometheus OTLP receiver + let exporter = opentelemetry_otlp::MetricExporter::builder() + .with_http() + .with_endpoint(endpoint) + .build() + .expect("Failed to create OTLP metric exporter"); + + // Build the meter provider with a periodic exporter + let provider = SdkMeterProvider::builder() + .with_resource(resource.clone()) + .with_periodic_exporter(exporter) + .build(); + + // Register as the global meter provider + global::set_meter_provider(provider.clone()); + provider +} + +/// Registers observable gauges that track database connection pool health. /// -/// This ensures all pending spans are flushed to the OTLP endpoint -/// before the application exits. +/// Creates three instruments: +/// - `db.pool.connections.active` — connections currently in use +/// - `db.pool.connections.idle` — connections waiting for work +/// - `db.pool.utilization` — percentage of pool capacity in use /// /// # Arguments -/// * `provider` - The tracer provider returned from `init_telemetry` -pub fn shutdown_telemetry(provider: SdkTracerProvider) { - if let Err(e) = provider.shutdown() { - eprintln!("Failed to shutdown tracer provider: {e}"); - } +/// * `meter` - The meter to register gauges on +/// * `pool` - The sqlx connection pool to observe +pub fn register_pool_metrics( + meter: &opentelemetry::metrics::Meter, + pool: sqlx::Pool<sqlx::Postgres>, +) { + // Active connections = total size minus idle connections + let pool_clone = pool.clone(); + meter + .u64_observable_gauge("db.pool.connections.active") + .with_description("Active connections in the pool") + .with_callback(move |observer| { + let active = pool_clone + .size() + .saturating_sub(pool_clone.num_idle() as u32); + observer.observe(u64::from(active), &[]); + }) + .build(); + + // Idle connections sitting in the pool waiting for work + let pool_clone = pool.clone(); + meter + .u64_observable_gauge("db.pool.connections.idle") + .with_description("Idle connections in the pool") + .with_callback(move |observer| { + observer.observe(pool_clone.num_idle() as u64, &[]); + }) + .build(); + + // Utilization = (active / max_size) * 100 + let pool_clone = pool.clone(); + meter + .f64_observable_gauge("db.pool.utilization") + .with_description("Pool utilization percentage") + .with_unit("%") + .with_callback(move |observer| { + let active = pool_clone + .size() + .saturating_sub(pool_clone.num_idle() as u32); + let max_size = pool_clone.options().get_max_connections(); + let utilization = if max_size > 0 { + (f64::from(active) / f64::from(max_size)) * 100.0 + } else { + 0.0 + }; + observer.observe(utilization, &[]); + }) + .build(); } diff --git a/distributed_observability/orders/Cargo.toml b/distributed_observability/orders/Cargo.toml @@ -57,3 +57,4 @@ tracing-opentelemetry = { workspace = true } axum-tracing-opentelemetry = { workspace = true } reqwest-middleware = { workspace = true } reqwest-tracing = { workspace = true } +axum-otel-metrics = { workspace = true } diff --git a/distributed_observability/orders/src/main.rs b/distributed_observability/orders/src/main.rs @@ -41,6 +41,7 @@ use axum::{ routing::{get, post, put}, Router, }; +use axum_otel_metrics::HttpMetricsLayerBuilder; use axum_tracing_opentelemetry::middleware::{OtelAxumLayer, OtelInResponseLayer}; use std::net::SocketAddr; use tower_http::cors::CorsLayer; @@ -86,7 +87,7 @@ async fn main() -> Result<()> { dotenvy::dotenv().ok(); // Initialize telemetry (tracing + OpenTelemetry) - let tracer_provider = telemetry::init_telemetry("orders"); + let _telemetry_guard = telemetry::init_telemetry("orders"); // Load configuration from config.toml let config = Config::load()?; @@ -103,6 +104,13 @@ async fn main() -> Result<()> { // Initialize database connection pool let db = Database::new(&config.database.url, config.database.max_connections).await?; + // Register observable gauges for connection pool health metrics + let meter = opentelemetry::global::meter("inventory-service"); + telemetry::register_pool_metrics(&meter, db.pool().clone()); + + // Build automatic HTTP RED metrics layer + let metrics = HttpMetricsLayerBuilder::new().build(); + // Create HTTP client with tracing middleware for automatic // span creation and trace context propagation let reqwest_client = reqwest::Client::builder() @@ -139,6 +147,8 @@ async fn main() -> Result<()> { "/orders/{order_uuid}/shipment/status", put(handlers::update_shipment_status), ) + // Automatic RED metrics (request rate, error rate, duration) + .layer(metrics) // Include trace context as header into the response .layer(OtelInResponseLayer::default()) // Start OpenTelemetry trace on incoming request @@ -168,8 +178,8 @@ async fn main() -> Result<()> { let listener = tokio::net::TcpListener::bind(addr).await?; axum::serve(listener, app).await?; - // Shutdown telemetry on exit (flush pending spans) - telemetry::shutdown_telemetry(tracer_provider); + // Telemetry is flushed and shut down by the `Drop` impl on `_telemetry` + // when this function returns. Ok(()) } diff --git a/distributed_observability/orders/src/telemetry.rs b/distributed_observability/orders/src/telemetry.rs @@ -1,54 +1,106 @@ //! Telemetry initialization for the Inventory service. //! -//! This module configures OpenTelemetry tracing with OTLP export to Jaeger. -//! It sets up a layered tracing subscriber that outputs both to stdout and -//! sends spans to the configured OTLP endpoint. +//! This module configures OpenTelemetry tracing and metrics pipelines. +//! Traces are exported via gRPC to Jaeger, and metrics are pushed via +//! OTLP HTTP to Prometheus. It sets up a layered tracing subscriber +//! that outputs both to stdout and sends spans to the configured endpoint. use opentelemetry::global; use opentelemetry::trace::TracerProvider as _; +use opentelemetry::KeyValue; use opentelemetry_otlp::WithExportConfig; -use opentelemetry_sdk::{propagation::TraceContextPropagator, trace::SdkTracerProvider, Resource}; +use opentelemetry_sdk::{ + metrics::SdkMeterProvider, propagation::TraceContextPropagator, trace::SdkTracerProvider, + Resource, +}; use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt, EnvFilter}; -/// Initializes the telemetry pipeline for the service. -/// -/// This function sets up: -/// - An OTLP exporter that sends traces to Jaeger (via OTEL_EXPORTER_OTLP_ENDPOINT) -/// - A tracing subscriber with environment-based filtering (RUST_LOG) -/// - Console output for local debugging +/// Holds both the tracer and meter providers, ensuring they are +/// shut down gracefully when the application exits. Dropping the guard +/// flushes pending telemetry and shuts down each provider in turn, so +/// callers only need to keep the value alive for the lifetime of the +/// application (e.g. `let _telemetry = telemetry::init_telemetry("inventory");`). +pub struct TelemetryGuard { + tracer_provider: SdkTracerProvider, + meter_provider: SdkMeterProvider, +} + +impl Drop for TelemetryGuard { + fn drop(&mut self) { + tracing::info!("Shutting down telemetry..."); + if let Err(e) = self.meter_provider.shutdown() { + eprintln!("Error shutting down meter provider: {:?}", e); + } + if let Err(e) = self.tracer_provider.shutdown() { + eprintln!("Error shutting down tracer provider: {:?}", e); + } + tracing::info!("Telemetry shutdown complete"); + } +} + +/// Initializes the full telemetry pipeline (tracing + metrics). /// /// # Arguments -/// * `service_name` - The name of this service, used to identify traces in Jaeger +/// * `service_name` - The name of this service, used to identify traces and metrics /// /// # Returns -/// * `SdkTracerProvider` - The tracer provider, which should be kept alive and -/// shut down gracefully when the application exits +/// * `TelemetryGuard` - Holds both providers; keep alive for the lifetime of the app /// /// # Panics -/// Panics if the OTLP exporter or tracing subscriber cannot be initialized. -pub fn init_telemetry(service_name: &str) -> SdkTracerProvider { +/// Panics if the OTLP exporters or tracing subscriber cannot be initialized. +pub fn init_telemetry(service_name: &str) -> TelemetryGuard { // Set up W3C Trace Context propagator for cross-service trace correlation global::set_text_map_propagator(TraceContextPropagator::new()); - // Get the OTLP endpoint from environment, defaulting to localhost for local dev - let otlp_endpoint = std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT") - .unwrap_or_else(|_| "http://localhost:4317".to_string()); + // Build a shared resource describing this service + let resource = Resource::builder() + .with_service_name(service_name.to_string()) + .with_attributes([ + KeyValue::new("service.version", env!("CARGO_PKG_VERSION")), + KeyValue::new( + "deployment.environment", + std::env::var("ENVIRONMENT").unwrap_or_else(|_| "development".into()), + ), + ]) + .build(); + + // Resolve separate endpoints for traces (gRPC) and metrics (HTTP) + let traces_endpoint = std::env::var("OTEL_EXPORTER_OTLP_TRACES_ENDPOINT") + .or_else(|_| std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT")) + .unwrap_or_else(|_| "http://localhost:4317".into()); + + let metrics_endpoint = std::env::var("OTEL_EXPORTER_OTLP_METRICS_ENDPOINT") + .or_else(|_| std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT")) + .unwrap_or_else(|_| "http://localhost:9090/api/v1/otlp/v1/metrics".into()); + + // Initialize both providers + let tracer_provider = init_tracer_provider(service_name, &traces_endpoint, &resource); + let meter_provider = init_meter_provider(&metrics_endpoint, &resource); + + TelemetryGuard { + tracer_provider, + meter_provider, + } +} +/// Creates the OTLP gRPC trace exporter and tracer provider, +/// then wires it into the global tracing subscriber. +fn init_tracer_provider( + service_name: &str, + endpoint: &str, + resource: &Resource, +) -> SdkTracerProvider { // Create the OTLP exporter configured to send traces via gRPC let exporter = opentelemetry_otlp::SpanExporter::builder() .with_tonic() - .with_endpoint(&otlp_endpoint) + .with_endpoint(endpoint) .build() - .expect("Failed to create OTLP exporter"); + .expect("Failed to create OTLP span exporter"); // Build the tracer provider with the exporter and service resource let provider = SdkTracerProvider::builder() .with_batch_exporter(exporter) - .with_resource( - Resource::builder_empty() - .with_service_name(service_name.to_string()) - .build(), - ) + .with_resource(resource.clone()) .build(); // Create a tracer from the provider for the OpenTelemetry layer @@ -78,15 +130,81 @@ pub fn init_telemetry(service_name: &str) -> SdkTracerProvider { provider } -/// Shuts down the telemetry pipeline gracefully. +/// Creates the OTLP HTTP metric exporter and meter provider, +/// then registers it as the global meter provider. +fn init_meter_provider(endpoint: &str, resource: &Resource) -> SdkMeterProvider { + // Create the OTLP HTTP exporter targeting Prometheus OTLP receiver + let exporter = opentelemetry_otlp::MetricExporter::builder() + .with_http() + .with_endpoint(endpoint) + .build() + .expect("Failed to create OTLP metric exporter"); + + // Build the meter provider with a periodic exporter + let provider = SdkMeterProvider::builder() + .with_resource(resource.clone()) + .with_periodic_exporter(exporter) + .build(); + + // Register as the global meter provider + global::set_meter_provider(provider.clone()); + provider +} + +/// Registers observable gauges that track database connection pool health. /// -/// This ensures all pending spans are flushed to the OTLP endpoint -/// before the application exits. +/// Creates three instruments: +/// - `db.pool.connections.active` — connections currently in use +/// - `db.pool.connections.idle` — connections waiting for work +/// - `db.pool.utilization` — percentage of pool capacity in use /// /// # Arguments -/// * `provider` - The tracer provider returned from `init_telemetry` -pub fn shutdown_telemetry(provider: SdkTracerProvider) { - if let Err(e) = provider.shutdown() { - eprintln!("Failed to shutdown tracer provider: {e}"); - } +/// * `meter` - The meter to register gauges on +/// * `pool` - The sqlx connection pool to observe +pub fn register_pool_metrics( + meter: &opentelemetry::metrics::Meter, + pool: sqlx::Pool<sqlx::Postgres>, +) { + // Active connections = total size minus idle connections + let pool_clone = pool.clone(); + meter + .u64_observable_gauge("db.pool.connections.active") + .with_description("Active connections in the pool") + .with_callback(move |observer| { + let active = pool_clone + .size() + .saturating_sub(pool_clone.num_idle() as u32); + observer.observe(u64::from(active), &[]); + }) + .build(); + + // Idle connections sitting in the pool waiting for work + let pool_clone = pool.clone(); + meter + .u64_observable_gauge("db.pool.connections.idle") + .with_description("Idle connections in the pool") + .with_callback(move |observer| { + observer.observe(pool_clone.num_idle() as u64, &[]); + }) + .build(); + + // Utilization = (active / max_size) * 100 + let pool_clone = pool.clone(); + meter + .f64_observable_gauge("db.pool.utilization") + .with_description("Pool utilization percentage") + .with_unit("%") + .with_callback(move |observer| { + let active = pool_clone + .size() + .saturating_sub(pool_clone.num_idle() as u32); + let max_size = pool_clone.options().get_max_connections(); + let utilization = if max_size > 0 { + (f64::from(active) / f64::from(max_size)) * 100.0 + } else { + 0.0 + }; + observer.observe(utilization, &[]); + }) + .build(); } diff --git a/distributed_observability/otelmart/Cargo.toml b/distributed_observability/otelmart/Cargo.toml @@ -69,5 +69,6 @@ tracing-opentelemetry = { workspace = true } axum-tracing-opentelemetry = { workspace = true } reqwest-middleware = { workspace = true } reqwest-tracing = { workspace = true } +axum-otel-metrics = { workspace = true } [dev-dependencies] diff --git a/distributed_observability/otelmart/src/main.rs b/distributed_observability/otelmart/src/main.rs @@ -14,6 +14,7 @@ use axum::{ routing::{get, post}, Router, }; +use axum_otel_metrics::HttpMetricsLayerBuilder; use axum_tracing_opentelemetry::middleware::{OtelAxumLayer, OtelInResponseLayer}; use std::net::SocketAddr; use std::path::PathBuf; @@ -91,7 +92,7 @@ async fn spa_fallback_handler( #[tokio::main] async fn main() -> Result<()> { // Initialize telemetry (tracing + OpenTelemetry) - let tracer_provider = telemetry::init_telemetry("otelmart"); + let _telemetry_guard = telemetry::init_telemetry("otelmart"); // Load configuration from config.toml let config = Config::load()?; @@ -110,6 +111,13 @@ async fn main() -> Result<()> { // Initialize database connection let db = Database::new(&config.database.url, config.database.max_connections).await?; + // Register observable gauges for connection pool health metrics + let meter = opentelemetry::global::meter("inventory-service"); + telemetry::register_pool_metrics(&meter, db.pool().clone()); + + // Build automatic HTTP RED metrics layer + let metrics = HttpMetricsLayerBuilder::new().build(); + // Create HTTP client with tracing middleware for automatic // span creation and trace context propagation let reqwest_client = reqwest::Client::builder() @@ -196,6 +204,8 @@ async fn main() -> Result<()> { let serve_dir = serve_dir.clone(); async move { spa_fallback_handler(uri, req, serve_dir, static_dir).await } }) + // Automatic RED metrics (request rate, error rate, duration) + .layer(metrics) // Include trace context as header into the response .layer(OtelInResponseLayer::default()) // Start OpenTelemetry trace on incoming request @@ -209,8 +219,8 @@ async fn main() -> Result<()> { let listener = tokio::net::TcpListener::bind(addr).await?; axum::serve(listener, app).await?; - // Shutdown telemetry on exit (flush pending spans) - telemetry::shutdown_telemetry(tracer_provider); + // Telemetry is flushed and shut down by the `Drop` impl on `_telemetry` + // when this function returns. Ok(()) } diff --git a/distributed_observability/otelmart/src/telemetry.rs b/distributed_observability/otelmart/src/telemetry.rs @@ -1,54 +1,106 @@ -//! Telemetry initialization for the OtelMart service. +//! Telemetry initialization for the Inventory service. //! -//! This module configures OpenTelemetry tracing with OTLP export to Jaeger. -//! It sets up a layered tracing subscriber that outputs both to stdout and -//! sends spans to the configured OTLP endpoint. +//! This module configures OpenTelemetry tracing and metrics pipelines. +//! Traces are exported via gRPC to Jaeger, and metrics are pushed via +//! OTLP HTTP to Prometheus. It sets up a layered tracing subscriber +//! that outputs both to stdout and sends spans to the configured endpoint. use opentelemetry::global; use opentelemetry::trace::TracerProvider as _; +use opentelemetry::KeyValue; use opentelemetry_otlp::WithExportConfig; -use opentelemetry_sdk::{propagation::TraceContextPropagator, trace::SdkTracerProvider, Resource}; +use opentelemetry_sdk::{ + metrics::SdkMeterProvider, propagation::TraceContextPropagator, trace::SdkTracerProvider, + Resource, +}; use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt, EnvFilter}; -/// Initializes the telemetry pipeline for the service. -/// -/// This function sets up: -/// - An OTLP exporter that sends traces to Jaeger (via OTEL_EXPORTER_OTLP_ENDPOINT) -/// - A tracing subscriber with environment-based filtering (RUST_LOG) -/// - Console output for local debugging +/// Holds both the tracer and meter providers, ensuring they are +/// shut down gracefully when the application exits. Dropping the guard +/// flushes pending telemetry and shuts down each provider in turn, so +/// callers only need to keep the value alive for the lifetime of the +/// application (e.g. `let _telemetry = telemetry::init_telemetry("inventory");`). +pub struct TelemetryGuard { + tracer_provider: SdkTracerProvider, + meter_provider: SdkMeterProvider, +} + +impl Drop for TelemetryGuard { + fn drop(&mut self) { + tracing::info!("Shutting down telemetry..."); + if let Err(e) = self.meter_provider.shutdown() { + eprintln!("Error shutting down meter provider: {:?}", e); + } + if let Err(e) = self.tracer_provider.shutdown() { + eprintln!("Error shutting down tracer provider: {:?}", e); + } + tracing::info!("Telemetry shutdown complete"); + } +} + +/// Initializes the full telemetry pipeline (tracing + metrics). /// /// # Arguments -/// * `service_name` - The name of this service, used to identify traces in Jaeger +/// * `service_name` - The name of this service, used to identify traces and metrics /// /// # Returns -/// * `SdkTracerProvider` - The tracer provider, which should be kept alive and -/// shut down gracefully when the application exits +/// * `TelemetryGuard` - Holds both providers; keep alive for the lifetime of the app /// /// # Panics -/// Panics if the OTLP exporter or tracing subscriber cannot be initialized. -pub fn init_telemetry(service_name: &str) -> SdkTracerProvider { +/// Panics if the OTLP exporters or tracing subscriber cannot be initialized. +pub fn init_telemetry(service_name: &str) -> TelemetryGuard { // Set up W3C Trace Context propagator for cross-service trace correlation global::set_text_map_propagator(TraceContextPropagator::new()); - // Get the OTLP endpoint from environment, defaulting to localhost for local dev - let otlp_endpoint = std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT") - .unwrap_or_else(|_| "http://localhost:4317".to_string()); + // Build a shared resource describing this service + let resource = Resource::builder() + .with_service_name(service_name.to_string()) + .with_attributes([ + KeyValue::new("service.version", env!("CARGO_PKG_VERSION")), + KeyValue::new( + "deployment.environment", + std::env::var("ENVIRONMENT").unwrap_or_else(|_| "development".into()), + ), + ]) + .build(); + + // Resolve separate endpoints for traces (gRPC) and metrics (HTTP) + let traces_endpoint = std::env::var("OTEL_EXPORTER_OTLP_TRACES_ENDPOINT") + .or_else(|_| std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT")) + .unwrap_or_else(|_| "http://localhost:4317".into()); + + let metrics_endpoint = std::env::var("OTEL_EXPORTER_OTLP_METRICS_ENDPOINT") + .or_else(|_| std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT")) + .unwrap_or_else(|_| "http://localhost:9090/api/v1/otlp/v1/metrics".into()); + + // Initialize both providers + let tracer_provider = init_tracer_provider(service_name, &traces_endpoint, &resource); + let meter_provider = init_meter_provider(&metrics_endpoint, &resource); + + TelemetryGuard { + tracer_provider, + meter_provider, + } +} +/// Creates the OTLP gRPC trace exporter and tracer provider, +/// then wires it into the global tracing subscriber. +fn init_tracer_provider( + service_name: &str, + endpoint: &str, + resource: &Resource, +) -> SdkTracerProvider { // Create the OTLP exporter configured to send traces via gRPC let exporter = opentelemetry_otlp::SpanExporter::builder() .with_tonic() - .with_endpoint(&otlp_endpoint) + .with_endpoint(endpoint) .build() - .expect("Failed to create OTLP exporter"); + .expect("Failed to create OTLP span exporter"); // Build the tracer provider with the exporter and service resource let provider = SdkTracerProvider::builder() .with_batch_exporter(exporter) - .with_resource( - Resource::builder_empty() - .with_service_name(service_name.to_string()) - .build(), - ) + .with_resource(resource.clone()) .build(); // Create a tracer from the provider for the OpenTelemetry layer @@ -78,15 +130,81 @@ pub fn init_telemetry(service_name: &str) -> SdkTracerProvider { provider } -/// Shuts down the telemetry pipeline gracefully. +/// Creates the OTLP HTTP metric exporter and meter provider, +/// then registers it as the global meter provider. +fn init_meter_provider(endpoint: &str, resource: &Resource) -> SdkMeterProvider { + // Create the OTLP HTTP exporter targeting Prometheus OTLP receiver + let exporter = opentelemetry_otlp::MetricExporter::builder() + .with_http() + .with_endpoint(endpoint) + .build() + .expect("Failed to create OTLP metric exporter"); + + // Build the meter provider with a periodic exporter + let provider = SdkMeterProvider::builder() + .with_resource(resource.clone()) + .with_periodic_exporter(exporter) + .build(); + + // Register as the global meter provider + global::set_meter_provider(provider.clone()); + provider +} + +/// Registers observable gauges that track database connection pool health. /// -/// This ensures all pending spans are flushed to the OTLP endpoint -/// before the application exits. +/// Creates three instruments: +/// - `db.pool.connections.active` — connections currently in use +/// - `db.pool.connections.idle` — connections waiting for work +/// - `db.pool.utilization` — percentage of pool capacity in use /// /// # Arguments -/// * `provider` - The tracer provider returned from `init_telemetry` -pub fn shutdown_telemetry(provider: SdkTracerProvider) { - if let Err(e) = provider.shutdown() { - eprintln!("Failed to shutdown tracer provider: {e}"); - } +/// * `meter` - The meter to register gauges on +/// * `pool` - The sqlx connection pool to observe +pub fn register_pool_metrics( + meter: &opentelemetry::metrics::Meter, + pool: sqlx::Pool<sqlx::Postgres>, +) { + // Active connections = total size minus idle connections + let pool_clone = pool.clone(); + meter + .u64_observable_gauge("db.pool.connections.active") + .with_description("Active connections in the pool") + .with_callback(move |observer| { + let active = pool_clone + .size() + .saturating_sub(pool_clone.num_idle() as u32); + observer.observe(u64::from(active), &[]); + }) + .build(); + + // Idle connections sitting in the pool waiting for work + let pool_clone = pool.clone(); + meter + .u64_observable_gauge("db.pool.connections.idle") + .with_description("Idle connections in the pool") + .with_callback(move |observer| { + observer.observe(pool_clone.num_idle() as u64, &[]); + }) + .build(); + + // Utilization = (active / max_size) * 100 + let pool_clone = pool.clone(); + meter + .f64_observable_gauge("db.pool.utilization") + .with_description("Pool utilization percentage") + .with_unit("%") + .with_callback(move |observer| { + let active = pool_clone + .size() + .saturating_sub(pool_clone.num_idle() as u32); + let max_size = pool_clone.options().get_max_connections(); + let utilization = if max_size > 0 { + (f64::from(active) / f64::from(max_size)) * 100.0 + } else { + 0.0 + }; + observer.observe(utilization, &[]); + }) + .build(); } diff --git a/distributed_observability/products/Cargo.toml b/distributed_observability/products/Cargo.toml @@ -54,3 +54,4 @@ tracing-opentelemetry = { workspace = true } axum-tracing-opentelemetry = { workspace = true } reqwest-middleware = { workspace = true } reqwest-tracing = { workspace = true } +axum-otel-metrics = { workspace = true } diff --git a/distributed_observability/products/src/main.rs b/distributed_observability/products/src/main.rs @@ -35,6 +35,7 @@ use axum::{ routing::{get, put}, Router, }; +use axum_otel_metrics::HttpMetricsLayerBuilder; use axum_tracing_opentelemetry::middleware::{OtelAxumLayer, OtelInResponseLayer}; use std::net::SocketAddr; use tower_http::cors::CorsLayer; @@ -49,7 +50,7 @@ async fn main() -> Result<()> { // This is optional - the service will work without it dotenvy::dotenv().ok(); - let tracer_provider = telemetry::init_telemetry("products"); + let _telemetry_guard = telemetry::init_telemetry("products"); // Load configuration from config.toml let config = Config::load()?; @@ -63,6 +64,13 @@ async fn main() -> Result<()> { // This sets the search path to the products schema for all connections let db = Database::new(&config.database.url, config.database.max_connections).await?; + // Register observable gauges for connection pool health metrics + let meter = opentelemetry::global::meter("inventory-service"); + telemetry::register_pool_metrics(&meter, db.pool().clone()); + + // Build automatic HTTP RED metrics layer + let metrics = HttpMetricsLayerBuilder::new().build(); + // Build the application router with all endpoints let app = Router::new() // Product endpoints @@ -76,6 +84,8 @@ async fn main() -> Result<()> { // Add database pool to application state // All handlers will have access to this via State extractor .with_state(db.pool().clone()) + // Automatic RED metrics (request rate, error rate, duration) + .layer(metrics) .layer(OtelInResponseLayer) .layer(OtelAxumLayer::default()) .layer(CorsLayer::permissive()); @@ -95,8 +105,8 @@ async fn main() -> Result<()> { let listener = tokio::net::TcpListener::bind(addr).await?; axum::serve(listener, app).await?; - // Graceful shutdown telemetry - telemetry::shutdown_telemetry(tracer_provider); + // Telemetry is flushed and shut down by the `Drop` impl on `_telemetry` + // when this function returns. Ok(()) } diff --git a/distributed_observability/products/src/telemetry.rs b/distributed_observability/products/src/telemetry.rs @@ -1,47 +1,126 @@ -use opentelemetry::{global::set_text_map_propagator, trace::TracerProvider}; -use opentelemetry_otlp::{SpanExporter, WithExportConfig}; +//! Telemetry initialization for the Inventory service. +//! +//! This module configures OpenTelemetry tracing and metrics pipelines. +//! Traces are exported via gRPC to Jaeger, and metrics are pushed via +//! OTLP HTTP to Prometheus. It sets up a layered tracing subscriber +//! that outputs both to stdout and sends spans to the configured endpoint. + +use opentelemetry::global; +use opentelemetry::trace::TracerProvider as _; +use opentelemetry::KeyValue; +use opentelemetry_otlp::WithExportConfig; use opentelemetry_sdk::{ - propagation::TraceContextPropagator, - trace::{RandomIdGenerator, SdkTracerProvider}, + metrics::SdkMeterProvider, propagation::TraceContextPropagator, trace::SdkTracerProvider, Resource, }; use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt, EnvFilter}; -pub fn init_telemetry(service_name: &str) -> SdkTracerProvider { - // Cross service trace correlation - set_text_map_propagator(TraceContextPropagator::new()); +/// Holds both the tracer and meter providers, ensuring they are +/// shut down gracefully when the application exits. Dropping the guard +/// flushes pending telemetry and shuts down each provider in turn, so +/// callers only need to keep the value alive for the lifetime of the +/// application (e.g. `let _telemetry = telemetry::init_telemetry("inventory");`). +pub struct TelemetryGuard { + tracer_provider: SdkTracerProvider, + meter_provider: SdkMeterProvider, +} - let otlp_endpoint = std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT") - .unwrap_or_else(|_| "http://localhost:4317".to_string()); +impl Drop for TelemetryGuard { + fn drop(&mut self) { + tracing::info!("Shutting down telemetry..."); + if let Err(e) = self.meter_provider.shutdown() { + eprintln!("Error shutting down meter provider: {:?}", e); + } + if let Err(e) = self.tracer_provider.shutdown() { + eprintln!("Error shutting down tracer provider: {:?}", e); + } + tracing::info!("Telemetry shutdown complete"); + } +} + +/// Initializes the full telemetry pipeline (tracing + metrics). +/// +/// # Arguments +/// * `service_name` - The name of this service, used to identify traces and metrics +/// +/// # Returns +/// * `TelemetryGuard` - Holds both providers; keep alive for the lifetime of the app +/// +/// # Panics +/// Panics if the OTLP exporters or tracing subscriber cannot be initialized. +pub fn init_telemetry(service_name: &str) -> TelemetryGuard { + // Set up W3C Trace Context propagator for cross-service trace correlation + global::set_text_map_propagator(TraceContextPropagator::new()); + + // Build a shared resource describing this service + let resource = Resource::builder() + .with_service_name(service_name.to_string()) + .with_attributes([ + KeyValue::new("service.version", env!("CARGO_PKG_VERSION")), + KeyValue::new( + "deployment.environment", + std::env::var("ENVIRONMENT").unwrap_or_else(|_| "development".into()), + ), + ]) + .build(); + + // Resolve separate endpoints for traces (gRPC) and metrics (HTTP) + let traces_endpoint = std::env::var("OTEL_EXPORTER_OTLP_TRACES_ENDPOINT") + .or_else(|_| std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT")) + .unwrap_or_else(|_| "http://localhost:4317".into()); + + let metrics_endpoint = std::env::var("OTEL_EXPORTER_OTLP_METRICS_ENDPOINT") + .or_else(|_| std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT")) + .unwrap_or_else(|_| "http://localhost:9090/api/v1/otlp/v1/metrics".into()); + + // Initialize both providers + let tracer_provider = init_tracer_provider(service_name, &traces_endpoint, &resource); + let meter_provider = init_meter_provider(&metrics_endpoint, &resource); + + TelemetryGuard { + tracer_provider, + meter_provider, + } +} - let exporter = SpanExporter::builder() +/// Creates the OTLP gRPC trace exporter and tracer provider, +/// then wires it into the global tracing subscriber. +fn init_tracer_provider( + service_name: &str, + endpoint: &str, + resource: &Resource, +) -> SdkTracerProvider { + // Create the OTLP exporter configured to send traces via gRPC + let exporter = opentelemetry_otlp::SpanExporter::builder() .with_tonic() - .with_endpoint(&otlp_endpoint) + .with_endpoint(endpoint) .build() - .expect("Failed to create OTLP exporter"); + .expect("Failed to create OTLP span exporter"); + // Build the tracer provider with the exporter and service resource let provider = SdkTracerProvider::builder() .with_batch_exporter(exporter) - .with_id_generator(RandomIdGenerator::default()) - .with_resource( - Resource::builder_empty() - .with_service_name(service_name.to_string()) - .build(), - ) + .with_resource(resource.clone()) .build(); + // Create a tracer from the provider for the OpenTelemetry layer let tracer = provider.tracer(service_name.to_string()); + // Create the OpenTelemetry tracing layer let otel_layer = tracing_opentelemetry::layer().with_tracer(tracer); + // Create an environment filter for log levels (defaults to "info") + // Include otel::tracing=trace to allow axum-tracing-opentelemetry spans through let env_filter = EnvFilter::try_from_default_env() .unwrap_or_else(|_| EnvFilter::new("info,otel::tracing=trace")); + // Create a formatting layer for console output let fmt_layer = tracing_subscriber::fmt::layer() .with_target(true) .with_thread_ids(false) .compact(); + // Combine all layers into a subscriber and set it as the global default tracing_subscriber::registry() .with(env_filter) .with(fmt_layer) @@ -51,8 +130,81 @@ pub fn init_telemetry(service_name: &str) -> SdkTracerProvider { provider } -pub fn shutdown_telemetry(provider: SdkTracerProvider) { - if let Err(e) = provider.shutdown() { - eprintln!("Failed to shutdown tracer provider: {e}"); - } +/// Creates the OTLP HTTP metric exporter and meter provider, +/// then registers it as the global meter provider. +fn init_meter_provider(endpoint: &str, resource: &Resource) -> SdkMeterProvider { + // Create the OTLP HTTP exporter targeting Prometheus OTLP receiver + let exporter = opentelemetry_otlp::MetricExporter::builder() + .with_http() + .with_endpoint(endpoint) + .build() + .expect("Failed to create OTLP metric exporter"); + + // Build the meter provider with a periodic exporter + let provider = SdkMeterProvider::builder() + .with_resource(resource.clone()) + .with_periodic_exporter(exporter) + .build(); + + // Register as the global meter provider + global::set_meter_provider(provider.clone()); + provider +} + +/// Registers observable gauges that track database connection pool health. +/// +/// Creates three instruments: +/// - `db.pool.connections.active` — connections currently in use +/// - `db.pool.connections.idle` — connections waiting for work +/// - `db.pool.utilization` — percentage of pool capacity in use +/// +/// # Arguments +/// * `meter` - The meter to register gauges on +/// * `pool` - The sqlx connection pool to observe +pub fn register_pool_metrics( + meter: &opentelemetry::metrics::Meter, + pool: sqlx::Pool<sqlx::Postgres>, +) { + // Active connections = total size minus idle connections + let pool_clone = pool.clone(); + meter + .u64_observable_gauge("db.pool.connections.active") + .with_description("Active connections in the pool") + .with_callback(move |observer| { + let active = pool_clone + .size() + .saturating_sub(pool_clone.num_idle() as u32); + observer.observe(u64::from(active), &[]); + }) + .build(); + + // Idle connections sitting in the pool waiting for work + let pool_clone = pool.clone(); + meter + .u64_observable_gauge("db.pool.connections.idle") + .with_description("Idle connections in the pool") + .with_callback(move |observer| { + observer.observe(pool_clone.num_idle() as u64, &[]); + }) + .build(); + + // Utilization = (active / max_size) * 100 + let pool_clone = pool.clone(); + meter + .f64_observable_gauge("db.pool.utilization") + .with_description("Pool utilization percentage") + .with_unit("%") + .with_callback(move |observer| { + let active = pool_clone + .size() + .saturating_sub(pool_clone.num_idle() as u32); + let max_size = pool_clone.options().get_max_connections(); + let utilization = if max_size > 0 { + (f64::from(active) / f64::from(max_size)) * 100.0 + } else { + 0.0 + }; + observer.observe(utilization, &[]); + }) + .build(); }