1use crate::error::{OTelSdkError, OTelSdkResult};
38use crate::resource::Resource;
39use crate::trace::Span;
40use crate::trace::{SpanData, SpanExporter};
41use opentelemetry::Context;
42#[cfg(feature = "experimental_metrics_bound_instruments")]
43use opentelemetry::KeyValue;
44use opentelemetry::{otel_debug, otel_error, otel_warn};
45use std::cmp::min;
46use std::sync::atomic::{AtomicUsize, Ordering};
47use std::sync::{Arc, Mutex};
48use std::{env, str::FromStr, time::Duration};
49
50use std::sync::atomic::AtomicBool;
51use std::thread;
52use std::time::Instant;
53
54pub const OTEL_BSP_SCHEDULE_DELAY: &str = "OTEL_BSP_SCHEDULE_DELAY";
57pub const OTEL_BSP_SCHEDULE_DELAY_DEFAULT: Duration = Duration::from_millis(5_000);
59pub const OTEL_BSP_MAX_QUEUE_SIZE: &str = "OTEL_BSP_MAX_QUEUE_SIZE";
62pub const OTEL_BSP_MAX_QUEUE_SIZE_DEFAULT: usize = 2_048;
64pub const OTEL_BSP_MAX_EXPORT_BATCH_SIZE: &str = "OTEL_BSP_MAX_EXPORT_BATCH_SIZE";
68pub const OTEL_BSP_MAX_EXPORT_BATCH_SIZE_DEFAULT: usize = 512;
70pub const OTEL_BSP_EXPORT_TIMEOUT: &str = "OTEL_BSP_EXPORT_TIMEOUT";
77pub const OTEL_BSP_EXPORT_TIMEOUT_DEFAULT: Duration = Duration::from_millis(30_000);
81pub(crate) const OTEL_BSP_MAX_CONCURRENT_EXPORTS: &str = "OTEL_BSP_MAX_CONCURRENT_EXPORTS";
82pub(crate) const OTEL_BSP_MAX_CONCURRENT_EXPORTS_DEFAULT: usize = 1;
84
85pub trait SpanProcessor: Send + Sync + std::fmt::Debug {
89 fn on_start(&self, span: &mut Span, cx: &Context);
93
94 fn on_end(&self, span: SpanData);
203 fn force_flush(&self) -> OTelSdkResult;
205 fn shutdown_with_timeout(&self, timeout: Duration) -> OTelSdkResult;
210 fn shutdown(&self) -> OTelSdkResult {
212 self.shutdown_with_timeout(Duration::from_secs(5))
213 }
214 fn set_resource(&mut self, _resource: &Resource) {}
216}
217
218#[derive(Debug)]
239pub struct SimpleSpanProcessor<T: SpanExporter> {
240 exporter: Mutex<T>,
241 is_shutdown: AtomicBool,
242
243 #[cfg(feature = "experimental_metrics_bound_instruments")]
248 processed_success: opentelemetry::metrics::BoundCounter<u64>,
249 #[cfg(feature = "experimental_metrics_bound_instruments")]
250 processed_after_shutdown: opentelemetry::metrics::BoundCounter<u64>,
251}
252
253impl<T: SpanExporter> SimpleSpanProcessor<T> {
254 pub fn new(exporter: T) -> Self {
256 #[cfg(feature = "experimental_metrics_bound_instruments")]
257 let (processed_success, processed_after_shutdown) = {
258 static INSTANCE_COUNTER: AtomicUsize = AtomicUsize::new(0);
259 let instance_id = INSTANCE_COUNTER.fetch_add(1, Ordering::Relaxed);
260 let component_name = format!("simple_span_processor/{instance_id}");
261
262 let meter = opentelemetry::global::meter("otel.sdk");
263 let counter = meter
264 .u64_counter("otel.sdk.processor.span.processed")
265 .with_description(
266 "The number of spans for which the processing has finished, \
267 either successful or failed.",
268 )
269 .with_unit("{span}")
270 .build();
271
272 let success_attrs = [
276 KeyValue::new("otel.component.type", "simple_span_processor"),
277 KeyValue::new("otel.component.name", component_name.clone()),
278 ];
279 let after_shutdown_attrs = [
280 KeyValue::new("error.type", "already_shutdown"),
281 KeyValue::new("otel.component.type", "simple_span_processor"),
282 KeyValue::new("otel.component.name", component_name),
283 ];
284
285 (
286 counter.bind(&success_attrs),
287 counter.bind(&after_shutdown_attrs),
288 )
289 };
290
291 Self {
292 exporter: Mutex::new(exporter),
293 is_shutdown: AtomicBool::new(false),
294 #[cfg(feature = "experimental_metrics_bound_instruments")]
295 processed_success,
296 #[cfg(feature = "experimental_metrics_bound_instruments")]
297 processed_after_shutdown,
298 }
299 }
300}
301
302impl<T: SpanExporter> SpanProcessor for SimpleSpanProcessor<T> {
303 fn on_start(&self, _span: &mut Span, _cx: &Context) {
304 }
306
307 fn on_end(&self, span: SpanData) {
308 if !span.span_context.is_sampled() {
309 return;
310 }
311
312 if self.is_shutdown.load(Ordering::Relaxed) {
314 #[cfg(feature = "experimental_metrics_bound_instruments")]
316 self.processed_after_shutdown.add(1);
317 otel_warn!(
318 name: "SimpleSpanProcessor.OnEnd.AfterShutdown",
319 message = "Spans are being emitted even after Shutdown. This indicates incorrect lifecycle management of TracerProvider in application. Spans will not be exported."
320 );
321 return;
322 }
323
324 let result = match self.exporter.lock() {
325 Ok(exporter) => {
326 #[cfg(feature = "experimental_metrics_bound_instruments")]
329 self.processed_success.add(1);
330 futures_executor::block_on(exporter.export(vec![span]))
331 }
332 Err(_) => Err(OTelSdkError::InternalFailure(
333 "SimpleSpanProcessor mutex poison".into(),
334 )),
335 };
336
337 if let Err(err) = result {
338 otel_debug!(
340 name: "SimpleProcessor.OnEnd.Error",
341 reason = format!("{:?}", err)
342 );
343 }
344 }
345
346 fn force_flush(&self) -> OTelSdkResult {
347 Ok(())
349 }
350
351 fn shutdown_with_timeout(&self, timeout: Duration) -> OTelSdkResult {
352 self.is_shutdown.store(true, Ordering::Relaxed);
353 if let Ok(exporter) = self.exporter.lock() {
354 exporter.shutdown_with_timeout(timeout)
355 } else {
356 Err(OTelSdkError::InternalFailure(
357 "SimpleSpanProcessor mutex poison at shutdown".into(),
358 ))
359 }
360 }
361
362 fn set_resource(&mut self, resource: &Resource) {
363 if let Ok(mut exporter) = self.exporter.lock() {
364 exporter.set_resource(resource);
365 }
366 }
367}
368
369use std::sync::mpsc::sync_channel;
435use std::sync::mpsc::Receiver;
436use std::sync::mpsc::RecvTimeoutError;
437use std::sync::mpsc::SyncSender;
438
439#[allow(clippy::large_enum_variant)]
441#[derive(Debug)]
442enum BatchMessage {
443 ExportSpan(Arc<AtomicBool>),
445 ForceFlush(SyncSender<OTelSdkResult>),
446 Shutdown(SyncSender<OTelSdkResult>),
447 SetResource(Arc<Resource>),
448}
449
450#[derive(Debug)]
486pub struct BatchSpanProcessor {
487 span_sender: SyncSender<SpanData>, message_sender: SyncSender<BatchMessage>, handle: Mutex<Option<thread::JoinHandle<()>>>,
490 forceflush_timeout: Duration,
491 export_span_message_sent: Arc<AtomicBool>,
492 current_batch_size: Arc<AtomicUsize>,
493 max_export_batch_size: usize,
494 dropped_spans_count: AtomicUsize,
495 max_queue_size: usize,
496
497 #[cfg(feature = "experimental_metrics_bound_instruments")]
503 processed_queue_full: opentelemetry::metrics::BoundCounter<u64>,
504 #[cfg(feature = "experimental_metrics_bound_instruments")]
505 processed_after_shutdown: opentelemetry::metrics::BoundCounter<u64>,
506}
507
508impl BatchSpanProcessor {
509 pub fn new<E>(
511 mut exporter: E,
512 config: BatchConfig,
513 ) -> Self
517 where
518 E: SpanExporter + Send + 'static,
519 {
520 let (span_sender, span_receiver) = sync_channel::<SpanData>(config.max_queue_size);
521 let (message_sender, message_receiver) = sync_channel::<BatchMessage>(64); let max_queue_size = config.max_queue_size;
523 let max_export_batch_size = config.max_export_batch_size;
524 let current_batch_size = Arc::new(AtomicUsize::new(0));
525 let current_batch_size_for_thread = current_batch_size.clone();
526
527 #[cfg(feature = "experimental_metrics_bound_instruments")]
531 let (processed_success, processed_queue_full, processed_after_shutdown) = {
532 static INSTANCE_COUNTER: AtomicUsize = AtomicUsize::new(0);
533 let instance_id = INSTANCE_COUNTER.fetch_add(1, Ordering::Relaxed);
534 let component_name = format!("batching_span_processor/{instance_id}");
535
536 let meter = opentelemetry::global::meter("otel.sdk");
537 let counter = meter
538 .u64_counter("otel.sdk.processor.span.processed")
539 .with_description(
540 "The number of spans for which the processing has finished, \
541 either successful or failed.",
542 )
543 .with_unit("{span}")
544 .build();
545
546 let success_attrs = [
550 KeyValue::new("otel.component.type", "batching_span_processor"),
551 KeyValue::new("otel.component.name", component_name.clone()),
552 ];
553 let queue_full_attrs = [
554 KeyValue::new("error.type", "queue_full"),
555 KeyValue::new("otel.component.type", "batching_span_processor"),
556 KeyValue::new("otel.component.name", component_name.clone()),
557 ];
558 let after_shutdown_attrs = [
559 KeyValue::new("error.type", "already_shutdown"),
560 KeyValue::new("otel.component.type", "batching_span_processor"),
561 KeyValue::new("otel.component.name", component_name),
562 ];
563
564 (
565 counter.bind(&success_attrs),
566 counter.bind(&queue_full_attrs),
567 counter.bind(&after_shutdown_attrs),
568 )
569 };
570
571 let handle = thread::Builder::new()
572 .name("OpenTelemetry.Traces.BatchProcessor".to_string())
573 .spawn(move || {
574 let _suppress_guard = Context::enter_telemetry_suppressed_scope();
575 otel_debug!(
576 name: "BatchSpanProcessor.ThreadStarted",
577 interval_in_millisecs = config.scheduled_delay.as_millis(),
578 max_export_batch_size = config.max_export_batch_size,
579 max_queue_size = config.max_queue_size,
580 );
581 let mut spans = Vec::with_capacity(config.max_export_batch_size);
582 let mut last_export_time = Instant::now();
583 let current_batch_size = current_batch_size_for_thread;
584
585 #[cfg(feature = "experimental_metrics_bound_instruments")]
588 let record_processed_success = move |count: u64| processed_success.add(count);
589 #[cfg(not(feature = "experimental_metrics_bound_instruments"))]
590 let record_processed_success = |_count: u64| {};
591 loop {
592 let remaining_time_option = config
593 .scheduled_delay
594 .checked_sub(last_export_time.elapsed());
595 let remaining_time = match remaining_time_option {
596 Some(remaining_time) => remaining_time,
597 None => config.scheduled_delay,
598 };
599 match message_receiver.recv_timeout(remaining_time) {
600 Ok(message) => match message {
601 BatchMessage::ExportSpan(export_span_message_sent) => {
602 export_span_message_sent.store(false, Ordering::Relaxed);
604 otel_debug!(
605 name: "BatchSpanProcessor.ExportingDueToBatchSize",
606 );
607 let _ = Self::get_spans_and_export(
608 &span_receiver,
609 &exporter,
610 &mut spans,
611 &mut last_export_time,
612 ¤t_batch_size,
613 &config,
614 &record_processed_success,
615 );
616 }
617 BatchMessage::ForceFlush(sender) => {
618 otel_debug!(name: "BatchSpanProcessor.ExportingDueToForceFlush");
619 let result = Self::get_spans_and_export(
620 &span_receiver,
621 &exporter,
622 &mut spans,
623 &mut last_export_time,
624 ¤t_batch_size,
625 &config,
626 &record_processed_success,
627 );
628 let _ = sender.send(result);
629 }
630 BatchMessage::Shutdown(sender) => {
631 otel_debug!(name: "BatchSpanProcessor.ExportingDueToShutdown");
632 let result = Self::get_spans_and_export(
633 &span_receiver,
634 &exporter,
635 &mut spans,
636 &mut last_export_time,
637 ¤t_batch_size,
638 &config,
639 &record_processed_success,
640 );
641 let _ = exporter.shutdown();
642 let _ = sender.send(result);
643
644 otel_debug!(
645 name: "BatchSpanProcessor.ThreadExiting",
646 reason = "ShutdownRequested"
647 );
648 break;
652 }
653 BatchMessage::SetResource(resource) => {
654 exporter.set_resource(&resource);
655 }
656 },
657 Err(RecvTimeoutError::Timeout) => {
658 otel_debug!(
659 name: "BatchSpanProcessor.ExportingDueToTimer",
660 );
661
662 let _ = Self::get_spans_and_export(
663 &span_receiver,
664 &exporter,
665 &mut spans,
666 &mut last_export_time,
667 ¤t_batch_size,
668 &config,
669 &record_processed_success,
670 );
671 }
672 Err(RecvTimeoutError::Disconnected) => {
673 otel_debug!(
676 name: "BatchSpanProcessor.ThreadExiting",
677 reason = "MessageSenderDisconnected"
678 );
679 break;
680 }
681 }
682 }
683 otel_debug!(
684 name: "BatchSpanProcessor.ThreadStopped"
685 );
686 })
687 .expect("Failed to spawn thread"); Self {
690 span_sender,
691 message_sender,
692 handle: Mutex::new(Some(handle)),
693 forceflush_timeout: Duration::from_secs(5), dropped_spans_count: AtomicUsize::new(0),
695 max_queue_size,
696 export_span_message_sent: Arc::new(AtomicBool::new(false)),
697 current_batch_size,
698 max_export_batch_size,
699 #[cfg(feature = "experimental_metrics_bound_instruments")]
700 processed_queue_full,
701 #[cfg(feature = "experimental_metrics_bound_instruments")]
702 processed_after_shutdown,
703 }
704 }
705
706 pub fn builder<E>(exporter: E) -> BatchSpanProcessorBuilder<E>
708 where
709 E: SpanExporter + Send + 'static,
710 {
711 BatchSpanProcessorBuilder {
712 exporter,
713 config: BatchConfig::default(),
714 }
715 }
716
717 #[inline]
721 fn get_spans_and_export<E, F>(
722 spans_receiver: &Receiver<SpanData>,
723 exporter: &E,
724 spans: &mut Vec<SpanData>,
725 last_export_time: &mut Instant,
726 current_batch_size: &AtomicUsize,
727 config: &BatchConfig,
728 record_processed_success: &F,
729 ) -> OTelSdkResult
730 where
731 E: SpanExporter + Send + Sync + 'static,
732 F: Fn(u64),
733 {
734 let target = current_batch_size.load(Ordering::Acquire); let mut result = OTelSdkResult::Ok(());
736 let mut total_exported_spans: usize = 0;
737
738 while target > 0 && total_exported_spans < target {
739 let batch_limit = config
740 .max_export_batch_size
741 .min(target - total_exported_spans);
742
743 while let Ok(span) = spans_receiver.try_recv() {
745 spans.push(span);
746 if spans.len() == batch_limit {
747 break;
748 }
749 }
750
751 let count_of_spans = spans.len(); if count_of_spans == 0 {
753 break;
754 }
755 total_exported_spans += count_of_spans;
756
757 record_processed_success(count_of_spans as u64);
760
761 result = Self::export_batch_sync(exporter, spans, last_export_time); current_batch_size.fetch_sub(count_of_spans, Ordering::AcqRel);
764 }
765 result
766 }
767
768 #[allow(clippy::vec_box)]
769 fn export_batch_sync<E>(
770 exporter: &E,
771 batch: &mut Vec<SpanData>,
772 last_export_time: &mut Instant,
773 ) -> OTelSdkResult
774 where
775 E: SpanExporter + ?Sized,
776 {
777 *last_export_time = Instant::now();
778
779 if batch.is_empty() {
780 return OTelSdkResult::Ok(());
781 }
782
783 let export = exporter.export(batch.split_off(0));
790 let export_result = futures_executor::block_on(export);
791
792 match export_result {
793 Ok(_) => OTelSdkResult::Ok(()),
794 Err(err) => {
795 otel_error!(
796 name: "BatchSpanProcessor.ExportError",
797 error = format!("{}", err)
798 );
799 OTelSdkResult::Err(err)
800 }
801 }
802 }
803}
804
805impl SpanProcessor for BatchSpanProcessor {
806 fn on_start(&self, _span: &mut Span, _cx: &Context) {
808 }
810
811 fn on_end(&self, span: SpanData) {
813 let previous_batch_size = self.current_batch_size.fetch_add(1, Ordering::AcqRel);
818 let result = self.span_sender.try_send(span);
819
820 match result {
822 Ok(_) => {
823 if previous_batch_size + 1 >= self.max_export_batch_size {
827 if !self.export_span_message_sent.load(Ordering::Relaxed) {
834 if !self.export_span_message_sent.swap(true, Ordering::Relaxed) {
844 match self.message_sender.try_send(BatchMessage::ExportSpan(
845 self.export_span_message_sent.clone(),
846 )) {
847 Ok(_) => {
848 }
850 Err(_err) => {
851 self.export_span_message_sent
855 .store(false, Ordering::Relaxed);
856 }
857 }
858 }
859 }
860 }
861 }
862 Err(std::sync::mpsc::TrySendError::Full(_)) => {
863 self.current_batch_size.fetch_sub(1, Ordering::AcqRel);
865 #[cfg(feature = "experimental_metrics_bound_instruments")]
867 self.processed_queue_full.add(1);
868 if self.dropped_spans_count.fetch_add(1, Ordering::Relaxed) == 0 {
871 otel_warn!(name: "BatchSpanProcessor.SpanDroppingStarted",
872 message = "BatchSpanProcessor dropped a Span due to queue full. No further log will be emitted for further drops until Shutdown. During Shutdown time, a log will be emitted with exact count of total spans dropped.");
873 }
874 }
875 Err(std::sync::mpsc::TrySendError::Disconnected(_)) => {
876 self.current_batch_size.fetch_sub(1, Ordering::AcqRel);
878 #[cfg(feature = "experimental_metrics_bound_instruments")]
880 self.processed_after_shutdown.add(1);
881 otel_warn!(
884 name: "BatchSpanProcessor.OnEnd.AfterShutdown",
885 message = "Spans are being emitted even after Shutdown. This indicates incorrect lifecycle management of TracerProvider in application. Spans will not be exported."
886 );
887 }
888 }
889 }
890
891 fn force_flush(&self) -> OTelSdkResult {
893 let (sender, receiver) = std::sync::mpsc::sync_channel(1);
894 match self
895 .message_sender
896 .try_send(BatchMessage::ForceFlush(sender))
897 {
898 Ok(_) => receiver
899 .recv_timeout(self.forceflush_timeout)
900 .map_err(|err| {
901 if err == std::sync::mpsc::RecvTimeoutError::Timeout {
902 OTelSdkError::Timeout(self.forceflush_timeout)
903 } else {
904 OTelSdkError::InternalFailure(format!("{err}"))
905 }
906 })?,
907 Err(std::sync::mpsc::TrySendError::Full(_)) => {
908 otel_debug!(
910 name: "BatchSpanProcessor.ForceFlush.ControlChannelFull",
911 message = "Control message to flush the worker thread could not be sent as the control channel is full. This can occur if user repeatedly calls force_flush/shutdown without finishing the previous call."
912 );
913 Err(OTelSdkError::InternalFailure("ForceFlush cannot be performed as Control channel is full. This can occur if user repeatedly calls force_flush/shutdown without finishing the previous call.".into()))
914 }
915 Err(std::sync::mpsc::TrySendError::Disconnected(_)) => {
916 otel_debug!(
919 name: "BatchSpanProcessor.ForceFlush.AlreadyShutdown",
920 message = "ForceFlush invoked after Shutdown. This will not perform Flush and indicates a incorrect lifecycle management in Application."
921 );
922
923 Err(OTelSdkError::AlreadyShutdown)
924 }
925 }
926 }
927
928 fn shutdown_with_timeout(&self, timeout: Duration) -> OTelSdkResult {
930 let dropped_spans = self.dropped_spans_count.load(Ordering::Relaxed);
931 let max_queue_size = self.max_queue_size;
932 if dropped_spans > 0 {
933 otel_warn!(
934 name: "BatchSpanProcessor.SpansDropped",
935 dropped_span_count = dropped_spans,
936 max_queue_size = max_queue_size,
937 message = "Spans were dropped due to a queue being full. The count represents the total count of spans dropped in the lifetime of this BatchSpanProcessor. Consider increasing the queue size and/or decrease delay between intervals."
938 );
939 }
940
941 let (sender, receiver) = std::sync::mpsc::sync_channel(1);
942 match self.message_sender.try_send(BatchMessage::Shutdown(sender)) {
943 Ok(_) => {
944 receiver
945 .recv_timeout(timeout)
946 .map(|_| {
947 if let Some(handle) = self.handle.lock().unwrap().take() {
950 handle.join().unwrap();
951 }
952 OTelSdkResult::Ok(())
953 })
954 .map_err(|err| match err {
955 std::sync::mpsc::RecvTimeoutError::Timeout => {
956 otel_error!(
957 name: "BatchSpanProcessor.Shutdown.Timeout",
958 message = "BatchSpanProcessor shutdown timing out."
959 );
960 OTelSdkError::Timeout(timeout)
961 }
962 _ => {
963 otel_error!(
964 name: "BatchSpanProcessor.Shutdown.Error",
965 error = format!("{}", err)
966 );
967 OTelSdkError::InternalFailure(format!("{err}"))
968 }
969 })?
970 }
971 Err(std::sync::mpsc::TrySendError::Full(_)) => {
972 otel_debug!(
974 name: "BatchSpanProcessor.Shutdown.ControlChannelFull",
975 message = "Control message to shutdown the worker thread could not be sent as the control channel is full. This can occur if user repeatedly calls force_flush/shutdown without finishing the previous call."
976 );
977 Err(OTelSdkError::InternalFailure("Shutdown cannot be performed as Control channel is full. This can occur if user repeatedly calls force_flush/shutdown without finishing the previous call.".into()))
978 }
979 Err(std::sync::mpsc::TrySendError::Disconnected(_)) => {
980 otel_debug!(
983 name: "BatchSpanProcessor.Shutdown.AlreadyShutdown",
984 message = "Shutdown is being invoked more than once. This is noop, but indicates a potential issue in the application's lifecycle management."
985 );
986
987 Err(OTelSdkError::AlreadyShutdown)
988 }
989 }
990 }
991
992 fn set_resource(&mut self, resource: &Resource) {
994 let resource = Arc::new(resource.clone());
995 let _ = self
996 .message_sender
997 .try_send(BatchMessage::SetResource(resource));
998 }
999}
1000
1001#[derive(Debug, Default)]
1003pub struct BatchSpanProcessorBuilder<E>
1004where
1005 E: SpanExporter + Send + 'static,
1006{
1007 exporter: E,
1008 config: BatchConfig,
1009}
1010
1011impl<E> BatchSpanProcessorBuilder<E>
1012where
1013 E: SpanExporter + Send + 'static,
1014{
1015 pub fn with_batch_config(self, config: BatchConfig) -> Self {
1017 BatchSpanProcessorBuilder { config, ..self }
1018 }
1019
1020 pub fn build(self) -> BatchSpanProcessor {
1022 BatchSpanProcessor::new(self.exporter, self.config)
1023 }
1024}
1025
1026#[derive(Debug)]
1029pub struct BatchConfig {
1030 pub(crate) max_queue_size: usize,
1033
1034 pub(crate) scheduled_delay: Duration,
1037
1038 #[allow(dead_code)]
1039 pub(crate) max_export_batch_size: usize,
1044
1045 #[allow(dead_code)]
1046 pub(crate) max_export_timeout: Duration,
1048
1049 #[allow(dead_code)]
1050 pub(crate) max_concurrent_exports: usize,
1051}
1052
1053impl Default for BatchConfig {
1054 fn default() -> Self {
1055 BatchConfigBuilder::default().build()
1056 }
1057}
1058
1059#[derive(Debug)]
1061pub struct BatchConfigBuilder {
1062 max_queue_size: usize,
1063 scheduled_delay: Duration,
1064 max_export_batch_size: usize,
1065 max_export_timeout: Duration,
1066 max_concurrent_exports: usize,
1067}
1068
1069impl Default for BatchConfigBuilder {
1070 fn default() -> Self {
1081 BatchConfigBuilder {
1082 max_queue_size: OTEL_BSP_MAX_QUEUE_SIZE_DEFAULT,
1083 scheduled_delay: OTEL_BSP_SCHEDULE_DELAY_DEFAULT,
1084 max_export_batch_size: OTEL_BSP_MAX_EXPORT_BATCH_SIZE_DEFAULT,
1085 max_export_timeout: OTEL_BSP_EXPORT_TIMEOUT_DEFAULT,
1086 max_concurrent_exports: OTEL_BSP_MAX_CONCURRENT_EXPORTS_DEFAULT,
1087 }
1088 .init_from_env_vars()
1089 }
1090}
1091
1092impl BatchConfigBuilder {
1093 pub fn with_max_queue_size(mut self, max_queue_size: usize) -> Self {
1102 self.max_queue_size = max_queue_size;
1103 self
1104 }
1105
1106 pub fn with_max_export_batch_size(mut self, max_export_batch_size: usize) -> Self {
1116 self.max_export_batch_size = max_export_batch_size;
1117 self
1118 }
1119
1120 #[cfg(feature = "experimental_trace_batch_span_processor_with_async_runtime")]
1121 pub fn with_max_concurrent_exports(mut self, max_concurrent_exports: usize) -> Self {
1138 self.max_concurrent_exports = max_concurrent_exports;
1139 self
1140 }
1141
1142 pub fn with_scheduled_delay(mut self, scheduled_delay: Duration) -> Self {
1150 self.scheduled_delay = scheduled_delay;
1151 self
1152 }
1153
1154 #[cfg(feature = "experimental_trace_batch_span_processor_with_async_runtime")]
1162 pub fn with_max_export_timeout(mut self, max_export_timeout: Duration) -> Self {
1163 self.max_export_timeout = max_export_timeout;
1164 self
1165 }
1166
1167 pub fn build(self) -> BatchConfig {
1170 let max_export_batch_size = min(self.max_export_batch_size, self.max_queue_size);
1173
1174 BatchConfig {
1175 max_queue_size: self.max_queue_size,
1176 scheduled_delay: self.scheduled_delay,
1177 max_export_timeout: self.max_export_timeout,
1178 max_concurrent_exports: self.max_concurrent_exports,
1179 max_export_batch_size,
1180 }
1181 }
1182
1183 fn init_from_env_vars(mut self) -> Self {
1184 if let Some(max_concurrent_exports) = env::var(OTEL_BSP_MAX_CONCURRENT_EXPORTS)
1185 .ok()
1186 .and_then(|max_concurrent_exports| usize::from_str(&max_concurrent_exports).ok())
1187 {
1188 self.max_concurrent_exports = max_concurrent_exports;
1189 }
1190
1191 if let Some(max_queue_size) = env::var(OTEL_BSP_MAX_QUEUE_SIZE)
1192 .ok()
1193 .and_then(|queue_size| usize::from_str(&queue_size).ok())
1194 {
1195 self.max_queue_size = max_queue_size;
1196 }
1197
1198 if let Some(scheduled_delay) = env::var(OTEL_BSP_SCHEDULE_DELAY)
1199 .ok()
1200 .and_then(|delay| u64::from_str(&delay).ok())
1201 {
1202 self.scheduled_delay = Duration::from_millis(scheduled_delay);
1203 }
1204
1205 if let Some(max_export_batch_size) = env::var(OTEL_BSP_MAX_EXPORT_BATCH_SIZE)
1206 .ok()
1207 .and_then(|batch_size| usize::from_str(&batch_size).ok())
1208 {
1209 self.max_export_batch_size = max_export_batch_size;
1210 }
1211
1212 if self.max_export_batch_size > self.max_queue_size {
1215 self.max_export_batch_size = self.max_queue_size;
1216 }
1217
1218 if let Some(max_export_timeout) = env::var(OTEL_BSP_EXPORT_TIMEOUT)
1219 .ok()
1220 .and_then(|timeout| u64::from_str(&timeout).ok())
1221 {
1222 self.max_export_timeout = Duration::from_millis(max_export_timeout);
1223 }
1224
1225 self
1226 }
1227}
1228
1229#[cfg(all(test, feature = "testing", feature = "trace"))]
1230mod tests {
1231 use super::{
1233 BatchSpanProcessor, SimpleSpanProcessor, SpanProcessor, OTEL_BSP_EXPORT_TIMEOUT,
1234 OTEL_BSP_MAX_EXPORT_BATCH_SIZE, OTEL_BSP_MAX_QUEUE_SIZE, OTEL_BSP_MAX_QUEUE_SIZE_DEFAULT,
1235 OTEL_BSP_SCHEDULE_DELAY, OTEL_BSP_SCHEDULE_DELAY_DEFAULT,
1236 };
1237 use crate::error::OTelSdkResult;
1238 use crate::testing::trace::new_test_export_span_data;
1239 use crate::trace::span_processor::{
1240 OTEL_BSP_EXPORT_TIMEOUT_DEFAULT, OTEL_BSP_MAX_CONCURRENT_EXPORTS,
1241 OTEL_BSP_MAX_CONCURRENT_EXPORTS_DEFAULT, OTEL_BSP_MAX_EXPORT_BATCH_SIZE_DEFAULT,
1242 };
1243 use crate::trace::InMemorySpanExporterBuilder;
1244 use crate::trace::{BatchConfig, BatchConfigBuilder, SpanEvents, SpanLinks};
1245 use crate::trace::{SpanData, SpanExporter};
1246 use opentelemetry::trace::{SpanContext, SpanId, SpanKind, Status};
1247 use std::fmt::Debug;
1248 use std::time::Duration;
1249
1250 #[test]
1251 fn simple_span_processor_on_end_calls_export() {
1252 let exporter = InMemorySpanExporterBuilder::new().build();
1253 let processor = SimpleSpanProcessor::new(exporter.clone());
1254 let span_data = new_test_export_span_data();
1255 processor.on_end(span_data.clone());
1256 assert_eq!(exporter.get_finished_spans().unwrap()[0], span_data);
1257 let _result = processor.shutdown();
1258 }
1259
1260 #[test]
1261 fn simple_span_processor_on_end_skips_export_if_not_sampled() {
1262 let exporter = InMemorySpanExporterBuilder::new().build();
1263 let processor = SimpleSpanProcessor::new(exporter.clone());
1264 let unsampled = SpanData {
1265 span_context: SpanContext::empty_context(),
1266 parent_span_id: SpanId::INVALID,
1267 parent_span_is_remote: false,
1268 span_kind: SpanKind::Internal,
1269 name: "opentelemetry".into(),
1270 start_time: opentelemetry::time::now(),
1271 end_time: opentelemetry::time::now(),
1272 attributes: Vec::new(),
1273 dropped_attributes_count: 0,
1274 events: SpanEvents::default(),
1275 links: SpanLinks::default(),
1276 status: Status::Unset,
1277 instrumentation_scope: Default::default(),
1278 };
1279 processor.on_end(unsampled);
1280 assert!(exporter.get_finished_spans().unwrap().is_empty());
1281 }
1282
1283 #[test]
1284 fn simple_span_processor_shutdown_calls_shutdown() {
1285 let exporter = InMemorySpanExporterBuilder::new().build();
1286 let processor = SimpleSpanProcessor::new(exporter.clone());
1287 let span_data = new_test_export_span_data();
1288 processor.on_end(span_data.clone());
1289 assert!(!exporter.get_finished_spans().unwrap().is_empty());
1290 let _result = processor.shutdown();
1291 assert!(exporter.get_finished_spans().unwrap().is_empty());
1293 }
1294
1295 #[test]
1296 fn test_default_const_values() {
1297 assert_eq!(OTEL_BSP_MAX_QUEUE_SIZE, "OTEL_BSP_MAX_QUEUE_SIZE");
1298 assert_eq!(OTEL_BSP_MAX_QUEUE_SIZE_DEFAULT, 2048);
1299 assert_eq!(OTEL_BSP_SCHEDULE_DELAY, "OTEL_BSP_SCHEDULE_DELAY");
1300 assert_eq!(OTEL_BSP_SCHEDULE_DELAY_DEFAULT.as_millis(), 5000);
1301 assert_eq!(
1302 OTEL_BSP_MAX_EXPORT_BATCH_SIZE,
1303 "OTEL_BSP_MAX_EXPORT_BATCH_SIZE"
1304 );
1305 assert_eq!(OTEL_BSP_MAX_EXPORT_BATCH_SIZE_DEFAULT, 512);
1306 assert_eq!(OTEL_BSP_EXPORT_TIMEOUT, "OTEL_BSP_EXPORT_TIMEOUT");
1307 assert_eq!(OTEL_BSP_EXPORT_TIMEOUT_DEFAULT.as_millis(), 30000);
1308 }
1309
1310 #[test]
1311 fn test_default_batch_config_adheres_to_specification() {
1312 let env_vars = vec![
1313 OTEL_BSP_SCHEDULE_DELAY,
1314 OTEL_BSP_EXPORT_TIMEOUT,
1315 OTEL_BSP_MAX_QUEUE_SIZE,
1316 OTEL_BSP_MAX_EXPORT_BATCH_SIZE,
1317 OTEL_BSP_MAX_CONCURRENT_EXPORTS,
1318 ];
1319
1320 let config = temp_env::with_vars_unset(env_vars, BatchConfig::default);
1321
1322 assert_eq!(
1323 config.max_concurrent_exports,
1324 OTEL_BSP_MAX_CONCURRENT_EXPORTS_DEFAULT
1325 );
1326 assert_eq!(config.scheduled_delay, OTEL_BSP_SCHEDULE_DELAY_DEFAULT);
1327 assert_eq!(config.max_export_timeout, OTEL_BSP_EXPORT_TIMEOUT_DEFAULT);
1328 assert_eq!(config.max_queue_size, OTEL_BSP_MAX_QUEUE_SIZE_DEFAULT);
1329 assert_eq!(
1330 config.max_export_batch_size,
1331 OTEL_BSP_MAX_EXPORT_BATCH_SIZE_DEFAULT
1332 );
1333 }
1334
1335 #[test]
1336 fn test_code_based_config_overrides_env_vars() {
1337 let env_vars = vec![
1338 (OTEL_BSP_EXPORT_TIMEOUT, Some("60000")),
1339 (OTEL_BSP_MAX_CONCURRENT_EXPORTS, Some("5")),
1340 (OTEL_BSP_MAX_EXPORT_BATCH_SIZE, Some("1024")),
1341 (OTEL_BSP_MAX_QUEUE_SIZE, Some("4096")),
1342 (OTEL_BSP_SCHEDULE_DELAY, Some("2000")),
1343 ];
1344
1345 temp_env::with_vars(env_vars, || {
1346 let config = BatchConfigBuilder::default()
1347 .with_max_export_batch_size(512)
1348 .with_max_queue_size(2048)
1349 .with_scheduled_delay(Duration::from_millis(1000));
1350 #[cfg(feature = "experimental_trace_batch_span_processor_with_async_runtime")]
1351 let config = {
1352 config
1353 .with_max_concurrent_exports(10)
1354 .with_max_export_timeout(Duration::from_millis(2000))
1355 };
1356 let config = config.build();
1357
1358 assert_eq!(config.max_export_batch_size, 512);
1359 assert_eq!(config.max_queue_size, 2048);
1360 assert_eq!(config.scheduled_delay, Duration::from_millis(1000));
1361 #[cfg(feature = "experimental_trace_batch_span_processor_with_async_runtime")]
1362 {
1363 assert_eq!(config.max_concurrent_exports, 10);
1364 assert_eq!(config.max_export_timeout, Duration::from_millis(2000));
1365 }
1366 });
1367 }
1368
1369 #[test]
1370 fn test_batch_config_configurable_by_env_vars() {
1371 let env_vars = vec![
1372 (OTEL_BSP_SCHEDULE_DELAY, Some("2000")),
1373 (OTEL_BSP_EXPORT_TIMEOUT, Some("60000")),
1374 (OTEL_BSP_MAX_QUEUE_SIZE, Some("4096")),
1375 (OTEL_BSP_MAX_EXPORT_BATCH_SIZE, Some("1024")),
1376 ];
1377
1378 let config = temp_env::with_vars(env_vars, BatchConfig::default);
1379
1380 assert_eq!(config.scheduled_delay, Duration::from_millis(2000));
1381 assert_eq!(config.max_export_timeout, Duration::from_millis(60000));
1382 assert_eq!(config.max_queue_size, 4096);
1383 assert_eq!(config.max_export_batch_size, 1024);
1384 }
1385
1386 #[test]
1387 fn test_batch_config_max_export_batch_size_validation() {
1388 let env_vars = vec![
1389 (OTEL_BSP_MAX_QUEUE_SIZE, Some("256")),
1390 (OTEL_BSP_MAX_EXPORT_BATCH_SIZE, Some("1024")),
1391 ];
1392
1393 let config = temp_env::with_vars(env_vars, BatchConfig::default);
1394
1395 assert_eq!(config.max_queue_size, 256);
1396 assert_eq!(config.max_export_batch_size, 256);
1397 assert_eq!(config.scheduled_delay, OTEL_BSP_SCHEDULE_DELAY_DEFAULT);
1398 assert_eq!(config.max_export_timeout, OTEL_BSP_EXPORT_TIMEOUT_DEFAULT);
1399 }
1400
1401 #[test]
1402 fn test_batch_config_with_fields() {
1403 let batch = BatchConfigBuilder::default()
1404 .with_max_export_batch_size(10)
1405 .with_scheduled_delay(Duration::from_millis(10))
1406 .with_max_queue_size(10);
1407 #[cfg(feature = "experimental_trace_batch_span_processor_with_async_runtime")]
1408 let batch = {
1409 batch
1410 .with_max_concurrent_exports(10)
1411 .with_max_export_timeout(Duration::from_millis(10))
1412 };
1413 let batch = batch.build();
1414 assert_eq!(batch.max_export_batch_size, 10);
1415 assert_eq!(batch.scheduled_delay, Duration::from_millis(10));
1416 assert_eq!(batch.max_queue_size, 10);
1417 #[cfg(feature = "experimental_trace_batch_span_processor_with_async_runtime")]
1418 {
1419 assert_eq!(batch.max_concurrent_exports, 10);
1420 assert_eq!(batch.max_export_timeout, Duration::from_millis(10));
1421 }
1422 }
1423
1424 fn create_test_span(name: &str) -> SpanData {
1426 SpanData {
1427 span_context: SpanContext::empty_context(),
1428 parent_span_id: SpanId::INVALID,
1429 parent_span_is_remote: false,
1430 span_kind: SpanKind::Internal,
1431 name: name.to_string().into(),
1432 start_time: opentelemetry::time::now(),
1433 end_time: opentelemetry::time::now(),
1434 attributes: Vec::new(),
1435 dropped_attributes_count: 0,
1436 events: SpanEvents::default(),
1437 links: SpanLinks::default(),
1438 status: Status::Unset,
1439 instrumentation_scope: Default::default(),
1440 }
1441 }
1442
1443 use crate::Resource;
1444 use opentelemetry::{Key, KeyValue, Value};
1445 use std::{
1446 sync::{
1447 atomic::{AtomicUsize, Ordering},
1448 Arc, Mutex,
1449 },
1450 time::Instant,
1451 };
1452
1453 #[derive(Debug)]
1455 struct MockSpanExporter {
1456 exported_spans: Arc<Mutex<Vec<SpanData>>>,
1457 exported_resource: Arc<Mutex<Option<Resource>>>,
1458 }
1459
1460 impl MockSpanExporter {
1461 fn new() -> Self {
1462 Self {
1463 exported_spans: Arc::new(Mutex::new(Vec::new())),
1464 exported_resource: Arc::new(Mutex::new(None)),
1465 }
1466 }
1467 }
1468
1469 impl SpanExporter for MockSpanExporter {
1470 async fn export(&self, batch: Vec<SpanData>) -> OTelSdkResult {
1471 let exported_spans = self.exported_spans.clone();
1472 exported_spans.lock().unwrap().extend(batch);
1473 Ok(())
1474 }
1475
1476 fn shutdown(&self) -> OTelSdkResult {
1477 Ok(())
1478 }
1479 fn set_resource(&mut self, resource: &Resource) {
1480 let mut exported_resource = self.exported_resource.lock().unwrap();
1481 *exported_resource = Some(resource.clone());
1482 }
1483 }
1484
1485 #[test]
1486 fn batchspanprocessor_handles_on_end() {
1487 let exporter = MockSpanExporter::new();
1488 let exporter_shared = exporter.exported_spans.clone();
1489 let config = BatchConfigBuilder::default()
1490 .with_max_queue_size(10)
1491 .with_max_export_batch_size(10)
1492 .with_scheduled_delay(Duration::from_secs(5))
1493 .build();
1494 let processor = BatchSpanProcessor::new(exporter, config);
1495
1496 let test_span = create_test_span("test_span");
1497 processor.on_end(test_span.clone());
1498
1499 std::thread::sleep(Duration::from_secs(6));
1501
1502 let exported_spans = exporter_shared.lock().unwrap();
1503 assert_eq!(exported_spans.len(), 1);
1504 assert_eq!(exported_spans[0].name, "test_span");
1505 }
1506
1507 #[test]
1508 fn batchspanprocessor_force_flush() {
1509 let exporter = MockSpanExporter::new();
1510 let exporter_shared = exporter.exported_spans.clone(); let config = BatchConfigBuilder::default()
1512 .with_max_queue_size(10)
1513 .with_max_export_batch_size(10)
1514 .with_scheduled_delay(Duration::from_secs(5))
1515 .build();
1516 let processor = BatchSpanProcessor::new(exporter, config);
1517
1518 let test_span = create_test_span("force_flush_span");
1520 processor.on_end(test_span.clone());
1521
1522 let flush_result = processor.force_flush();
1524 assert!(flush_result.is_ok(), "Force flush failed unexpectedly");
1525
1526 let exported_spans = exporter_shared.lock().unwrap();
1528 assert_eq!(
1529 exported_spans.len(),
1530 1,
1531 "Unexpected number of exported spans"
1532 );
1533 assert_eq!(exported_spans[0].name, "force_flush_span");
1534 }
1535
1536 #[test]
1537 fn batchspanprocessor_does_not_overdrain_unaccounted_spans() {
1538 let exporter = MockSpanExporter::new();
1539 let exported_spans = exporter.exported_spans.clone();
1540 let (sender, receiver) = std::sync::mpsc::sync_channel(4);
1541 let current_batch_size = AtomicUsize::new(1);
1542 let config = BatchConfigBuilder::default()
1543 .with_max_queue_size(4)
1544 .with_max_export_batch_size(4)
1545 .build();
1546 let mut spans = Vec::with_capacity(config.max_export_batch_size);
1547 let mut last_export_time = Instant::now();
1548
1549 sender.send(create_test_span("counted")).unwrap();
1550 sender.send(create_test_span("unaccounted")).unwrap();
1551
1552 let result = BatchSpanProcessor::get_spans_and_export(
1553 &receiver,
1554 &exporter,
1555 &mut spans,
1556 &mut last_export_time,
1557 ¤t_batch_size,
1558 &config,
1559 &|_count: u64| {},
1560 );
1561
1562 assert!(result.is_ok(), "export should succeed");
1563 assert_eq!(
1564 current_batch_size.load(Ordering::Relaxed),
1565 0,
1566 "helper should only subtract the counted span"
1567 );
1568 assert_eq!(
1569 exported_spans.lock().unwrap().len(),
1570 1,
1571 "helper should export at most the target batch size snapshot"
1572 );
1573 assert!(
1574 receiver.try_recv().is_ok(),
1575 "one span should remain queued for a later export cycle"
1576 );
1577 }
1578
1579 #[test]
1580 fn batchspanprocessor_drain_handles_counted_but_not_yet_enqueued_spans() {
1581 let exporter = MockSpanExporter::new();
1587 let exported_spans = exporter.exported_spans.clone();
1588 let (sender, receiver) = std::sync::mpsc::sync_channel(4);
1589 let current_batch_size = AtomicUsize::new(2);
1591 let config = BatchConfigBuilder::default()
1592 .with_max_queue_size(4)
1593 .with_max_export_batch_size(4)
1594 .build();
1595 let mut spans = Vec::with_capacity(config.max_export_batch_size);
1596 let mut last_export_time = Instant::now();
1597
1598 sender.send(create_test_span("landed")).unwrap();
1599
1600 let result = BatchSpanProcessor::get_spans_and_export(
1601 &receiver,
1602 &exporter,
1603 &mut spans,
1604 &mut last_export_time,
1605 ¤t_batch_size,
1606 &config,
1607 &|_count: u64| {},
1608 );
1609
1610 assert!(result.is_ok(), "export should succeed");
1611 assert_eq!(
1612 exported_spans.lock().unwrap().len(),
1613 1,
1614 "only the span that landed in the channel can be exported"
1615 );
1616 assert_eq!(
1617 current_batch_size.load(Ordering::Relaxed),
1618 1,
1619 "the count of the not-yet-enqueued span must survive the drain"
1620 );
1621
1622 sender.send(create_test_span("late")).unwrap();
1625 let result = BatchSpanProcessor::get_spans_and_export(
1626 &receiver,
1627 &exporter,
1628 &mut spans,
1629 &mut last_export_time,
1630 ¤t_batch_size,
1631 &config,
1632 &|_count: u64| {},
1633 );
1634
1635 assert!(result.is_ok(), "export should succeed");
1636 assert_eq!(
1637 exported_spans.lock().unwrap().len(),
1638 2,
1639 "the late span should be exported on the next cycle"
1640 );
1641 assert_eq!(
1642 current_batch_size.load(Ordering::Relaxed),
1643 0,
1644 "counter should settle to zero once everything is exported"
1645 );
1646 }
1647
1648 #[derive(Debug)]
1649 struct BlockingExporter {
1650 exported_count: Arc<AtomicUsize>,
1651 export_started: std::sync::mpsc::SyncSender<()>,
1652 release: Arc<Mutex<std::sync::mpsc::Receiver<()>>>,
1653 }
1654
1655 impl SpanExporter for BlockingExporter {
1656 async fn export(&self, batch: Vec<SpanData>) -> OTelSdkResult {
1657 let _ = self.export_started.try_send(());
1658 let _ = self.release.lock().unwrap().recv();
1660 self.exported_count.fetch_add(batch.len(), Ordering::SeqCst);
1661 Ok(())
1662 }
1663 }
1664
1665 #[test]
1666 fn batchspanprocessor_on_end_reverts_count_when_queue_full() {
1667 let (started_sender, started_receiver) = std::sync::mpsc::sync_channel(8);
1668 let (release_sender, release_receiver) = std::sync::mpsc::sync_channel(8);
1669 let exported_count = Arc::new(AtomicUsize::new(0));
1670 let exporter = BlockingExporter {
1671 exported_count: exported_count.clone(),
1672 export_started: started_sender,
1673 release: Arc::new(Mutex::new(release_receiver)),
1674 };
1675 let config = BatchConfigBuilder::default()
1676 .with_max_queue_size(4)
1677 .with_max_export_batch_size(4)
1678 .with_scheduled_delay(Duration::from_secs(60))
1679 .build();
1680 let processor = BatchSpanProcessor::new(exporter, config);
1681
1682 for _ in 0..4 {
1685 processor.on_end(create_test_span("first_batch"));
1686 }
1687 started_receiver
1688 .recv_timeout(Duration::from_secs(5))
1689 .expect("worker should start exporting the first batch");
1690
1691 for _ in 0..4 {
1694 processor.on_end(create_test_span("second_batch"));
1695 }
1696 for _ in 0..2 {
1697 processor.on_end(create_test_span("overflow"));
1698 }
1699
1700 assert_eq!(processor.dropped_spans_count.load(Ordering::Relaxed), 2);
1701 assert_eq!(
1704 processor.current_batch_size.load(Ordering::Relaxed),
1705 8,
1706 "dropped spans must not remain counted as pending"
1707 );
1708
1709 release_sender.send(()).unwrap();
1711 release_sender.send(()).unwrap();
1712 let flush_result = processor.force_flush();
1713 assert!(flush_result.is_ok(), "force flush failed unexpectedly");
1714
1715 assert_eq!(
1716 exported_count.load(Ordering::SeqCst),
1717 8,
1718 "all spans that entered the queue must be exported"
1719 );
1720 assert_eq!(
1721 processor.current_batch_size.load(Ordering::Relaxed),
1722 0,
1723 "counter should settle to zero; a leftover value indicates the \
1724 queue-full path did not revert its increment"
1725 );
1726 }
1727
1728 #[derive(Debug)]
1731 struct CountingSpanExporter {
1732 count: Arc<AtomicUsize>,
1733 }
1734
1735 impl SpanExporter for CountingSpanExporter {
1736 async fn export(&self, batch: Vec<SpanData>) -> OTelSdkResult {
1737 self.count.fetch_add(batch.len(), Ordering::SeqCst);
1738 std::thread::sleep(Duration::from_millis(20));
1740 Ok(())
1741 }
1742 }
1743
1744 #[test]
1749 fn batchspanprocessor_all_spans_accounted_for() {
1750 let count = Arc::new(AtomicUsize::new(0));
1751 let exporter = CountingSpanExporter {
1752 count: count.clone(),
1753 };
1754
1755 let config = BatchConfigBuilder::default()
1756 .with_max_queue_size(2048)
1757 .with_max_export_batch_size(512)
1758 .with_scheduled_delay(Duration::from_millis(5))
1759 .build();
1760
1761 let processor = BatchSpanProcessor::new(exporter, config);
1762
1763 let total_spans_per_thread = 10_000;
1764 let num_threads = 4;
1765 let total_spans_to_emit = total_spans_per_thread * num_threads;
1766
1767 std::thread::scope(|s| {
1768 for _ in 0..num_threads {
1769 s.spawn(|| {
1770 for _ in 0..total_spans_per_thread {
1771 processor.on_end(create_test_span("stress test span"));
1772 }
1773 });
1774 }
1775 });
1776
1777 processor.shutdown().unwrap();
1779
1780 let spans_received = count.load(Ordering::SeqCst);
1781 let spans_dropped = processor.dropped_spans_count.load(Ordering::SeqCst);
1782
1783 assert_eq!(
1785 spans_received + spans_dropped,
1786 total_spans_to_emit,
1787 "Spans unaccounted for! Received: {spans_received}, Dropped: {spans_dropped}, Total emitted: {total_spans_to_emit}"
1788 );
1789 }
1790
1791 #[test]
1792 fn batchspanprocessor_shutdown() {
1793 let exporter = InMemorySpanExporterBuilder::new()
1795 .keep_records_on_shutdown()
1796 .build();
1797 let processor = BatchSpanProcessor::new(exporter.clone(), BatchConfig::default());
1798
1799 let record = create_test_span("test_span");
1800
1801 processor.on_end(record);
1802 processor.force_flush().unwrap();
1803 processor.shutdown().unwrap();
1804
1805 processor.on_end(create_test_span("after_shutdown_span"));
1807
1808 assert_eq!(1, exporter.get_finished_spans().unwrap().len());
1809 assert!(exporter.is_shutdown_called());
1810 }
1811
1812 #[test]
1813 fn batchspanprocessor_handles_dropped_spans() {
1814 #[derive(Debug)]
1817 struct SlowExporter {
1818 exported_count: Arc<std::sync::atomic::AtomicUsize>,
1819 }
1820
1821 impl SpanExporter for SlowExporter {
1822 async fn export(&self, batch: Vec<SpanData>) -> OTelSdkResult {
1823 std::thread::sleep(Duration::from_millis(50));
1825 self.exported_count
1826 .fetch_add(batch.len(), Ordering::Relaxed);
1827 Ok(())
1828 }
1829
1830 fn shutdown(&self) -> OTelSdkResult {
1831 Ok(())
1832 }
1833
1834 fn set_resource(&mut self, _resource: &Resource) {}
1835 }
1836
1837 let exported_count = Arc::new(std::sync::atomic::AtomicUsize::new(0));
1838 let exporter = SlowExporter {
1839 exported_count: exported_count.clone(),
1840 };
1841
1842 let max_queue_size = 10;
1843 let config = BatchConfigBuilder::default()
1844 .with_max_queue_size(max_queue_size)
1845 .with_max_export_batch_size(5)
1846 .with_scheduled_delay(Duration::from_millis(10))
1847 .build();
1848 let processor = BatchSpanProcessor::new(exporter, config);
1849
1850 let total_spans_to_send = 100;
1852 for i in 0..total_spans_to_send {
1853 let span = create_test_span(&format!("span_{}", i));
1854 processor.on_end(span);
1855 }
1856
1857 let _ = processor.force_flush();
1859
1860 let dropped = processor.dropped_spans_count.load(Ordering::Relaxed);
1861 let exported = exported_count.load(Ordering::Relaxed);
1862
1863 assert_eq!(
1865 dropped + exported,
1866 total_spans_to_send,
1867 "dropped ({}) + exported ({}) should equal total sent ({})",
1868 dropped,
1869 exported,
1870 total_spans_to_send
1871 );
1872
1873 assert!(
1875 dropped > 0,
1876 "Expected some spans to be dropped due to full queue. Exported: {}",
1877 exported
1878 );
1879 }
1880
1881 #[test]
1882 fn batchspanprocessor_sync_ignores_max_concurrent_exports() {
1883 #[derive(Debug)]
1884 struct TrackingExporter {
1885 active: Arc<AtomicUsize>,
1886 max_inflight: Arc<AtomicUsize>,
1887 export_calls: Arc<AtomicUsize>,
1888 delay: Duration,
1889 }
1890
1891 impl SpanExporter for TrackingExporter {
1892 async fn export(&self, _batch: Vec<SpanData>) -> OTelSdkResult {
1893 self.export_calls.fetch_add(1, Ordering::SeqCst);
1894 let inflight = self.active.fetch_add(1, Ordering::SeqCst) + 1;
1895 self.max_inflight.fetch_max(inflight, Ordering::SeqCst);
1896
1897 std::thread::sleep(self.delay);
1898 self.active.fetch_sub(1, Ordering::SeqCst);
1899 Ok(())
1900 }
1901 }
1902
1903 let active = Arc::new(AtomicUsize::new(0));
1904 let max_inflight = Arc::new(AtomicUsize::new(0));
1905 let export_calls = Arc::new(AtomicUsize::new(0));
1906 let exporter = TrackingExporter {
1907 active: active.clone(),
1908 max_inflight: max_inflight.clone(),
1909 export_calls: export_calls.clone(),
1910 delay: Duration::from_millis(50),
1911 };
1912
1913 let config = BatchConfig {
1914 max_export_batch_size: 1,
1915 max_queue_size: 16,
1916 scheduled_delay: Duration::from_secs(3600),
1917 max_export_timeout: Duration::from_secs(5),
1918 max_concurrent_exports: 4,
1919 };
1920
1921 let processor = BatchSpanProcessor::new(exporter, config);
1922
1923 processor.on_end(new_test_export_span_data());
1924 processor.on_end(new_test_export_span_data());
1925 processor.on_end(new_test_export_span_data());
1926
1927 processor.force_flush().expect("force flush failed");
1928 processor.shutdown().expect("shutdown failed");
1929
1930 assert_eq!(
1931 export_calls.load(Ordering::SeqCst),
1932 3,
1933 "expected three exports for three spans with max_export_batch_size=1"
1934 );
1935 assert_eq!(
1936 max_inflight.load(Ordering::SeqCst),
1937 1,
1938 "sync BatchSpanProcessor should export serially regardless of max_concurrent_exports"
1939 );
1940 }
1941
1942 #[test]
1943 fn validate_span_attributes_exported_correctly() {
1944 let exporter = MockSpanExporter::new();
1945 let exporter_shared = exporter.exported_spans.clone();
1946 let config = BatchConfigBuilder::default().build();
1947 let processor = BatchSpanProcessor::new(exporter, config);
1948
1949 let mut span_data = create_test_span("attribute_validation");
1951 span_data.attributes = vec![
1952 KeyValue::new("key1", "value1"),
1953 KeyValue::new("key2", "value2"),
1954 ];
1955 processor.on_end(span_data.clone());
1956
1957 let _ = processor.force_flush();
1959
1960 let exported_spans = exporter_shared.lock().unwrap();
1962 assert_eq!(exported_spans.len(), 1);
1963 let exported_span = &exported_spans[0];
1964 assert!(exported_span
1965 .attributes
1966 .contains(&KeyValue::new("key1", "value1")));
1967 assert!(exported_span
1968 .attributes
1969 .contains(&KeyValue::new("key2", "value2")));
1970 }
1971
1972 #[test]
1973 fn batchspanprocessor_sets_and_exports_with_resource() {
1974 let exporter = MockSpanExporter::new();
1975 let exporter_shared = exporter.exported_spans.clone();
1976 let resource_shared = exporter.exported_resource.clone();
1977 let config = BatchConfigBuilder::default().build();
1978 let mut processor = BatchSpanProcessor::new(exporter, config);
1979
1980 let resource = Resource::builder_empty()
1982 .with_attributes(vec![KeyValue::new("service.name", "test_service")])
1983 .build();
1984 processor.set_resource(&resource);
1985
1986 let test_span = create_test_span("resource_test");
1988 processor.on_end(test_span.clone());
1989
1990 let _ = processor.force_flush();
1992
1993 let exported_spans = exporter_shared.lock().unwrap();
1995 assert_eq!(exported_spans.len(), 1);
1996
1997 let exported_resource = resource_shared.lock().unwrap();
1999 assert!(exported_resource.is_some());
2000 assert_eq!(
2001 exported_resource
2002 .as_ref()
2003 .unwrap()
2004 .get(&Key::new("service.name")),
2005 Some(Value::from("test_service"))
2006 );
2007 }
2008
2009 #[tokio::test(flavor = "current_thread")]
2010 async fn test_batch_processor_current_thread_runtime() {
2011 let exporter = MockSpanExporter::new();
2012 let exporter_shared = exporter.exported_spans.clone();
2013
2014 let config = BatchConfigBuilder::default()
2015 .with_max_queue_size(5)
2016 .with_max_export_batch_size(3)
2017 .build();
2018
2019 let processor = BatchSpanProcessor::new(exporter, config);
2020
2021 for _ in 0..4 {
2022 let span = new_test_export_span_data();
2023 processor.on_end(span);
2024 }
2025
2026 processor.force_flush().unwrap();
2027
2028 let exported_spans = exporter_shared.lock().unwrap();
2029 assert_eq!(exported_spans.len(), 4);
2030 }
2031
2032 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
2033 async fn test_batch_processor_multi_thread_count_1_runtime() {
2034 let exporter = MockSpanExporter::new();
2035 let exporter_shared = exporter.exported_spans.clone();
2036
2037 let config = BatchConfigBuilder::default()
2038 .with_max_queue_size(5)
2039 .with_max_export_batch_size(3)
2040 .build();
2041
2042 let processor = BatchSpanProcessor::new(exporter, config);
2043
2044 for _ in 0..4 {
2045 let span = new_test_export_span_data();
2046 processor.on_end(span);
2047 }
2048
2049 processor.force_flush().unwrap();
2050
2051 let exported_spans = exporter_shared.lock().unwrap();
2052 assert_eq!(exported_spans.len(), 4);
2053 }
2054
2055 #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
2056 async fn test_batch_processor_multi_thread() {
2057 let exporter = MockSpanExporter::new();
2058 let exporter_shared = exporter.exported_spans.clone();
2059
2060 let config = BatchConfigBuilder::default()
2061 .with_max_queue_size(20)
2062 .with_max_export_batch_size(5)
2063 .build();
2064
2065 let processor = Arc::new(BatchSpanProcessor::new(exporter, config));
2067
2068 let mut handles = vec![];
2069 for _ in 0..10 {
2070 let processor_clone = Arc::clone(&processor);
2071 let handle = tokio::spawn(async move {
2072 let span = new_test_export_span_data();
2073 processor_clone.on_end(span);
2074 });
2075 handles.push(handle);
2076 }
2077
2078 for handle in handles {
2079 handle.await.unwrap();
2080 }
2081
2082 processor.force_flush().unwrap();
2083
2084 let exported_spans = exporter_shared.lock().unwrap();
2086 assert_eq!(exported_spans.len(), 10);
2087 }
2088
2089 #[cfg(feature = "experimental_metrics_bound_instruments")]
2093 fn sum_processed_spans(
2094 metric_exporter: &crate::metrics::InMemoryMetricExporter,
2095 error_type: Option<&str>,
2096 ) -> u64 {
2097 use crate::metrics::data::{AggregatedMetrics, MetricData};
2098
2099 let metrics = metric_exporter.get_finished_metrics().unwrap();
2100 let mut total: u64 = 0;
2101 for rm in &metrics {
2102 for sm in &rm.scope_metrics {
2103 for metric in &sm.metrics {
2104 if metric.name == "otel.sdk.processor.span.processed" {
2105 if let AggregatedMetrics::U64(MetricData::Sum(sum)) = &metric.data {
2106 for dp in sum.data_points() {
2107 let dp_error_type = dp
2108 .attributes()
2109 .find(|kv| kv.key.as_str() == "error.type")
2110 .map(|kv| kv.value.as_str().to_string());
2111 let matches = match error_type {
2112 Some(expected) => dp_error_type.as_deref() == Some(expected),
2113 None => dp_error_type.is_none(),
2114 };
2115 if matches {
2116 total += dp.value();
2117 }
2118 }
2119 }
2120 }
2121 }
2122 }
2123 }
2124 total
2125 }
2126
2127 #[cfg(feature = "experimental_metrics_bound_instruments")]
2128 mod self_obs {
2129 use super::*;
2130
2131 #[cfg(feature = "experimental_metrics_bound_instruments")]
2137 #[test]
2138 #[ignore]
2139 fn self_diagnostics_counter_records_success() {
2140 use crate::metrics::{InMemoryMetricExporter, SdkMeterProvider};
2141
2142 let metric_exporter = InMemoryMetricExporter::default();
2143 let meter_provider = SdkMeterProvider::builder()
2144 .with_periodic_exporter(metric_exporter.clone())
2145 .build();
2146 opentelemetry::global::set_meter_provider(meter_provider.clone());
2147
2148 let span_exporter = InMemorySpanExporterBuilder::new().build();
2149 let config = BatchConfigBuilder::default()
2150 .with_max_queue_size(256)
2151 .with_max_export_batch_size(64)
2152 .with_scheduled_delay(Duration::from_secs(60))
2153 .build();
2154 let processor = BatchSpanProcessor::new(span_exporter, config);
2155
2156 for _ in 0..10 {
2157 processor.on_end(create_test_span("success"));
2158 }
2159
2160 processor.force_flush().unwrap();
2163 meter_provider.force_flush().unwrap();
2164
2165 let processed = sum_processed_spans(&metric_exporter, None);
2166 assert_eq!(
2167 processed, 10,
2168 "expected 10 processed spans, got {processed}"
2169 );
2170
2171 processor.shutdown().unwrap();
2172 meter_provider.shutdown().unwrap();
2173 }
2174
2175 #[cfg(feature = "experimental_metrics_bound_instruments")]
2182 #[test]
2183 #[ignore]
2184 fn self_diagnostics_counter_records_queue_full_drops() {
2185 use crate::metrics::{InMemoryMetricExporter, SdkMeterProvider};
2186
2187 let metric_exporter = InMemoryMetricExporter::default();
2188 let meter_provider = SdkMeterProvider::builder()
2189 .with_periodic_exporter(metric_exporter.clone())
2190 .build();
2191 opentelemetry::global::set_meter_provider(meter_provider.clone());
2192
2193 let (started_sender, started_receiver) = std::sync::mpsc::sync_channel(8);
2194 let (release_sender, release_receiver) = std::sync::mpsc::sync_channel(8);
2195 let exported_count = Arc::new(AtomicUsize::new(0));
2196 let exporter = BlockingExporter {
2197 exported_count: exported_count.clone(),
2198 export_started: started_sender,
2199 release: Arc::new(Mutex::new(release_receiver)),
2200 };
2201 let config = BatchConfigBuilder::default()
2202 .with_max_queue_size(4)
2203 .with_max_export_batch_size(4)
2204 .with_scheduled_delay(Duration::from_secs(60))
2205 .build();
2206 let processor = BatchSpanProcessor::new(exporter, config);
2207
2208 for _ in 0..4 {
2211 processor.on_end(create_test_span("first_batch"));
2212 }
2213 started_receiver
2214 .recv_timeout(Duration::from_secs(5))
2215 .expect("worker should start exporting the first batch");
2216
2217 for _ in 0..4 {
2220 processor.on_end(create_test_span("second_batch"));
2221 }
2222 for _ in 0..2 {
2223 processor.on_end(create_test_span("overflow"));
2224 }
2225
2226 release_sender.send(()).unwrap();
2228 release_sender.send(()).unwrap();
2229 processor.force_flush().unwrap();
2230 meter_provider.force_flush().unwrap();
2231
2232 let queue_full = sum_processed_spans(&metric_exporter, Some("queue_full"));
2233 assert_eq!(
2234 queue_full, 2,
2235 "expected 2 queue_full drops, got {queue_full}"
2236 );
2237
2238 processor.shutdown().unwrap();
2239 meter_provider.shutdown().unwrap();
2240 }
2241
2242 #[cfg(feature = "experimental_metrics_bound_instruments")]
2248 #[test]
2249 #[ignore]
2250 fn self_diagnostics_counter_records_already_shutdown_drops() {
2251 use crate::metrics::{InMemoryMetricExporter, SdkMeterProvider};
2252
2253 let metric_exporter = InMemoryMetricExporter::default();
2254 let meter_provider = SdkMeterProvider::builder()
2255 .with_periodic_exporter(metric_exporter.clone())
2256 .build();
2257 opentelemetry::global::set_meter_provider(meter_provider.clone());
2258
2259 let span_exporter = InMemorySpanExporterBuilder::new().build();
2260 let processor = BatchSpanProcessor::new(span_exporter, BatchConfig::default());
2261
2262 processor.shutdown().unwrap();
2265
2266 for _ in 0..7 {
2267 processor.on_end(create_test_span("after_shutdown"));
2268 }
2269
2270 meter_provider.force_flush().unwrap();
2271
2272 let already_shutdown = sum_processed_spans(&metric_exporter, Some("already_shutdown"));
2273 assert_eq!(
2274 already_shutdown, 7,
2275 "expected 7 already_shutdown drops, got {already_shutdown}"
2276 );
2277
2278 meter_provider.shutdown().unwrap();
2279 }
2280
2281 #[cfg(feature = "experimental_metrics_bound_instruments")]
2287 #[test]
2288 #[ignore]
2289 fn simple_self_diagnostics_counter_records_success() {
2290 use crate::metrics::{InMemoryMetricExporter, SdkMeterProvider};
2291
2292 let metric_exporter = InMemoryMetricExporter::default();
2293 let meter_provider = SdkMeterProvider::builder()
2294 .with_periodic_exporter(metric_exporter.clone())
2295 .build();
2296 opentelemetry::global::set_meter_provider(meter_provider.clone());
2297
2298 let span_exporter = InMemorySpanExporterBuilder::new().build();
2299 let processor = SimpleSpanProcessor::new(span_exporter);
2300
2301 for _ in 0..10 {
2302 processor.on_end(new_test_export_span_data());
2303 }
2304
2305 meter_provider.force_flush().unwrap();
2306
2307 let processed = sum_processed_spans(&metric_exporter, None);
2308 assert_eq!(
2309 processed, 10,
2310 "expected 10 processed spans, got {processed}"
2311 );
2312
2313 processor.shutdown().unwrap();
2314 meter_provider.shutdown().unwrap();
2315 }
2316
2317 #[cfg(feature = "experimental_metrics_bound_instruments")]
2323 #[test]
2324 #[ignore]
2325 fn simple_self_diagnostics_counter_records_already_shutdown_drops() {
2326 use crate::metrics::{InMemoryMetricExporter, SdkMeterProvider};
2327
2328 let metric_exporter = InMemoryMetricExporter::default();
2329 let meter_provider = SdkMeterProvider::builder()
2330 .with_periodic_exporter(metric_exporter.clone())
2331 .build();
2332 opentelemetry::global::set_meter_provider(meter_provider.clone());
2333
2334 let span_exporter = InMemorySpanExporterBuilder::new().build();
2335 let processor = SimpleSpanProcessor::new(span_exporter);
2336
2337 processor.shutdown().unwrap();
2340
2341 for _ in 0..7 {
2342 processor.on_end(new_test_export_span_data());
2343 }
2344
2345 meter_provider.force_flush().unwrap();
2346
2347 let already_shutdown = sum_processed_spans(&metric_exporter, Some("already_shutdown"));
2348 assert_eq!(
2349 already_shutdown, 7,
2350 "expected 7 already_shutdown drops, got {already_shutdown}"
2351 );
2352 let success = sum_processed_spans(&metric_exporter, None);
2353 assert_eq!(
2354 success, 0,
2355 "post-shutdown spans must not be counted as success, got {success}"
2356 );
2357
2358 meter_provider.shutdown().unwrap();
2359 }
2360 }
2361}