exercises

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

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 }