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 #[cfg(feature = "grpc-mtls")]
427 collector_spiffe_id: None,
428 propagated_span_fields: crate::log_bridge::PROPAGATED_SPAN_FIELDS,
429 #[cfg(feature = "profiling")]
430 pyroscope_endpoint: None,
431 }
432 }
433
434 pub fn from_env() -> TelemetryBuilder {
444 TelemetryBuilder {
445 service_name: None,
446 service_version: None,
447 deployment_environment: None,
448 sampler: None,
449 metrics: true,
450 logs: false,
451 protocol: None,
452 default_endpoint: None,
453 max_export_batch_size: None,
454 metric_export_interval: None,
455 export_timeout: None,
456 shutdown_timeout: DEFAULT_SHUTDOWN_TIMEOUT,
457 log_filter: None,
458 log_format: LogFormat::default(),
459 extra_layers: Vec::new(),
460 extra_metric_readers: Vec::new(),
461 runtime_metrics: true,
462 #[cfg(feature = "grpc-mtls")]
463 mtls: None,
464 #[cfg(feature = "grpc-mtls")]
465 mtls_source: None,
466 #[cfg(feature = "grpc-mtls")]
467 collector_spiffe_id: None,
468 propagated_span_fields: crate::log_bridge::PROPAGATED_SPAN_FIELDS,
469 #[cfg(feature = "profiling")]
470 pyroscope_endpoint: None,
471 }
472 }
473}
474
475#[must_use = "a TelemetryBuilder does nothing until .init() is called"]
493pub struct TelemetryBuilder {
494 service_name: Option<String>,
495 service_version: Option<String>,
496 deployment_environment: Option<String>,
497 sampler: Option<TraceSampler>,
498 metrics: bool,
499 logs: bool,
500 protocol: Option<ExportProtocol>,
501 max_export_batch_size: Option<usize>,
502 metric_export_interval: Option<Duration>,
503 export_timeout: Option<Duration>,
504 shutdown_timeout: Duration,
505 log_filter: Option<String>,
506 log_format: LogFormat,
507 extra_layers: Vec<
508 Box<dyn tracing_subscriber::Layer<tracing_subscriber::Registry> + Send + Sync + 'static>,
509 >,
510 extra_metric_readers: Vec<MeterProviderInstaller>,
511 runtime_metrics: bool,
512 #[cfg(feature = "grpc-mtls")]
513 mtls: Option<MtlsMaterial>,
514 #[cfg(feature = "grpc-mtls")]
515 mtls_source: Option<std::sync::Arc<dyn CertSource>>,
516 #[cfg(feature = "grpc-mtls")]
517 collector_spiffe_id: Option<String>,
518 propagated_span_fields: &'static [&'static str],
519 #[cfg(feature = "profiling")]
520 pyroscope_endpoint: Option<String>,
521 default_endpoint: Option<String>,
522}
523
524type MeterProviderInstaller =
529 Box<dyn FnOnce(MeterProviderBuilder) -> MeterProviderBuilder + Send + Sync>;
530
531impl TelemetryBuilder {
532 pub fn with_log_filter(mut self, directive: impl Into<String>) -> Self {
537 self.log_filter = Some(directive.into());
538 self
539 }
540
541 pub fn with_log_format(mut self, format: LogFormat) -> Self {
543 self.log_format = format;
544 self
545 }
546
547 pub fn with_version(mut self, version: &str) -> Self {
549 self.service_version = Some(version.to_string());
550 self
551 }
552
553 pub fn with_environment(mut self, environment: &str) -> Self {
555 self.deployment_environment = Some(environment.to_string());
556 self
557 }
558
559 #[cfg(feature = "grpc-mtls")]
576 pub fn with_mtls(mut self, material: MtlsMaterial) -> Self {
577 self.mtls = Some(material);
578 self.mtls_source = None;
579 self.protocol = Some(ExportProtocol::Grpc);
580 self
581 }
582
583 #[cfg(feature = "grpc-mtls")]
598 pub fn with_mtls_source(mut self, source: std::sync::Arc<dyn CertSource>) -> Self {
599 self.mtls_source = Some(source);
600 self.mtls = None;
601 self.protocol = Some(ExportProtocol::Grpc);
602 self
603 }
604
605 #[cfg(feature = "grpc-mtls")]
618 pub fn with_collector_spiffe_id(mut self, id: impl Into<String>) -> Self {
619 self.collector_spiffe_id = Some(id.into());
620 self
621 }
622
623 pub fn with_sampler(mut self, sampler: TraceSampler) -> Self {
626 self.sampler = Some(sampler);
627 self
628 }
629
630 pub fn with_metrics(mut self, enabled: bool) -> Self {
632 self.metrics = enabled;
633 self
634 }
635
636 pub fn with_runtime_metrics(mut self, enabled: bool) -> Self {
652 self.runtime_metrics = enabled;
653 self
654 }
655
656 pub fn with_default_endpoint(mut self, endpoint: impl Into<String>) -> Self {
661 self.default_endpoint = Some(endpoint.into());
662 self
663 }
664
665 pub fn with_protocol(mut self, protocol: ExportProtocol) -> Self {
669 self.protocol = Some(protocol);
670 self
671 }
672
673 pub fn with_max_export_batch_size(mut self, size: usize) -> Self {
678 self.max_export_batch_size = Some(size);
679 self
680 }
681
682 pub fn with_metric_export_interval(mut self, interval: Duration) -> Self {
687 self.metric_export_interval = Some(interval);
688 self
689 }
690
691 pub fn with_logs(mut self, enabled: bool) -> Self {
698 self.logs = enabled;
699 self
700 }
701
702 pub fn with_propagated_span_fields(mut self, fields: &'static [&'static str]) -> Self {
715 self.propagated_span_fields = fields;
716 self
717 }
718
719 pub fn with_export_timeout(mut self, timeout: Duration) -> Self {
723 self.export_timeout = Some(timeout);
724 self
725 }
726
727 pub fn with_shutdown_timeout(mut self, timeout: Duration) -> Self {
734 self.shutdown_timeout = timeout;
735 self
736 }
737
738 #[cfg(feature = "profiling")]
756 pub fn with_profiling(mut self, endpoint: &str) -> Self {
757 self.pyroscope_endpoint = Some(endpoint.to_string());
758 self
759 }
760
761 pub fn with_meter_provider_setup<F>(mut self, setup: F) -> Self
818 where
819 F: FnOnce(MeterProviderBuilder) -> MeterProviderBuilder + Send + Sync + 'static,
820 {
821 self.extra_metric_readers.push(Box::new(setup));
822 self
823 }
824
825 pub fn with_layer<L>(mut self, layer: L) -> Self
826 where
827 L: tracing_subscriber::Layer<tracing_subscriber::Registry> + Send + Sync + 'static,
828 {
829 self.extra_layers.push(Box::new(layer));
830 self
831 }
832
833 #[cfg_attr(not(feature = "grpc-mtls"), allow(unused_mut))]
848 pub fn init(mut self) -> Result<TelemetryHandles, Box<dyn Error>> {
849 let log_filter = match self.log_filter.as_deref() {
850 Some(directive) => tracing_subscriber::EnvFilter::try_new(directive)?,
851 None => tracing_subscriber::EnvFilter::from_default_env(),
852 };
853
854 if let Some(interval) = self.metric_export_interval
855 && interval.is_zero()
856 {
857 return Err("metric_export_interval must be greater than zero".into());
858 }
859
860 let protocol = self.protocol.or_else(protocol_from_env).unwrap_or({
861 #[cfg(feature = "grpc")]
862 {
863 ExportProtocol::Grpc
864 }
865 #[cfg(all(not(feature = "grpc"), feature = "http"))]
866 {
867 ExportProtocol::HttpProtobuf
868 }
869 });
870
871 let default_endpoint = match protocol {
872 #[cfg(feature = "grpc")]
873 ExportProtocol::Grpc => "http://localhost:4317",
874 #[cfg(feature = "http")]
875 ExportProtocol::HttpProtobuf => "http://localhost:4318",
876 };
877 let endpoint = resolve_endpoint(
878 std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT").ok(),
879 self.default_endpoint.as_deref(),
880 default_endpoint,
881 );
882
883 let export_timeout = self.export_timeout.or_else(timeout_from_env);
885
886 let service_name = self.service_name.unwrap_or_else(|| {
888 std::env::var("OTEL_SERVICE_NAME").unwrap_or_else(|_| "unknown_service".to_string())
889 });
890
891 let resource = build_resource(
892 &service_name,
893 self.service_version.as_deref(),
894 self.deployment_environment.as_deref(),
895 );
896
897 let sampler = match self.sampler {
898 Some(s) => s,
899 None => sampler_from_env()?.unwrap_or(TraceSampler::AlwaysOn),
900 };
901
902 #[cfg(feature = "grpc-mtls")]
903 let mtls_transport = MtlsTransport::resolve(
904 self.mtls.take(),
905 self.mtls_source.take(),
906 &endpoint,
907 export_timeout,
908 self.collector_spiffe_id.as_deref(),
909 )?;
910
911 let trace_exporter = build_span_exporter(
913 protocol,
914 &endpoint,
915 export_timeout,
916 #[cfg(feature = "grpc-mtls")]
917 mtls_transport.as_ref(),
918 )?;
919
920 let batch_processor = if let Some(size) = self.max_export_batch_size {
921 BatchSpanProcessor::builder(trace_exporter)
922 .with_batch_config(
923 BatchConfigBuilder::default()
924 .with_max_export_batch_size(size)
925 .build(),
926 )
927 .build()
928 } else {
929 BatchSpanProcessor::builder(trace_exporter).build()
930 };
931
932 let tracer_provider = SdkTracerProvider::builder()
933 .with_resource(resource.clone())
934 .with_sampler(sampler.into_sdk_sampler())
935 .with_span_processor(batch_processor)
936 .build();
937
938 opentelemetry::global::set_tracer_provider(tracer_provider.clone());
939
940 let propagator = TextMapCompositePropagator::new(vec![
942 Box::new(TraceContextPropagator::new()),
943 Box::new(BaggagePropagator::new()),
944 ]);
945 opentelemetry::global::set_text_map_propagator(propagator);
946
947 let meter_provider = if self.metrics {
949 let metric_exporter = build_metric_exporter(
950 protocol,
951 &endpoint,
952 export_timeout,
953 #[cfg(feature = "grpc-mtls")]
954 mtls_transport.as_ref(),
955 )?;
956
957 let periodic_reader = if let Some(interval) = self.metric_export_interval {
958 PeriodicReader::builder(metric_exporter)
959 .with_interval(interval)
960 .build()
961 } else {
962 PeriodicReader::builder(metric_exporter).build()
963 };
964
965 let mut mp_builder = SdkMeterProvider::builder()
966 .with_resource(resource.clone())
967 .with_reader(periodic_reader);
968 for installer in self.extra_metric_readers {
969 mp_builder = installer(mp_builder);
970 }
971 let mp = mp_builder.build();
972
973 opentelemetry::global::set_meter_provider(mp.clone());
974
975 if self.runtime_metrics {
979 crate::runtime_metrics::install();
980 }
981
982 Some(mp)
983 } else {
984 None
985 };
986
987 let logger_provider = if self.logs {
989 let log_exporter = build_log_exporter(
990 protocol,
991 &endpoint,
992 export_timeout,
993 #[cfg(feature = "grpc-mtls")]
994 mtls_transport.as_ref(),
995 )?;
996
997 let lp = SdkLoggerProvider::builder()
998 .with_resource(resource)
999 .with_batch_exporter(log_exporter)
1000 .build();
1001
1002 Some(lp)
1003 } else {
1004 None
1005 };
1006
1007 #[cfg(feature = "profiling")]
1009 let profiling_handle = if let Some(ref endpoint) = self.pyroscope_endpoint {
1010 let identity = profiling::ProfilingIdentity {
1015 host_name: hostname::get()
1016 .ok()
1017 .and_then(|h| h.into_string().ok())
1018 .filter(|h| !h.is_empty()),
1019 deployment_environment: self.deployment_environment.clone(),
1020 service_version: self.service_version.clone(),
1021 };
1022 profiling::start_pyroscope_bridge(&service_name, endpoint, &identity)?
1023 } else {
1024 None
1025 };
1026 #[cfg(not(feature = "profiling"))]
1027 let _profiling_handle: Option<()> = None;
1028
1029 let extra = if self.extra_layers.is_empty() {
1035 None
1036 } else {
1037 Some(self.extra_layers)
1038 };
1039
1040 macro_rules! install_subscriber {
1041 ($fmt_layer:expr) => {{
1042 let otel_layer = tracing_opentelemetry::layer()
1046 .with_tracer(tracing_bridge_tracer(&tracer_provider));
1047 let registry = tracing_subscriber::registry()
1050 .with(extra)
1051 .with(crate::export_backoff::ExportFailureBackoff::default())
1052 .with(log_filter)
1053 .with($fmt_layer)
1054 .with(otel_layer);
1055
1056 #[cfg(feature = "profiling-bridge-pyroscope-rs")]
1059 #[allow(deprecated)]
1060 let registry = registry.with(crate::profiling::ProfilingTagLayer);
1061
1062 if let Some(lp) = &logger_provider {
1063 if let Err(e) = registry
1064 .with(crate::log_bridge::SpanAwareLogBridge::new(
1065 lp,
1066 self.propagated_span_fields,
1067 ))
1068 .try_init()
1069 {
1070 eprintln!(
1071 "otel-bootstrap: global tracing subscriber already installed — \
1072 OTLP log records will NOT be exported to the collector: {e}"
1073 );
1074 }
1075 } else if let Err(e) = registry.try_init() {
1076 eprintln!(
1077 "otel-bootstrap: global tracing subscriber already installed — \
1078 OTLP telemetry will NOT be exported to the collector: {e}"
1079 );
1080 }
1081 }};
1082 }
1083
1084 match self.log_format {
1085 LogFormat::Pretty => install_subscriber!(tracing_subscriber::fmt::layer()),
1086 LogFormat::Json => install_subscriber!(tracing_subscriber::fmt::layer().json()),
1087 }
1088
1089 let boot_owner = boot::Timeline::global()
1091 .attach(&tracer_provider, &service_name)
1092 .is_some();
1093
1094 Ok(TelemetryHandles {
1095 tracer_provider,
1096 meter_provider,
1097 logger_provider,
1098 shutdown_timeout: self.shutdown_timeout,
1099 boot_owner,
1100 #[cfg(feature = "profiling")]
1101 profiling_handle,
1102 })
1103 }
1104}
1105
1106pub fn init_telemetry(service_name: &str) -> Result<TelemetryHandles, Box<dyn Error>> {
1120 Telemetry::builder(service_name).init()
1121}
1122
1123pub fn init_telemetry_with_sampler(
1140 service_name: &str,
1141 sampler: Option<TraceSampler>,
1142) -> Result<TelemetryHandles, Box<dyn Error>> {
1143 let builder = Telemetry::builder(service_name);
1144 match sampler {
1145 Some(s) => builder.with_sampler(s),
1146 None => builder, }
1148 .init()
1149}
1150
1151fn timeout_from_env() -> Option<Duration> {
1153 let ms = std::env::var("OTEL_EXPORTER_OTLP_TIMEOUT").ok()?;
1154 let ms: u64 = ms.trim().parse().ok()?;
1155 Some(Duration::from_millis(ms))
1156}
1157
1158#[cfg(feature = "grpc-mtls")]
1162enum MtlsTransport {
1163 Snapshot(MtlsMaterial),
1164 Live(tonic::transport::Channel),
1165}
1166
1167#[cfg(feature = "grpc-mtls")]
1168impl MtlsTransport {
1169 fn resolve(
1170 material: Option<MtlsMaterial>,
1171 source: Option<std::sync::Arc<dyn CertSource>>,
1172 endpoint: &str,
1173 timeout: Option<Duration>,
1174 expected_id: Option<&str>,
1175 ) -> Result<Option<Self>, Box<dyn Error>> {
1176 let source = match (source, material.as_ref(), expected_id) {
1179 (Some(source), _, _) => Some(source),
1180 (None, Some(material), Some(_)) => {
1181 Some(std::sync::Arc::new(StaticCertSource::new(material.clone()))
1182 as std::sync::Arc<dyn CertSource>)
1183 }
1184 _ => None,
1185 };
1186 match (source, material) {
1187 (Some(source), _) => {
1188 let channel = rotating_mtls::channel(endpoint, source, timeout, expected_id)
1189 .map_err(|e| -> Box<dyn Error> { e.to_string().into() })?;
1190 Ok(Some(Self::Live(channel)))
1191 }
1192 (None, Some(material)) => Ok(Some(Self::Snapshot(material))),
1193 (None, None) => Ok(None),
1194 }
1195 }
1196
1197 fn apply<B: opentelemetry_otlp::WithTonicConfig>(&self, builder: B) -> B {
1198 match self {
1199 Self::Snapshot(material) => builder.with_tls_config(build_tls_config(material)),
1200 Self::Live(channel) => builder.with_channel(channel.clone()),
1201 }
1202 }
1203}
1204
1205#[cfg(feature = "grpc-mtls")]
1212fn build_tls_config(material: &MtlsMaterial) -> tonic::transport::ClientTlsConfig {
1213 use tonic::transport::{Certificate, ClientTlsConfig, Identity};
1214 ClientTlsConfig::new()
1215 .ca_certificate(Certificate::from_pem(&material.trust_bundle_pem))
1216 .identity(Identity::from_pem(
1217 &material.client_cert_chain_pem,
1218 &material.client_key_pem,
1219 ))
1220}
1221
1222fn resolve_endpoint(
1225 configured: Option<String>,
1226 runtime_default: Option<&str>,
1227 fallback: &str,
1228) -> String {
1229 configured
1230 .or_else(|| runtime_default.map(str::to_owned))
1231 .unwrap_or_else(|| fallback.to_owned())
1232}
1233
1234fn build_span_exporter(
1235 protocol: ExportProtocol,
1236 endpoint: &str,
1237 timeout: Option<Duration>,
1238 #[cfg(feature = "grpc-mtls")] mtls: Option<&MtlsTransport>,
1239) -> Result<opentelemetry_otlp::SpanExporter, Box<dyn Error>> {
1240 match protocol {
1241 #[cfg(feature = "grpc")]
1242 ExportProtocol::Grpc => {
1243 let mut b = opentelemetry_otlp::SpanExporter::builder()
1244 .with_tonic()
1245 .with_endpoint(endpoint);
1246 if let Some(t) = timeout {
1247 b = b.with_timeout(t);
1248 }
1249 #[cfg(feature = "grpc-mtls")]
1250 if let Some(m) = mtls {
1251 b = m.apply(b);
1252 }
1253 Ok(b.build()?)
1254 }
1255 #[cfg(feature = "http")]
1256 ExportProtocol::HttpProtobuf => {
1257 let mut b = opentelemetry_otlp::SpanExporter::builder()
1258 .with_http()
1259 .with_endpoint(endpoint);
1260 if let Some(t) = timeout {
1261 b = b.with_timeout(t);
1262 }
1263 Ok(b.build()?)
1264 }
1265 }
1266}
1267
1268fn build_metric_exporter(
1269 protocol: ExportProtocol,
1270 endpoint: &str,
1271 timeout: Option<Duration>,
1272 #[cfg(feature = "grpc-mtls")] mtls: Option<&MtlsTransport>,
1273) -> Result<opentelemetry_otlp::MetricExporter, Box<dyn Error>> {
1274 match protocol {
1275 #[cfg(feature = "grpc")]
1276 ExportProtocol::Grpc => {
1277 let mut b = opentelemetry_otlp::MetricExporter::builder()
1278 .with_tonic()
1279 .with_endpoint(endpoint);
1280 if let Some(t) = timeout {
1281 b = b.with_timeout(t);
1282 }
1283 #[cfg(feature = "grpc-mtls")]
1284 if let Some(m) = mtls {
1285 b = m.apply(b);
1286 }
1287 Ok(b.build()?)
1288 }
1289 #[cfg(feature = "http")]
1290 ExportProtocol::HttpProtobuf => {
1291 let mut b = opentelemetry_otlp::MetricExporter::builder()
1292 .with_http()
1293 .with_endpoint(endpoint);
1294 if let Some(t) = timeout {
1295 b = b.with_timeout(t);
1296 }
1297 Ok(b.build()?)
1298 }
1299 }
1300}
1301
1302fn build_log_exporter(
1303 protocol: ExportProtocol,
1304 endpoint: &str,
1305 timeout: Option<Duration>,
1306 #[cfg(feature = "grpc-mtls")] mtls: Option<&MtlsTransport>,
1307) -> Result<opentelemetry_otlp::LogExporter, Box<dyn Error>> {
1308 match protocol {
1309 #[cfg(feature = "grpc")]
1310 ExportProtocol::Grpc => {
1311 let mut b = opentelemetry_otlp::LogExporter::builder()
1312 .with_tonic()
1313 .with_endpoint(endpoint);
1314 if let Some(t) = timeout {
1315 b = b.with_timeout(t);
1316 }
1317 #[cfg(feature = "grpc-mtls")]
1318 if let Some(m) = mtls {
1319 b = m.apply(b);
1320 }
1321 Ok(b.build()?)
1322 }
1323 #[cfg(feature = "http")]
1324 ExportProtocol::HttpProtobuf => {
1325 let mut b = opentelemetry_otlp::LogExporter::builder()
1326 .with_http()
1327 .with_endpoint(endpoint);
1328 if let Some(t) = timeout {
1329 b = b.with_timeout(t);
1330 }
1331 Ok(b.build()?)
1332 }
1333 }
1334}
1335
1336pub fn build_resource(
1351 service_name: &str,
1352 service_version: Option<&str>,
1353 deployment_environment: Option<&str>,
1354) -> Resource {
1355 let hostname = hostname::get()
1356 .ok()
1357 .and_then(|h| h.into_string().ok())
1358 .unwrap_or_default();
1359
1360 let mut builder = Resource::builder()
1361 .with_service_name(service_name.to_string())
1362 .with_attributes([
1363 KeyValue::new(HOST_NAME, hostname),
1364 KeyValue::new(PROCESS_PID, std::process::id() as i64),
1365 ]);
1366
1367 if let Some(version) = service_version {
1368 builder = builder.with_attribute(KeyValue::new(SERVICE_VERSION, version.to_string()));
1369 }
1370
1371 if let Some(env) = deployment_environment {
1372 builder =
1373 builder.with_attribute(KeyValue::new(DEPLOYMENT_ENVIRONMENT_NAME, env.to_string()));
1374 }
1375
1376 builder.build()
1377}
1378
1379#[cfg(feature = "axum")]
1397pub fn axum_layer() -> axum_middleware::OtelTraceLayer {
1398 axum_middleware::OtelTraceLayer
1399}
1400
1401#[cfg(feature = "axum")]
1432pub fn span_enricher_layer<T>() -> axum_middleware::SpanEnricherLayer<T>
1433where
1434 T: span_enrichment::EnrichSpan + Clone + Send + Sync + 'static,
1435{
1436 axum_middleware::SpanEnricherLayer::default()
1437}
1438
1439#[cfg(feature = "tonic-tracing")]
1461pub fn grpc_client_layer() -> grpc_middleware::GrpcClientTraceLayer {
1462 grpc_middleware::GrpcClientTraceLayer
1463}
1464
1465#[cfg(feature = "tonic-tracing")]
1480pub fn grpc_server_layer() -> grpc_middleware::GrpcServerTraceLayer {
1481 grpc_middleware::GrpcServerTraceLayer
1482}
1483
1484#[cfg(test)]
1485mod tests {
1486 use super::*;
1487
1488 #[test]
1490 fn runtime_metrics_can_be_disabled() {
1491 assert!(
1492 Telemetry::builder("rm-default").runtime_metrics,
1493 "runtime metrics are on by default"
1494 );
1495 assert!(
1496 !Telemetry::builder("rm-off")
1497 .with_runtime_metrics(false)
1498 .runtime_metrics
1499 );
1500 }
1501
1502 #[tokio::test]
1510 async fn shutdown_absorbs_provider_errors() {
1511 let handles = TelemetryHandles {
1512 tracer_provider: SdkTracerProvider::builder().build(),
1513 meter_provider: Some(SdkMeterProvider::builder().build()),
1514 logger_provider: Some(SdkLoggerProvider::builder().build()),
1515 shutdown_timeout: DEFAULT_SHUTDOWN_TIMEOUT,
1516 boot_owner: false,
1517 #[cfg(feature = "profiling")]
1518 profiling_handle: None,
1519 };
1520
1521 handles.shutdown().expect("first shutdown succeeds");
1522 handles
1523 .shutdown()
1524 .expect("second shutdown absorbs the already-shut-down errors");
1525 }
1526 use opentelemetry::trace::{Span as _, Tracer as _};
1527 use std::sync::Mutex;
1528
1529 static ENV_LOCK: Mutex<()> = Mutex::new(());
1530
1531 #[test]
1532 fn tracing_bridge_uses_sdk_tracer() {
1533 let provider = SdkTracerProvider::builder().build();
1534 let tracer = tracing_bridge_tracer(&provider);
1535 let span = tracer.start("bridge-regression");
1536
1537 assert!(span.span_context().is_valid());
1538
1539 provider.shutdown().expect("provider shutdown");
1540 }
1541
1542 #[test]
1543 fn resource_contains_all_attributes_when_provided() {
1544 let resource = build_resource("test-svc", Some("1.2.3"), Some("staging"));
1545
1546 assert_eq!(
1547 resource.get(&opentelemetry::Key::new("service.name")),
1548 Some(opentelemetry::Value::from("test-svc")),
1549 );
1550 assert_eq!(
1551 resource.get(&opentelemetry::Key::new(SERVICE_VERSION)),
1552 Some(opentelemetry::Value::from("1.2.3")),
1553 );
1554 assert_eq!(
1555 resource.get(&opentelemetry::Key::new(DEPLOYMENT_ENVIRONMENT_NAME)),
1556 Some(opentelemetry::Value::from("staging")),
1557 );
1558 assert!(resource.get(&opentelemetry::Key::new(HOST_NAME)).is_some());
1559 assert!(
1560 resource
1561 .get(&opentelemetry::Key::new(PROCESS_PID))
1562 .is_some()
1563 );
1564 }
1565
1566 #[test]
1567 fn resource_graceful_when_optional_values_omitted() {
1568 let resource = build_resource("test-svc", None, None);
1569
1570 assert_eq!(
1571 resource.get(&opentelemetry::Key::new("service.name")),
1572 Some(opentelemetry::Value::from("test-svc")),
1573 );
1574 assert!(
1575 resource
1576 .get(&opentelemetry::Key::new(SERVICE_VERSION))
1577 .is_none()
1578 );
1579 assert!(
1580 resource
1581 .get(&opentelemetry::Key::new(DEPLOYMENT_ENVIRONMENT_NAME))
1582 .is_none()
1583 );
1584 assert!(resource.get(&opentelemetry::Key::new(HOST_NAME)).is_some());
1586 assert!(
1587 resource
1588 .get(&opentelemetry::Key::new(PROCESS_PID))
1589 .is_some()
1590 );
1591 }
1592
1593 #[test]
1594 fn trace_sampler_ratio_converts_to_sdk() {
1595 let sampler = TraceSampler::TraceIdRatio(0.5);
1596 let sdk = sampler.into_sdk_sampler();
1597 assert_eq!(format!("{sdk:?}"), "TraceIdRatioBased(0.5)");
1598 }
1599
1600 #[test]
1601 fn trace_sampler_parent_based_converts_to_sdk() {
1602 let sampler = TraceSampler::ParentBased(Box::new(TraceSampler::TraceIdRatio(0.25)));
1603 let sdk = sampler.into_sdk_sampler();
1604 let debug = format!("{sdk:?}");
1605 assert!(debug.contains("ParentBased"));
1606 assert!(debug.contains("0.25"));
1607 }
1608
1609 unsafe fn set_env(key: &str, val: &str) {
1611 unsafe {
1612 std::env::set_var(key, val);
1613 }
1614 }
1615
1616 unsafe fn remove_env(key: &str) {
1617 unsafe {
1618 std::env::remove_var(key);
1619 }
1620 }
1621
1622 #[test]
1623 fn sampler_from_env_reads_traceidratio() {
1624 let _lock = ENV_LOCK.lock().unwrap();
1625 unsafe {
1626 set_env("OTEL_TRACES_SAMPLER", "traceidratio");
1627 set_env("OTEL_TRACES_SAMPLER_ARG", "0.42");
1628 }
1629
1630 let sampler = sampler_from_env()
1631 .expect("should not error")
1632 .expect("should return Some");
1633 assert!(
1634 matches!(sampler, TraceSampler::TraceIdRatio(r) if (r - 0.42).abs() < f64::EPSILON)
1635 );
1636
1637 unsafe {
1638 remove_env("OTEL_TRACES_SAMPLER");
1639 remove_env("OTEL_TRACES_SAMPLER_ARG");
1640 }
1641 }
1642
1643 #[test]
1644 fn sampler_from_env_returns_none_when_unset() {
1645 let _lock = ENV_LOCK.lock().unwrap();
1646 unsafe {
1647 remove_env("OTEL_TRACES_SAMPLER");
1648 }
1649 assert!(sampler_from_env().expect("should not error").is_none());
1650 }
1651
1652 #[test]
1653 fn sampler_from_env_reads_parentbased_traceidratio() {
1654 let _lock = ENV_LOCK.lock().unwrap();
1655 unsafe {
1656 set_env("OTEL_TRACES_SAMPLER", "parentbased_traceidratio");
1657 set_env("OTEL_TRACES_SAMPLER_ARG", "0.1");
1658 }
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::TraceIdRatio(r) if (r - 0.1).abs() < f64::EPSILON))
1665 );
1666
1667 unsafe {
1668 remove_env("OTEL_TRACES_SAMPLER");
1669 remove_env("OTEL_TRACES_SAMPLER_ARG");
1670 }
1671 }
1672
1673 #[test]
1674 fn sampler_from_env_parentbased_always_on() {
1675 let _lock = ENV_LOCK.lock().unwrap();
1676 unsafe {
1677 set_env("OTEL_TRACES_SAMPLER", "parentbased_always_on");
1678 }
1679 let sampler = sampler_from_env()
1680 .expect("should not error")
1681 .expect("should return Some");
1682 assert!(
1683 matches!(sampler, TraceSampler::ParentBased(inner) if matches!(*inner, TraceSampler::AlwaysOn))
1684 );
1685 unsafe {
1686 remove_env("OTEL_TRACES_SAMPLER");
1687 }
1688 }
1689
1690 #[test]
1691 fn sampler_from_env_parentbased_always_off() {
1692 let _lock = ENV_LOCK.lock().unwrap();
1693 unsafe {
1694 set_env("OTEL_TRACES_SAMPLER", "parentbased_always_off");
1695 }
1696 let sampler = sampler_from_env()
1697 .expect("should not error")
1698 .expect("should return Some");
1699 assert!(
1700 matches!(sampler, TraceSampler::ParentBased(inner) if matches!(*inner, TraceSampler::AlwaysOff))
1701 );
1702 unsafe {
1703 remove_env("OTEL_TRACES_SAMPLER");
1704 }
1705 }
1706
1707 #[test]
1708 fn sampler_from_env_always_on() {
1709 let _lock = ENV_LOCK.lock().unwrap();
1710 unsafe {
1711 set_env("OTEL_TRACES_SAMPLER", "always_on");
1712 }
1713 let sampler = sampler_from_env()
1714 .expect("should not error")
1715 .expect("should return Some");
1716 assert!(matches!(sampler, TraceSampler::AlwaysOn));
1717 unsafe {
1718 remove_env("OTEL_TRACES_SAMPLER");
1719 }
1720 }
1721
1722 #[test]
1723 fn sampler_from_env_always_off() {
1724 let _lock = ENV_LOCK.lock().unwrap();
1725 unsafe {
1726 set_env("OTEL_TRACES_SAMPLER", "always_off");
1727 }
1728 let sampler = sampler_from_env()
1729 .expect("should not error")
1730 .expect("should return Some");
1731 assert!(matches!(sampler, TraceSampler::AlwaysOff));
1732 unsafe {
1733 remove_env("OTEL_TRACES_SAMPLER");
1734 }
1735 }
1736
1737 #[test]
1738 fn sampler_from_env_unknown_returns_error() {
1739 let _lock = ENV_LOCK.lock().unwrap();
1740 unsafe {
1741 set_env("OTEL_TRACES_SAMPLER", "unknown_sampler");
1742 }
1743 let err = sampler_from_env().expect_err("unknown sampler should produce an error");
1744 assert!(
1745 err.to_string().contains("unknown_sampler"),
1746 "error message should include the unknown name, got: {err}"
1747 );
1748 unsafe {
1749 remove_env("OTEL_TRACES_SAMPLER");
1750 }
1751 }
1752
1753 #[test]
1754 fn trace_sampler_always_on_converts_to_sdk() {
1755 let sdk = TraceSampler::AlwaysOn.into_sdk_sampler();
1756 assert_eq!(format!("{sdk:?}"), "AlwaysOn");
1757 }
1758
1759 #[test]
1760 fn trace_sampler_always_off_converts_to_sdk() {
1761 let sdk = TraceSampler::AlwaysOff.into_sdk_sampler();
1762 assert_eq!(format!("{sdk:?}"), "AlwaysOff");
1763 }
1764
1765 #[test]
1766 fn builder_has_sensible_defaults() {
1767 let builder = Telemetry::builder("test-svc");
1768 assert_eq!(builder.service_name.as_deref(), Some("test-svc"));
1769 assert!(builder.service_version.is_none());
1770 assert!(builder.deployment_environment.is_none());
1771 assert!(builder.sampler.is_none());
1772 assert!(builder.metrics);
1773 assert!(!builder.logs);
1774 assert!(builder.protocol.is_none());
1775 assert!(builder.max_export_batch_size.is_none());
1776 assert!(builder.metric_export_interval.is_none());
1777 assert!(builder.export_timeout.is_none());
1778 }
1779
1780 #[test]
1781 fn from_env_builder_has_no_service_name() {
1782 let builder = Telemetry::from_env();
1783 assert!(builder.service_name.is_none());
1784 }
1785
1786 #[test]
1787 fn with_export_timeout_stores_value() {
1788 let timeout = Duration::from_secs(5);
1789 let builder = Telemetry::builder("test-svc").with_export_timeout(timeout);
1790 assert_eq!(builder.export_timeout, Some(timeout));
1791 }
1792
1793 #[test]
1794 fn timeout_from_env_reads_milliseconds() {
1795 let _lock = ENV_LOCK.lock().unwrap();
1796 unsafe {
1797 set_env("OTEL_EXPORTER_OTLP_TIMEOUT", "5000");
1798 }
1799 let t = timeout_from_env();
1800 assert_eq!(t, Some(Duration::from_millis(5000)));
1801 unsafe {
1802 remove_env("OTEL_EXPORTER_OTLP_TIMEOUT");
1803 }
1804 }
1805
1806 #[test]
1807 fn timeout_from_env_returns_none_when_unset() {
1808 let _lock = ENV_LOCK.lock().unwrap();
1809 unsafe {
1810 remove_env("OTEL_EXPORTER_OTLP_TIMEOUT");
1811 }
1812 assert_eq!(timeout_from_env(), None);
1813 }
1814
1815 #[test]
1816 fn service_name_from_env_used_when_none_given() {
1817 let builder = Telemetry::from_env();
1818 assert!(builder.service_name.is_none());
1819 }
1820
1821 #[test]
1822 fn explicit_service_name_overrides_env_var() {
1823 let builder = Telemetry::builder("explicit-svc");
1824 assert_eq!(builder.service_name.as_deref(), Some("explicit-svc"));
1825 }
1826
1827 #[test]
1828 fn from_env_builder_service_name_is_none() {
1829 let builder = Telemetry::from_env();
1830 assert!(builder.service_name.is_none());
1831 }
1832
1833 #[test]
1834 fn init_returns_error_for_unknown_otel_traces_sampler() {
1835 let _lock = ENV_LOCK.lock().unwrap();
1836 unsafe {
1837 set_env("OTEL_TRACES_SAMPLER", "not_a_real_sampler");
1838 }
1839 let result = Telemetry::builder("test-svc").with_metrics(false).init();
1840 let err = result
1841 .err()
1842 .expect("unknown sampler env var should cause init to fail");
1843 assert!(
1844 err.to_string().contains("not_a_real_sampler"),
1845 "error should name the unknown sampler, got: {err}"
1846 );
1847 unsafe {
1848 remove_env("OTEL_TRACES_SAMPLER");
1849 }
1850 }
1851
1852 #[test]
1853 fn with_max_export_batch_size_stores_value() {
1854 let builder = Telemetry::builder("test-svc").with_max_export_batch_size(1024);
1855 assert_eq!(builder.max_export_batch_size, Some(1024));
1856 }
1857
1858 #[test]
1859 fn with_metric_export_interval_stores_value() {
1860 let interval = Duration::from_secs(30);
1861 let builder = Telemetry::builder("test-svc").with_metric_export_interval(interval);
1862 assert_eq!(builder.metric_export_interval, Some(interval));
1863 }
1864
1865 #[test]
1866 fn init_rejects_zero_metric_export_interval() {
1867 let err = Telemetry::builder("test-svc")
1868 .with_metric_export_interval(Duration::ZERO)
1869 .with_metrics(false)
1870 .init()
1871 .err()
1872 .expect("expected error for zero interval");
1873 assert!(
1874 err.to_string().contains("metric_export_interval"),
1875 "error message should mention metric_export_interval, got: {err}"
1876 );
1877 }
1878
1879 #[test]
1880 fn builder_with_custom_values() {
1881 let builder = Telemetry::builder("test-svc")
1882 .with_version("2.0.0")
1883 .with_environment("production")
1884 .with_sampler(TraceSampler::TraceIdRatio(0.5))
1885 .with_metrics(false);
1886
1887 assert_eq!(builder.service_name.as_deref(), Some("test-svc"));
1888 assert_eq!(builder.service_version.as_deref(), Some("2.0.0"));
1889 assert_eq!(
1890 builder.deployment_environment.as_deref(),
1891 Some("production")
1892 );
1893 assert!(
1894 matches!(builder.sampler, Some(TraceSampler::TraceIdRatio(r)) if (r - 0.5).abs() < f64::EPSILON)
1895 );
1896 assert!(!builder.metrics);
1897 }
1898
1899 #[test]
1900 fn builder_stores_programmatic_log_configuration() {
1901 let builder = Telemetry::builder("test-svc")
1902 .with_log_filter("info,opentelemetry_sdk=warn")
1903 .with_log_format(LogFormat::Json);
1904
1905 assert_eq!(
1906 builder.log_filter.as_deref(),
1907 Some("info,opentelemetry_sdk=warn")
1908 );
1909 assert_eq!(builder.log_format, LogFormat::Json);
1910 }
1911
1912 #[test]
1913 fn init_rejects_invalid_programmatic_log_filter_before_provider_setup() {
1914 let setup_ran = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
1915 let setup_ran_in_closure = std::sync::Arc::clone(&setup_ran);
1916
1917 let error = Telemetry::builder("test-svc")
1918 .with_log_filter("[")
1919 .with_meter_provider_setup(move |builder| {
1920 setup_ran_in_closure.store(true, std::sync::atomic::Ordering::SeqCst);
1921 builder
1922 })
1923 .init()
1924 .err()
1925 .expect("invalid filter must fail initialization");
1926
1927 assert!(error.to_string().contains("invalid filter directive"));
1928 assert!(!setup_ran.load(std::sync::atomic::Ordering::SeqCst));
1929 }
1930
1931 #[test]
1932 fn builder_with_default_endpoint() {
1933 let builder = Telemetry::builder("svc").with_default_endpoint("http://otel-collector:4317");
1934 assert_eq!(
1935 builder.default_endpoint.as_deref(),
1936 Some("http://otel-collector:4317")
1937 );
1938 }
1939
1940 #[test]
1941 fn a_configured_endpoint_wins_over_the_runtime_default() {
1942 assert_eq!(
1943 resolve_endpoint(
1944 Some("http://c:4317".into()),
1945 Some("http://d:4317"),
1946 "http://localhost:4317"
1947 ),
1948 "http://c:4317"
1949 );
1950 assert_eq!(
1951 resolve_endpoint(None, Some("http://d:4317"), "http://localhost:4317"),
1952 "http://d:4317"
1953 );
1954 assert_eq!(
1955 resolve_endpoint(None, None, "http://localhost:4317"),
1956 "http://localhost:4317"
1957 );
1958 }
1959
1960 #[test]
1961 #[cfg(feature = "grpc")]
1962 fn builder_with_protocol_grpc() {
1963 let builder = Telemetry::builder("test-svc").with_protocol(ExportProtocol::Grpc);
1964 assert_eq!(builder.protocol, Some(ExportProtocol::Grpc));
1965 }
1966
1967 #[test]
1968 #[cfg(feature = "http")]
1969 fn builder_with_protocol_http() {
1970 let builder = Telemetry::builder("test-svc").with_protocol(ExportProtocol::HttpProtobuf);
1971 assert_eq!(builder.protocol, Some(ExportProtocol::HttpProtobuf));
1972 }
1973
1974 #[test]
1975 #[cfg(feature = "grpc")]
1976 fn protocol_from_env_reads_grpc() {
1977 let _lock = ENV_LOCK.lock().unwrap();
1978 unsafe {
1979 set_env("OTEL_EXPORTER_OTLP_PROTOCOL", "grpc");
1980 }
1981 assert_eq!(protocol_from_env(), Some(ExportProtocol::Grpc));
1982 unsafe {
1983 remove_env("OTEL_EXPORTER_OTLP_PROTOCOL");
1984 }
1985 }
1986
1987 #[test]
1988 #[cfg(feature = "http")]
1989 fn protocol_from_env_reads_http_protobuf() {
1990 let _lock = ENV_LOCK.lock().unwrap();
1991 unsafe {
1992 set_env("OTEL_EXPORTER_OTLP_PROTOCOL", "http/protobuf");
1993 }
1994 assert_eq!(protocol_from_env(), Some(ExportProtocol::HttpProtobuf));
1995 unsafe {
1996 remove_env("OTEL_EXPORTER_OTLP_PROTOCOL");
1997 }
1998 }
1999
2000 #[test]
2001 fn protocol_from_env_returns_none_when_unset() {
2002 let _lock = ENV_LOCK.lock().unwrap();
2003 unsafe {
2004 remove_env("OTEL_EXPORTER_OTLP_PROTOCOL");
2005 }
2006 assert_eq!(protocol_from_env(), None);
2007 }
2008
2009 #[test]
2010 fn protocol_from_env_returns_none_for_unknown() {
2011 let _lock = ENV_LOCK.lock().unwrap();
2012 unsafe {
2013 set_env("OTEL_EXPORTER_OTLP_PROTOCOL", "websocket");
2014 }
2015 assert_eq!(protocol_from_env(), None);
2016 unsafe {
2017 remove_env("OTEL_EXPORTER_OTLP_PROTOCOL");
2018 }
2019 }
2020
2021 #[test]
2022 fn builder_is_send_and_sync() {
2023 fn assert_send_sync<T: Send + Sync>() {}
2024 assert_send_sync::<TelemetryBuilder>();
2025 }
2026
2027 #[test]
2028 fn with_shutdown_timeout_stores_value() {
2029 let timeout = Duration::from_secs(10);
2030 let builder = Telemetry::builder("test-svc").with_shutdown_timeout(timeout);
2031 assert_eq!(builder.shutdown_timeout, timeout);
2032 }
2033
2034 #[test]
2035 fn default_shutdown_timeout_is_five_seconds() {
2036 let builder = Telemetry::builder("test-svc");
2037 assert_eq!(builder.shutdown_timeout, Duration::from_secs(5));
2038 }
2039
2040 #[cfg(feature = "testing")]
2048 #[test]
2049 fn drop_completes_within_shutdown_timeout() {
2050 let mut handles = crate::Telemetry::testing("drop-timeout-test");
2052 handles.shutdown_timeout = Duration::from_millis(100);
2054
2055 let start = std::time::Instant::now();
2056 drop(handles);
2057 let elapsed = start.elapsed();
2058
2059 assert!(
2061 elapsed < Duration::from_millis(500),
2062 "drop took {elapsed:?}, expected < 500 ms"
2063 );
2064 }
2065}