telemetry.rs (7978B)
1 //! Telemetry initialization for the Inventory service. 2 //! 3 //! This module configures OpenTelemetry tracing and metrics pipelines. 4 //! Traces are exported via gRPC to Jaeger, and metrics are pushed via 5 //! OTLP HTTP to Prometheus. It sets up a layered tracing subscriber 6 //! that outputs both to stdout and sends spans to the configured endpoint. 7 8 use opentelemetry::global; 9 use opentelemetry::trace::TracerProvider as _; 10 use opentelemetry::KeyValue; 11 use opentelemetry_otlp::WithExportConfig; 12 use opentelemetry_sdk::{ 13 metrics::SdkMeterProvider, propagation::TraceContextPropagator, trace::SdkTracerProvider, 14 Resource, 15 }; 16 use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt, EnvFilter}; 17 18 /// Holds both the tracer and meter providers, ensuring they are 19 /// shut down gracefully when the application exits. Dropping the guard 20 /// flushes pending telemetry and shuts down each provider in turn, so 21 /// callers only need to keep the value alive for the lifetime of the 22 /// application (e.g. `let _telemetry = telemetry::init_telemetry("inventory");`). 23 pub struct TelemetryGuard { 24 tracer_provider: SdkTracerProvider, 25 meter_provider: SdkMeterProvider, 26 } 27 28 impl Drop for TelemetryGuard { 29 fn drop(&mut self) { 30 tracing::info!("Shutting down telemetry..."); 31 if let Err(e) = self.meter_provider.shutdown() { 32 eprintln!("Error shutting down meter provider: {:?}", e); 33 } 34 if let Err(e) = self.tracer_provider.shutdown() { 35 eprintln!("Error shutting down tracer provider: {:?}", e); 36 } 37 tracing::info!("Telemetry shutdown complete"); 38 } 39 } 40 41 /// Initializes the full telemetry pipeline (tracing + metrics). 42 /// 43 /// # Arguments 44 /// * `service_name` - The name of this service, used to identify traces and metrics 45 /// 46 /// # Returns 47 /// * `TelemetryGuard` - Holds both providers; keep alive for the lifetime of the app 48 /// 49 /// # Panics 50 /// Panics if the OTLP exporters or tracing subscriber cannot be initialized. 51 pub fn init_telemetry(service_name: &str) -> TelemetryGuard { 52 // Set up W3C Trace Context propagator for cross-service trace correlation 53 global::set_text_map_propagator(TraceContextPropagator::new()); 54 55 // Build a shared resource describing this service 56 let resource = Resource::builder() 57 .with_service_name(service_name.to_string()) 58 .with_attributes([ 59 KeyValue::new("service.version", env!("CARGO_PKG_VERSION")), 60 KeyValue::new( 61 "deployment.environment", 62 std::env::var("ENVIRONMENT").unwrap_or_else(|_| "development".into()), 63 ), 64 ]) 65 .build(); 66 67 // Resolve separate endpoints for traces (gRPC) and metrics (HTTP) 68 let traces_endpoint = std::env::var("OTEL_EXPORTER_OTLP_TRACES_ENDPOINT") 69 .or_else(|_| std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT")) 70 .unwrap_or_else(|_| "http://localhost:9090/api/v1/otlp/v1/metrics".into()); 71 72 let metrics_endpoint = std::env::var("OTEL_EXPORTER_OTLP_METRICS_ENDPOINT") 73 .or_else(|_| std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT")) 74 .unwrap_or_else(|_| "http://localhost:4317".into()); 75 76 // Initialize both providers 77 let tracer_provider = init_tracer_provider(service_name, &traces_endpoint, &resource); 78 let meter_provider = init_meter_provider(&metrics_endpoint, &resource); 79 80 TelemetryGuard { 81 tracer_provider, 82 meter_provider, 83 } 84 } 85 86 /// Creates the OTLP gRPC trace exporter and tracer provider, 87 /// then wires it into the global tracing subscriber. 88 fn init_tracer_provider( 89 service_name: &str, 90 endpoint: &str, 91 resource: &Resource, 92 ) -> SdkTracerProvider { 93 // Create the OTLP exporter configured to send traces via gRPC 94 let exporter = opentelemetry_otlp::SpanExporter::builder() 95 .with_tonic() 96 .with_endpoint(endpoint) 97 .build() 98 .expect("Failed to create OTLP span exporter"); 99 100 // Build the tracer provider with the exporter and service resource 101 let provider = SdkTracerProvider::builder() 102 .with_batch_exporter(exporter) 103 .with_resource(resource.clone()) 104 .build(); 105 106 // Create a tracer from the provider for the OpenTelemetry layer 107 let tracer = provider.tracer(service_name.to_string()); 108 109 // Create the OpenTelemetry tracing layer 110 let otel_layer = tracing_opentelemetry::layer().with_tracer(tracer); 111 112 // Create an environment filter for log levels (defaults to "info") 113 // Include otel::tracing=trace to allow axum-tracing-opentelemetry spans through 114 let env_filter = EnvFilter::try_from_default_env() 115 .unwrap_or_else(|_| EnvFilter::new("info,otel::tracing=trace")); 116 117 // Create a formatting layer for console output 118 let fmt_layer = tracing_subscriber::fmt::layer() 119 .with_target(true) 120 .with_thread_ids(false) 121 .compact(); 122 123 // Combine all layers into a subscriber and set it as the global default 124 tracing_subscriber::registry() 125 .with(env_filter) 126 .with(fmt_layer) 127 .with(otel_layer) 128 .init(); 129 130 provider 131 } 132 133 /// Creates the OTLP HTTP metric exporter and meter provider, 134 /// then registers it as the global meter provider. 135 fn init_meter_provider(endpoint: &str, resource: &Resource) -> SdkMeterProvider { 136 // Create the OTLP HTTP exporter targeting Prometheus OTLP receiver 137 let exporter = opentelemetry_otlp::MetricExporter::builder() 138 .with_http() 139 .with_endpoint(endpoint) 140 .build() 141 .expect("Failed to create OTLP metric exporter"); 142 143 // Build the meter provider with a periodic exporter 144 let provider = SdkMeterProvider::builder() 145 .with_resource(resource.clone()) 146 .with_periodic_exporter(exporter) 147 .build(); 148 149 // Register as the global meter provider 150 global::set_meter_provider(provider.clone()); 151 provider 152 } 153 154 /// Registers observable gauges that track database connection pool health. 155 /// 156 /// Creates three instruments: 157 /// - `db.pool.connections.active` — connections currently in use 158 /// - `db.pool.connections.idle` — connections waiting for work 159 /// - `db.pool.utilization` — percentage of pool capacity in use 160 /// 161 /// # Arguments 162 /// * `meter` - The meter to register gauges on 163 /// * `pool` - The sqlx connection pool to observe 164 pub fn register_pool_metrics( 165 meter: &opentelemetry::metrics::Meter, 166 pool: sqlx::Pool<sqlx::Postgres>, 167 ) { 168 // Active connections = total size minus idle connections 169 let pool_clone = pool.clone(); 170 meter 171 .u64_observable_gauge("db.pool.connections.active") 172 .with_description("Active connections in the pool") 173 .with_callback(move |observer| { 174 let active = pool_clone 175 .size() 176 .saturating_sub(pool_clone.num_idle() as u32); 177 observer.observe(u64::from(active), &[]); 178 }) 179 .build(); 180 181 // Idle connections sitting in the pool waiting for work 182 let pool_clone = pool.clone(); 183 meter 184 .u64_observable_gauge("db.pool.connections.idle") 185 .with_description("Idle connections in the pool") 186 .with_callback(move |observer| { 187 observer.observe(pool_clone.num_idle() as u64, &[]); 188 }) 189 .build(); 190 191 // Utilization = (active / max_size) * 100 192 let pool_clone = pool.clone(); 193 meter 194 .f64_observable_gauge("db.pool.utilization") 195 .with_description("Pool utilization percentage") 196 .with_unit("%") 197 .with_callback(move |observer| { 198 let active = pool_clone 199 .size() 200 .saturating_sub(pool_clone.num_idle() as u32); 201 let max_size = pool_clone.options().get_max_connections(); 202 let utilization = if max_size > 0 { 203 (f64::from(active) / f64::from(max_size)) * 100.0 204 } else { 205 0.0 206 }; 207 observer.observe(utilization, &[]); 208 }) 209 .build(); 210 }