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