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 #[error("error loading GRPC config")]
27 GRPCError(#[from] SlimConfigError),
28
29 #[error("error building exporter")]
30 OpenTelemetryInitError(#[from] ExporterBuildError),
31
32 #[error("error parsing filter directives")]
34 FilterParseError(#[from] tracing_subscriber::filter::ParseError),
35
36 #[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
60impl 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 pub fn with_enabled(mut self, enabled: bool) -> Self {
106 self.enabled = enabled;
107 self
108 }
109
110 pub fn with_grpc_config(mut self, grpc_config: ClientConfig) -> Self {
120 self.grpc = grpc_config;
121 self
122 }
123
124 pub fn with_service_name(mut self, service_name: String) -> Self {
134 self.service_name = service_name;
135 self
136 }
137
138 pub fn with_service_version(mut self, service_version: String) -> Self {
148 self.service_version = service_version;
149 self
150 }
151
152 pub fn with_environment(mut self, environment: String) -> Self {
162 self.environment = environment;
163 self
164 }
165
166 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 pub fn enabled(&self) -> bool {
186 self.enabled
187 }
188
189 pub fn grpc_config(&self) -> &ClientConfig {
195 &self.grpc
196 }
197
198 pub fn service_name(&self) -> &str {
204 &self.service_name
205 }
206
207 pub fn service_version(&self) -> &str {
213 &self.service_version
214 }
215
216 pub fn environment(&self) -> &str {
222 &self.environment
223 }
224
225 pub fn metrics_interval_secs(&self) -> u64 {
231 self.metrics_interval_secs
232 }
233}
234
235impl 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 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 }
298
299fn 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, }
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 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 let level_filter = if let Ok(env_value) = std::env::var("RUST_LOG") {
412 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 let builder =
442 EnvFilter::builder().with_default_directive(resolve_level(self.log_level()).into());
443
444 let filter_string = if is_default {
445 self.filters
447 .iter()
448 .map(|m| {
449 let module = m.split('=').next().unwrap_or(m);
452 format!("{module}={}", self.log_level())
453 })
454 .collect::<Vec<_>>()
455 .join(",")
456 } else {
457 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 let endpoint = self.opentelemetry.grpc.endpoint.clone();
480
481 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 let exporter = opentelemetry_otlp::SpanExporter::builder()
495 .with_tonic()
496 .with_endpoint(&endpoint)
497 .build()?;
498
499 let tracer_provider = SdkTracerProvider::builder()
500 .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 global::set_meter_provider(meter_provider.clone());
532
533 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 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 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#[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 let config = TracingConfiguration::default();
645 let guard = config.setup_tracing_subscriber().unwrap();
646 drop(guard); }
648}