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 deployment_environment: None,
408 sampler: None,
409 metrics: true,
410 logs: false,
411 protocol: None,
412 default_endpoint: None,
413 max_export_batch_size: None,
414 metric_export_interval: None,
415 export_timeout: None,
416 shutdown_timeout: DEFAULT_SHUTDOWN_TIMEOUT,
417 log_filter: None,
418 log_format: LogFormat::default(),
419 extra_layers: Vec::new(),
420 extra_metric_readers: Vec::new(),
421 runtime_metrics: true,
422 #[cfg(feature = "grpc-mtls")]
423 mtls: None,
424 #[cfg(feature = "grpc-mtls")]
425 mtls_source: None,
426 propagated_span_fields: crate::log_bridge::PROPAGATED_SPAN_FIELDS,
427 #[cfg(feature = "profiling")]
428 pyroscope_endpoint: None,
429 }
430 }
431
432 pub fn from_env() -> TelemetryBuilder {
442 TelemetryBuilder {
443 service_name: None,
444 service_version: None,
445 deployment_environment: None,
446 sampler: None,
447 metrics: true,
448 logs: false,
449 protocol: None,
450 default_endpoint: None,
451 max_export_batch_size: None,
452 metric_export_interval: None,
453 export_timeout: None,
454 shutdown_timeout: DEFAULT_SHUTDOWN_TIMEOUT,
455 log_filter: None,
456 log_format: LogFormat::default(),
457 extra_layers: Vec::new(),
458 extra_metric_readers: Vec::new(),
459 runtime_metrics: true,
460 #[cfg(feature = "grpc-mtls")]
461 mtls: None,
462 #[cfg(feature = "grpc-mtls")]
463 mtls_source: None,
464 propagated_span_fields: crate::log_bridge::PROPAGATED_SPAN_FIELDS,
465 #[cfg(feature = "profiling")]
466 pyroscope_endpoint: None,
467 }
468 }
469}
470
471#[must_use = "a TelemetryBuilder does nothing until .init() is called"]
489pub struct TelemetryBuilder {
490 service_name: Option<String>,
491 service_version: Option<String>,
492 deployment_environment: Option<String>,
493 sampler: Option<TraceSampler>,
494 metrics: bool,
495 logs: bool,
496 protocol: Option<ExportProtocol>,
497 max_export_batch_size: Option<usize>,
498 metric_export_interval: Option<Duration>,
499 export_timeout: Option<Duration>,
500 shutdown_timeout: Duration,
501 log_filter: Option<String>,
502 log_format: LogFormat,
503 extra_layers: Vec<
504 Box<dyn tracing_subscriber::Layer<tracing_subscriber::Registry> + Send + Sync + 'static>,
505 >,
506 extra_metric_readers: Vec<MeterProviderInstaller>,
507 runtime_metrics: bool,
508 #[cfg(feature = "grpc-mtls")]
509 mtls: Option<MtlsMaterial>,
510 #[cfg(feature = "grpc-mtls")]
511 mtls_source: Option<std::sync::Arc<dyn CertSource>>,
512 propagated_span_fields: &'static [&'static str],
513 #[cfg(feature = "profiling")]
514 pyroscope_endpoint: Option<String>,
515 default_endpoint: Option<String>,
516}
517
518type MeterProviderInstaller =
523 Box<dyn FnOnce(MeterProviderBuilder) -> MeterProviderBuilder + Send + Sync>;
524
525impl TelemetryBuilder {
526 pub fn with_log_filter(mut self, directive: impl Into<String>) -> Self {
531 self.log_filter = Some(directive.into());
532 self
533 }
534
535 pub fn with_log_format(mut self, format: LogFormat) -> Self {
537 self.log_format = format;
538 self
539 }
540
541 pub fn with_version(mut self, version: &str) -> Self {
543 self.service_version = Some(version.to_string());
544 self
545 }
546
547 pub fn with_environment(mut self, environment: &str) -> Self {
549 self.deployment_environment = Some(environment.to_string());
550 self
551 }
552
553 #[cfg(feature = "grpc-mtls")]
570 pub fn with_mtls(mut self, material: MtlsMaterial) -> Self {
571 self.mtls = Some(material);
572 self.mtls_source = None;
573 self.protocol = Some(ExportProtocol::Grpc);
574 self
575 }
576
577 #[cfg(feature = "grpc-mtls")]
592 pub fn with_mtls_source(mut self, source: std::sync::Arc<dyn CertSource>) -> Self {
593 self.mtls_source = Some(source);
594 self.mtls = None;
595 self.protocol = Some(ExportProtocol::Grpc);
596 self
597 }
598
599 pub fn with_sampler(mut self, sampler: TraceSampler) -> Self {
602 self.sampler = Some(sampler);
603 self
604 }
605
606 pub fn with_metrics(mut self, enabled: bool) -> Self {
608 self.metrics = enabled;
609 self
610 }
611
612 pub fn with_runtime_metrics(mut self, enabled: bool) -> Self {
628 self.runtime_metrics = enabled;
629 self
630 }
631
632 pub fn with_default_endpoint(mut self, endpoint: impl Into<String>) -> Self {
637 self.default_endpoint = Some(endpoint.into());
638 self
639 }
640
641 pub fn with_protocol(mut self, protocol: ExportProtocol) -> Self {
645 self.protocol = Some(protocol);
646 self
647 }
648
649 pub fn with_max_export_batch_size(mut self, size: usize) -> Self {
654 self.max_export_batch_size = Some(size);
655 self
656 }
657
658 pub fn with_metric_export_interval(mut self, interval: Duration) -> Self {
663 self.metric_export_interval = Some(interval);
664 self
665 }
666
667 pub fn with_logs(mut self, enabled: bool) -> Self {
674 self.logs = enabled;
675 self
676 }
677
678 pub fn with_propagated_span_fields(mut self, fields: &'static [&'static str]) -> Self {
691 self.propagated_span_fields = fields;
692 self
693 }
694
695 pub fn with_export_timeout(mut self, timeout: Duration) -> Self {
699 self.export_timeout = Some(timeout);
700 self
701 }
702
703 pub fn with_shutdown_timeout(mut self, timeout: Duration) -> Self {
710 self.shutdown_timeout = timeout;
711 self
712 }
713
714 #[cfg(feature = "profiling")]
732 pub fn with_profiling(mut self, endpoint: &str) -> Self {
733 self.pyroscope_endpoint = Some(endpoint.to_string());
734 self
735 }
736
737 pub fn with_meter_provider_setup<F>(mut self, setup: F) -> Self
794 where
795 F: FnOnce(MeterProviderBuilder) -> MeterProviderBuilder + Send + Sync + 'static,
796 {
797 self.extra_metric_readers.push(Box::new(setup));
798 self
799 }
800
801 pub fn with_layer<L>(mut self, layer: L) -> Self
802 where
803 L: tracing_subscriber::Layer<tracing_subscriber::Registry> + Send + Sync + 'static,
804 {
805 self.extra_layers.push(Box::new(layer));
806 self
807 }
808
809 #[cfg_attr(not(feature = "grpc-mtls"), allow(unused_mut))]
824 pub fn init(mut self) -> Result<TelemetryHandles, Box<dyn Error>> {
825 let log_filter = match self.log_filter.as_deref() {
826 Some(directive) => tracing_subscriber::EnvFilter::try_new(directive)?,
827 None => tracing_subscriber::EnvFilter::from_default_env(),
828 };
829
830 if let Some(interval) = self.metric_export_interval
831 && interval.is_zero()
832 {
833 return Err("metric_export_interval must be greater than zero".into());
834 }
835
836 let protocol = self.protocol.or_else(protocol_from_env).unwrap_or({
837 #[cfg(feature = "grpc")]
838 {
839 ExportProtocol::Grpc
840 }
841 #[cfg(all(not(feature = "grpc"), feature = "http"))]
842 {
843 ExportProtocol::HttpProtobuf
844 }
845 });
846
847 let default_endpoint = match protocol {
848 #[cfg(feature = "grpc")]
849 ExportProtocol::Grpc => "http://localhost:4317",
850 #[cfg(feature = "http")]
851 ExportProtocol::HttpProtobuf => "http://localhost:4318",
852 };
853 let endpoint = resolve_endpoint(
854 std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT").ok(),
855 self.default_endpoint.as_deref(),
856 default_endpoint,
857 );
858
859 let export_timeout = self.export_timeout.or_else(timeout_from_env);
861
862 let service_name = self.service_name.unwrap_or_else(|| {
864 std::env::var("OTEL_SERVICE_NAME").unwrap_or_else(|_| "unknown_service".to_string())
865 });
866
867 let resource = build_resource(
868 &service_name,
869 self.service_version.as_deref(),
870 self.deployment_environment.as_deref(),
871 );
872
873 let sampler = match self.sampler {
874 Some(s) => s,
875 None => sampler_from_env()?.unwrap_or(TraceSampler::AlwaysOn),
876 };
877
878 #[cfg(feature = "grpc-mtls")]
879 let mtls_transport = MtlsTransport::resolve(
880 self.mtls.take(),
881 self.mtls_source.take(),
882 &endpoint,
883 export_timeout,
884 )?;
885
886 let trace_exporter = build_span_exporter(
888 protocol,
889 &endpoint,
890 export_timeout,
891 #[cfg(feature = "grpc-mtls")]
892 mtls_transport.as_ref(),
893 )?;
894
895 let batch_processor = if let Some(size) = self.max_export_batch_size {
896 BatchSpanProcessor::builder(trace_exporter)
897 .with_batch_config(
898 BatchConfigBuilder::default()
899 .with_max_export_batch_size(size)
900 .build(),
901 )
902 .build()
903 } else {
904 BatchSpanProcessor::builder(trace_exporter).build()
905 };
906
907 let tracer_provider = SdkTracerProvider::builder()
908 .with_resource(resource.clone())
909 .with_sampler(sampler.into_sdk_sampler())
910 .with_span_processor(batch_processor)
911 .build();
912
913 opentelemetry::global::set_tracer_provider(tracer_provider.clone());
914
915 let propagator = TextMapCompositePropagator::new(vec![
917 Box::new(TraceContextPropagator::new()),
918 Box::new(BaggagePropagator::new()),
919 ]);
920 opentelemetry::global::set_text_map_propagator(propagator);
921
922 let meter_provider = if self.metrics {
924 let metric_exporter = build_metric_exporter(
925 protocol,
926 &endpoint,
927 export_timeout,
928 #[cfg(feature = "grpc-mtls")]
929 mtls_transport.as_ref(),
930 )?;
931
932 let periodic_reader = if let Some(interval) = self.metric_export_interval {
933 PeriodicReader::builder(metric_exporter)
934 .with_interval(interval)
935 .build()
936 } else {
937 PeriodicReader::builder(metric_exporter).build()
938 };
939
940 let mut mp_builder = SdkMeterProvider::builder()
941 .with_resource(resource.clone())
942 .with_reader(periodic_reader);
943 for installer in self.extra_metric_readers {
944 mp_builder = installer(mp_builder);
945 }
946 let mp = mp_builder.build();
947
948 opentelemetry::global::set_meter_provider(mp.clone());
949
950 if self.runtime_metrics {
954 crate::runtime_metrics::install();
955 }
956
957 Some(mp)
958 } else {
959 None
960 };
961
962 let logger_provider = if self.logs {
964 let log_exporter = build_log_exporter(
965 protocol,
966 &endpoint,
967 export_timeout,
968 #[cfg(feature = "grpc-mtls")]
969 mtls_transport.as_ref(),
970 )?;
971
972 let lp = SdkLoggerProvider::builder()
973 .with_resource(resource)
974 .with_batch_exporter(log_exporter)
975 .build();
976
977 Some(lp)
978 } else {
979 None
980 };
981
982 #[cfg(feature = "profiling")]
984 let profiling_handle = if let Some(ref endpoint) = self.pyroscope_endpoint {
985 let identity = profiling::ProfilingIdentity {
990 host_name: hostname::get()
991 .ok()
992 .and_then(|h| h.into_string().ok())
993 .filter(|h| !h.is_empty()),
994 deployment_environment: self.deployment_environment.clone(),
995 service_version: self.service_version.clone(),
996 };
997 profiling::start_pyroscope_bridge(&service_name, endpoint, &identity)?
998 } else {
999 None
1000 };
1001 #[cfg(not(feature = "profiling"))]
1002 let _profiling_handle: Option<()> = None;
1003
1004 let extra = if self.extra_layers.is_empty() {
1010 None
1011 } else {
1012 Some(self.extra_layers)
1013 };
1014
1015 macro_rules! install_subscriber {
1016 ($fmt_layer:expr) => {{
1017 let otel_layer = tracing_opentelemetry::layer()
1021 .with_tracer(tracing_bridge_tracer(&tracer_provider));
1022 let registry = tracing_subscriber::registry()
1025 .with(extra)
1026 .with(crate::export_backoff::ExportFailureBackoff::default())
1027 .with(log_filter)
1028 .with($fmt_layer)
1029 .with(otel_layer);
1030
1031 #[cfg(feature = "profiling-bridge-pyroscope-rs")]
1034 #[allow(deprecated)]
1035 let registry = registry.with(crate::profiling::ProfilingTagLayer);
1036
1037 if let Some(lp) = &logger_provider {
1038 if let Err(e) = registry
1039 .with(crate::log_bridge::SpanAwareLogBridge::new(
1040 lp,
1041 self.propagated_span_fields,
1042 ))
1043 .try_init()
1044 {
1045 eprintln!(
1046 "otel-bootstrap: global tracing subscriber already installed — \
1047 OTLP log records will NOT be exported to the collector: {e}"
1048 );
1049 }
1050 } else if let Err(e) = registry.try_init() {
1051 eprintln!(
1052 "otel-bootstrap: global tracing subscriber already installed — \
1053 OTLP telemetry will NOT be exported to the collector: {e}"
1054 );
1055 }
1056 }};
1057 }
1058
1059 match self.log_format {
1060 LogFormat::Pretty => install_subscriber!(tracing_subscriber::fmt::layer()),
1061 LogFormat::Json => install_subscriber!(tracing_subscriber::fmt::layer().json()),
1062 }
1063
1064 let boot_owner = boot::Timeline::global()
1066 .attach(&tracer_provider, &service_name)
1067 .is_some();
1068
1069 Ok(TelemetryHandles {
1070 tracer_provider,
1071 meter_provider,
1072 logger_provider,
1073 shutdown_timeout: self.shutdown_timeout,
1074 boot_owner,
1075 #[cfg(feature = "profiling")]
1076 profiling_handle,
1077 })
1078 }
1079}
1080
1081pub fn init_telemetry(service_name: &str) -> Result<TelemetryHandles, Box<dyn Error>> {
1095 Telemetry::builder(service_name).init()
1096}
1097
1098pub fn init_telemetry_with_sampler(
1115 service_name: &str,
1116 sampler: Option<TraceSampler>,
1117) -> Result<TelemetryHandles, Box<dyn Error>> {
1118 let builder = Telemetry::builder(service_name);
1119 match sampler {
1120 Some(s) => builder.with_sampler(s),
1121 None => builder, }
1123 .init()
1124}
1125
1126fn timeout_from_env() -> Option<Duration> {
1128 let ms = std::env::var("OTEL_EXPORTER_OTLP_TIMEOUT").ok()?;
1129 let ms: u64 = ms.trim().parse().ok()?;
1130 Some(Duration::from_millis(ms))
1131}
1132
1133#[cfg(feature = "grpc-mtls")]
1137enum MtlsTransport {
1138 Snapshot(MtlsMaterial),
1139 Live(tonic::transport::Channel),
1140}
1141
1142#[cfg(feature = "grpc-mtls")]
1143impl MtlsTransport {
1144 fn resolve(
1145 material: Option<MtlsMaterial>,
1146 source: Option<std::sync::Arc<dyn CertSource>>,
1147 endpoint: &str,
1148 timeout: Option<Duration>,
1149 ) -> Result<Option<Self>, Box<dyn Error>> {
1150 match (source, material) {
1151 (Some(source), _) => {
1152 let channel = rotating_mtls::channel(endpoint, source, timeout)
1153 .map_err(|e| -> Box<dyn Error> { e.to_string().into() })?;
1154 Ok(Some(Self::Live(channel)))
1155 }
1156 (None, Some(material)) => Ok(Some(Self::Snapshot(material))),
1157 (None, None) => Ok(None),
1158 }
1159 }
1160
1161 fn apply<B: opentelemetry_otlp::WithTonicConfig>(&self, builder: B) -> B {
1162 match self {
1163 Self::Snapshot(material) => builder.with_tls_config(build_tls_config(material)),
1164 Self::Live(channel) => builder.with_channel(channel.clone()),
1165 }
1166 }
1167}
1168
1169#[cfg(feature = "grpc-mtls")]
1176fn build_tls_config(material: &MtlsMaterial) -> tonic::transport::ClientTlsConfig {
1177 use tonic::transport::{Certificate, ClientTlsConfig, Identity};
1178 ClientTlsConfig::new()
1179 .ca_certificate(Certificate::from_pem(&material.trust_bundle_pem))
1180 .identity(Identity::from_pem(
1181 &material.client_cert_chain_pem,
1182 &material.client_key_pem,
1183 ))
1184}
1185
1186fn resolve_endpoint(
1189 configured: Option<String>,
1190 runtime_default: Option<&str>,
1191 fallback: &str,
1192) -> String {
1193 configured
1194 .or_else(|| runtime_default.map(str::to_owned))
1195 .unwrap_or_else(|| fallback.to_owned())
1196}
1197
1198fn build_span_exporter(
1199 protocol: ExportProtocol,
1200 endpoint: &str,
1201 timeout: Option<Duration>,
1202 #[cfg(feature = "grpc-mtls")] mtls: Option<&MtlsTransport>,
1203) -> Result<opentelemetry_otlp::SpanExporter, Box<dyn Error>> {
1204 match protocol {
1205 #[cfg(feature = "grpc")]
1206 ExportProtocol::Grpc => {
1207 let mut b = opentelemetry_otlp::SpanExporter::builder()
1208 .with_tonic()
1209 .with_endpoint(endpoint);
1210 if let Some(t) = timeout {
1211 b = b.with_timeout(t);
1212 }
1213 #[cfg(feature = "grpc-mtls")]
1214 if let Some(m) = mtls {
1215 b = m.apply(b);
1216 }
1217 Ok(b.build()?)
1218 }
1219 #[cfg(feature = "http")]
1220 ExportProtocol::HttpProtobuf => {
1221 let mut b = opentelemetry_otlp::SpanExporter::builder()
1222 .with_http()
1223 .with_endpoint(endpoint);
1224 if let Some(t) = timeout {
1225 b = b.with_timeout(t);
1226 }
1227 Ok(b.build()?)
1228 }
1229 }
1230}
1231
1232fn build_metric_exporter(
1233 protocol: ExportProtocol,
1234 endpoint: &str,
1235 timeout: Option<Duration>,
1236 #[cfg(feature = "grpc-mtls")] mtls: Option<&MtlsTransport>,
1237) -> Result<opentelemetry_otlp::MetricExporter, Box<dyn Error>> {
1238 match protocol {
1239 #[cfg(feature = "grpc")]
1240 ExportProtocol::Grpc => {
1241 let mut b = opentelemetry_otlp::MetricExporter::builder()
1242 .with_tonic()
1243 .with_endpoint(endpoint);
1244 if let Some(t) = timeout {
1245 b = b.with_timeout(t);
1246 }
1247 #[cfg(feature = "grpc-mtls")]
1248 if let Some(m) = mtls {
1249 b = m.apply(b);
1250 }
1251 Ok(b.build()?)
1252 }
1253 #[cfg(feature = "http")]
1254 ExportProtocol::HttpProtobuf => {
1255 let mut b = opentelemetry_otlp::MetricExporter::builder()
1256 .with_http()
1257 .with_endpoint(endpoint);
1258 if let Some(t) = timeout {
1259 b = b.with_timeout(t);
1260 }
1261 Ok(b.build()?)
1262 }
1263 }
1264}
1265
1266fn build_log_exporter(
1267 protocol: ExportProtocol,
1268 endpoint: &str,
1269 timeout: Option<Duration>,
1270 #[cfg(feature = "grpc-mtls")] mtls: Option<&MtlsTransport>,
1271) -> Result<opentelemetry_otlp::LogExporter, Box<dyn Error>> {
1272 match protocol {
1273 #[cfg(feature = "grpc")]
1274 ExportProtocol::Grpc => {
1275 let mut b = opentelemetry_otlp::LogExporter::builder()
1276 .with_tonic()
1277 .with_endpoint(endpoint);
1278 if let Some(t) = timeout {
1279 b = b.with_timeout(t);
1280 }
1281 #[cfg(feature = "grpc-mtls")]
1282 if let Some(m) = mtls {
1283 b = m.apply(b);
1284 }
1285 Ok(b.build()?)
1286 }
1287 #[cfg(feature = "http")]
1288 ExportProtocol::HttpProtobuf => {
1289 let mut b = opentelemetry_otlp::LogExporter::builder()
1290 .with_http()
1291 .with_endpoint(endpoint);
1292 if let Some(t) = timeout {
1293 b = b.with_timeout(t);
1294 }
1295 Ok(b.build()?)
1296 }
1297 }
1298}
1299
1300pub fn build_resource(
1315 service_name: &str,
1316 service_version: Option<&str>,
1317 deployment_environment: Option<&str>,
1318) -> Resource {
1319 let hostname = hostname::get()
1320 .ok()
1321 .and_then(|h| h.into_string().ok())
1322 .unwrap_or_default();
1323
1324 let mut builder = Resource::builder()
1325 .with_service_name(service_name.to_string())
1326 .with_attributes([
1327 KeyValue::new(HOST_NAME, hostname),
1328 KeyValue::new(PROCESS_PID, std::process::id() as i64),
1329 ]);
1330
1331 if let Some(version) = service_version {
1332 builder = builder.with_attribute(KeyValue::new(SERVICE_VERSION, version.to_string()));
1333 }
1334
1335 if let Some(env) = deployment_environment {
1336 builder =
1337 builder.with_attribute(KeyValue::new(DEPLOYMENT_ENVIRONMENT_NAME, env.to_string()));
1338 }
1339
1340 builder.build()
1341}
1342
1343#[cfg(feature = "axum")]
1361pub fn axum_layer() -> axum_middleware::OtelTraceLayer {
1362 axum_middleware::OtelTraceLayer
1363}
1364
1365#[cfg(feature = "axum")]
1396pub fn span_enricher_layer<T>() -> axum_middleware::SpanEnricherLayer<T>
1397where
1398 T: span_enrichment::EnrichSpan + Clone + Send + Sync + 'static,
1399{
1400 axum_middleware::SpanEnricherLayer::default()
1401}
1402
1403#[cfg(feature = "tonic-tracing")]
1425pub fn grpc_client_layer() -> grpc_middleware::GrpcClientTraceLayer {
1426 grpc_middleware::GrpcClientTraceLayer
1427}
1428
1429#[cfg(feature = "tonic-tracing")]
1444pub fn grpc_server_layer() -> grpc_middleware::GrpcServerTraceLayer {
1445 grpc_middleware::GrpcServerTraceLayer
1446}
1447
1448#[cfg(test)]
1449mod tests {
1450 use super::*;
1451
1452 #[test]
1454 fn runtime_metrics_can_be_disabled() {
1455 assert!(
1456 Telemetry::builder("rm-default").runtime_metrics,
1457 "runtime metrics are on by default"
1458 );
1459 assert!(
1460 !Telemetry::builder("rm-off")
1461 .with_runtime_metrics(false)
1462 .runtime_metrics
1463 );
1464 }
1465
1466 #[tokio::test]
1474 async fn shutdown_absorbs_provider_errors() {
1475 let handles = TelemetryHandles {
1476 tracer_provider: SdkTracerProvider::builder().build(),
1477 meter_provider: Some(SdkMeterProvider::builder().build()),
1478 logger_provider: Some(SdkLoggerProvider::builder().build()),
1479 shutdown_timeout: DEFAULT_SHUTDOWN_TIMEOUT,
1480 boot_owner: false,
1481 #[cfg(feature = "profiling")]
1482 profiling_handle: None,
1483 };
1484
1485 handles.shutdown().expect("first shutdown succeeds");
1486 handles
1487 .shutdown()
1488 .expect("second shutdown absorbs the already-shut-down errors");
1489 }
1490 use opentelemetry::trace::{Span as _, Tracer as _};
1491 use std::sync::Mutex;
1492
1493 static ENV_LOCK: Mutex<()> = Mutex::new(());
1494
1495 #[test]
1496 fn tracing_bridge_uses_sdk_tracer() {
1497 let provider = SdkTracerProvider::builder().build();
1498 let tracer = tracing_bridge_tracer(&provider);
1499 let span = tracer.start("bridge-regression");
1500
1501 assert!(span.span_context().is_valid());
1502
1503 provider.shutdown().expect("provider shutdown");
1504 }
1505
1506 #[test]
1507 fn resource_contains_all_attributes_when_provided() {
1508 let resource = build_resource("test-svc", Some("1.2.3"), Some("staging"));
1509
1510 assert_eq!(
1511 resource.get(&opentelemetry::Key::new("service.name")),
1512 Some(opentelemetry::Value::from("test-svc")),
1513 );
1514 assert_eq!(
1515 resource.get(&opentelemetry::Key::new(SERVICE_VERSION)),
1516 Some(opentelemetry::Value::from("1.2.3")),
1517 );
1518 assert_eq!(
1519 resource.get(&opentelemetry::Key::new(DEPLOYMENT_ENVIRONMENT_NAME)),
1520 Some(opentelemetry::Value::from("staging")),
1521 );
1522 assert!(resource.get(&opentelemetry::Key::new(HOST_NAME)).is_some());
1523 assert!(
1524 resource
1525 .get(&opentelemetry::Key::new(PROCESS_PID))
1526 .is_some()
1527 );
1528 }
1529
1530 #[test]
1531 fn resource_graceful_when_optional_values_omitted() {
1532 let resource = build_resource("test-svc", None, None);
1533
1534 assert_eq!(
1535 resource.get(&opentelemetry::Key::new("service.name")),
1536 Some(opentelemetry::Value::from("test-svc")),
1537 );
1538 assert!(
1539 resource
1540 .get(&opentelemetry::Key::new(SERVICE_VERSION))
1541 .is_none()
1542 );
1543 assert!(
1544 resource
1545 .get(&opentelemetry::Key::new(DEPLOYMENT_ENVIRONMENT_NAME))
1546 .is_none()
1547 );
1548 assert!(resource.get(&opentelemetry::Key::new(HOST_NAME)).is_some());
1550 assert!(
1551 resource
1552 .get(&opentelemetry::Key::new(PROCESS_PID))
1553 .is_some()
1554 );
1555 }
1556
1557 #[test]
1558 fn trace_sampler_ratio_converts_to_sdk() {
1559 let sampler = TraceSampler::TraceIdRatio(0.5);
1560 let sdk = sampler.into_sdk_sampler();
1561 assert_eq!(format!("{sdk:?}"), "TraceIdRatioBased(0.5)");
1562 }
1563
1564 #[test]
1565 fn trace_sampler_parent_based_converts_to_sdk() {
1566 let sampler = TraceSampler::ParentBased(Box::new(TraceSampler::TraceIdRatio(0.25)));
1567 let sdk = sampler.into_sdk_sampler();
1568 let debug = format!("{sdk:?}");
1569 assert!(debug.contains("ParentBased"));
1570 assert!(debug.contains("0.25"));
1571 }
1572
1573 unsafe fn set_env(key: &str, val: &str) {
1575 unsafe {
1576 std::env::set_var(key, val);
1577 }
1578 }
1579
1580 unsafe fn remove_env(key: &str) {
1581 unsafe {
1582 std::env::remove_var(key);
1583 }
1584 }
1585
1586 #[test]
1587 fn sampler_from_env_reads_traceidratio() {
1588 let _lock = ENV_LOCK.lock().unwrap();
1589 unsafe {
1590 set_env("OTEL_TRACES_SAMPLER", "traceidratio");
1591 set_env("OTEL_TRACES_SAMPLER_ARG", "0.42");
1592 }
1593
1594 let sampler = sampler_from_env()
1595 .expect("should not error")
1596 .expect("should return Some");
1597 assert!(
1598 matches!(sampler, TraceSampler::TraceIdRatio(r) if (r - 0.42).abs() < f64::EPSILON)
1599 );
1600
1601 unsafe {
1602 remove_env("OTEL_TRACES_SAMPLER");
1603 remove_env("OTEL_TRACES_SAMPLER_ARG");
1604 }
1605 }
1606
1607 #[test]
1608 fn sampler_from_env_returns_none_when_unset() {
1609 let _lock = ENV_LOCK.lock().unwrap();
1610 unsafe {
1611 remove_env("OTEL_TRACES_SAMPLER");
1612 }
1613 assert!(sampler_from_env().expect("should not error").is_none());
1614 }
1615
1616 #[test]
1617 fn sampler_from_env_reads_parentbased_traceidratio() {
1618 let _lock = ENV_LOCK.lock().unwrap();
1619 unsafe {
1620 set_env("OTEL_TRACES_SAMPLER", "parentbased_traceidratio");
1621 set_env("OTEL_TRACES_SAMPLER_ARG", "0.1");
1622 }
1623
1624 let sampler = sampler_from_env()
1625 .expect("should not error")
1626 .expect("should return Some");
1627 assert!(
1628 matches!(sampler, TraceSampler::ParentBased(inner) if matches!(*inner, TraceSampler::TraceIdRatio(r) if (r - 0.1).abs() < f64::EPSILON))
1629 );
1630
1631 unsafe {
1632 remove_env("OTEL_TRACES_SAMPLER");
1633 remove_env("OTEL_TRACES_SAMPLER_ARG");
1634 }
1635 }
1636
1637 #[test]
1638 fn sampler_from_env_parentbased_always_on() {
1639 let _lock = ENV_LOCK.lock().unwrap();
1640 unsafe {
1641 set_env("OTEL_TRACES_SAMPLER", "parentbased_always_on");
1642 }
1643 let sampler = sampler_from_env()
1644 .expect("should not error")
1645 .expect("should return Some");
1646 assert!(
1647 matches!(sampler, TraceSampler::ParentBased(inner) if matches!(*inner, TraceSampler::AlwaysOn))
1648 );
1649 unsafe {
1650 remove_env("OTEL_TRACES_SAMPLER");
1651 }
1652 }
1653
1654 #[test]
1655 fn sampler_from_env_parentbased_always_off() {
1656 let _lock = ENV_LOCK.lock().unwrap();
1657 unsafe {
1658 set_env("OTEL_TRACES_SAMPLER", "parentbased_always_off");
1659 }
1660 let sampler = sampler_from_env()
1661 .expect("should not error")
1662 .expect("should return Some");
1663 assert!(
1664 matches!(sampler, TraceSampler::ParentBased(inner) if matches!(*inner, TraceSampler::AlwaysOff))
1665 );
1666 unsafe {
1667 remove_env("OTEL_TRACES_SAMPLER");
1668 }
1669 }
1670
1671 #[test]
1672 fn sampler_from_env_always_on() {
1673 let _lock = ENV_LOCK.lock().unwrap();
1674 unsafe {
1675 set_env("OTEL_TRACES_SAMPLER", "always_on");
1676 }
1677 let sampler = sampler_from_env()
1678 .expect("should not error")
1679 .expect("should return Some");
1680 assert!(matches!(sampler, TraceSampler::AlwaysOn));
1681 unsafe {
1682 remove_env("OTEL_TRACES_SAMPLER");
1683 }
1684 }
1685
1686 #[test]
1687 fn sampler_from_env_always_off() {
1688 let _lock = ENV_LOCK.lock().unwrap();
1689 unsafe {
1690 set_env("OTEL_TRACES_SAMPLER", "always_off");
1691 }
1692 let sampler = sampler_from_env()
1693 .expect("should not error")
1694 .expect("should return Some");
1695 assert!(matches!(sampler, TraceSampler::AlwaysOff));
1696 unsafe {
1697 remove_env("OTEL_TRACES_SAMPLER");
1698 }
1699 }
1700
1701 #[test]
1702 fn sampler_from_env_unknown_returns_error() {
1703 let _lock = ENV_LOCK.lock().unwrap();
1704 unsafe {
1705 set_env("OTEL_TRACES_SAMPLER", "unknown_sampler");
1706 }
1707 let err = sampler_from_env().expect_err("unknown sampler should produce an error");
1708 assert!(
1709 err.to_string().contains("unknown_sampler"),
1710 "error message should include the unknown name, got: {err}"
1711 );
1712 unsafe {
1713 remove_env("OTEL_TRACES_SAMPLER");
1714 }
1715 }
1716
1717 #[test]
1718 fn trace_sampler_always_on_converts_to_sdk() {
1719 let sdk = TraceSampler::AlwaysOn.into_sdk_sampler();
1720 assert_eq!(format!("{sdk:?}"), "AlwaysOn");
1721 }
1722
1723 #[test]
1724 fn trace_sampler_always_off_converts_to_sdk() {
1725 let sdk = TraceSampler::AlwaysOff.into_sdk_sampler();
1726 assert_eq!(format!("{sdk:?}"), "AlwaysOff");
1727 }
1728
1729 #[test]
1730 fn builder_has_sensible_defaults() {
1731 let builder = Telemetry::builder("test-svc");
1732 assert_eq!(builder.service_name.as_deref(), Some("test-svc"));
1733 assert!(builder.service_version.is_none());
1734 assert!(builder.deployment_environment.is_none());
1735 assert!(builder.sampler.is_none());
1736 assert!(builder.metrics);
1737 assert!(!builder.logs);
1738 assert!(builder.protocol.is_none());
1739 assert!(builder.max_export_batch_size.is_none());
1740 assert!(builder.metric_export_interval.is_none());
1741 assert!(builder.export_timeout.is_none());
1742 }
1743
1744 #[test]
1745 fn from_env_builder_has_no_service_name() {
1746 let builder = Telemetry::from_env();
1747 assert!(builder.service_name.is_none());
1748 }
1749
1750 #[test]
1751 fn with_export_timeout_stores_value() {
1752 let timeout = Duration::from_secs(5);
1753 let builder = Telemetry::builder("test-svc").with_export_timeout(timeout);
1754 assert_eq!(builder.export_timeout, Some(timeout));
1755 }
1756
1757 #[test]
1758 fn timeout_from_env_reads_milliseconds() {
1759 let _lock = ENV_LOCK.lock().unwrap();
1760 unsafe {
1761 set_env("OTEL_EXPORTER_OTLP_TIMEOUT", "5000");
1762 }
1763 let t = timeout_from_env();
1764 assert_eq!(t, Some(Duration::from_millis(5000)));
1765 unsafe {
1766 remove_env("OTEL_EXPORTER_OTLP_TIMEOUT");
1767 }
1768 }
1769
1770 #[test]
1771 fn timeout_from_env_returns_none_when_unset() {
1772 let _lock = ENV_LOCK.lock().unwrap();
1773 unsafe {
1774 remove_env("OTEL_EXPORTER_OTLP_TIMEOUT");
1775 }
1776 assert_eq!(timeout_from_env(), None);
1777 }
1778
1779 #[test]
1780 fn service_name_from_env_used_when_none_given() {
1781 let builder = Telemetry::from_env();
1782 assert!(builder.service_name.is_none());
1783 }
1784
1785 #[test]
1786 fn explicit_service_name_overrides_env_var() {
1787 let builder = Telemetry::builder("explicit-svc");
1788 assert_eq!(builder.service_name.as_deref(), Some("explicit-svc"));
1789 }
1790
1791 #[test]
1792 fn from_env_builder_service_name_is_none() {
1793 let builder = Telemetry::from_env();
1794 assert!(builder.service_name.is_none());
1795 }
1796
1797 #[test]
1798 fn init_returns_error_for_unknown_otel_traces_sampler() {
1799 let _lock = ENV_LOCK.lock().unwrap();
1800 unsafe {
1801 set_env("OTEL_TRACES_SAMPLER", "not_a_real_sampler");
1802 }
1803 let result = Telemetry::builder("test-svc").with_metrics(false).init();
1804 let err = result
1805 .err()
1806 .expect("unknown sampler env var should cause init to fail");
1807 assert!(
1808 err.to_string().contains("not_a_real_sampler"),
1809 "error should name the unknown sampler, got: {err}"
1810 );
1811 unsafe {
1812 remove_env("OTEL_TRACES_SAMPLER");
1813 }
1814 }
1815
1816 #[test]
1817 fn with_max_export_batch_size_stores_value() {
1818 let builder = Telemetry::builder("test-svc").with_max_export_batch_size(1024);
1819 assert_eq!(builder.max_export_batch_size, Some(1024));
1820 }
1821
1822 #[test]
1823 fn with_metric_export_interval_stores_value() {
1824 let interval = Duration::from_secs(30);
1825 let builder = Telemetry::builder("test-svc").with_metric_export_interval(interval);
1826 assert_eq!(builder.metric_export_interval, Some(interval));
1827 }
1828
1829 #[test]
1830 fn init_rejects_zero_metric_export_interval() {
1831 let err = Telemetry::builder("test-svc")
1832 .with_metric_export_interval(Duration::ZERO)
1833 .with_metrics(false)
1834 .init()
1835 .err()
1836 .expect("expected error for zero interval");
1837 assert!(
1838 err.to_string().contains("metric_export_interval"),
1839 "error message should mention metric_export_interval, got: {err}"
1840 );
1841 }
1842
1843 #[test]
1844 fn builder_with_custom_values() {
1845 let builder = Telemetry::builder("test-svc")
1846 .with_version("2.0.0")
1847 .with_environment("production")
1848 .with_sampler(TraceSampler::TraceIdRatio(0.5))
1849 .with_metrics(false);
1850
1851 assert_eq!(builder.service_name.as_deref(), Some("test-svc"));
1852 assert_eq!(builder.service_version.as_deref(), Some("2.0.0"));
1853 assert_eq!(
1854 builder.deployment_environment.as_deref(),
1855 Some("production")
1856 );
1857 assert!(
1858 matches!(builder.sampler, Some(TraceSampler::TraceIdRatio(r)) if (r - 0.5).abs() < f64::EPSILON)
1859 );
1860 assert!(!builder.metrics);
1861 }
1862
1863 #[test]
1864 fn builder_stores_programmatic_log_configuration() {
1865 let builder = Telemetry::builder("test-svc")
1866 .with_log_filter("info,opentelemetry_sdk=warn")
1867 .with_log_format(LogFormat::Json);
1868
1869 assert_eq!(
1870 builder.log_filter.as_deref(),
1871 Some("info,opentelemetry_sdk=warn")
1872 );
1873 assert_eq!(builder.log_format, LogFormat::Json);
1874 }
1875
1876 #[test]
1877 fn init_rejects_invalid_programmatic_log_filter_before_provider_setup() {
1878 let setup_ran = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
1879 let setup_ran_in_closure = std::sync::Arc::clone(&setup_ran);
1880
1881 let error = Telemetry::builder("test-svc")
1882 .with_log_filter("[")
1883 .with_meter_provider_setup(move |builder| {
1884 setup_ran_in_closure.store(true, std::sync::atomic::Ordering::SeqCst);
1885 builder
1886 })
1887 .init()
1888 .err()
1889 .expect("invalid filter must fail initialization");
1890
1891 assert!(error.to_string().contains("invalid filter directive"));
1892 assert!(!setup_ran.load(std::sync::atomic::Ordering::SeqCst));
1893 }
1894
1895 #[test]
1896 fn builder_with_default_endpoint() {
1897 let builder = Telemetry::builder("svc").with_default_endpoint("http://otel-collector:4317");
1898 assert_eq!(
1899 builder.default_endpoint.as_deref(),
1900 Some("http://otel-collector:4317")
1901 );
1902 }
1903
1904 #[test]
1905 fn a_configured_endpoint_wins_over_the_runtime_default() {
1906 assert_eq!(
1907 resolve_endpoint(
1908 Some("http://c:4317".into()),
1909 Some("http://d:4317"),
1910 "http://localhost:4317"
1911 ),
1912 "http://c:4317"
1913 );
1914 assert_eq!(
1915 resolve_endpoint(None, Some("http://d:4317"), "http://localhost:4317"),
1916 "http://d:4317"
1917 );
1918 assert_eq!(
1919 resolve_endpoint(None, None, "http://localhost:4317"),
1920 "http://localhost:4317"
1921 );
1922 }
1923
1924 #[test]
1925 #[cfg(feature = "grpc")]
1926 fn builder_with_protocol_grpc() {
1927 let builder = Telemetry::builder("test-svc").with_protocol(ExportProtocol::Grpc);
1928 assert_eq!(builder.protocol, Some(ExportProtocol::Grpc));
1929 }
1930
1931 #[test]
1932 #[cfg(feature = "http")]
1933 fn builder_with_protocol_http() {
1934 let builder = Telemetry::builder("test-svc").with_protocol(ExportProtocol::HttpProtobuf);
1935 assert_eq!(builder.protocol, Some(ExportProtocol::HttpProtobuf));
1936 }
1937
1938 #[test]
1939 #[cfg(feature = "grpc")]
1940 fn protocol_from_env_reads_grpc() {
1941 let _lock = ENV_LOCK.lock().unwrap();
1942 unsafe {
1943 set_env("OTEL_EXPORTER_OTLP_PROTOCOL", "grpc");
1944 }
1945 assert_eq!(protocol_from_env(), Some(ExportProtocol::Grpc));
1946 unsafe {
1947 remove_env("OTEL_EXPORTER_OTLP_PROTOCOL");
1948 }
1949 }
1950
1951 #[test]
1952 #[cfg(feature = "http")]
1953 fn protocol_from_env_reads_http_protobuf() {
1954 let _lock = ENV_LOCK.lock().unwrap();
1955 unsafe {
1956 set_env("OTEL_EXPORTER_OTLP_PROTOCOL", "http/protobuf");
1957 }
1958 assert_eq!(protocol_from_env(), Some(ExportProtocol::HttpProtobuf));
1959 unsafe {
1960 remove_env("OTEL_EXPORTER_OTLP_PROTOCOL");
1961 }
1962 }
1963
1964 #[test]
1965 fn protocol_from_env_returns_none_when_unset() {
1966 let _lock = ENV_LOCK.lock().unwrap();
1967 unsafe {
1968 remove_env("OTEL_EXPORTER_OTLP_PROTOCOL");
1969 }
1970 assert_eq!(protocol_from_env(), None);
1971 }
1972
1973 #[test]
1974 fn protocol_from_env_returns_none_for_unknown() {
1975 let _lock = ENV_LOCK.lock().unwrap();
1976 unsafe {
1977 set_env("OTEL_EXPORTER_OTLP_PROTOCOL", "websocket");
1978 }
1979 assert_eq!(protocol_from_env(), None);
1980 unsafe {
1981 remove_env("OTEL_EXPORTER_OTLP_PROTOCOL");
1982 }
1983 }
1984
1985 #[test]
1986 fn builder_is_send_and_sync() {
1987 fn assert_send_sync<T: Send + Sync>() {}
1988 assert_send_sync::<TelemetryBuilder>();
1989 }
1990
1991 #[test]
1992 fn with_shutdown_timeout_stores_value() {
1993 let timeout = Duration::from_secs(10);
1994 let builder = Telemetry::builder("test-svc").with_shutdown_timeout(timeout);
1995 assert_eq!(builder.shutdown_timeout, timeout);
1996 }
1997
1998 #[test]
1999 fn default_shutdown_timeout_is_five_seconds() {
2000 let builder = Telemetry::builder("test-svc");
2001 assert_eq!(builder.shutdown_timeout, Duration::from_secs(5));
2002 }
2003
2004 #[cfg(feature = "testing")]
2012 #[test]
2013 fn drop_completes_within_shutdown_timeout() {
2014 let mut handles = crate::Telemetry::testing("drop-timeout-test");
2016 handles.shutdown_timeout = Duration::from_millis(100);
2018
2019 let start = std::time::Instant::now();
2020 drop(handles);
2021 let elapsed = start.elapsed();
2022
2023 assert!(
2025 elapsed < Duration::from_millis(500),
2026 "drop took {elapsed:?}, expected < 500 ms"
2027 );
2028 }
2029}