exercises

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

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 }