Skip to main content

kmp_observability/
lib.rs

1mod buffered_quality_metrics_observer;
2mod embedded_telemetry_guard;
3#[cfg(feature = "otel")]
4pub mod metrics;
5mod otlp_query_adapter;
6#[cfg(feature = "otel")]
7pub mod quality_observers;
8mod quality_telemetry_observation;
9
10#[cfg(feature = "otel")]
11use opentelemetry::trace::TracerProvider as _;
12#[cfg(feature = "otel")]
13use opentelemetry_otlp::WithExportConfig;
14#[cfg(feature = "otel")]
15use opentelemetry_otlp::WithTonicConfig;
16#[cfg(feature = "otel")]
17use opentelemetry_otlp::tonic_types::transport::{Certificate, ClientTlsConfig, Identity};
18#[cfg(feature = "otel")]
19use opentelemetry_sdk::metrics::SdkMeterProvider;
20#[cfg(feature = "otel")]
21use opentelemetry_sdk::trace::SdkTracerProvider;
22#[cfg(feature = "otel")]
23use tracing_subscriber::{EnvFilter, fmt, layer::SubscriberExt, util::SubscriberInitExt};
24
25pub use buffered_quality_metrics_observer::BufferedQualityMetricsObserver;
26pub use embedded_telemetry_guard::EmbeddedTelemetryGuard;
27#[cfg(feature = "otel")]
28pub use metrics::KernelMetrics;
29pub use otlp_query_adapter::{
30    OtlpMetricsQueryClient, OtlpObservabilityQueryAdapter, OtlpQueryResponse,
31};
32pub use quality_telemetry_observation::QualityTelemetryObservation;
33
34/// Resources returned by `init_observability` for lifecycle management.
35#[cfg(feature = "otel")]
36pub struct ObservabilityGuard {
37    trace_provider: Option<SdkTracerProvider>,
38    meter_provider: Option<SdkMeterProvider>,
39    pub metrics: KernelMetrics,
40}
41
42/// Initializes structured logging with optional OpenTelemetry trace and metric export.
43///
44/// ## Environment variables
45///
46/// - `RUST_LOG`: log level filter (default: `info`)
47/// - `KMP_LOG_FORMAT`: `json` | `pretty` | compact (default)
48/// - `OTEL_EXPORTER_OTLP_ENDPOINT`: OTLP endpoint for metric export and, when
49///   enabled, trace export (e.g. `http://localhost:4317`)
50/// - `OTEL_TRACES_EXPORTER`: standard OTel traces exporter selector. When set
51///   to `none`, trace export is disabled even if OTLP metrics remain enabled.
52///   When unset, traces follow the OTLP endpoint configuration.
53#[cfg(feature = "otel")]
54pub fn init_observability(service_name: &str) -> ObservabilityGuard {
55    let env_filter = EnvFilter::try_from_default_env().unwrap_or_else(|_| EnvFilter::new("info"));
56    let log_format = std::env::var("KMP_LOG_FORMAT").unwrap_or_default();
57
58    let trace_provider = init_otel_tracer(service_name);
59    let otel_layer = trace_provider.as_ref().map(|provider| {
60        let tracer = provider.tracer(service_name.to_string());
61        tracing_opentelemetry::layer().with_tracer(tracer)
62    });
63
64    let meter_provider = metrics::init_otel_metrics(service_name);
65    if let Some(ref provider) = meter_provider {
66        opentelemetry::global::set_meter_provider(provider.clone());
67    }
68    let meter = opentelemetry::global::meter("kmp");
69    let kernel_metrics = KernelMetrics::new(&meter);
70
71    match log_format.as_str() {
72        "json" => {
73            tracing_subscriber::registry()
74                .with(env_filter)
75                .with(otel_layer)
76                .with(
77                    fmt::layer()
78                        .json()
79                        .with_target(true)
80                        .with_thread_ids(false)
81                        .with_file(false)
82                        .with_line_number(false),
83                )
84                .init();
85        }
86        "pretty" => {
87            tracing_subscriber::registry()
88                .with(env_filter)
89                .with(otel_layer)
90                .with(fmt::layer().pretty())
91                .init();
92        }
93        _ => {
94            tracing_subscriber::registry()
95                .with(env_filter)
96                .with(otel_layer)
97                .with(fmt::layer().compact())
98                .init();
99        }
100    }
101
102    tracing::info!(service = service_name, "observability initialized");
103
104    ObservabilityGuard {
105        trace_provider,
106        meter_provider,
107        metrics: kernel_metrics,
108    }
109}
110
111#[cfg(feature = "otel")]
112fn init_otel_tracer(service_name: &str) -> Option<SdkTracerProvider> {
113    let endpoint = std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT").ok()?;
114    if !traces_export_enabled(
115        std::env::var("OTEL_TRACES_EXPORTER").ok().as_deref(),
116        Some(endpoint.as_str()),
117    ) {
118        return None;
119    }
120    let endpoint = endpoint.trim().to_string();
121
122    let mut builder = opentelemetry_otlp::SpanExporter::builder()
123        .with_tonic()
124        .with_endpoint(endpoint);
125    if let Some(tls_config) = build_otlp_tls_config() {
126        builder = builder.with_tls_config(tls_config);
127    }
128    let exporter = builder.build().ok()?;
129
130    let provider = SdkTracerProvider::builder()
131        .with_resource(
132            opentelemetry_sdk::Resource::builder()
133                .with_service_name(service_name.to_string())
134                .build(),
135        )
136        .with_batch_exporter(exporter)
137        .build();
138
139    Some(provider)
140}
141
142#[cfg(feature = "otel")]
143fn traces_export_enabled(traces_exporter: Option<&str>, endpoint: Option<&str>) -> bool {
144    let Some(endpoint) = endpoint.map(str::trim).filter(|value| !value.is_empty()) else {
145        return false;
146    };
147    let _ = endpoint;
148
149    !matches!(
150        traces_exporter
151            .map(str::trim)
152            .filter(|value| !value.is_empty())
153            .map(|value| value.to_ascii_lowercase()),
154        Some(value) if value == "none"
155    )
156}
157
158/// Build a TLS config for OTLP exporters from environment variables.
159///
160/// - `OTEL_EXPORTER_OTLP_CA_PATH` — CA certificate for server verification
161/// - `OTEL_EXPORTER_OTLP_CERT_PATH` — client certificate for mTLS
162/// - `OTEL_EXPORTER_OTLP_KEY_PATH` — client key for mTLS
163///
164/// Returns `None` if no TLS variables are set (plaintext mode).
165#[cfg(feature = "otel")]
166pub(crate) fn build_otlp_tls_config() -> Option<ClientTlsConfig> {
167    let ca_path = std::env::var("OTEL_EXPORTER_OTLP_CA_PATH").ok();
168    let cert_path = std::env::var("OTEL_EXPORTER_OTLP_CERT_PATH").ok();
169    let key_path = std::env::var("OTEL_EXPORTER_OTLP_KEY_PATH").ok();
170
171    if ca_path.is_none() && cert_path.is_none() && key_path.is_none() {
172        return None;
173    }
174
175    let mut tls_config = ClientTlsConfig::new();
176
177    if let Some(ca_path) = ca_path {
178        let ca_pem = std::fs::read(ca_path.trim()).ok()?;
179        tls_config = tls_config.ca_certificate(Certificate::from_pem(ca_pem));
180    }
181
182    if let (Some(cert_path), Some(key_path)) = (cert_path, key_path) {
183        let cert_pem = std::fs::read(cert_path.trim()).ok()?;
184        let key_pem = std::fs::read(key_path.trim()).ok()?;
185        tls_config = tls_config.identity(Identity::from_pem(cert_pem, key_pem));
186    }
187
188    Some(tls_config)
189}
190
191/// Shuts down the OpenTelemetry providers, flushing pending data.
192#[cfg(feature = "otel")]
193pub fn shutdown_observability(guard: ObservabilityGuard) {
194    if let Some(provider) = guard.trace_provider
195        && let Err(error) = provider.shutdown()
196    {
197        tracing::warn!(%error, "opentelemetry trace shutdown failed");
198    }
199    if let Some(provider) = guard.meter_provider
200        && let Err(error) = provider.shutdown()
201    {
202        tracing::warn!(%error, "opentelemetry metrics shutdown failed");
203    }
204}
205
206#[cfg(all(test, feature = "otel"))]
207mod tests {
208    use super::*;
209
210    #[test]
211    fn kernel_metrics_instruments_are_constructible() {
212        let meter = opentelemetry::global::meter("test");
213        let metrics = KernelMetrics::new(&meter);
214        // Verify instruments exist and can record without panic
215        metrics.rpc_duration.record(0.1, &[]);
216        metrics.bundle_nodes.record(5, &[]);
217        metrics.bundle_relationships.record(3, &[]);
218        metrics.bundle_details.record(2, &[]);
219        metrics.rendered_tokens.record(100, &[]);
220        metrics.truncation_total.add(1, &[]);
221        metrics.projection_lag.record(0.05, &[]);
222    }
223
224    #[test]
225    fn traces_export_is_disabled_when_exporter_is_none() {
226        assert!(!traces_export_enabled(
227            Some("none"),
228            Some("https://collector:4317")
229        ));
230    }
231
232    #[test]
233    fn traces_export_is_disabled_without_endpoint() {
234        assert!(!traces_export_enabled(Some("otlp"), None));
235        assert!(!traces_export_enabled(None, Some("   ")));
236    }
237
238    #[test]
239    fn traces_export_defaults_to_enabled_when_endpoint_exists() {
240        assert!(traces_export_enabled(None, Some("https://collector:4317")));
241        assert!(traces_export_enabled(
242            Some("otlp"),
243            Some("https://collector:4317")
244        ));
245    }
246}