telemetry.rs (10263B)
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_appender_tracing::layer::OpenTelemetryTracingBridge; 12 use opentelemetry_otlp::WithExportConfig; 13 use opentelemetry_sdk::{ 14 logs::SdkLoggerProvider, metrics::SdkMeterProvider, propagation::TraceContextPropagator, 15 trace::SdkTracerProvider, Resource, 16 }; 17 use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt, EnvFilter, Layer}; 18 19 /// Holds the tracer, meter and logger providers, ensuring they are 20 /// shut down gracefully when the application exits. Dropping the guard 21 /// flushes pending telemetry and shuts down each provider in turn, so 22 /// callers only need to keep the value alive for the lifetime of the 23 /// application (e.g. `let _telemetry = telemetry::init_telemetry("inventory");`). 24 pub struct TelemetryGuard { 25 tracer_provider: SdkTracerProvider, 26 meter_provider: SdkMeterProvider, 27 logger_provider: SdkLoggerProvider, 28 } 29 30 impl Drop for TelemetryGuard { 31 fn drop(&mut self) { 32 tracing::info!("Shutting down telemetry..."); 33 34 // Logger first — prevents new log records from being generated 35 // during the shutdown of other providers. 36 if let Err(e) = self.logger_provider.shutdown() { 37 eprintln!("Error shutting down logger provider: {:?}", e); 38 } 39 if let Err(e) = self.meter_provider.shutdown() { 40 eprintln!("Error shutting down meter provider: {:?}", e); 41 } 42 if let Err(e) = self.tracer_provider.shutdown() { 43 eprintln!("Error shutting down tracer provider: {:?}", e); 44 } 45 46 tracing::info!("Telemetry shutdown complete"); 47 } 48 } 49 50 /// Initializes the full telemetry pipeline (tracing + metrics). 51 /// 52 /// # Arguments 53 /// * `service_name` - The name of this service, used to identify traces and metrics 54 /// 55 /// # Returns 56 /// * `TelemetryGuard` - Holds both providers; keep alive for the lifetime of the app 57 /// 58 /// # Panics 59 /// Panics if the OTLP exporters or tracing subscriber cannot be initialized. 60 pub fn init_telemetry(service_name: &str) -> TelemetryGuard { 61 // Set up W3C Trace Context propagator for cross-service trace correlation 62 global::set_text_map_propagator(TraceContextPropagator::new()); 63 64 // Build a shared resource describing this service 65 let resource = Resource::builder() 66 .with_service_name(service_name.to_string()) 67 .with_attributes([ 68 KeyValue::new("service.version", env!("CARGO_PKG_VERSION")), 69 KeyValue::new( 70 "deployment.environment", 71 std::env::var("ENVIRONMENT").unwrap_or_else(|_| "development".into()), 72 ), 73 ]) 74 .build(); 75 76 // Resolve separate endpoints for traces (gRPC) and metrics (HTTP) 77 // Traces → Jaeger (gRPC on port 4317) 78 let traces_endpoint = std::env::var("OTEL_EXPORTER_OTLP_TRACES_ENDPOINT") 79 .or_else(|_| std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT")) 80 .unwrap_or_else(|_| "http://localhost:9090/api/v1/otlp/v1/metrics".into()); 81 // Metrics → Prometheus 3.0 (HTTP on port 9090) 82 let metrics_endpoint = std::env::var("OTEL_EXPORTER_OTLP_METRICS_ENDPOINT") 83 .or_else(|_| std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT")) 84 .unwrap_or_else(|_| "http://localhost:4317".into()); 85 // Logs → Loki (HTTP on port 3100) 86 let logs_endpoint = std::env::var("OTEL_EXPORTER_OTLP_LOGS_ENDPOINT") 87 .unwrap_or_else(|_| "http://localhost:3100/otlp/v1/logs".into()); 88 89 // Initialize all three providers 90 let tracer_provider = init_tracer_provider(&traces_endpoint, &resource); 91 let meter_provider = init_meter_provider(&metrics_endpoint, &resource); 92 let logger_provider = init_logger_provider(&logs_endpoint, &resource); 93 94 // Build the subscriber layers 95 // OpenTelemetry traces layer — converts tracing spans to OTel spans for Jaeger 96 let tracer = tracer_provider.tracer(service_name.to_string()); 97 let otel_trace_layer = tracing_opentelemetry::layer().with_tracer(tracer); 98 // OpenTelemetry logs layer — bridges tracing events to OTel log records for Loki 99 let otel_logs_layer = OpenTelemetryTracingBridge::new(&logger_provider); 100 101 // Environment filter for log levels 102 // Include opentelemetry=off to prevent the SDK's internal diagnostics 103 // from feeding back into the OpenTelemetryTracingBridge (infinite recursion) 104 let env_filter = EnvFilter::try_from_default_env() 105 .unwrap_or_else(|_| EnvFilter::new("info,opentelemetry=off")); 106 107 // Console output layer: JSON in production, compact in development 108 let is_production = std::env::var("ENVIRONMENT") 109 .map(|e| e == "production") 110 .unwrap_or(false); 111 112 let fmt_layer = if is_production { 113 tracing_subscriber::fmt::layer() 114 .json() 115 .with_target(true) 116 .with_thread_ids(false) 117 .with_file(false) 118 .with_line_number(false) 119 .boxed() 120 } else { 121 tracing_subscriber::fmt::layer() 122 .with_target(true) 123 .with_thread_ids(false) 124 .compact() 125 .boxed() 126 }; 127 128 // Combine all layers into a subscriber and set it as the global default 129 // Layer ordering: env_filter at the top filters all downstream layers 130 tracing_subscriber::registry() 131 .with(env_filter) 132 .with(fmt_layer) 133 .with(otel_trace_layer) 134 .with(otel_logs_layer) 135 .init(); 136 137 TelemetryGuard { 138 tracer_provider, 139 meter_provider, 140 logger_provider, 141 } 142 } 143 144 /// Creates the OTLP gRPC trace exporter and tracer provider, 145 /// then wires it into the global tracing subscriber. 146 fn init_tracer_provider(endpoint: &str, resource: &Resource) -> SdkTracerProvider { 147 // Create the OTLP exporter configured to send traces via gRPC 148 let exporter = opentelemetry_otlp::SpanExporter::builder() 149 .with_tonic() 150 .with_endpoint(endpoint) 151 .build() 152 .expect("Failed to create OTLP span exporter"); 153 154 // Build the tracer provider with the exporter and service resource 155 let provider = SdkTracerProvider::builder() 156 .with_batch_exporter(exporter) 157 .with_resource(resource.clone()) 158 .build(); 159 160 // NOTE: the global subscriber is assembled once in `init_telemetry`, which 161 // combines the trace layer with the logs bridge. Setting it here too would 162 // panic with "a global default trace dispatcher has already been set". 163 provider 164 } 165 166 /// Creates the OTLP HTTP metric exporter and meter provider, 167 /// then registers it as the global meter provider. 168 fn init_meter_provider(endpoint: &str, resource: &Resource) -> SdkMeterProvider { 169 // Create the OTLP HTTP exporter targeting Prometheus OTLP receiver 170 let exporter = opentelemetry_otlp::MetricExporter::builder() 171 .with_http() 172 .with_endpoint(endpoint) 173 .build() 174 .expect("Failed to create OTLP metric exporter"); 175 176 // Build the meter provider with a periodic exporter 177 let provider = SdkMeterProvider::builder() 178 .with_resource(resource.clone()) 179 .with_periodic_exporter(exporter) 180 .build(); 181 182 // Register as the global meter provider 183 global::set_meter_provider(provider.clone()); 184 provider 185 } 186 187 /// Initialize the logger provider with OTLP HTTP exporter. 188 /// Exports log records to Loki via its native OTLP endpoint. 189 fn init_logger_provider(endpoint: &str, resource: &Resource) -> SdkLoggerProvider { 190 let exporter = opentelemetry_otlp::LogExporter::builder() 191 .with_http() 192 .with_endpoint(endpoint) 193 .build() 194 .expect("Failed to create OTLP log exporter"); 195 196 SdkLoggerProvider::builder() 197 .with_resource(resource.clone()) 198 .with_batch_exporter(exporter) 199 .build() 200 } 201 202 /// Registers observable gauges that track database connection pool health. 203 /// 204 /// Creates three instruments: 205 /// - `db.pool.connections.active` — connections currently in use 206 /// - `db.pool.connections.idle` — connections waiting for work 207 /// - `db.pool.utilization` — percentage of pool capacity in use 208 /// 209 /// # Arguments 210 /// * `meter` - The meter to register gauges on 211 /// * `pool` - The sqlx connection pool to observe 212 pub fn register_pool_metrics( 213 meter: &opentelemetry::metrics::Meter, 214 pool: sqlx::Pool<sqlx::Postgres>, 215 ) { 216 // Active connections = total size minus idle connections 217 let pool_clone = pool.clone(); 218 meter 219 .u64_observable_gauge("db.pool.connections.active") 220 .with_description("Active connections in the pool") 221 .with_callback(move |observer| { 222 let active = pool_clone 223 .size() 224 .saturating_sub(pool_clone.num_idle() as u32); 225 observer.observe(u64::from(active), &[]); 226 }) 227 .build(); 228 229 // Idle connections sitting in the pool waiting for work 230 let pool_clone = pool.clone(); 231 meter 232 .u64_observable_gauge("db.pool.connections.idle") 233 .with_description("Idle connections in the pool") 234 .with_callback(move |observer| { 235 observer.observe(pool_clone.num_idle() as u64, &[]); 236 }) 237 .build(); 238 239 // Utilization = (active / max_size) * 100 240 let pool_clone = pool.clone(); 241 meter 242 .f64_observable_gauge("db.pool.utilization") 243 .with_description("Pool utilization percentage") 244 .with_unit("%") 245 .with_callback(move |observer| { 246 let active = pool_clone 247 .size() 248 .saturating_sub(pool_clone.num_idle() as u32); 249 let max_size = pool_clone.options().get_max_connections(); 250 let utilization = if max_size > 0 { 251 (f64::from(active) / f64::from(max_size)) * 100.0 252 } else { 253 0.0 254 }; 255 observer.observe(utilization, &[]); 256 }) 257 .build(); 258 }