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;
40
41pub mod log_bridge;
42pub mod span_enrichment;
43
44pub use log_bridge::{
45 PROPAGATED_SPAN_FIELDS, SpanLogAttrs, record_span_log_attr, record_span_log_attr_on,
46};
47
48use opentelemetry::KeyValue;
49use opentelemetry::propagation::TextMapCompositePropagator;
50use opentelemetry_otlp::WithExportConfig;
51use opentelemetry_sdk::{
52 Resource,
53 logs::SdkLoggerProvider,
54 metrics::{MeterProviderBuilder, PeriodicReader, SdkMeterProvider},
55 propagation::{BaggagePropagator, TraceContextPropagator},
56 trace::{BatchConfigBuilder, BatchSpanProcessor, Sampler, SdkTracerProvider},
57};
58use opentelemetry_semantic_conventions::attribute::{
59 DEPLOYMENT_ENVIRONMENT_NAME, HOST_NAME, PROCESS_PID, SERVICE_VERSION,
60};
61use std::error::Error;
62use std::time::Duration;
63use tracing_subscriber::layer::SubscriberExt;
64use tracing_subscriber::util::SubscriberInitExt;
65
66#[derive(Debug, Clone)]
81pub enum TraceSampler {
82 AlwaysOn,
84 AlwaysOff,
86 TraceIdRatio(f64),
88 ParentBased(Box<TraceSampler>),
91}
92
93impl TraceSampler {
94 fn into_sdk_sampler(self) -> Sampler {
96 match self {
97 TraceSampler::AlwaysOn => Sampler::AlwaysOn,
98 TraceSampler::AlwaysOff => Sampler::AlwaysOff,
99 TraceSampler::TraceIdRatio(r) => Sampler::TraceIdRatioBased(r),
100 TraceSampler::ParentBased(inner) => {
101 Sampler::ParentBased(Box::new(inner.into_sdk_sampler()))
102 }
103 }
104 }
105}
106
107fn sampler_from_env() -> Result<Option<TraceSampler>, Box<dyn Error>> {
115 let name = match std::env::var("OTEL_TRACES_SAMPLER") {
116 Ok(v) => v,
117 Err(_) => return Ok(None),
118 };
119 let arg = std::env::var("OTEL_TRACES_SAMPLER_ARG").ok();
120 let sampler = match name.as_str() {
121 "always_on" => TraceSampler::AlwaysOn,
122 "always_off" => TraceSampler::AlwaysOff,
123 "traceidratio" => {
124 let ratio = arg
125 .as_deref()
126 .unwrap_or("1.0")
127 .parse::<f64>()
128 .unwrap_or(1.0);
129 TraceSampler::TraceIdRatio(ratio)
130 }
131 "parentbased_always_on" => TraceSampler::ParentBased(Box::new(TraceSampler::AlwaysOn)),
132 "parentbased_always_off" => TraceSampler::ParentBased(Box::new(TraceSampler::AlwaysOff)),
133 "parentbased_traceidratio" => {
134 let ratio = arg
135 .as_deref()
136 .unwrap_or("1.0")
137 .parse::<f64>()
138 .unwrap_or(1.0);
139 TraceSampler::ParentBased(Box::new(TraceSampler::TraceIdRatio(ratio)))
140 }
141 unknown => {
142 return Err(format!(
143 "OTEL_TRACES_SAMPLER: unrecognised sampler name '{unknown}'. \
144 Valid values: always_on, always_off, traceidratio, \
145 parentbased_always_on, parentbased_always_off, parentbased_traceidratio"
146 )
147 .into());
148 }
149 };
150 Ok(Some(sampler))
151}
152
153const DEFAULT_SHUTDOWN_TIMEOUT: Duration = Duration::from_secs(5);
155
156pub struct TelemetryHandles {
177 pub tracer_provider: SdkTracerProvider,
178 pub meter_provider: Option<SdkMeterProvider>,
179 pub logger_provider: Option<SdkLoggerProvider>,
180 shutdown_timeout: Duration,
181 #[cfg(feature = "profiling")]
182 pub profiling_handle: Option<profiling::ProfilingHandle>,
183}
184
185impl TelemetryHandles {
186 pub fn shutdown(&self) -> Result<(), Box<dyn Error>> {
199 self.tracer_provider.shutdown()?;
200 if let Some(mp) = &self.meter_provider {
201 mp.shutdown()?;
202 }
203 if let Some(lp) = &self.logger_provider {
204 lp.shutdown()?;
205 }
206 Ok(())
207 }
208}
209
210impl Drop for TelemetryHandles {
211 fn drop(&mut self) {
212 let tracer_provider = self.tracer_provider.clone();
213 let meter_provider = self.meter_provider.clone();
214 let logger_provider = self.logger_provider.clone();
215 let timeout = self.shutdown_timeout;
216
217 let (tx, rx) = std::sync::mpsc::channel();
218 std::thread::spawn(move || {
219 if let Err(e) = tracer_provider.shutdown() {
220 tracing::warn!("tracer provider shutdown error: {e}");
221 }
222 if let Some(mp) = meter_provider
223 && let Err(e) = mp.shutdown()
224 {
225 tracing::warn!("meter provider shutdown error: {e}");
226 }
227 if let Some(lp) = logger_provider
228 && let Err(e) = lp.shutdown()
229 {
230 tracing::warn!("logger provider shutdown error: {e}");
231 }
232 let _ = tx.send(());
233 });
234
235 if rx.recv_timeout(timeout).is_err() {
236 tracing::warn!(
237 "telemetry shutdown did not complete within {timeout:?}; \
238 some spans/metrics may not have been exported"
239 );
240 }
241 }
242}
243
244#[derive(Debug, Clone, Copy, PartialEq, Eq)]
267pub enum ExportProtocol {
268 #[cfg(feature = "grpc")]
270 Grpc,
271 #[cfg(feature = "http")]
273 HttpProtobuf,
274}
275
276#[cfg(feature = "grpc-mtls")]
286#[derive(Clone)]
287pub struct MtlsMaterial {
288 pub client_cert_chain_pem: Vec<u8>,
290 pub client_key_pem: Vec<u8>,
292 pub trust_bundle_pem: Vec<u8>,
294}
295
296#[cfg(feature = "grpc-mtls")]
297impl std::fmt::Debug for MtlsMaterial {
298 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
299 f.debug_struct("MtlsMaterial")
300 .field("client_cert_chain_pem", &"<redacted>")
301 .field("client_key_pem", &"<redacted>")
302 .field("trust_bundle_pem", &"<redacted>")
303 .finish()
304 }
305}
306
307fn protocol_from_env() -> Option<ExportProtocol> {
309 let val = std::env::var("OTEL_EXPORTER_OTLP_PROTOCOL").ok()?;
310 match val.trim() {
311 #[cfg(feature = "grpc")]
312 "grpc" => Some(ExportProtocol::Grpc),
313 #[cfg(feature = "http")]
314 "http/protobuf" => Some(ExportProtocol::HttpProtobuf),
315 _ => None,
316 }
317}
318
319pub struct Telemetry;
335
336impl Telemetry {
337 pub fn builder(service_name: &str) -> TelemetryBuilder {
341 TelemetryBuilder {
342 service_name: Some(service_name.to_string()),
343 service_version: None,
344 deployment_environment: None,
345 sampler: None,
346 metrics: true,
347 logs: false,
348 protocol: None,
349 max_export_batch_size: None,
350 metric_export_interval: None,
351 export_timeout: None,
352 shutdown_timeout: DEFAULT_SHUTDOWN_TIMEOUT,
353 extra_layers: Vec::new(),
354 extra_metric_readers: Vec::new(),
355 #[cfg(feature = "grpc-mtls")]
356 mtls: None,
357 propagated_span_fields: crate::log_bridge::PROPAGATED_SPAN_FIELDS,
358 #[cfg(feature = "profiling")]
359 pyroscope_endpoint: None,
360 }
361 }
362
363 pub fn from_env() -> TelemetryBuilder {
373 TelemetryBuilder {
374 service_name: None,
375 service_version: None,
376 deployment_environment: None,
377 sampler: None,
378 metrics: true,
379 logs: false,
380 protocol: None,
381 max_export_batch_size: None,
382 metric_export_interval: None,
383 export_timeout: None,
384 shutdown_timeout: DEFAULT_SHUTDOWN_TIMEOUT,
385 extra_layers: Vec::new(),
386 extra_metric_readers: Vec::new(),
387 #[cfg(feature = "grpc-mtls")]
388 mtls: None,
389 propagated_span_fields: crate::log_bridge::PROPAGATED_SPAN_FIELDS,
390 #[cfg(feature = "profiling")]
391 pyroscope_endpoint: None,
392 }
393 }
394}
395
396#[must_use = "a TelemetryBuilder does nothing until .init() is called"]
414pub struct TelemetryBuilder {
415 service_name: Option<String>,
416 service_version: Option<String>,
417 deployment_environment: Option<String>,
418 sampler: Option<TraceSampler>,
419 metrics: bool,
420 logs: bool,
421 protocol: Option<ExportProtocol>,
422 max_export_batch_size: Option<usize>,
423 metric_export_interval: Option<Duration>,
424 export_timeout: Option<Duration>,
425 shutdown_timeout: Duration,
426 extra_layers: Vec<
427 Box<dyn tracing_subscriber::Layer<tracing_subscriber::Registry> + Send + Sync + 'static>,
428 >,
429 extra_metric_readers: Vec<MeterProviderInstaller>,
430 #[cfg(feature = "grpc-mtls")]
431 mtls: Option<MtlsMaterial>,
432 propagated_span_fields: &'static [&'static str],
433 #[cfg(feature = "profiling")]
434 pyroscope_endpoint: Option<String>,
435}
436
437type MeterProviderInstaller =
442 Box<dyn FnOnce(MeterProviderBuilder) -> MeterProviderBuilder + Send + Sync>;
443
444impl TelemetryBuilder {
445 pub fn with_version(mut self, version: &str) -> Self {
447 self.service_version = Some(version.to_string());
448 self
449 }
450
451 pub fn with_environment(mut self, environment: &str) -> Self {
453 self.deployment_environment = Some(environment.to_string());
454 self
455 }
456
457 #[cfg(feature = "grpc-mtls")]
482 pub fn with_mtls(mut self, material: MtlsMaterial) -> Self {
483 self.mtls = Some(material);
484 self.protocol = Some(ExportProtocol::Grpc);
485 self
486 }
487
488 pub fn with_sampler(mut self, sampler: TraceSampler) -> Self {
491 self.sampler = Some(sampler);
492 self
493 }
494
495 pub fn with_metrics(mut self, enabled: bool) -> Self {
497 self.metrics = enabled;
498 self
499 }
500
501 pub fn with_protocol(mut self, protocol: ExportProtocol) -> Self {
505 self.protocol = Some(protocol);
506 self
507 }
508
509 pub fn with_max_export_batch_size(mut self, size: usize) -> Self {
514 self.max_export_batch_size = Some(size);
515 self
516 }
517
518 pub fn with_metric_export_interval(mut self, interval: Duration) -> Self {
523 self.metric_export_interval = Some(interval);
524 self
525 }
526
527 pub fn with_logs(mut self, enabled: bool) -> Self {
534 self.logs = enabled;
535 self
536 }
537
538 pub fn with_propagated_span_fields(mut self, fields: &'static [&'static str]) -> Self {
551 self.propagated_span_fields = fields;
552 self
553 }
554
555 pub fn with_export_timeout(mut self, timeout: Duration) -> Self {
559 self.export_timeout = Some(timeout);
560 self
561 }
562
563 pub fn with_shutdown_timeout(mut self, timeout: Duration) -> Self {
570 self.shutdown_timeout = timeout;
571 self
572 }
573
574 #[cfg(feature = "profiling")]
592 pub fn with_profiling(mut self, endpoint: &str) -> Self {
593 self.pyroscope_endpoint = Some(endpoint.to_string());
594 self
595 }
596
597 pub fn with_meter_provider_setup<F>(mut self, setup: F) -> Self
654 where
655 F: FnOnce(MeterProviderBuilder) -> MeterProviderBuilder + Send + Sync + 'static,
656 {
657 self.extra_metric_readers.push(Box::new(setup));
658 self
659 }
660
661 pub fn with_layer<L>(mut self, layer: L) -> Self
662 where
663 L: tracing_subscriber::Layer<tracing_subscriber::Registry> + Send + Sync + 'static,
664 {
665 self.extra_layers.push(Box::new(layer));
666 self
667 }
668
669 pub fn init(self) -> Result<TelemetryHandles, Box<dyn Error>> {
684 if let Some(interval) = self.metric_export_interval
685 && interval.is_zero()
686 {
687 return Err("metric_export_interval must be greater than zero".into());
688 }
689
690 let protocol = self.protocol.or_else(protocol_from_env).unwrap_or({
691 #[cfg(feature = "grpc")]
692 {
693 ExportProtocol::Grpc
694 }
695 #[cfg(all(not(feature = "grpc"), feature = "http"))]
696 {
697 ExportProtocol::HttpProtobuf
698 }
699 });
700
701 let default_endpoint = match protocol {
702 #[cfg(feature = "grpc")]
703 ExportProtocol::Grpc => "http://localhost:4317",
704 #[cfg(feature = "http")]
705 ExportProtocol::HttpProtobuf => "http://localhost:4318",
706 };
707 let endpoint = std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT")
708 .unwrap_or_else(|_| default_endpoint.to_string());
709
710 let export_timeout = self.export_timeout.or_else(timeout_from_env);
712
713 let service_name = self.service_name.unwrap_or_else(|| {
715 std::env::var("OTEL_SERVICE_NAME").unwrap_or_else(|_| "unknown_service".to_string())
716 });
717
718 let resource = build_resource(
719 &service_name,
720 self.service_version.as_deref(),
721 self.deployment_environment.as_deref(),
722 );
723
724 let sampler = match self.sampler {
725 Some(s) => s,
726 None => sampler_from_env()?.unwrap_or(TraceSampler::AlwaysOn),
727 };
728
729 let trace_exporter = build_span_exporter(
731 protocol,
732 &endpoint,
733 export_timeout,
734 #[cfg(feature = "grpc-mtls")]
735 self.mtls.as_ref(),
736 )?;
737
738 let batch_processor = if let Some(size) = self.max_export_batch_size {
739 BatchSpanProcessor::builder(trace_exporter)
740 .with_batch_config(
741 BatchConfigBuilder::default()
742 .with_max_export_batch_size(size)
743 .build(),
744 )
745 .build()
746 } else {
747 BatchSpanProcessor::builder(trace_exporter).build()
748 };
749
750 let tracer_provider = SdkTracerProvider::builder()
751 .with_resource(resource.clone())
752 .with_sampler(sampler.into_sdk_sampler())
753 .with_span_processor(batch_processor)
754 .build();
755
756 opentelemetry::global::set_tracer_provider(tracer_provider.clone());
757
758 let propagator = TextMapCompositePropagator::new(vec![
760 Box::new(TraceContextPropagator::new()),
761 Box::new(BaggagePropagator::new()),
762 ]);
763 opentelemetry::global::set_text_map_propagator(propagator);
764
765 let meter_provider = if self.metrics {
767 let metric_exporter = build_metric_exporter(
768 protocol,
769 &endpoint,
770 export_timeout,
771 #[cfg(feature = "grpc-mtls")]
772 self.mtls.as_ref(),
773 )?;
774
775 let periodic_reader = if let Some(interval) = self.metric_export_interval {
776 PeriodicReader::builder(metric_exporter)
777 .with_interval(interval)
778 .build()
779 } else {
780 PeriodicReader::builder(metric_exporter).build()
781 };
782
783 let mut mp_builder = SdkMeterProvider::builder()
784 .with_resource(resource.clone())
785 .with_reader(periodic_reader);
786 for installer in self.extra_metric_readers {
787 mp_builder = installer(mp_builder);
788 }
789 let mp = mp_builder.build();
790
791 opentelemetry::global::set_meter_provider(mp.clone());
792
793 Some(mp)
794 } else {
795 None
796 };
797
798 let logger_provider = if self.logs {
800 let log_exporter = build_log_exporter(
801 protocol,
802 &endpoint,
803 export_timeout,
804 #[cfg(feature = "grpc-mtls")]
805 self.mtls.as_ref(),
806 )?;
807
808 let lp = SdkLoggerProvider::builder()
809 .with_resource(resource)
810 .with_batch_exporter(log_exporter)
811 .build();
812
813 Some(lp)
814 } else {
815 None
816 };
817
818 #[cfg(feature = "profiling")]
820 let profiling_handle = if let Some(ref endpoint) = self.pyroscope_endpoint {
821 profiling::start_pyroscope_bridge(&service_name, endpoint)?
822 } else {
823 None
824 };
825 #[cfg(not(feature = "profiling"))]
826 let _profiling_handle: Option<()> = None;
827
828 let otel_layer = tracing_opentelemetry::layer();
830
831 let extra = if self.extra_layers.is_empty() {
836 None
837 } else {
838 Some(self.extra_layers)
839 };
840
841 let registry = tracing_subscriber::registry()
842 .with(extra)
843 .with(tracing_subscriber::EnvFilter::from_default_env())
844 .with(tracing_subscriber::fmt::layer())
845 .with(otel_layer);
846
847 #[cfg(feature = "profiling-bridge-pyroscope-rs")]
848 let registry = registry.with(crate::profiling::ProfilingTagLayer);
849
850 if let Some(lp) = &logger_provider {
851 if let Err(e) = registry
852 .with(crate::log_bridge::SpanAwareLogBridge::new(
853 lp,
854 self.propagated_span_fields,
855 ))
856 .try_init()
857 {
858 eprintln!(
859 "otel-bootstrap: global tracing subscriber already installed — \
860 OTLP log records will NOT be exported to the collector: {e}"
861 );
862 }
863 } else if let Err(e) = registry.try_init() {
864 eprintln!(
865 "otel-bootstrap: global tracing subscriber already installed — \
866 OTLP telemetry will NOT be exported to the collector: {e}"
867 );
868 }
869
870 Ok(TelemetryHandles {
871 tracer_provider,
872 meter_provider,
873 logger_provider,
874 shutdown_timeout: self.shutdown_timeout,
875 #[cfg(feature = "profiling")]
876 profiling_handle,
877 })
878 }
879}
880
881pub fn init_telemetry(service_name: &str) -> Result<TelemetryHandles, Box<dyn Error>> {
895 Telemetry::builder(service_name).init()
896}
897
898pub fn init_telemetry_with_sampler(
915 service_name: &str,
916 sampler: Option<TraceSampler>,
917) -> Result<TelemetryHandles, Box<dyn Error>> {
918 let builder = Telemetry::builder(service_name);
919 match sampler {
920 Some(s) => builder.with_sampler(s),
921 None => builder, }
923 .init()
924}
925
926fn timeout_from_env() -> Option<Duration> {
928 let ms = std::env::var("OTEL_EXPORTER_OTLP_TIMEOUT").ok()?;
929 let ms: u64 = ms.trim().parse().ok()?;
930 Some(Duration::from_millis(ms))
931}
932
933#[cfg(feature = "grpc-mtls")]
940fn build_tls_config(material: &MtlsMaterial) -> tonic::transport::ClientTlsConfig {
941 use tonic::transport::{Certificate, ClientTlsConfig, Identity};
942 ClientTlsConfig::new()
943 .ca_certificate(Certificate::from_pem(&material.trust_bundle_pem))
944 .identity(Identity::from_pem(
945 &material.client_cert_chain_pem,
946 &material.client_key_pem,
947 ))
948}
949
950fn build_span_exporter(
951 protocol: ExportProtocol,
952 endpoint: &str,
953 timeout: Option<Duration>,
954 #[cfg(feature = "grpc-mtls")] mtls: Option<&MtlsMaterial>,
955) -> Result<opentelemetry_otlp::SpanExporter, Box<dyn Error>> {
956 match protocol {
957 #[cfg(feature = "grpc")]
958 ExportProtocol::Grpc => {
959 let mut b = opentelemetry_otlp::SpanExporter::builder()
960 .with_tonic()
961 .with_endpoint(endpoint);
962 if let Some(t) = timeout {
963 b = b.with_timeout(t);
964 }
965 #[cfg(feature = "grpc-mtls")]
966 if let Some(m) = mtls {
967 use opentelemetry_otlp::WithTonicConfig as _;
968 b = b.with_tls_config(build_tls_config(m));
969 }
970 Ok(b.build()?)
971 }
972 #[cfg(feature = "http")]
973 ExportProtocol::HttpProtobuf => {
974 let mut b = opentelemetry_otlp::SpanExporter::builder()
975 .with_http()
976 .with_endpoint(endpoint);
977 if let Some(t) = timeout {
978 b = b.with_timeout(t);
979 }
980 Ok(b.build()?)
981 }
982 }
983}
984
985fn build_metric_exporter(
986 protocol: ExportProtocol,
987 endpoint: &str,
988 timeout: Option<Duration>,
989 #[cfg(feature = "grpc-mtls")] mtls: Option<&MtlsMaterial>,
990) -> Result<opentelemetry_otlp::MetricExporter, Box<dyn Error>> {
991 match protocol {
992 #[cfg(feature = "grpc")]
993 ExportProtocol::Grpc => {
994 let mut b = opentelemetry_otlp::MetricExporter::builder()
995 .with_tonic()
996 .with_endpoint(endpoint);
997 if let Some(t) = timeout {
998 b = b.with_timeout(t);
999 }
1000 #[cfg(feature = "grpc-mtls")]
1001 if let Some(m) = mtls {
1002 use opentelemetry_otlp::WithTonicConfig as _;
1003 b = b.with_tls_config(build_tls_config(m));
1004 }
1005 Ok(b.build()?)
1006 }
1007 #[cfg(feature = "http")]
1008 ExportProtocol::HttpProtobuf => {
1009 let mut b = opentelemetry_otlp::MetricExporter::builder()
1010 .with_http()
1011 .with_endpoint(endpoint);
1012 if let Some(t) = timeout {
1013 b = b.with_timeout(t);
1014 }
1015 Ok(b.build()?)
1016 }
1017 }
1018}
1019
1020fn build_log_exporter(
1021 protocol: ExportProtocol,
1022 endpoint: &str,
1023 timeout: Option<Duration>,
1024 #[cfg(feature = "grpc-mtls")] mtls: Option<&MtlsMaterial>,
1025) -> Result<opentelemetry_otlp::LogExporter, Box<dyn Error>> {
1026 match protocol {
1027 #[cfg(feature = "grpc")]
1028 ExportProtocol::Grpc => {
1029 let mut b = opentelemetry_otlp::LogExporter::builder()
1030 .with_tonic()
1031 .with_endpoint(endpoint);
1032 if let Some(t) = timeout {
1033 b = b.with_timeout(t);
1034 }
1035 #[cfg(feature = "grpc-mtls")]
1036 if let Some(m) = mtls {
1037 use opentelemetry_otlp::WithTonicConfig as _;
1038 b = b.with_tls_config(build_tls_config(m));
1039 }
1040 Ok(b.build()?)
1041 }
1042 #[cfg(feature = "http")]
1043 ExportProtocol::HttpProtobuf => {
1044 let mut b = opentelemetry_otlp::LogExporter::builder()
1045 .with_http()
1046 .with_endpoint(endpoint);
1047 if let Some(t) = timeout {
1048 b = b.with_timeout(t);
1049 }
1050 Ok(b.build()?)
1051 }
1052 }
1053}
1054
1055pub fn build_resource(
1070 service_name: &str,
1071 service_version: Option<&str>,
1072 deployment_environment: Option<&str>,
1073) -> Resource {
1074 let hostname = hostname::get()
1075 .ok()
1076 .and_then(|h| h.into_string().ok())
1077 .unwrap_or_default();
1078
1079 let mut builder = Resource::builder()
1080 .with_service_name(service_name.to_string())
1081 .with_attributes([
1082 KeyValue::new(HOST_NAME, hostname),
1083 KeyValue::new(PROCESS_PID, std::process::id() as i64),
1084 ]);
1085
1086 if let Some(version) = service_version {
1087 builder = builder.with_attribute(KeyValue::new(SERVICE_VERSION, version.to_string()));
1088 }
1089
1090 if let Some(env) = deployment_environment {
1091 builder =
1092 builder.with_attribute(KeyValue::new(DEPLOYMENT_ENVIRONMENT_NAME, env.to_string()));
1093 }
1094
1095 builder.build()
1096}
1097
1098#[cfg(feature = "axum")]
1116pub fn axum_layer() -> axum_middleware::OtelTraceLayer {
1117 axum_middleware::OtelTraceLayer
1118}
1119
1120#[cfg(feature = "axum")]
1151pub fn span_enricher_layer<T>() -> axum_middleware::SpanEnricherLayer<T>
1152where
1153 T: span_enrichment::EnrichSpan + Clone + Send + Sync + 'static,
1154{
1155 axum_middleware::SpanEnricherLayer::default()
1156}
1157
1158#[cfg(feature = "tonic-tracing")]
1180pub fn grpc_client_layer() -> grpc_middleware::GrpcClientTraceLayer {
1181 grpc_middleware::GrpcClientTraceLayer
1182}
1183
1184#[cfg(feature = "tonic-tracing")]
1199pub fn grpc_server_layer() -> grpc_middleware::GrpcServerTraceLayer {
1200 grpc_middleware::GrpcServerTraceLayer
1201}
1202
1203#[cfg(test)]
1204mod tests {
1205 use super::*;
1206 use std::sync::Mutex;
1207
1208 static ENV_LOCK: Mutex<()> = Mutex::new(());
1209
1210 #[test]
1211 fn resource_contains_all_attributes_when_provided() {
1212 let resource = build_resource("test-svc", Some("1.2.3"), Some("staging"));
1213
1214 assert_eq!(
1215 resource.get(&opentelemetry::Key::new("service.name")),
1216 Some(opentelemetry::Value::from("test-svc")),
1217 );
1218 assert_eq!(
1219 resource.get(&opentelemetry::Key::new(SERVICE_VERSION)),
1220 Some(opentelemetry::Value::from("1.2.3")),
1221 );
1222 assert_eq!(
1223 resource.get(&opentelemetry::Key::new(DEPLOYMENT_ENVIRONMENT_NAME)),
1224 Some(opentelemetry::Value::from("staging")),
1225 );
1226 assert!(resource.get(&opentelemetry::Key::new(HOST_NAME)).is_some());
1227 assert!(
1228 resource
1229 .get(&opentelemetry::Key::new(PROCESS_PID))
1230 .is_some()
1231 );
1232 }
1233
1234 #[test]
1235 fn resource_graceful_when_optional_values_omitted() {
1236 let resource = build_resource("test-svc", None, None);
1237
1238 assert_eq!(
1239 resource.get(&opentelemetry::Key::new("service.name")),
1240 Some(opentelemetry::Value::from("test-svc")),
1241 );
1242 assert!(
1243 resource
1244 .get(&opentelemetry::Key::new(SERVICE_VERSION))
1245 .is_none()
1246 );
1247 assert!(
1248 resource
1249 .get(&opentelemetry::Key::new(DEPLOYMENT_ENVIRONMENT_NAME))
1250 .is_none()
1251 );
1252 assert!(resource.get(&opentelemetry::Key::new(HOST_NAME)).is_some());
1254 assert!(
1255 resource
1256 .get(&opentelemetry::Key::new(PROCESS_PID))
1257 .is_some()
1258 );
1259 }
1260
1261 #[test]
1262 fn trace_sampler_ratio_converts_to_sdk() {
1263 let sampler = TraceSampler::TraceIdRatio(0.5);
1264 let sdk = sampler.into_sdk_sampler();
1265 assert_eq!(format!("{sdk:?}"), "TraceIdRatioBased(0.5)");
1266 }
1267
1268 #[test]
1269 fn trace_sampler_parent_based_converts_to_sdk() {
1270 let sampler = TraceSampler::ParentBased(Box::new(TraceSampler::TraceIdRatio(0.25)));
1271 let sdk = sampler.into_sdk_sampler();
1272 let debug = format!("{sdk:?}");
1273 assert!(debug.contains("ParentBased"));
1274 assert!(debug.contains("0.25"));
1275 }
1276
1277 unsafe fn set_env(key: &str, val: &str) {
1279 unsafe {
1280 std::env::set_var(key, val);
1281 }
1282 }
1283
1284 unsafe fn remove_env(key: &str) {
1285 unsafe {
1286 std::env::remove_var(key);
1287 }
1288 }
1289
1290 #[test]
1291 fn sampler_from_env_reads_traceidratio() {
1292 let _lock = ENV_LOCK.lock().unwrap();
1293 unsafe {
1294 set_env("OTEL_TRACES_SAMPLER", "traceidratio");
1295 set_env("OTEL_TRACES_SAMPLER_ARG", "0.42");
1296 }
1297
1298 let sampler = sampler_from_env()
1299 .expect("should not error")
1300 .expect("should return Some");
1301 assert!(
1302 matches!(sampler, TraceSampler::TraceIdRatio(r) if (r - 0.42).abs() < f64::EPSILON)
1303 );
1304
1305 unsafe {
1306 remove_env("OTEL_TRACES_SAMPLER");
1307 remove_env("OTEL_TRACES_SAMPLER_ARG");
1308 }
1309 }
1310
1311 #[test]
1312 fn sampler_from_env_returns_none_when_unset() {
1313 let _lock = ENV_LOCK.lock().unwrap();
1314 unsafe {
1315 remove_env("OTEL_TRACES_SAMPLER");
1316 }
1317 assert!(sampler_from_env().expect("should not error").is_none());
1318 }
1319
1320 #[test]
1321 fn sampler_from_env_reads_parentbased_traceidratio() {
1322 let _lock = ENV_LOCK.lock().unwrap();
1323 unsafe {
1324 set_env("OTEL_TRACES_SAMPLER", "parentbased_traceidratio");
1325 set_env("OTEL_TRACES_SAMPLER_ARG", "0.1");
1326 }
1327
1328 let sampler = sampler_from_env()
1329 .expect("should not error")
1330 .expect("should return Some");
1331 assert!(
1332 matches!(sampler, TraceSampler::ParentBased(inner) if matches!(*inner, TraceSampler::TraceIdRatio(r) if (r - 0.1).abs() < f64::EPSILON))
1333 );
1334
1335 unsafe {
1336 remove_env("OTEL_TRACES_SAMPLER");
1337 remove_env("OTEL_TRACES_SAMPLER_ARG");
1338 }
1339 }
1340
1341 #[test]
1342 fn sampler_from_env_parentbased_always_on() {
1343 let _lock = ENV_LOCK.lock().unwrap();
1344 unsafe {
1345 set_env("OTEL_TRACES_SAMPLER", "parentbased_always_on");
1346 }
1347 let sampler = sampler_from_env()
1348 .expect("should not error")
1349 .expect("should return Some");
1350 assert!(
1351 matches!(sampler, TraceSampler::ParentBased(inner) if matches!(*inner, TraceSampler::AlwaysOn))
1352 );
1353 unsafe {
1354 remove_env("OTEL_TRACES_SAMPLER");
1355 }
1356 }
1357
1358 #[test]
1359 fn sampler_from_env_parentbased_always_off() {
1360 let _lock = ENV_LOCK.lock().unwrap();
1361 unsafe {
1362 set_env("OTEL_TRACES_SAMPLER", "parentbased_always_off");
1363 }
1364 let sampler = sampler_from_env()
1365 .expect("should not error")
1366 .expect("should return Some");
1367 assert!(
1368 matches!(sampler, TraceSampler::ParentBased(inner) if matches!(*inner, TraceSampler::AlwaysOff))
1369 );
1370 unsafe {
1371 remove_env("OTEL_TRACES_SAMPLER");
1372 }
1373 }
1374
1375 #[test]
1376 fn sampler_from_env_always_on() {
1377 let _lock = ENV_LOCK.lock().unwrap();
1378 unsafe {
1379 set_env("OTEL_TRACES_SAMPLER", "always_on");
1380 }
1381 let sampler = sampler_from_env()
1382 .expect("should not error")
1383 .expect("should return Some");
1384 assert!(matches!(sampler, TraceSampler::AlwaysOn));
1385 unsafe {
1386 remove_env("OTEL_TRACES_SAMPLER");
1387 }
1388 }
1389
1390 #[test]
1391 fn sampler_from_env_always_off() {
1392 let _lock = ENV_LOCK.lock().unwrap();
1393 unsafe {
1394 set_env("OTEL_TRACES_SAMPLER", "always_off");
1395 }
1396 let sampler = sampler_from_env()
1397 .expect("should not error")
1398 .expect("should return Some");
1399 assert!(matches!(sampler, TraceSampler::AlwaysOff));
1400 unsafe {
1401 remove_env("OTEL_TRACES_SAMPLER");
1402 }
1403 }
1404
1405 #[test]
1406 fn sampler_from_env_unknown_returns_error() {
1407 let _lock = ENV_LOCK.lock().unwrap();
1408 unsafe {
1409 set_env("OTEL_TRACES_SAMPLER", "unknown_sampler");
1410 }
1411 let err = sampler_from_env().expect_err("unknown sampler should produce an error");
1412 assert!(
1413 err.to_string().contains("unknown_sampler"),
1414 "error message should include the unknown name, got: {err}"
1415 );
1416 unsafe {
1417 remove_env("OTEL_TRACES_SAMPLER");
1418 }
1419 }
1420
1421 #[test]
1422 fn trace_sampler_always_on_converts_to_sdk() {
1423 let sdk = TraceSampler::AlwaysOn.into_sdk_sampler();
1424 assert_eq!(format!("{sdk:?}"), "AlwaysOn");
1425 }
1426
1427 #[test]
1428 fn trace_sampler_always_off_converts_to_sdk() {
1429 let sdk = TraceSampler::AlwaysOff.into_sdk_sampler();
1430 assert_eq!(format!("{sdk:?}"), "AlwaysOff");
1431 }
1432
1433 #[test]
1434 fn builder_has_sensible_defaults() {
1435 let builder = Telemetry::builder("test-svc");
1436 assert_eq!(builder.service_name.as_deref(), Some("test-svc"));
1437 assert!(builder.service_version.is_none());
1438 assert!(builder.deployment_environment.is_none());
1439 assert!(builder.sampler.is_none());
1440 assert!(builder.metrics);
1441 assert!(!builder.logs);
1442 assert!(builder.protocol.is_none());
1443 assert!(builder.max_export_batch_size.is_none());
1444 assert!(builder.metric_export_interval.is_none());
1445 assert!(builder.export_timeout.is_none());
1446 }
1447
1448 #[test]
1449 fn from_env_builder_has_no_service_name() {
1450 let builder = Telemetry::from_env();
1451 assert!(builder.service_name.is_none());
1452 }
1453
1454 #[test]
1455 fn with_export_timeout_stores_value() {
1456 let timeout = Duration::from_secs(5);
1457 let builder = Telemetry::builder("test-svc").with_export_timeout(timeout);
1458 assert_eq!(builder.export_timeout, Some(timeout));
1459 }
1460
1461 #[test]
1462 fn timeout_from_env_reads_milliseconds() {
1463 let _lock = ENV_LOCK.lock().unwrap();
1464 unsafe {
1465 set_env("OTEL_EXPORTER_OTLP_TIMEOUT", "5000");
1466 }
1467 let t = timeout_from_env();
1468 assert_eq!(t, Some(Duration::from_millis(5000)));
1469 unsafe {
1470 remove_env("OTEL_EXPORTER_OTLP_TIMEOUT");
1471 }
1472 }
1473
1474 #[test]
1475 fn timeout_from_env_returns_none_when_unset() {
1476 let _lock = ENV_LOCK.lock().unwrap();
1477 unsafe {
1478 remove_env("OTEL_EXPORTER_OTLP_TIMEOUT");
1479 }
1480 assert_eq!(timeout_from_env(), None);
1481 }
1482
1483 #[test]
1484 fn service_name_from_env_used_when_none_given() {
1485 let builder = Telemetry::from_env();
1486 assert!(builder.service_name.is_none());
1487 }
1488
1489 #[test]
1490 fn explicit_service_name_overrides_env_var() {
1491 let builder = Telemetry::builder("explicit-svc");
1492 assert_eq!(builder.service_name.as_deref(), Some("explicit-svc"));
1493 }
1494
1495 #[test]
1496 fn from_env_builder_service_name_is_none() {
1497 let builder = Telemetry::from_env();
1498 assert!(builder.service_name.is_none());
1499 }
1500
1501 #[test]
1502 fn init_returns_error_for_unknown_otel_traces_sampler() {
1503 let _lock = ENV_LOCK.lock().unwrap();
1504 unsafe {
1505 set_env("OTEL_TRACES_SAMPLER", "not_a_real_sampler");
1506 }
1507 let result = Telemetry::builder("test-svc").with_metrics(false).init();
1508 let err = result
1509 .err()
1510 .expect("unknown sampler env var should cause init to fail");
1511 assert!(
1512 err.to_string().contains("not_a_real_sampler"),
1513 "error should name the unknown sampler, got: {err}"
1514 );
1515 unsafe {
1516 remove_env("OTEL_TRACES_SAMPLER");
1517 }
1518 }
1519
1520 #[test]
1521 fn with_max_export_batch_size_stores_value() {
1522 let builder = Telemetry::builder("test-svc").with_max_export_batch_size(1024);
1523 assert_eq!(builder.max_export_batch_size, Some(1024));
1524 }
1525
1526 #[test]
1527 fn with_metric_export_interval_stores_value() {
1528 let interval = Duration::from_secs(30);
1529 let builder = Telemetry::builder("test-svc").with_metric_export_interval(interval);
1530 assert_eq!(builder.metric_export_interval, Some(interval));
1531 }
1532
1533 #[test]
1534 fn init_rejects_zero_metric_export_interval() {
1535 let err = Telemetry::builder("test-svc")
1536 .with_metric_export_interval(Duration::ZERO)
1537 .with_metrics(false)
1538 .init()
1539 .err()
1540 .expect("expected error for zero interval");
1541 assert!(
1542 err.to_string().contains("metric_export_interval"),
1543 "error message should mention metric_export_interval, got: {err}"
1544 );
1545 }
1546
1547 #[test]
1548 fn builder_with_custom_values() {
1549 let builder = Telemetry::builder("test-svc")
1550 .with_version("2.0.0")
1551 .with_environment("production")
1552 .with_sampler(TraceSampler::TraceIdRatio(0.5))
1553 .with_metrics(false);
1554
1555 assert_eq!(builder.service_name.as_deref(), Some("test-svc"));
1556 assert_eq!(builder.service_version.as_deref(), Some("2.0.0"));
1557 assert_eq!(
1558 builder.deployment_environment.as_deref(),
1559 Some("production")
1560 );
1561 assert!(
1562 matches!(builder.sampler, Some(TraceSampler::TraceIdRatio(r)) if (r - 0.5).abs() < f64::EPSILON)
1563 );
1564 assert!(!builder.metrics);
1565 }
1566
1567 #[test]
1568 #[cfg(feature = "grpc")]
1569 fn builder_with_protocol_grpc() {
1570 let builder = Telemetry::builder("test-svc").with_protocol(ExportProtocol::Grpc);
1571 assert_eq!(builder.protocol, Some(ExportProtocol::Grpc));
1572 }
1573
1574 #[test]
1575 #[cfg(feature = "http")]
1576 fn builder_with_protocol_http() {
1577 let builder = Telemetry::builder("test-svc").with_protocol(ExportProtocol::HttpProtobuf);
1578 assert_eq!(builder.protocol, Some(ExportProtocol::HttpProtobuf));
1579 }
1580
1581 #[test]
1582 #[cfg(feature = "grpc")]
1583 fn protocol_from_env_reads_grpc() {
1584 let _lock = ENV_LOCK.lock().unwrap();
1585 unsafe {
1586 set_env("OTEL_EXPORTER_OTLP_PROTOCOL", "grpc");
1587 }
1588 assert_eq!(protocol_from_env(), Some(ExportProtocol::Grpc));
1589 unsafe {
1590 remove_env("OTEL_EXPORTER_OTLP_PROTOCOL");
1591 }
1592 }
1593
1594 #[test]
1595 #[cfg(feature = "http")]
1596 fn protocol_from_env_reads_http_protobuf() {
1597 let _lock = ENV_LOCK.lock().unwrap();
1598 unsafe {
1599 set_env("OTEL_EXPORTER_OTLP_PROTOCOL", "http/protobuf");
1600 }
1601 assert_eq!(protocol_from_env(), Some(ExportProtocol::HttpProtobuf));
1602 unsafe {
1603 remove_env("OTEL_EXPORTER_OTLP_PROTOCOL");
1604 }
1605 }
1606
1607 #[test]
1608 fn protocol_from_env_returns_none_when_unset() {
1609 let _lock = ENV_LOCK.lock().unwrap();
1610 unsafe {
1611 remove_env("OTEL_EXPORTER_OTLP_PROTOCOL");
1612 }
1613 assert_eq!(protocol_from_env(), None);
1614 }
1615
1616 #[test]
1617 fn protocol_from_env_returns_none_for_unknown() {
1618 let _lock = ENV_LOCK.lock().unwrap();
1619 unsafe {
1620 set_env("OTEL_EXPORTER_OTLP_PROTOCOL", "websocket");
1621 }
1622 assert_eq!(protocol_from_env(), None);
1623 unsafe {
1624 remove_env("OTEL_EXPORTER_OTLP_PROTOCOL");
1625 }
1626 }
1627
1628 #[test]
1629 fn builder_is_send_and_sync() {
1630 fn assert_send_sync<T: Send + Sync>() {}
1631 assert_send_sync::<TelemetryBuilder>();
1632 }
1633
1634 #[test]
1635 fn with_shutdown_timeout_stores_value() {
1636 let timeout = Duration::from_secs(10);
1637 let builder = Telemetry::builder("test-svc").with_shutdown_timeout(timeout);
1638 assert_eq!(builder.shutdown_timeout, timeout);
1639 }
1640
1641 #[test]
1642 fn default_shutdown_timeout_is_five_seconds() {
1643 let builder = Telemetry::builder("test-svc");
1644 assert_eq!(builder.shutdown_timeout, Duration::from_secs(5));
1645 }
1646
1647 #[cfg(feature = "testing")]
1655 #[test]
1656 fn drop_completes_within_shutdown_timeout() {
1657 let mut handles = crate::Telemetry::testing("drop-timeout-test");
1659 handles.shutdown_timeout = Duration::from_millis(100);
1661
1662 let start = std::time::Instant::now();
1663 drop(handles);
1664 let elapsed = start.elapsed();
1665
1666 assert!(
1668 elapsed < Duration::from_millis(500),
1669 "drop took {elapsed:?}, expected < 500 ms"
1670 );
1671 }
1672}