Skip to main content

slim_tracing/
native.rs

1use opentelemetry::{KeyValue, global, trace::TracerProvider as _};
2use opentelemetry_otlp::{ExporterBuildError, WithExportConfig};
3use opentelemetry_sdk::{
4    Resource,
5    metrics::{MeterProviderBuilder, PeriodicReader, SdkMeterProvider},
6    trace::{RandomIdGenerator, Sampler, SdkTracerProvider},
7};
8use opentelemetry_semantic_conventions::attribute::{
9    DEPLOYMENT_ENVIRONMENT_NAME, SERVICE_NAME, SERVICE_VERSION,
10};
11use serde::Deserialize;
12use thiserror::Error;
13use tracing::Level;
14use tracing_opentelemetry::{MetricsLayer, OpenTelemetryLayer};
15use tracing_subscriber::{EnvFilter, Layer, fmt, layer::SubscriberExt, util::SubscriberInitExt};
16
17use slim_config::{
18    client::ClientConfig, errors::ConfigError as SlimConfigError, tls::client::TlsClientConfig,
19};
20
21const OTEL_EXPORTER_OTLP_ENDPOINT: &str = "http://localhost:4317";
22
23#[derive(Error, Debug)]
24pub enum ConfigError {
25    // gRPC / remote configuration
26    #[error("error loading GRPC config")]
27    GRPCError(#[from] SlimConfigError),
28
29    #[error("error building exporter")]
30    OpenTelemetryInitError(#[from] ExporterBuildError),
31
32    // Filter parsing / directives
33    #[error("error parsing filter directives")]
34    FilterParseError(#[from] tracing_subscriber::filter::ParseError),
35
36    // Tracing subscriber initialization
37    #[error("error setting up tracing subscriber")]
38    TracingSetupError(#[from] tracing_subscriber::util::TryInitError),
39}
40
41#[derive(Clone, Debug, Deserialize)]
42#[serde(deny_unknown_fields)]
43pub struct TracingConfiguration {
44    #[serde(default = "default_log_level")]
45    log_level: String,
46
47    #[serde(default = "default_display_thread_names")]
48    display_thread_names: bool,
49
50    #[serde(default = "default_display_thread_ids")]
51    display_thread_ids: bool,
52
53    #[serde(default = "default_filter")]
54    filters: Vec<String>,
55
56    #[serde(default)]
57    opentelemetry: OpenTelemetryConfig,
58}
59
60// default implementation for TracingConfiguration
61impl Default for TracingConfiguration {
62    fn default() -> Self {
63        TracingConfiguration {
64            log_level: default_log_level(),
65            display_thread_names: default_display_thread_names(),
66            display_thread_ids: default_display_thread_ids(),
67            filters: default_filter(),
68            opentelemetry: OpenTelemetryConfig::default(),
69        }
70    }
71}
72
73#[derive(Clone, Debug, Deserialize)]
74#[serde(deny_unknown_fields)]
75pub struct OpenTelemetryConfig {
76    #[serde(default)]
77    enabled: bool,
78
79    #[serde(default)]
80    grpc: ClientConfig,
81
82    #[serde(default = "default_service_name")]
83    service_name: String,
84
85    #[serde(default = "default_service_version")]
86    service_version: String,
87
88    #[serde(default = "default_environment")]
89    environment: String,
90
91    #[serde(default = "default_metrics_interval")]
92    metrics_interval_secs: u64,
93}
94
95impl OpenTelemetryConfig {
96    /// Sets whether OpenTelemetry tracing and metrics are enabled.
97    ///
98    /// # Arguments
99    ///
100    /// * `enabled` - A boolean indicating whether OpenTelemetry should be enabled
101    ///
102    /// # Returns
103    ///
104    /// Returns `self` for method chaining
105    pub fn with_enabled(mut self, enabled: bool) -> Self {
106        self.enabled = enabled;
107        self
108    }
109
110    /// Sets the gRPC configuration for OpenTelemetry export.
111    ///
112    /// # Arguments
113    ///
114    /// * `grpc_config` - The gRPC client configuration to use for OpenTelemetry export
115    ///
116    /// # Returns
117    ///
118    /// Returns `self` for method chaining
119    pub fn with_grpc_config(mut self, grpc_config: ClientConfig) -> Self {
120        self.grpc = grpc_config;
121        self
122    }
123
124    /// Sets the service name for OpenTelemetry traces and metrics.
125    ///
126    /// # Arguments
127    ///
128    /// * `service_name` - The name of the service to be used in OpenTelemetry data
129    ///
130    /// # Returns
131    ///
132    /// Returns `self` for method chaining
133    pub fn with_service_name(mut self, service_name: String) -> Self {
134        self.service_name = service_name;
135        self
136    }
137
138    /// Sets the service version for OpenTelemetry traces and metrics.
139    ///
140    /// # Arguments
141    ///
142    /// * `service_version` - The version of the service to be used in OpenTelemetry data
143    ///
144    /// # Returns
145    ///
146    /// Returns `self` for method chaining
147    pub fn with_service_version(mut self, service_version: String) -> Self {
148        self.service_version = service_version;
149        self
150    }
151
152    /// Sets the deployment environment for OpenTelemetry traces and metrics.
153    ///
154    /// # Arguments
155    ///
156    /// * `environment` - The deployment environment (e.g., "development", "production")
157    ///
158    /// # Returns
159    ///
160    /// Returns `self` for method chaining
161    pub fn with_environment(mut self, environment: String) -> Self {
162        self.environment = environment;
163        self
164    }
165
166    /// Sets the interval in seconds between metric exports.
167    ///
168    /// # Arguments
169    ///
170    /// * `metrics_interval_secs` - The interval in seconds between metric exports
171    ///
172    /// # Returns
173    ///
174    /// Returns `self` for method chaining
175    pub fn with_metrics_interval_secs(mut self, metrics_interval_secs: u64) -> Self {
176        self.metrics_interval_secs = metrics_interval_secs;
177        self
178    }
179
180    /// Returns whether OpenTelemetry tracing and metrics are enabled.
181    ///
182    /// # Returns
183    ///
184    /// Returns a boolean indicating whether OpenTelemetry is enabled
185    pub fn enabled(&self) -> bool {
186        self.enabled
187    }
188
189    /// Returns the gRPC configuration for OpenTelemetry export.
190    ///
191    /// # Returns
192    ///
193    /// Returns a reference to the gRPC client configuration
194    pub fn grpc_config(&self) -> &ClientConfig {
195        &self.grpc
196    }
197
198    /// Returns the service name used in OpenTelemetry data.
199    ///
200    /// # Returns
201    ///
202    /// Returns a reference to the service name string
203    pub fn service_name(&self) -> &str {
204        &self.service_name
205    }
206
207    /// Returns the service version used in OpenTelemetry data.
208    ///
209    /// # Returns
210    ///
211    /// Returns a reference to the service version string
212    pub fn service_version(&self) -> &str {
213        &self.service_version
214    }
215
216    /// Returns the deployment environment used in OpenTelemetry data.
217    ///
218    /// # Returns
219    ///
220    /// Returns a reference to the environment string
221    pub fn environment(&self) -> &str {
222        &self.environment
223    }
224
225    /// Returns the interval in seconds between metric exports.
226    ///
227    /// # Returns
228    ///
229    /// Returns the metrics interval in seconds
230    pub fn metrics_interval_secs(&self) -> u64 {
231        self.metrics_interval_secs
232    }
233}
234
235// default implementation for OpenTelemetryConfig
236impl Default for OpenTelemetryConfig {
237    fn default() -> Self {
238        OpenTelemetryConfig {
239            enabled: false,
240            grpc: ClientConfig::with_endpoint(OTEL_EXPORTER_OTLP_ENDPOINT)
241                .with_tls_setting(TlsClientConfig::new().with_insecure(true)),
242            service_name: default_service_name(),
243            service_version: default_service_version(),
244            environment: default_environment(),
245            metrics_interval_secs: default_metrics_interval(),
246        }
247    }
248}
249
250fn default_log_level() -> String {
251    "info".to_string()
252}
253
254fn default_display_thread_names() -> bool {
255    true
256}
257
258fn default_display_thread_ids() -> bool {
259    false
260}
261
262fn default_filter() -> Vec<String> {
263    // Only module names here. Their effective level will be the configured `log_level`.
264    vec![
265        "slim_datapath".to_string(),
266        "slim_service".to_string(),
267        "slim_controller".to_string(),
268        "slim_auth".to_string(),
269        "slim_config".to_string(),
270        "slim_mls".to_string(),
271        "slim_session".to_string(),
272        "slim_signal".to_string(),
273        "slim_tracing".to_string(),
274        "_slim_bindings".to_string(),
275        "slim_testing".to_string(),
276        "slim".to_string(),
277        "slim_examples".to_string(),
278        "sdk_mock".to_string(),
279        "client".to_string(),
280    ]
281}
282
283fn default_service_name() -> String {
284    "slim-data-plane".to_string()
285}
286
287fn default_service_version() -> String {
288    "v0.1.0".to_string()
289}
290
291fn default_environment() -> String {
292    "development".to_string()
293}
294
295fn default_metrics_interval() -> u64 {
296    30 // default to 30 seconds
297}
298
299// function to convert string tracing level to tracing::Level
300fn resolve_level(level: &str) -> tracing::Level {
301    let level = level.to_lowercase();
302    match level.as_str() {
303        "trace" => Level::TRACE,
304        "debug" => Level::DEBUG,
305        "info" => Level::INFO,
306        "warn" => Level::WARN,
307        "error" => Level::ERROR,
308        _ => Level::INFO, // default level
309    }
310}
311
312pub struct OtelGuard {
313    tracer_provider: Option<SdkTracerProvider>,
314    meter_provider: Option<SdkMeterProvider>,
315}
316
317impl Drop for OtelGuard {
318    fn drop(&mut self) {
319        if let Some(tracer) = self.tracer_provider.take()
320            && let Err(err) = tracer.shutdown()
321        {
322            eprintln!("Error shutting down tracer provider: {err:?}");
323        }
324
325        if let Some(meter) = self.meter_provider.take()
326            && let Err(err) = meter.shutdown()
327        {
328            eprintln!("Error shutting down meter provider: {err:?}");
329        }
330    }
331}
332
333impl TracingConfiguration {
334    pub fn with_log_level(self, log_level: String) -> Self {
335        TracingConfiguration { log_level, ..self }
336    }
337
338    pub fn with_display_thread_names(self, display_thread_names: bool) -> Self {
339        TracingConfiguration {
340            display_thread_names,
341            ..self
342        }
343    }
344
345    pub fn with_display_thread_ids(self, display_thread_ids: bool) -> Self {
346        TracingConfiguration {
347            display_thread_ids,
348            ..self
349        }
350    }
351
352    pub fn with_filter(self, filter: Vec<String>) -> Self {
353        TracingConfiguration {
354            filters: filter,
355            ..self
356        }
357    }
358
359    pub fn with_opentelemetry_config(mut self, config: OpenTelemetryConfig) -> Self {
360        self.opentelemetry = config;
361        self
362    }
363
364    pub fn enable_opentelemetry(mut self) -> Self {
365        self.opentelemetry.enabled = true;
366        self
367    }
368
369    pub fn with_metrics_interval(mut self, interval_secs: u64) -> Self {
370        self.opentelemetry.metrics_interval_secs = interval_secs;
371        self
372    }
373
374    pub fn log_level(&self) -> &str {
375        &self.log_level
376    }
377
378    pub fn display_thread_names(&self) -> bool {
379        self.display_thread_names
380    }
381
382    pub fn display_thread_ids(&self) -> bool {
383        self.display_thread_ids
384    }
385
386    pub fn filter(&self) -> &Vec<String> {
387        &self.filters
388    }
389
390    /// Set up a subscriber
391    pub fn setup_tracing_subscriber(&self) -> Result<OtelGuard, ConfigError> {
392        let fmt_layer = fmt::layer()
393            .with_thread_ids(self.display_thread_ids)
394            .with_thread_names(self.display_thread_names)
395            .with_line_number(true)
396            .with_filter(tracing_subscriber::filter::filter_fn(
397                |metadata: &tracing::Metadata| {
398                    !metadata
399                        .fields()
400                        .iter()
401                        .any(|field| field.name() == "telemetry")
402                },
403            ));
404
405        // Build the EnvFilter with correct precedence:
406        // 1. Environment variable (RUST_LOG) overrides everything (both modules & levels)
407        // 2. User-provided filter (if differs from the default) overrides default (modules & levels)
408        // 3. Default filter modules use the configured `log_level`
409        //
410        // Additionally, environment variable has highest priority.
411        let level_filter = if let Ok(env_value) = std::env::var("RUST_LOG") {
412            // Highest priority: environment.
413            // If env_value has no global directive (a bare level), and consists only of module=level
414            // directives, then append a global "off" so that unspecified modules are silenced.
415            // Examples:
416            //   slim=debug                  -> slim=debug,off
417            //   slim=debug,slim_auth=trace  -> slim=debug,slim_auth=trace,off
418            //   debug                       -> debug            (keep global)
419            //   info,slim=debug             -> info,slim=debug  (keep global)
420            let needs_global_off = {
421                let tokens: Vec<&str> = env_value
422                    .split(',')
423                    .map(|t| t.trim())
424                    .filter(|t| !t.is_empty())
425                    .collect();
426                let bare_level_present = tokens
427                    .iter()
428                    .any(|t| matches!(*t, "trace" | "debug" | "info" | "warn" | "error" | "off"));
429                !bare_level_present
430            };
431            let augmented = if needs_global_off {
432                format!("{env_value},off")
433            } else {
434                env_value
435            };
436            EnvFilter::new(augmented)
437        } else {
438            let is_default = self.filters == default_filter();
439
440            // Always set a fallback directive using the configured log_level.
441            let builder =
442                EnvFilter::builder().with_default_directive(resolve_level(self.log_level()).into());
443
444            let filter_string = if is_default {
445                // Apply the configured log_level to each default module.
446                self.filters
447                    .iter()
448                    .map(|m| {
449                        // In case a module accidentally already contains a level (e.g. "foo=debug"),
450                        // keep only the part before '=' to enforce overriding with `log_level`.
451                        let module = m.split('=').next().unwrap_or(m);
452                        format!("{module}={}", self.log_level())
453                    })
454                    .collect::<Vec<_>>()
455                    .join(",")
456            } else {
457                // Custom filter provided: treat entries as authoritative.
458                // They may include levels (module=level) or just modules.
459                // For entries without explicit level, append the configured log_level.
460                self.filters
461                    .iter()
462                    .map(|d| {
463                        if d.contains('=') {
464                            d.clone()
465                        } else {
466                            format!("{d}={}", self.log_level())
467                        }
468                    })
469                    .collect::<Vec<_>>()
470                    .join(",")
471            };
472
473            builder.parse_lossy(filter_string)
474        };
475
476        if self.opentelemetry.enabled {
477            // TODO(msardara): derive a tonic channel directly when opentelemetry-otlp
478            // upgrades to tonic version 0.13.0
479            let endpoint = self.opentelemetry.grpc.endpoint.clone();
480
481            // resource
482            let resource = Resource::builder()
483                .with_attributes([
484                    KeyValue::new(SERVICE_NAME, self.opentelemetry.service_name.clone()),
485                    KeyValue::new(SERVICE_VERSION, self.opentelemetry.service_version.clone()),
486                    KeyValue::new(
487                        DEPLOYMENT_ENVIRONMENT_NAME,
488                        self.opentelemetry.environment.clone(),
489                    ),
490                ])
491                .build();
492
493            // init tracer provider
494            let exporter = opentelemetry_otlp::SpanExporter::builder()
495                .with_tonic()
496                .with_endpoint(&endpoint)
497                .build()?;
498
499            let tracer_provider = SdkTracerProvider::builder()
500                // TODO(zkacsand): customize sampling strategy
501                .with_sampler(Sampler::ParentBased(Box::new(Sampler::TraceIdRatioBased(
502                    1.0,
503                ))))
504                .with_id_generator(RandomIdGenerator::default())
505                .with_resource(resource.clone())
506                .with_batch_exporter(exporter)
507                .build();
508
509            let exporter = opentelemetry_otlp::MetricExporter::builder()
510                .with_tonic()
511                .with_endpoint(&endpoint)
512                .with_temporality(opentelemetry_sdk::metrics::Temporality::default())
513                .build()?;
514
515            let reader = PeriodicReader::builder(exporter)
516                .with_interval(std::time::Duration::from_secs(
517                    self.opentelemetry.metrics_interval_secs,
518                ))
519                .build();
520
521            let stdout_reader =
522                PeriodicReader::builder(opentelemetry_stdout::MetricExporter::default()).build();
523
524            let meter_provider = MeterProviderBuilder::default()
525                .with_resource(resource.clone())
526                .with_reader(reader)
527                .with_reader(stdout_reader)
528                .build();
529
530            // set global meter provider
531            global::set_meter_provider(meter_provider.clone());
532
533            // Sst up the trace context propagator
534            let propagator = opentelemetry_sdk::propagation::TraceContextPropagator::new();
535            global::set_text_map_propagator(propagator);
536
537            let tracer = tracer_provider.tracer("tracing-otel-subscriber");
538
539            // Construct the subscriber with OpenTelemetry
540            tracing_subscriber::registry()
541                .with(level_filter)
542                .with(fmt_layer)
543                .with(MetricsLayer::new(meter_provider.clone()))
544                .with(OpenTelemetryLayer::new(tracer))
545                .try_init()?;
546
547            Ok(OtelGuard {
548                tracer_provider: Some(tracer_provider),
549                meter_provider: Some(meter_provider),
550            })
551        } else {
552            // Basic subscriber without OpenTelemetry
553            tracing_subscriber::registry()
554                .with(level_filter)
555                .with(fmt_layer)
556                .try_init()?;
557
558            Ok(OtelGuard {
559                tracer_provider: None,
560                meter_provider: None,
561            })
562        }
563    }
564}
565
566// tests
567#[cfg(test)]
568mod tests {
569    use super::*;
570
571    #[test]
572    fn test_default_tracing_configuration() {
573        let config = TracingConfiguration::default();
574        assert_eq!(config.log_level, default_log_level());
575        assert_eq!(config.display_thread_names, default_display_thread_names());
576        assert_eq!(config.display_thread_ids, default_display_thread_ids());
577        assert_eq!(config.filters, default_filter());
578    }
579
580    #[test]
581    fn test_resolve_level() {
582        assert_eq!(resolve_level("trace"), Level::TRACE);
583        assert_eq!(resolve_level("debug"), Level::DEBUG);
584        assert_eq!(resolve_level("info"), Level::INFO);
585        assert_eq!(resolve_level("warn"), Level::WARN);
586        assert_eq!(resolve_level("error"), Level::ERROR);
587        assert_eq!(resolve_level("invalid"), Level::INFO);
588    }
589
590    #[test]
591    fn test_tracing_configuration_builder_methods() {
592        let config = TracingConfiguration::default()
593            .with_log_level("debug".to_string())
594            .with_display_thread_names(false)
595            .with_display_thread_ids(true)
596            .with_filter(vec!["debug".to_string()]);
597
598        assert_eq!(config.log_level(), "debug");
599        assert!(!config.display_thread_names());
600        assert!(config.display_thread_ids());
601        assert_eq!(config.filter(), &vec!["debug".to_string()]);
602    }
603
604    #[test]
605    fn test_opentelemetry_config_default() {
606        let config = OpenTelemetryConfig::default();
607        assert!(!config.enabled());
608        assert_eq!(config.service_name(), default_service_name());
609        assert_eq!(config.grpc_config().endpoint, OTEL_EXPORTER_OTLP_ENDPOINT);
610        assert_eq!(config.service_version(), default_service_version());
611        assert_eq!(config.environment(), default_environment());
612        assert_eq!(config.metrics_interval_secs(), default_metrics_interval());
613    }
614
615    #[test]
616    fn test_tracing_configuration_with_opentelemetry() {
617        let otel_config = OpenTelemetryConfig::default()
618            .with_enabled(true)
619            .with_service_name("test-service".to_string())
620            .with_service_version("1.0.0".to_string());
621
622        let config = TracingConfiguration::default().with_opentelemetry_config(otel_config);
623
624        assert!(config.opentelemetry.enabled());
625        assert_eq!(config.opentelemetry.service_name(), "test-service");
626        assert_eq!(config.opentelemetry.service_version(), "1.0.0");
627    }
628
629    #[test]
630    fn test_enable_opentelemetry() {
631        let config = TracingConfiguration::default().enable_opentelemetry();
632        assert!(config.opentelemetry.enabled());
633    }
634
635    #[test]
636    fn test_with_metrics_interval() {
637        let config = TracingConfiguration::default().with_metrics_interval(60);
638        assert_eq!(config.opentelemetry.metrics_interval_secs(), 60);
639    }
640
641    #[test]
642    fn test_otel_guard_drop() {
643        // This test verifies that OtelGuard can be created and dropped without panicking
644        let config = TracingConfiguration::default();
645        let guard = config.setup_tracing_subscriber().unwrap();
646        drop(guard); // Should not panic
647    }
648}