Skip to main content

kmp_observability/
lib.rs

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