1#[cfg(not(any(feature = "grpc", feature = "http")))]
27compile_error!("at least one transport feature must be enabled: `grpc` or `http`");
28
29#[cfg(feature = "testing")]
30pub mod testing;
31
32#[cfg(feature = "axum")]
33pub mod axum_middleware;
34
35#[cfg(feature = "tonic-tracing")]
36pub mod grpc_middleware;
37
38#[cfg(feature = "profiling")]
39pub mod profiling;
40mod runtime_metrics;
41
42pub mod boot;
43pub mod export_backoff;
44pub mod instrumented_port;
45#[cfg(feature = "grpc-mtls")]
46mod rotating_mtls;
47#[cfg(feature = "grpc-mtls")]
48pub use rotating_mtls::{CertSource, StaticCertSource};
49pub mod log_bridge;
50pub mod span_enrichment;
51pub mod spanned;
52
53pub use instrumented_port::{Instrumented, InstrumentedArc};
54pub use log_bridge::{
55 PROPAGATED_SPAN_FIELDS, SpanLogAttrs, record_span_log_attr, record_span_log_attr_on,
56};
57pub use spanned::{Spanned, in_span};
58
59use opentelemetry::KeyValue;
60use opentelemetry::propagation::TextMapCompositePropagator;
61use opentelemetry_otlp::WithExportConfig;
62use opentelemetry_sdk::{
63 Resource,
64 logs::SdkLoggerProvider,
65 metrics::{MeterProviderBuilder, PeriodicReader, SdkMeterProvider},
66 propagation::{BaggagePropagator, TraceContextPropagator},
67 trace::{BatchConfigBuilder, BatchSpanProcessor, Sampler, SdkTracer, SdkTracerProvider},
68};
69use opentelemetry_semantic_conventions::attribute::{
70 DEPLOYMENT_ENVIRONMENT_NAME, HOST_NAME, PROCESS_PID, SERVICE_VERSION,
71};
72use std::error::Error;
73use std::time::Duration;
74use tracing_subscriber::layer::SubscriberExt;
75use tracing_subscriber::util::SubscriberInitExt;
76
77fn tracing_bridge_tracer(provider: &SdkTracerProvider) -> SdkTracer {
78 use opentelemetry::trace::TracerProvider as _;
79
80 provider.tracer(env!("CARGO_PKG_NAME"))
81}
82
83#[derive(Debug, Clone)]
98pub enum TraceSampler {
99 AlwaysOn,
101 AlwaysOff,
103 TraceIdRatio(f64),
105 ParentBased(Box<TraceSampler>),
108}
109
110#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
112pub enum LogFormat {
113 #[default]
115 Pretty,
116 Json,
118}
119
120impl TraceSampler {
121 fn into_sdk_sampler(self) -> Sampler {
123 match self {
124 TraceSampler::AlwaysOn => Sampler::AlwaysOn,
125 TraceSampler::AlwaysOff => Sampler::AlwaysOff,
126 TraceSampler::TraceIdRatio(r) => Sampler::TraceIdRatioBased(r),
127 TraceSampler::ParentBased(inner) => {
128 Sampler::ParentBased(Box::new(inner.into_sdk_sampler()))
129 }
130 }
131 }
132}
133
134fn sampler_from_env() -> Result<Option<TraceSampler>, Box<dyn Error>> {
142 let name = match std::env::var("OTEL_TRACES_SAMPLER") {
143 Ok(v) => v,
144 Err(_) => return Ok(None),
145 };
146 let arg = std::env::var("OTEL_TRACES_SAMPLER_ARG").ok();
147 let sampler = match name.as_str() {
148 "always_on" => TraceSampler::AlwaysOn,
149 "always_off" => TraceSampler::AlwaysOff,
150 "traceidratio" => {
151 let ratio = arg
152 .as_deref()
153 .unwrap_or("1.0")
154 .parse::<f64>()
155 .unwrap_or(1.0);
156 TraceSampler::TraceIdRatio(ratio)
157 }
158 "parentbased_always_on" => TraceSampler::ParentBased(Box::new(TraceSampler::AlwaysOn)),
159 "parentbased_always_off" => TraceSampler::ParentBased(Box::new(TraceSampler::AlwaysOff)),
160 "parentbased_traceidratio" => {
161 let ratio = arg
162 .as_deref()
163 .unwrap_or("1.0")
164 .parse::<f64>()
165 .unwrap_or(1.0);
166 TraceSampler::ParentBased(Box::new(TraceSampler::TraceIdRatio(ratio)))
167 }
168 unknown => {
169 return Err(format!(
170 "OTEL_TRACES_SAMPLER: unrecognised sampler name '{unknown}'. \
171 Valid values: always_on, always_off, traceidratio, \
172 parentbased_always_on, parentbased_always_off, parentbased_traceidratio"
173 )
174 .into());
175 }
176 };
177 Ok(Some(sampler))
178}
179
180const DEFAULT_SHUTDOWN_TIMEOUT: Duration = Duration::from_secs(5);
182
183pub struct TelemetryHandles {
204 pub tracer_provider: SdkTracerProvider,
205 pub meter_provider: Option<SdkMeterProvider>,
206 pub logger_provider: Option<SdkLoggerProvider>,
207 shutdown_timeout: Duration,
208 boot_owner: bool,
211 #[cfg(feature = "profiling")]
212 pub profiling_handle: Option<profiling::ProfilingHandle>,
213}
214
215impl TelemetryHandles {
216 pub fn shutdown(&self) -> Result<(), Box<dyn Error>> {
247 self.close_boot();
248 if let Err(e) = self.tracer_provider.shutdown() {
249 tracing::warn!("tracer provider shutdown error: {e}");
250 }
251 if let Some(mp) = &self.meter_provider
252 && let Err(e) = mp.shutdown()
253 {
254 tracing::warn!("meter provider shutdown error: {e}");
255 }
256 if let Some(lp) = &self.logger_provider
257 && let Err(e) = lp.shutdown()
258 {
259 tracing::warn!("logger provider shutdown error: {e}");
260 }
261 Ok(())
262 }
263
264 fn close_boot(&self) {
266 if self.boot_owner {
267 boot::Timeline::global().close();
268 }
269 }
270}
271
272impl Drop for TelemetryHandles {
273 fn drop(&mut self) {
274 self.close_boot();
275 let tracer_provider = self.tracer_provider.clone();
276 let meter_provider = self.meter_provider.clone();
277 let logger_provider = self.logger_provider.clone();
278 let timeout = self.shutdown_timeout;
279
280 let (tx, rx) = std::sync::mpsc::channel();
281 std::thread::spawn(move || {
282 if let Err(e) = tracer_provider.shutdown() {
283 tracing::warn!("tracer provider shutdown error: {e}");
284 }
285 if let Some(mp) = meter_provider
286 && let Err(e) = mp.shutdown()
287 {
288 tracing::warn!("meter provider shutdown error: {e}");
289 }
290 if let Some(lp) = logger_provider
291 && let Err(e) = lp.shutdown()
292 {
293 tracing::warn!("logger provider shutdown error: {e}");
294 }
295 let _ = tx.send(());
296 });
297
298 if rx.recv_timeout(timeout).is_err() {
299 tracing::warn!(
300 "telemetry shutdown did not complete within {timeout:?}; \
301 some spans/metrics may not have been exported"
302 );
303 }
304 }
305}
306
307#[derive(Debug, Clone, Copy, PartialEq, Eq)]
330pub enum ExportProtocol {
331 #[cfg(feature = "grpc")]
333 Grpc,
334 #[cfg(feature = "http")]
336 HttpProtobuf,
337}
338
339#[cfg(feature = "grpc-mtls")]
349#[derive(Clone)]
350pub struct MtlsMaterial {
351 pub client_cert_chain_pem: Vec<u8>,
353 pub client_key_pem: Vec<u8>,
355 pub trust_bundle_pem: Vec<u8>,
357}
358
359#[cfg(feature = "grpc-mtls")]
360impl std::fmt::Debug for MtlsMaterial {
361 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
362 f.debug_struct("MtlsMaterial")
363 .field("client_cert_chain_pem", &"<redacted>")
364 .field("client_key_pem", &"<redacted>")
365 .field("trust_bundle_pem", &"<redacted>")
366 .finish()
367 }
368}
369
370fn protocol_from_env() -> Option<ExportProtocol> {
372 let val = std::env::var("OTEL_EXPORTER_OTLP_PROTOCOL").ok()?;
373 match val.trim() {
374 #[cfg(feature = "grpc")]
375 "grpc" => Some(ExportProtocol::Grpc),
376 #[cfg(feature = "http")]
377 "http/protobuf" => Some(ExportProtocol::HttpProtobuf),
378 _ => None,
379 }
380}
381
382pub struct Telemetry;
398
399impl Telemetry {
400 pub fn builder(service_name: &str) -> TelemetryBuilder {
404 TelemetryBuilder {
405 service_name: Some(service_name.to_string()),
406 service_version: None,
407 service_namespace: None,
408 deployment_environment: None,
409 sampler: None,
410 metrics: true,
411 logs: false,
412 protocol: None,
413 default_endpoint: None,
414 max_export_batch_size: None,
415 metric_export_interval: None,
416 export_timeout: None,
417 shutdown_timeout: DEFAULT_SHUTDOWN_TIMEOUT,
418 log_filter: None,
419 log_format: LogFormat::default(),
420 extra_layers: Vec::new(),
421 extra_metric_readers: Vec::new(),
422 runtime_metrics: true,
423 #[cfg(feature = "grpc-mtls")]
424 mtls: None,
425 #[cfg(feature = "grpc-mtls")]
426 mtls_source: None,
427 #[cfg(feature = "grpc-mtls")]
428 collector_spiffe_id: None,
429 propagated_span_fields: crate::log_bridge::PROPAGATED_SPAN_FIELDS,
430 #[cfg(feature = "profiling")]
431 pyroscope_endpoint: None,
432 }
433 }
434
435 pub fn from_env() -> TelemetryBuilder {
445 TelemetryBuilder {
446 service_name: None,
447 service_version: None,
448 service_namespace: None,
449 deployment_environment: None,
450 sampler: None,
451 metrics: true,
452 logs: false,
453 protocol: None,
454 default_endpoint: None,
455 max_export_batch_size: None,
456 metric_export_interval: None,
457 export_timeout: None,
458 shutdown_timeout: DEFAULT_SHUTDOWN_TIMEOUT,
459 log_filter: None,
460 log_format: LogFormat::default(),
461 extra_layers: Vec::new(),
462 extra_metric_readers: Vec::new(),
463 runtime_metrics: true,
464 #[cfg(feature = "grpc-mtls")]
465 mtls: None,
466 #[cfg(feature = "grpc-mtls")]
467 mtls_source: None,
468 #[cfg(feature = "grpc-mtls")]
469 collector_spiffe_id: None,
470 propagated_span_fields: crate::log_bridge::PROPAGATED_SPAN_FIELDS,
471 #[cfg(feature = "profiling")]
472 pyroscope_endpoint: None,
473 }
474 }
475}
476
477#[must_use = "a TelemetryBuilder does nothing until .init() is called"]
495pub struct TelemetryBuilder {
496 service_name: Option<String>,
497 service_version: Option<String>,
498 service_namespace: Option<String>,
499 deployment_environment: Option<String>,
500 sampler: Option<TraceSampler>,
501 metrics: bool,
502 logs: bool,
503 protocol: Option<ExportProtocol>,
504 max_export_batch_size: Option<usize>,
505 metric_export_interval: Option<Duration>,
506 export_timeout: Option<Duration>,
507 shutdown_timeout: Duration,
508 log_filter: Option<String>,
509 log_format: LogFormat,
510 extra_layers: Vec<
511 Box<dyn tracing_subscriber::Layer<tracing_subscriber::Registry> + Send + Sync + 'static>,
512 >,
513 extra_metric_readers: Vec<MeterProviderInstaller>,
514 runtime_metrics: bool,
515 #[cfg(feature = "grpc-mtls")]
516 mtls: Option<MtlsMaterial>,
517 #[cfg(feature = "grpc-mtls")]
518 mtls_source: Option<std::sync::Arc<dyn CertSource>>,
519 #[cfg(feature = "grpc-mtls")]
520 collector_spiffe_id: Option<String>,
521 propagated_span_fields: &'static [&'static str],
522 #[cfg(feature = "profiling")]
523 pyroscope_endpoint: Option<String>,
524 default_endpoint: Option<String>,
525}
526
527type MeterProviderInstaller =
532 Box<dyn FnOnce(MeterProviderBuilder) -> MeterProviderBuilder + Send + Sync>;
533
534impl TelemetryBuilder {
535 pub fn with_log_filter(mut self, directive: impl Into<String>) -> Self {
540 self.log_filter = Some(directive.into());
541 self
542 }
543
544 pub fn with_log_format(mut self, format: LogFormat) -> Self {
546 self.log_format = format;
547 self
548 }
549
550 pub fn with_version(mut self, version: &str) -> Self {
552 self.service_version = Some(version.to_string());
553 self
554 }
555
556 pub fn with_service_namespace(mut self, namespace: &str) -> Self {
559 self.service_namespace = Some(namespace.to_string()).filter(|n| !n.is_empty());
560 self
561 }
562
563 pub fn with_environment(mut self, environment: &str) -> Self {
565 self.deployment_environment = Some(environment.to_string());
566 self
567 }
568
569 #[cfg(feature = "grpc-mtls")]
586 pub fn with_mtls(mut self, material: MtlsMaterial) -> Self {
587 self.mtls = Some(material);
588 self.mtls_source = None;
589 self.protocol = Some(ExportProtocol::Grpc);
590 self
591 }
592
593 #[cfg(feature = "grpc-mtls")]
608 pub fn with_mtls_source(mut self, source: std::sync::Arc<dyn CertSource>) -> Self {
609 self.mtls_source = Some(source);
610 self.mtls = None;
611 self.protocol = Some(ExportProtocol::Grpc);
612 self
613 }
614
615 #[cfg(feature = "grpc-mtls")]
628 pub fn with_collector_spiffe_id(mut self, id: impl Into<String>) -> Self {
629 self.collector_spiffe_id = Some(id.into());
630 self
631 }
632
633 pub fn with_sampler(mut self, sampler: TraceSampler) -> Self {
636 self.sampler = Some(sampler);
637 self
638 }
639
640 pub fn with_metrics(mut self, enabled: bool) -> Self {
642 self.metrics = enabled;
643 self
644 }
645
646 pub fn with_runtime_metrics(mut self, enabled: bool) -> Self {
662 self.runtime_metrics = enabled;
663 self
664 }
665
666 pub fn with_default_endpoint(mut self, endpoint: impl Into<String>) -> Self {
671 self.default_endpoint = Some(endpoint.into());
672 self
673 }
674
675 pub fn with_protocol(mut self, protocol: ExportProtocol) -> Self {
679 self.protocol = Some(protocol);
680 self
681 }
682
683 pub fn with_max_export_batch_size(mut self, size: usize) -> Self {
688 self.max_export_batch_size = Some(size);
689 self
690 }
691
692 pub fn with_metric_export_interval(mut self, interval: Duration) -> Self {
697 self.metric_export_interval = Some(interval);
698 self
699 }
700
701 pub fn with_logs(mut self, enabled: bool) -> Self {
708 self.logs = enabled;
709 self
710 }
711
712 pub fn with_propagated_span_fields(mut self, fields: &'static [&'static str]) -> Self {
725 self.propagated_span_fields = fields;
726 self
727 }
728
729 pub fn with_export_timeout(mut self, timeout: Duration) -> Self {
733 self.export_timeout = Some(timeout);
734 self
735 }
736
737 pub fn with_shutdown_timeout(mut self, timeout: Duration) -> Self {
744 self.shutdown_timeout = timeout;
745 self
746 }
747
748 #[cfg(feature = "profiling")]
766 pub fn with_profiling(mut self, endpoint: &str) -> Self {
767 self.pyroscope_endpoint = Some(endpoint.to_string());
768 self
769 }
770
771 pub fn with_meter_provider_setup<F>(mut self, setup: F) -> Self
828 where
829 F: FnOnce(MeterProviderBuilder) -> MeterProviderBuilder + Send + Sync + 'static,
830 {
831 self.extra_metric_readers.push(Box::new(setup));
832 self
833 }
834
835 pub fn with_layer<L>(mut self, layer: L) -> Self
836 where
837 L: tracing_subscriber::Layer<tracing_subscriber::Registry> + Send + Sync + 'static,
838 {
839 self.extra_layers.push(Box::new(layer));
840 self
841 }
842
843 #[cfg_attr(not(feature = "grpc-mtls"), allow(unused_mut))]
858 pub fn init(mut self) -> Result<TelemetryHandles, Box<dyn Error>> {
859 let log_filter = match self.log_filter.as_deref() {
860 Some(directive) => tracing_subscriber::EnvFilter::try_new(directive)?,
861 None => tracing_subscriber::EnvFilter::from_default_env(),
862 };
863
864 if let Some(interval) = self.metric_export_interval
865 && interval.is_zero()
866 {
867 return Err("metric_export_interval must be greater than zero".into());
868 }
869
870 let protocol = self.protocol.or_else(protocol_from_env).unwrap_or({
871 #[cfg(feature = "grpc")]
872 {
873 ExportProtocol::Grpc
874 }
875 #[cfg(all(not(feature = "grpc"), feature = "http"))]
876 {
877 ExportProtocol::HttpProtobuf
878 }
879 });
880
881 let default_endpoint = match protocol {
882 #[cfg(feature = "grpc")]
883 ExportProtocol::Grpc => "http://localhost:4317",
884 #[cfg(feature = "http")]
885 ExportProtocol::HttpProtobuf => "http://localhost:4318",
886 };
887 let endpoint = resolve_endpoint(
888 std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT").ok(),
889 self.default_endpoint.as_deref(),
890 default_endpoint,
891 );
892
893 let export_timeout = self.export_timeout.or_else(timeout_from_env);
895
896 let service_name = self.service_name.unwrap_or_else(|| {
898 std::env::var("OTEL_SERVICE_NAME").unwrap_or_else(|_| "unknown_service".to_string())
899 });
900
901 let resource = build_resource_in(
902 &service_name,
903 self.service_namespace.as_deref(),
904 self.service_version.as_deref(),
905 self.deployment_environment.as_deref(),
906 );
907
908 let sampler = match self.sampler {
909 Some(s) => s,
910 None => sampler_from_env()?.unwrap_or(TraceSampler::AlwaysOn),
911 };
912
913 #[cfg(feature = "grpc-mtls")]
914 let mtls_transport = MtlsTransport::resolve(
915 self.mtls.take(),
916 self.mtls_source.take(),
917 &endpoint,
918 export_timeout,
919 self.collector_spiffe_id.as_deref(),
920 )?;
921
922 let trace_exporter = build_span_exporter(
924 protocol,
925 &endpoint,
926 export_timeout,
927 #[cfg(feature = "grpc-mtls")]
928 mtls_transport.as_ref(),
929 )?;
930
931 let batch_processor = if let Some(size) = self.max_export_batch_size {
932 BatchSpanProcessor::builder(trace_exporter)
933 .with_batch_config(
934 BatchConfigBuilder::default()
935 .with_max_export_batch_size(size)
936 .build(),
937 )
938 .build()
939 } else {
940 BatchSpanProcessor::builder(trace_exporter).build()
941 };
942
943 let tracer_provider = SdkTracerProvider::builder()
944 .with_resource(resource.clone())
945 .with_sampler(sampler.into_sdk_sampler())
946 .with_span_processor(batch_processor)
947 .build();
948
949 opentelemetry::global::set_tracer_provider(tracer_provider.clone());
950
951 let propagator = TextMapCompositePropagator::new(vec![
953 Box::new(TraceContextPropagator::new()),
954 Box::new(BaggagePropagator::new()),
955 ]);
956 opentelemetry::global::set_text_map_propagator(propagator);
957
958 let meter_provider = if self.metrics {
960 let metric_exporter = build_metric_exporter(
961 protocol,
962 &endpoint,
963 export_timeout,
964 #[cfg(feature = "grpc-mtls")]
965 mtls_transport.as_ref(),
966 )?;
967
968 let periodic_reader = if let Some(interval) = self.metric_export_interval {
969 PeriodicReader::builder(metric_exporter)
970 .with_interval(interval)
971 .build()
972 } else {
973 PeriodicReader::builder(metric_exporter).build()
974 };
975
976 let mut mp_builder = SdkMeterProvider::builder()
977 .with_resource(resource.clone())
978 .with_reader(periodic_reader);
979 for installer in self.extra_metric_readers {
980 mp_builder = installer(mp_builder);
981 }
982 let mp = mp_builder.build();
983
984 opentelemetry::global::set_meter_provider(mp.clone());
985
986 if self.runtime_metrics {
990 crate::runtime_metrics::install();
991 }
992
993 Some(mp)
994 } else {
995 None
996 };
997
998 let logger_provider = if self.logs {
1000 let log_exporter = build_log_exporter(
1001 protocol,
1002 &endpoint,
1003 export_timeout,
1004 #[cfg(feature = "grpc-mtls")]
1005 mtls_transport.as_ref(),
1006 )?;
1007
1008 let lp = SdkLoggerProvider::builder()
1009 .with_resource(resource)
1010 .with_batch_exporter(log_exporter)
1011 .build();
1012
1013 Some(lp)
1014 } else {
1015 None
1016 };
1017
1018 #[cfg(feature = "profiling")]
1020 let profiling_handle = if let Some(ref endpoint) = self.pyroscope_endpoint {
1021 let identity = profiling::ProfilingIdentity {
1026 host_name: hostname::get()
1027 .ok()
1028 .and_then(|h| h.into_string().ok())
1029 .filter(|h| !h.is_empty()),
1030 deployment_environment: self.deployment_environment.clone(),
1031 service_version: self.service_version.clone(),
1032 };
1033 profiling::start_pyroscope_bridge(&service_name, endpoint, &identity)?
1034 } else {
1035 None
1036 };
1037 #[cfg(not(feature = "profiling"))]
1038 let _profiling_handle: Option<()> = None;
1039
1040 let extra = if self.extra_layers.is_empty() {
1046 None
1047 } else {
1048 Some(self.extra_layers)
1049 };
1050
1051 macro_rules! install_subscriber {
1052 ($fmt_layer:expr) => {{
1053 let otel_layer = tracing_opentelemetry::layer()
1057 .with_tracer(tracing_bridge_tracer(&tracer_provider));
1058 let registry = tracing_subscriber::registry()
1061 .with(extra)
1062 .with(crate::export_backoff::ExportFailureBackoff::default())
1063 .with(log_filter)
1064 .with($fmt_layer)
1065 .with(otel_layer);
1066
1067 #[cfg(feature = "profiling-bridge-pyroscope-rs")]
1070 #[allow(deprecated)]
1071 let registry = registry.with(crate::profiling::ProfilingTagLayer);
1072
1073 if let Some(lp) = &logger_provider {
1074 if let Err(e) = registry
1075 .with(crate::log_bridge::SpanAwareLogBridge::new(
1076 lp,
1077 self.propagated_span_fields,
1078 ))
1079 .try_init()
1080 {
1081 eprintln!(
1082 "otel-bootstrap: global tracing subscriber already installed — \
1083 OTLP log records will NOT be exported to the collector: {e}"
1084 );
1085 }
1086 } else if let Err(e) = registry.try_init() {
1087 eprintln!(
1088 "otel-bootstrap: global tracing subscriber already installed — \
1089 OTLP telemetry will NOT be exported to the collector: {e}"
1090 );
1091 }
1092 }};
1093 }
1094
1095 match self.log_format {
1096 LogFormat::Pretty => install_subscriber!(tracing_subscriber::fmt::layer()),
1097 LogFormat::Json => install_subscriber!(tracing_subscriber::fmt::layer().json()),
1098 }
1099
1100 let boot_owner = boot::Timeline::global()
1102 .attach(&tracer_provider, &service_name)
1103 .is_some();
1104
1105 Ok(TelemetryHandles {
1106 tracer_provider,
1107 meter_provider,
1108 logger_provider,
1109 shutdown_timeout: self.shutdown_timeout,
1110 boot_owner,
1111 #[cfg(feature = "profiling")]
1112 profiling_handle,
1113 })
1114 }
1115}
1116
1117pub fn init_telemetry(service_name: &str) -> Result<TelemetryHandles, Box<dyn Error>> {
1131 Telemetry::builder(service_name).init()
1132}
1133
1134pub fn init_telemetry_with_sampler(
1151 service_name: &str,
1152 sampler: Option<TraceSampler>,
1153) -> Result<TelemetryHandles, Box<dyn Error>> {
1154 let builder = Telemetry::builder(service_name);
1155 match sampler {
1156 Some(s) => builder.with_sampler(s),
1157 None => builder, }
1159 .init()
1160}
1161
1162fn timeout_from_env() -> Option<Duration> {
1164 let ms = std::env::var("OTEL_EXPORTER_OTLP_TIMEOUT").ok()?;
1165 let ms: u64 = ms.trim().parse().ok()?;
1166 Some(Duration::from_millis(ms))
1167}
1168
1169#[cfg(feature = "grpc-mtls")]
1173enum MtlsTransport {
1174 Snapshot(MtlsMaterial),
1175 Live(tonic::transport::Channel),
1176}
1177
1178#[cfg(feature = "grpc-mtls")]
1179impl MtlsTransport {
1180 fn resolve(
1181 material: Option<MtlsMaterial>,
1182 source: Option<std::sync::Arc<dyn CertSource>>,
1183 endpoint: &str,
1184 timeout: Option<Duration>,
1185 expected_id: Option<&str>,
1186 ) -> Result<Option<Self>, Box<dyn Error>> {
1187 let source = match (source, material.as_ref(), expected_id) {
1190 (Some(source), _, _) => Some(source),
1191 (None, Some(material), Some(_)) => {
1192 Some(std::sync::Arc::new(StaticCertSource::new(material.clone()))
1193 as std::sync::Arc<dyn CertSource>)
1194 }
1195 _ => None,
1196 };
1197 match (source, material) {
1198 (Some(source), _) => {
1199 let channel = rotating_mtls::channel(endpoint, source, timeout, expected_id)
1200 .map_err(|e| -> Box<dyn Error> { e.to_string().into() })?;
1201 Ok(Some(Self::Live(channel)))
1202 }
1203 (None, Some(material)) => Ok(Some(Self::Snapshot(material))),
1204 (None, None) => Ok(None),
1205 }
1206 }
1207
1208 fn apply<B: opentelemetry_otlp::WithTonicConfig>(&self, builder: B) -> B {
1209 match self {
1210 Self::Snapshot(material) => builder.with_tls_config(build_tls_config(material)),
1211 Self::Live(channel) => builder.with_channel(channel.clone()),
1212 }
1213 }
1214}
1215
1216#[cfg(feature = "grpc-mtls")]
1223fn build_tls_config(material: &MtlsMaterial) -> tonic::transport::ClientTlsConfig {
1224 use tonic::transport::{Certificate, ClientTlsConfig, Identity};
1225 ClientTlsConfig::new()
1226 .ca_certificate(Certificate::from_pem(&material.trust_bundle_pem))
1227 .identity(Identity::from_pem(
1228 &material.client_cert_chain_pem,
1229 &material.client_key_pem,
1230 ))
1231}
1232
1233fn resolve_endpoint(
1236 configured: Option<String>,
1237 runtime_default: Option<&str>,
1238 fallback: &str,
1239) -> String {
1240 configured
1241 .or_else(|| runtime_default.map(str::to_owned))
1242 .unwrap_or_else(|| fallback.to_owned())
1243}
1244
1245fn build_span_exporter(
1246 protocol: ExportProtocol,
1247 endpoint: &str,
1248 timeout: Option<Duration>,
1249 #[cfg(feature = "grpc-mtls")] mtls: Option<&MtlsTransport>,
1250) -> Result<opentelemetry_otlp::SpanExporter, Box<dyn Error>> {
1251 match protocol {
1252 #[cfg(feature = "grpc")]
1253 ExportProtocol::Grpc => {
1254 let mut b = opentelemetry_otlp::SpanExporter::builder()
1255 .with_tonic()
1256 .with_endpoint(endpoint);
1257 if let Some(t) = timeout {
1258 b = b.with_timeout(t);
1259 }
1260 #[cfg(feature = "grpc-mtls")]
1261 if let Some(m) = mtls {
1262 b = m.apply(b);
1263 }
1264 Ok(b.build()?)
1265 }
1266 #[cfg(feature = "http")]
1267 ExportProtocol::HttpProtobuf => {
1268 let mut b = opentelemetry_otlp::SpanExporter::builder()
1269 .with_http()
1270 .with_endpoint(endpoint);
1271 if let Some(t) = timeout {
1272 b = b.with_timeout(t);
1273 }
1274 Ok(b.build()?)
1275 }
1276 }
1277}
1278
1279fn build_metric_exporter(
1280 protocol: ExportProtocol,
1281 endpoint: &str,
1282 timeout: Option<Duration>,
1283 #[cfg(feature = "grpc-mtls")] mtls: Option<&MtlsTransport>,
1284) -> Result<opentelemetry_otlp::MetricExporter, Box<dyn Error>> {
1285 match protocol {
1286 #[cfg(feature = "grpc")]
1287 ExportProtocol::Grpc => {
1288 let mut b = opentelemetry_otlp::MetricExporter::builder()
1289 .with_tonic()
1290 .with_endpoint(endpoint);
1291 if let Some(t) = timeout {
1292 b = b.with_timeout(t);
1293 }
1294 #[cfg(feature = "grpc-mtls")]
1295 if let Some(m) = mtls {
1296 b = m.apply(b);
1297 }
1298 Ok(b.build()?)
1299 }
1300 #[cfg(feature = "http")]
1301 ExportProtocol::HttpProtobuf => {
1302 let mut b = opentelemetry_otlp::MetricExporter::builder()
1303 .with_http()
1304 .with_endpoint(endpoint);
1305 if let Some(t) = timeout {
1306 b = b.with_timeout(t);
1307 }
1308 Ok(b.build()?)
1309 }
1310 }
1311}
1312
1313fn build_log_exporter(
1314 protocol: ExportProtocol,
1315 endpoint: &str,
1316 timeout: Option<Duration>,
1317 #[cfg(feature = "grpc-mtls")] mtls: Option<&MtlsTransport>,
1318) -> Result<opentelemetry_otlp::LogExporter, Box<dyn Error>> {
1319 match protocol {
1320 #[cfg(feature = "grpc")]
1321 ExportProtocol::Grpc => {
1322 let mut b = opentelemetry_otlp::LogExporter::builder()
1323 .with_tonic()
1324 .with_endpoint(endpoint);
1325 if let Some(t) = timeout {
1326 b = b.with_timeout(t);
1327 }
1328 #[cfg(feature = "grpc-mtls")]
1329 if let Some(m) = mtls {
1330 b = m.apply(b);
1331 }
1332 Ok(b.build()?)
1333 }
1334 #[cfg(feature = "http")]
1335 ExportProtocol::HttpProtobuf => {
1336 let mut b = opentelemetry_otlp::LogExporter::builder()
1337 .with_http()
1338 .with_endpoint(endpoint);
1339 if let Some(t) = timeout {
1340 b = b.with_timeout(t);
1341 }
1342 Ok(b.build()?)
1343 }
1344 }
1345}
1346
1347pub fn build_resource(
1362 service_name: &str,
1363 service_version: Option<&str>,
1364 deployment_environment: Option<&str>,
1365) -> Resource {
1366 build_resource_in(service_name, None, service_version, deployment_environment)
1367}
1368
1369fn build_resource_in(
1370 service_name: &str,
1371 service_namespace: Option<&str>,
1372 service_version: Option<&str>,
1373 deployment_environment: Option<&str>,
1374) -> Resource {
1375 let hostname = hostname::get()
1376 .ok()
1377 .and_then(|h| h.into_string().ok())
1378 .unwrap_or_default();
1379
1380 let mut builder = Resource::builder()
1381 .with_service_name(service_name.to_string())
1382 .with_attributes([
1383 KeyValue::new(HOST_NAME, hostname),
1384 KeyValue::new(PROCESS_PID, std::process::id() as i64),
1385 ]);
1386
1387 if let Some(namespace) = service_namespace {
1388 builder = builder.with_attribute(KeyValue::new("service.namespace", namespace.to_string()));
1389 }
1390
1391 if let Some(version) = service_version {
1392 builder = builder.with_attribute(KeyValue::new(SERVICE_VERSION, version.to_string()));
1393 }
1394
1395 if let Some(env) = deployment_environment {
1396 builder =
1397 builder.with_attribute(KeyValue::new(DEPLOYMENT_ENVIRONMENT_NAME, env.to_string()));
1398 }
1399
1400 builder.build()
1401}
1402
1403#[cfg(feature = "axum")]
1421pub fn axum_layer() -> axum_middleware::OtelTraceLayer {
1422 axum_middleware::OtelTraceLayer
1423}
1424
1425#[cfg(feature = "axum")]
1456pub fn span_enricher_layer<T>() -> axum_middleware::SpanEnricherLayer<T>
1457where
1458 T: span_enrichment::EnrichSpan + Clone + Send + Sync + 'static,
1459{
1460 axum_middleware::SpanEnricherLayer::default()
1461}
1462
1463#[cfg(feature = "tonic-tracing")]
1485pub fn grpc_client_layer() -> grpc_middleware::GrpcClientTraceLayer {
1486 grpc_middleware::GrpcClientTraceLayer
1487}
1488
1489#[cfg(feature = "tonic-tracing")]
1504pub fn grpc_server_layer() -> grpc_middleware::GrpcServerTraceLayer {
1505 grpc_middleware::GrpcServerTraceLayer
1506}
1507
1508#[cfg(test)]
1509mod tests {
1510 use super::*;
1511
1512 #[test]
1514 fn runtime_metrics_can_be_disabled() {
1515 assert!(
1516 Telemetry::builder("rm-default").runtime_metrics,
1517 "runtime metrics are on by default"
1518 );
1519 assert!(
1520 !Telemetry::builder("rm-off")
1521 .with_runtime_metrics(false)
1522 .runtime_metrics
1523 );
1524 }
1525
1526 #[tokio::test]
1534 async fn shutdown_absorbs_provider_errors() {
1535 let handles = TelemetryHandles {
1536 tracer_provider: SdkTracerProvider::builder().build(),
1537 meter_provider: Some(SdkMeterProvider::builder().build()),
1538 logger_provider: Some(SdkLoggerProvider::builder().build()),
1539 shutdown_timeout: DEFAULT_SHUTDOWN_TIMEOUT,
1540 boot_owner: false,
1541 #[cfg(feature = "profiling")]
1542 profiling_handle: None,
1543 };
1544
1545 handles.shutdown().expect("first shutdown succeeds");
1546 handles
1547 .shutdown()
1548 .expect("second shutdown absorbs the already-shut-down errors");
1549 }
1550 use opentelemetry::trace::{Span as _, Tracer as _};
1551 use std::sync::Mutex;
1552
1553 static ENV_LOCK: Mutex<()> = Mutex::new(());
1554
1555 #[test]
1556 fn tracing_bridge_uses_sdk_tracer() {
1557 let provider = SdkTracerProvider::builder().build();
1558 let tracer = tracing_bridge_tracer(&provider);
1559 let span = tracer.start("bridge-regression");
1560
1561 assert!(span.span_context().is_valid());
1562
1563 provider.shutdown().expect("provider shutdown");
1564 }
1565
1566 #[test]
1567 fn resource_contains_all_attributes_when_provided() {
1568 let resource = build_resource("test-svc", Some("1.2.3"), Some("staging"));
1569
1570 assert_eq!(
1571 resource.get(&opentelemetry::Key::new("service.name")),
1572 Some(opentelemetry::Value::from("test-svc")),
1573 );
1574 assert_eq!(
1575 resource.get(&opentelemetry::Key::new(SERVICE_VERSION)),
1576 Some(opentelemetry::Value::from("1.2.3")),
1577 );
1578 assert_eq!(
1579 resource.get(&opentelemetry::Key::new(DEPLOYMENT_ENVIRONMENT_NAME)),
1580 Some(opentelemetry::Value::from("staging")),
1581 );
1582 assert!(resource.get(&opentelemetry::Key::new(HOST_NAME)).is_some());
1583 assert!(
1584 resource
1585 .get(&opentelemetry::Key::new(PROCESS_PID))
1586 .is_some()
1587 );
1588 }
1589
1590 #[test]
1591 fn resource_carries_the_service_namespace_only_when_set() {
1592 let key = opentelemetry::Key::new("service.namespace");
1593 let with = build_resource_in("svc", Some("is.brefwiz.internal"), None, None);
1594 assert_eq!(
1595 with.get(&key),
1596 Some(opentelemetry::Value::from("is.brefwiz.internal")),
1597 );
1598 assert!(build_resource("svc", None, None).get(&key).is_none());
1599 }
1600
1601 #[test]
1602 fn builder_namespace_is_absent_unless_non_empty() {
1603 let set = Telemetry::builder("svc").with_service_namespace("realm.example");
1604 assert_eq!(set.service_namespace.as_deref(), Some("realm.example"));
1605 let empty = Telemetry::builder("svc").with_service_namespace("");
1606 assert!(empty.service_namespace.is_none());
1607 assert!(Telemetry::builder("svc").service_namespace.is_none());
1608 }
1609
1610 #[test]
1611 fn resource_graceful_when_optional_values_omitted() {
1612 let resource = build_resource("test-svc", None, None);
1613
1614 assert_eq!(
1615 resource.get(&opentelemetry::Key::new("service.name")),
1616 Some(opentelemetry::Value::from("test-svc")),
1617 );
1618 assert!(
1619 resource
1620 .get(&opentelemetry::Key::new(SERVICE_VERSION))
1621 .is_none()
1622 );
1623 assert!(
1624 resource
1625 .get(&opentelemetry::Key::new(DEPLOYMENT_ENVIRONMENT_NAME))
1626 .is_none()
1627 );
1628 assert!(resource.get(&opentelemetry::Key::new(HOST_NAME)).is_some());
1630 assert!(
1631 resource
1632 .get(&opentelemetry::Key::new(PROCESS_PID))
1633 .is_some()
1634 );
1635 }
1636
1637 #[test]
1638 fn trace_sampler_ratio_converts_to_sdk() {
1639 let sampler = TraceSampler::TraceIdRatio(0.5);
1640 let sdk = sampler.into_sdk_sampler();
1641 assert_eq!(format!("{sdk:?}"), "TraceIdRatioBased(0.5)");
1642 }
1643
1644 #[test]
1645 fn trace_sampler_parent_based_converts_to_sdk() {
1646 let sampler = TraceSampler::ParentBased(Box::new(TraceSampler::TraceIdRatio(0.25)));
1647 let sdk = sampler.into_sdk_sampler();
1648 let debug = format!("{sdk:?}");
1649 assert!(debug.contains("ParentBased"));
1650 assert!(debug.contains("0.25"));
1651 }
1652
1653 unsafe fn set_env(key: &str, val: &str) {
1655 unsafe {
1656 std::env::set_var(key, val);
1657 }
1658 }
1659
1660 unsafe fn remove_env(key: &str) {
1661 unsafe {
1662 std::env::remove_var(key);
1663 }
1664 }
1665
1666 #[test]
1667 fn sampler_from_env_reads_traceidratio() {
1668 let _lock = ENV_LOCK.lock().unwrap();
1669 unsafe {
1670 set_env("OTEL_TRACES_SAMPLER", "traceidratio");
1671 set_env("OTEL_TRACES_SAMPLER_ARG", "0.42");
1672 }
1673
1674 let sampler = sampler_from_env()
1675 .expect("should not error")
1676 .expect("should return Some");
1677 assert!(
1678 matches!(sampler, TraceSampler::TraceIdRatio(r) if (r - 0.42).abs() < f64::EPSILON)
1679 );
1680
1681 unsafe {
1682 remove_env("OTEL_TRACES_SAMPLER");
1683 remove_env("OTEL_TRACES_SAMPLER_ARG");
1684 }
1685 }
1686
1687 #[test]
1688 fn sampler_from_env_returns_none_when_unset() {
1689 let _lock = ENV_LOCK.lock().unwrap();
1690 unsafe {
1691 remove_env("OTEL_TRACES_SAMPLER");
1692 }
1693 assert!(sampler_from_env().expect("should not error").is_none());
1694 }
1695
1696 #[test]
1697 fn sampler_from_env_reads_parentbased_traceidratio() {
1698 let _lock = ENV_LOCK.lock().unwrap();
1699 unsafe {
1700 set_env("OTEL_TRACES_SAMPLER", "parentbased_traceidratio");
1701 set_env("OTEL_TRACES_SAMPLER_ARG", "0.1");
1702 }
1703
1704 let sampler = sampler_from_env()
1705 .expect("should not error")
1706 .expect("should return Some");
1707 assert!(
1708 matches!(sampler, TraceSampler::ParentBased(inner) if matches!(*inner, TraceSampler::TraceIdRatio(r) if (r - 0.1).abs() < f64::EPSILON))
1709 );
1710
1711 unsafe {
1712 remove_env("OTEL_TRACES_SAMPLER");
1713 remove_env("OTEL_TRACES_SAMPLER_ARG");
1714 }
1715 }
1716
1717 #[test]
1718 fn sampler_from_env_parentbased_always_on() {
1719 let _lock = ENV_LOCK.lock().unwrap();
1720 unsafe {
1721 set_env("OTEL_TRACES_SAMPLER", "parentbased_always_on");
1722 }
1723 let sampler = sampler_from_env()
1724 .expect("should not error")
1725 .expect("should return Some");
1726 assert!(
1727 matches!(sampler, TraceSampler::ParentBased(inner) if matches!(*inner, TraceSampler::AlwaysOn))
1728 );
1729 unsafe {
1730 remove_env("OTEL_TRACES_SAMPLER");
1731 }
1732 }
1733
1734 #[test]
1735 fn sampler_from_env_parentbased_always_off() {
1736 let _lock = ENV_LOCK.lock().unwrap();
1737 unsafe {
1738 set_env("OTEL_TRACES_SAMPLER", "parentbased_always_off");
1739 }
1740 let sampler = sampler_from_env()
1741 .expect("should not error")
1742 .expect("should return Some");
1743 assert!(
1744 matches!(sampler, TraceSampler::ParentBased(inner) if matches!(*inner, TraceSampler::AlwaysOff))
1745 );
1746 unsafe {
1747 remove_env("OTEL_TRACES_SAMPLER");
1748 }
1749 }
1750
1751 #[test]
1752 fn sampler_from_env_always_on() {
1753 let _lock = ENV_LOCK.lock().unwrap();
1754 unsafe {
1755 set_env("OTEL_TRACES_SAMPLER", "always_on");
1756 }
1757 let sampler = sampler_from_env()
1758 .expect("should not error")
1759 .expect("should return Some");
1760 assert!(matches!(sampler, TraceSampler::AlwaysOn));
1761 unsafe {
1762 remove_env("OTEL_TRACES_SAMPLER");
1763 }
1764 }
1765
1766 #[test]
1767 fn sampler_from_env_always_off() {
1768 let _lock = ENV_LOCK.lock().unwrap();
1769 unsafe {
1770 set_env("OTEL_TRACES_SAMPLER", "always_off");
1771 }
1772 let sampler = sampler_from_env()
1773 .expect("should not error")
1774 .expect("should return Some");
1775 assert!(matches!(sampler, TraceSampler::AlwaysOff));
1776 unsafe {
1777 remove_env("OTEL_TRACES_SAMPLER");
1778 }
1779 }
1780
1781 #[test]
1782 fn sampler_from_env_unknown_returns_error() {
1783 let _lock = ENV_LOCK.lock().unwrap();
1784 unsafe {
1785 set_env("OTEL_TRACES_SAMPLER", "unknown_sampler");
1786 }
1787 let err = sampler_from_env().expect_err("unknown sampler should produce an error");
1788 assert!(
1789 err.to_string().contains("unknown_sampler"),
1790 "error message should include the unknown name, got: {err}"
1791 );
1792 unsafe {
1793 remove_env("OTEL_TRACES_SAMPLER");
1794 }
1795 }
1796
1797 #[test]
1798 fn trace_sampler_always_on_converts_to_sdk() {
1799 let sdk = TraceSampler::AlwaysOn.into_sdk_sampler();
1800 assert_eq!(format!("{sdk:?}"), "AlwaysOn");
1801 }
1802
1803 #[test]
1804 fn trace_sampler_always_off_converts_to_sdk() {
1805 let sdk = TraceSampler::AlwaysOff.into_sdk_sampler();
1806 assert_eq!(format!("{sdk:?}"), "AlwaysOff");
1807 }
1808
1809 #[test]
1810 fn builder_has_sensible_defaults() {
1811 let builder = Telemetry::builder("test-svc");
1812 assert_eq!(builder.service_name.as_deref(), Some("test-svc"));
1813 assert!(builder.service_version.is_none());
1814 assert!(builder.deployment_environment.is_none());
1815 assert!(builder.sampler.is_none());
1816 assert!(builder.metrics);
1817 assert!(!builder.logs);
1818 assert!(builder.protocol.is_none());
1819 assert!(builder.max_export_batch_size.is_none());
1820 assert!(builder.metric_export_interval.is_none());
1821 assert!(builder.export_timeout.is_none());
1822 }
1823
1824 #[test]
1825 fn from_env_builder_has_no_service_name() {
1826 let builder = Telemetry::from_env();
1827 assert!(builder.service_name.is_none());
1828 }
1829
1830 #[test]
1831 fn with_export_timeout_stores_value() {
1832 let timeout = Duration::from_secs(5);
1833 let builder = Telemetry::builder("test-svc").with_export_timeout(timeout);
1834 assert_eq!(builder.export_timeout, Some(timeout));
1835 }
1836
1837 #[test]
1838 fn timeout_from_env_reads_milliseconds() {
1839 let _lock = ENV_LOCK.lock().unwrap();
1840 unsafe {
1841 set_env("OTEL_EXPORTER_OTLP_TIMEOUT", "5000");
1842 }
1843 let t = timeout_from_env();
1844 assert_eq!(t, Some(Duration::from_millis(5000)));
1845 unsafe {
1846 remove_env("OTEL_EXPORTER_OTLP_TIMEOUT");
1847 }
1848 }
1849
1850 #[test]
1851 fn timeout_from_env_returns_none_when_unset() {
1852 let _lock = ENV_LOCK.lock().unwrap();
1853 unsafe {
1854 remove_env("OTEL_EXPORTER_OTLP_TIMEOUT");
1855 }
1856 assert_eq!(timeout_from_env(), None);
1857 }
1858
1859 #[test]
1860 fn service_name_from_env_used_when_none_given() {
1861 let builder = Telemetry::from_env();
1862 assert!(builder.service_name.is_none());
1863 }
1864
1865 #[test]
1866 fn explicit_service_name_overrides_env_var() {
1867 let builder = Telemetry::builder("explicit-svc");
1868 assert_eq!(builder.service_name.as_deref(), Some("explicit-svc"));
1869 }
1870
1871 #[test]
1872 fn from_env_builder_service_name_is_none() {
1873 let builder = Telemetry::from_env();
1874 assert!(builder.service_name.is_none());
1875 }
1876
1877 #[test]
1878 fn init_returns_error_for_unknown_otel_traces_sampler() {
1879 let _lock = ENV_LOCK.lock().unwrap();
1880 unsafe {
1881 set_env("OTEL_TRACES_SAMPLER", "not_a_real_sampler");
1882 }
1883 let result = Telemetry::builder("test-svc").with_metrics(false).init();
1884 let err = result
1885 .err()
1886 .expect("unknown sampler env var should cause init to fail");
1887 assert!(
1888 err.to_string().contains("not_a_real_sampler"),
1889 "error should name the unknown sampler, got: {err}"
1890 );
1891 unsafe {
1892 remove_env("OTEL_TRACES_SAMPLER");
1893 }
1894 }
1895
1896 #[test]
1897 fn with_max_export_batch_size_stores_value() {
1898 let builder = Telemetry::builder("test-svc").with_max_export_batch_size(1024);
1899 assert_eq!(builder.max_export_batch_size, Some(1024));
1900 }
1901
1902 #[test]
1903 fn with_metric_export_interval_stores_value() {
1904 let interval = Duration::from_secs(30);
1905 let builder = Telemetry::builder("test-svc").with_metric_export_interval(interval);
1906 assert_eq!(builder.metric_export_interval, Some(interval));
1907 }
1908
1909 #[test]
1910 fn init_rejects_zero_metric_export_interval() {
1911 let err = Telemetry::builder("test-svc")
1912 .with_metric_export_interval(Duration::ZERO)
1913 .with_metrics(false)
1914 .init()
1915 .err()
1916 .expect("expected error for zero interval");
1917 assert!(
1918 err.to_string().contains("metric_export_interval"),
1919 "error message should mention metric_export_interval, got: {err}"
1920 );
1921 }
1922
1923 #[test]
1924 fn builder_with_custom_values() {
1925 let builder = Telemetry::builder("test-svc")
1926 .with_version("2.0.0")
1927 .with_environment("production")
1928 .with_sampler(TraceSampler::TraceIdRatio(0.5))
1929 .with_metrics(false);
1930
1931 assert_eq!(builder.service_name.as_deref(), Some("test-svc"));
1932 assert_eq!(builder.service_version.as_deref(), Some("2.0.0"));
1933 assert_eq!(
1934 builder.deployment_environment.as_deref(),
1935 Some("production")
1936 );
1937 assert!(
1938 matches!(builder.sampler, Some(TraceSampler::TraceIdRatio(r)) if (r - 0.5).abs() < f64::EPSILON)
1939 );
1940 assert!(!builder.metrics);
1941 }
1942
1943 #[test]
1944 fn builder_stores_programmatic_log_configuration() {
1945 let builder = Telemetry::builder("test-svc")
1946 .with_log_filter("info,opentelemetry_sdk=warn")
1947 .with_log_format(LogFormat::Json);
1948
1949 assert_eq!(
1950 builder.log_filter.as_deref(),
1951 Some("info,opentelemetry_sdk=warn")
1952 );
1953 assert_eq!(builder.log_format, LogFormat::Json);
1954 }
1955
1956 #[test]
1957 fn init_rejects_invalid_programmatic_log_filter_before_provider_setup() {
1958 let setup_ran = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
1959 let setup_ran_in_closure = std::sync::Arc::clone(&setup_ran);
1960
1961 let error = Telemetry::builder("test-svc")
1962 .with_log_filter("[")
1963 .with_meter_provider_setup(move |builder| {
1964 setup_ran_in_closure.store(true, std::sync::atomic::Ordering::SeqCst);
1965 builder
1966 })
1967 .init()
1968 .err()
1969 .expect("invalid filter must fail initialization");
1970
1971 assert!(error.to_string().contains("invalid filter directive"));
1972 assert!(!setup_ran.load(std::sync::atomic::Ordering::SeqCst));
1973 }
1974
1975 #[test]
1976 fn builder_with_default_endpoint() {
1977 let builder = Telemetry::builder("svc").with_default_endpoint("http://otel-collector:4317");
1978 assert_eq!(
1979 builder.default_endpoint.as_deref(),
1980 Some("http://otel-collector:4317")
1981 );
1982 }
1983
1984 #[test]
1985 fn a_configured_endpoint_wins_over_the_runtime_default() {
1986 assert_eq!(
1987 resolve_endpoint(
1988 Some("http://c:4317".into()),
1989 Some("http://d:4317"),
1990 "http://localhost:4317"
1991 ),
1992 "http://c:4317"
1993 );
1994 assert_eq!(
1995 resolve_endpoint(None, Some("http://d:4317"), "http://localhost:4317"),
1996 "http://d:4317"
1997 );
1998 assert_eq!(
1999 resolve_endpoint(None, None, "http://localhost:4317"),
2000 "http://localhost:4317"
2001 );
2002 }
2003
2004 #[test]
2005 #[cfg(feature = "grpc")]
2006 fn builder_with_protocol_grpc() {
2007 let builder = Telemetry::builder("test-svc").with_protocol(ExportProtocol::Grpc);
2008 assert_eq!(builder.protocol, Some(ExportProtocol::Grpc));
2009 }
2010
2011 #[test]
2012 #[cfg(feature = "http")]
2013 fn builder_with_protocol_http() {
2014 let builder = Telemetry::builder("test-svc").with_protocol(ExportProtocol::HttpProtobuf);
2015 assert_eq!(builder.protocol, Some(ExportProtocol::HttpProtobuf));
2016 }
2017
2018 #[test]
2019 #[cfg(feature = "grpc")]
2020 fn protocol_from_env_reads_grpc() {
2021 let _lock = ENV_LOCK.lock().unwrap();
2022 unsafe {
2023 set_env("OTEL_EXPORTER_OTLP_PROTOCOL", "grpc");
2024 }
2025 assert_eq!(protocol_from_env(), Some(ExportProtocol::Grpc));
2026 unsafe {
2027 remove_env("OTEL_EXPORTER_OTLP_PROTOCOL");
2028 }
2029 }
2030
2031 #[test]
2032 #[cfg(feature = "http")]
2033 fn protocol_from_env_reads_http_protobuf() {
2034 let _lock = ENV_LOCK.lock().unwrap();
2035 unsafe {
2036 set_env("OTEL_EXPORTER_OTLP_PROTOCOL", "http/protobuf");
2037 }
2038 assert_eq!(protocol_from_env(), Some(ExportProtocol::HttpProtobuf));
2039 unsafe {
2040 remove_env("OTEL_EXPORTER_OTLP_PROTOCOL");
2041 }
2042 }
2043
2044 #[test]
2045 fn protocol_from_env_returns_none_when_unset() {
2046 let _lock = ENV_LOCK.lock().unwrap();
2047 unsafe {
2048 remove_env("OTEL_EXPORTER_OTLP_PROTOCOL");
2049 }
2050 assert_eq!(protocol_from_env(), None);
2051 }
2052
2053 #[test]
2054 fn protocol_from_env_returns_none_for_unknown() {
2055 let _lock = ENV_LOCK.lock().unwrap();
2056 unsafe {
2057 set_env("OTEL_EXPORTER_OTLP_PROTOCOL", "websocket");
2058 }
2059 assert_eq!(protocol_from_env(), None);
2060 unsafe {
2061 remove_env("OTEL_EXPORTER_OTLP_PROTOCOL");
2062 }
2063 }
2064
2065 #[test]
2066 fn builder_is_send_and_sync() {
2067 fn assert_send_sync<T: Send + Sync>() {}
2068 assert_send_sync::<TelemetryBuilder>();
2069 }
2070
2071 #[test]
2072 fn with_shutdown_timeout_stores_value() {
2073 let timeout = Duration::from_secs(10);
2074 let builder = Telemetry::builder("test-svc").with_shutdown_timeout(timeout);
2075 assert_eq!(builder.shutdown_timeout, timeout);
2076 }
2077
2078 #[test]
2079 fn default_shutdown_timeout_is_five_seconds() {
2080 let builder = Telemetry::builder("test-svc");
2081 assert_eq!(builder.shutdown_timeout, Duration::from_secs(5));
2082 }
2083
2084 #[cfg(feature = "testing")]
2092 #[test]
2093 fn drop_completes_within_shutdown_timeout() {
2094 let mut handles = crate::Telemetry::testing("drop-timeout-test");
2096 handles.shutdown_timeout = Duration::from_millis(100);
2098
2099 let start = std::time::Instant::now();
2100 drop(handles);
2101 let elapsed = start.elapsed();
2102
2103 assert!(
2105 elapsed < Duration::from_millis(500),
2106 "drop took {elapsed:?}, expected < 500 ms"
2107 );
2108 }
2109}