Skip to main content

opentelemetry_sdk/trace/
span_processor.rs

1//! # OpenTelemetry Span Processor Interface
2//!
3//! Span processor is an interface which allows hooks for span start and end method
4//! invocations. The span processors are invoked only when
5//! [`is_recording`] is true.
6//!
7//! Built-in span processors are responsible for batching and conversion of spans to
8//! exportable representation and passing batches to exporters.
9//!
10//! Span processors can be registered directly on SDK [`TracerProvider`] and they are
11//! invoked in the same order as they were registered.
12//!
13//! All `Tracer` instances created by a `TracerProvider` share the same span processors.
14//! Changes to this collection reflect in all `Tracer` instances.
15//!
16//! The following diagram shows `SpanProcessor`'s relationship to other components
17//! in the SDK:
18//!
19//! ```ascii
20//!   +-----+--------------+   +-----------------------+   +-------------------+
21//!   |     |              |   |                       |   |                   |
22//!   |     |              |   | (Batch)SpanProcessor  |   |    SpanExporter   |
23//!   |     |              +---> (Simple)SpanProcessor +--->  (OTLPExporter)   |
24//!   |     |              |   |                       |   |                   |
25//!   | SDK | Tracer.span()|   +-----------------------+   +-------------------+
26//!   |     | Span.end()   |
27//!   |     |              |
28//!   |     |              |
29//!   |     |              |
30//!   |     |              |
31//!   +-----+--------------+
32//! ```
33//!
34//! [`is_recording`]: opentelemetry::trace::Span::is_recording()
35//! [`TracerProvider`]: opentelemetry::trace::TracerProvider
36
37use 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
54/// Environment variable for configuring the delay interval (in milliseconds)
55/// between two consecutive exports for the [`BatchSpanProcessor`].
56pub const OTEL_BSP_SCHEDULE_DELAY: &str = "OTEL_BSP_SCHEDULE_DELAY";
57/// Default delay interval between two consecutive exports.
58pub const OTEL_BSP_SCHEDULE_DELAY_DEFAULT: Duration = Duration::from_millis(5_000);
59/// Environment variable for configuring the maximum queue size for the
60/// [`BatchSpanProcessor`].
61pub const OTEL_BSP_MAX_QUEUE_SIZE: &str = "OTEL_BSP_MAX_QUEUE_SIZE";
62/// Default maximum queue size
63pub const OTEL_BSP_MAX_QUEUE_SIZE_DEFAULT: usize = 2_048;
64/// Environment variable for configuring the maximum batch size for the
65/// [`BatchSpanProcessor`], must be less than or equal to
66/// `OTEL_BSP_MAX_QUEUE_SIZE`.
67pub const OTEL_BSP_MAX_EXPORT_BATCH_SIZE: &str = "OTEL_BSP_MAX_EXPORT_BATCH_SIZE";
68/// Default maximum batch size
69pub const OTEL_BSP_MAX_EXPORT_BATCH_SIZE_DEFAULT: usize = 512;
70/// Environment variable for configuring the maximum allowed time to export
71/// data.
72///
73/// This value is honored by
74/// `span_processor_with_async_runtime::BatchSpanProcessor`. The thread-based
75/// [`BatchSpanProcessor`] ignores this setting.
76pub const OTEL_BSP_EXPORT_TIMEOUT: &str = "OTEL_BSP_EXPORT_TIMEOUT";
77/// Default maximum allowed time to export data.
78///
79/// See [`OTEL_BSP_EXPORT_TIMEOUT`] for which processors honor this value.
80pub 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";
82/// Default max concurrent exports for BSP
83pub(crate) const OTEL_BSP_MAX_CONCURRENT_EXPORTS_DEFAULT: usize = 1;
84
85/// `SpanProcessor` is an interface which allows hooks for span start and end
86/// method invocations. The span processors are invoked only when is_recording
87/// is true.
88pub trait SpanProcessor: Send + Sync + std::fmt::Debug {
89    /// `on_start` is called when a `Span` is started.  This method is called
90    /// synchronously on the thread that started the span, therefore it should
91    /// not block or throw exceptions.
92    fn on_start(&self, span: &mut Span, cx: &Context);
93
94    /// `on_end` is called after a `Span` is ended (i.e., the end timestamp is
95    /// already set). This method is called synchronously within the `Span::end`
96    /// API, therefore it should not block or throw an exception.
97    ///
98    /// # Accessing Context
99    ///
100    /// **Important**: Do not rely on [`Context::current()`] in `on_end`. When `on_end`
101    /// is called during span cleanup, `Context::current()` returns whatever context
102    /// happens to be active at that moment, which is typically unrelated to the span
103    /// being ended. Contexts can be activated in any order and are not necessarily
104    /// hierarchical.
105    ///
106    /// **Best Practice**: Extract any needed context information in [`on_start`]
107    /// and store it as span attributes. This ensures the information is available
108    /// in the [`SpanData`] passed to `on_end`.
109    ///
110    /// # Example
111    ///
112    /// ```rust,ignore
113    /// impl SpanProcessor for MyProcessor {
114    ///     fn on_start(&self, span: &mut Span, cx: &Context) {
115    ///         // Extract baggage and store as span attribute
116    ///         if let Some(value) = cx.baggage().get("my-key") {
117    ///             span.set_attribute(KeyValue::new("my-key", value.to_string()));
118    ///         }
119    ///     }
120    ///
121    ///     fn on_end(&self, span: SpanData) {
122    ///         // Access the attribute stored in on_start
123    ///         let my_value = span.attributes.iter()
124    ///             .find(|kv| kv.key.as_str() == "my-key");
125    ///     }
126    /// }
127    /// ```
128    ///
129    /// # Filtering completed spans
130    ///
131    /// **Warning:** Filtering individual spans can produce incomplete or broken
132    /// traces, such as an exported child whose parent was discarded. This does
133    /// not coordinate filtering across spans or services. For coordinated
134    /// decisions based on completed spans, prefer [tail-based sampling] in the
135    /// OpenTelemetry Collector or another telemetry pipeline. All spans in a
136    /// trace must reach the same tail-sampling instance; it cannot recover spans
137    /// already discarded by SDK sampling or filtering.
138    ///
139    /// [tail-based sampling]: https://github.com/open-telemetry/opentelemetry-collector-contrib/tree/main/processor/tailsamplingprocessor
140    ///
141    /// If SDK processor filtering fits your requirements and you accept this
142    /// tradeoff, wrap another processor and delegate only the spans that satisfy
143    /// your condition. This example uses an attribute, but the
144    /// condition can use any information in [`SpanData`] or other processor state.
145    /// Register only the wrapper with [`SdkTracerProvider`](crate::trace::SdkTracerProvider), since separately
146    /// registered processors receive spans independently.
147    ///
148    /// ```rust
149    /// use opentelemetry::{Context, Value};
150    /// use opentelemetry_sdk::{
151    ///     error::OTelSdkResult,
152    ///     trace::{Span, SpanData, SpanProcessor},
153    ///     Resource,
154    /// };
155    /// use std::time::Duration;
156    ///
157    /// #[derive(Debug)]
158    /// struct FilteringSpanProcessor<P> {
159    ///     next: P,
160    /// }
161    ///
162    /// impl<P> FilteringSpanProcessor<P> {
163    ///     fn new(next: P) -> Self {
164    ///         Self { next }
165    ///     }
166    /// }
167    ///
168    /// impl<P: SpanProcessor> SpanProcessor for FilteringSpanProcessor<P> {
169    ///     fn on_start(&self, span: &mut Span, cx: &Context) {
170    ///         self.next.on_start(span, cx);
171    ///     }
172    ///
173    ///     fn on_end(&self, span: SpanData) {
174    ///         let should_drop = span.attributes.iter().any(|attribute| {
175    ///             attribute.key.as_str() == "example.drop"
176    ///                 && attribute.value == Value::Bool(true)
177    ///         });
178    ///
179    ///         if !should_drop {
180    ///             self.next.on_end(span);
181    ///         }
182    ///     }
183    ///
184    ///     fn force_flush(&self) -> OTelSdkResult {
185    ///         self.next.force_flush()
186    ///     }
187    ///
188    ///     fn shutdown_with_timeout(&self, timeout: Duration) -> OTelSdkResult {
189    ///         self.next.shutdown_with_timeout(timeout)
190    ///     }
191    ///
192    ///     fn set_resource(&mut self, resource: &Resource) {
193    ///         self.next.set_resource(resource);
194    ///     }
195    /// }
196    /// ```
197    ///
198    /// [`on_start`]: SpanProcessor::on_start
199    /// [`Context::current()`]: opentelemetry::Context::current
200    ///
201    /// TODO - This method should take reference to `SpanData`
202    fn on_end(&self, span: SpanData);
203    /// Force the spans lying in the cache to be exported.
204    fn force_flush(&self) -> OTelSdkResult;
205    /// Shuts down the processor. Called when SDK is shut down. This is an
206    /// opportunity for processors to do any cleanup required.
207    ///
208    /// Implementation should make sure shutdown can be called multiple times.
209    fn shutdown_with_timeout(&self, timeout: Duration) -> OTelSdkResult;
210    /// shutdown the processor with a default timeout.
211    fn shutdown(&self) -> OTelSdkResult {
212        self.shutdown_with_timeout(Duration::from_secs(5))
213    }
214    /// Set the resource for the span processor.
215    fn set_resource(&mut self, _resource: &Resource) {}
216}
217
218/// A [SpanProcessor] that passes finished spans to the configured
219/// `SpanExporter`, as soon as they are finished, without any batching. This is
220/// typically useful for debugging and testing. For scenarios requiring higher
221/// performance/throughput, consider using [BatchSpanProcessor].
222/// Spans are exported synchronously
223/// in the same thread that emits the log record.
224/// When using this processor with the OTLP Exporter, the following exporter
225/// features are supported:
226/// - `grpc-tonic`: This requires TracerProvider to be created within a tokio
227///   runtime. Spans can be emitted from any thread, including tokio runtime
228///   threads.
229/// - `reqwest-blocking-client`: TracerProvider may be created anywhere, but
230///   spans must be emitted from a non-tokio runtime thread.
231/// - `reqwest-client`: TracerProvider may be created anywhere, but spans must be
232///   emitted from a tokio runtime thread.
233///
234/// The OTLP HTTP exporter chooses its default HTTP client from enabled crate
235/// features. That choice is not processor-aware. If you enable async HTTP
236/// clients such as `reqwest-client` or `hyper-client`, ensure this processor is
237/// only used from a thread where those clients can run.
238#[derive(Debug)]
239pub struct SimpleSpanProcessor<T: SpanExporter> {
240    exporter: Mutex<T>,
241    is_shutdown: AtomicBool,
242
243    // Self-diagnostics: otel.sdk.processor.span.processed counter, gated behind
244    // experimental_metrics_bound_instruments. The SimpleSpanProcessor exports
245    // each span synchronously and has no queue, so the only processor-side drop
246    // is `already_shutdown`.
247    #[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    /// Create a new [SimpleSpanProcessor] using the provided exporter.
255    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            // Attribute values follow the OTel semantic conventions for SDK metrics:
273            // https://github.com/open-telemetry/semantic-conventions/blob/main/docs/otel/sdk-metrics.md#metric-otelsdkprocessorspanprocessed
274            // https://github.com/open-telemetry/semantic-conventions/blob/main/docs/registry/attributes/otel.md#otel-component-attributes
275            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        // Ignored
305    }
306
307    fn on_end(&self, span: SpanData) {
308        if !span.span_context.is_sampled() {
309            return;
310        }
311
312        // noop after shutdown
313        if self.is_shutdown.load(Ordering::Relaxed) {
314            // Record the post-shutdown drop in self-diagnostics before returning.
315            #[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                // Count the span as processed right before submitting it to the
327                // exporter, independent of the export outcome, per semconv.
328                #[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            // TODO: check error type, and log `error` only if the error is user-actionable, else log `debug`
339            otel_debug!(
340                name: "SimpleProcessor.OnEnd.Error",
341                reason = format!("{:?}", err)
342            );
343        }
344    }
345
346    fn force_flush(&self) -> OTelSdkResult {
347        // Nothing to flush for simple span processor.
348        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
369/// The `BatchSpanProcessor` collects finished spans in a buffer and exports them
370/// in batches to the configured `SpanExporter`. This processor is ideal for
371/// high-throughput environments, as it minimizes the overhead of exporting spans
372/// individually. It uses a **dedicated background thread** to manage and export spans
373/// asynchronously, ensuring that the application's main execution flow is not blocked.
374///
375/// When using this processor with the OTLP Exporter, the following exporter
376/// features are supported:
377/// - `grpc-tonic`: This requires `TracerProvider` to be created within a tokio
378///   runtime.
379/// - `reqwest-blocking-client`: Works with a regular `main` or `tokio::main`.
380///
381/// In other words, async HTTP clients like `reqwest-client` and `hyper-client`
382/// are not supported by this default processor. The OTLP HTTP exporter chooses
383/// its default HTTP client from enabled crate features and cannot tell which
384/// processor will drive it. If your dependency graph enables async HTTP client
385/// features, either pass an explicit blocking client for this processor or use
386/// the experimental async-runtime batch span processor.
387///
388/// # Example
389///
390/// This example demonstrates how to configure and use the `BatchSpanProcessor`
391/// with a custom configuration. Note that a dedicated thread is used internally
392/// to manage the export process.
393///
394/// ```rust
395/// # #[cfg(feature = "testing")]
396/// # {
397/// use opentelemetry::global;
398/// use opentelemetry_sdk::trace::{
399///     BatchSpanProcessor, BatchConfigBuilder, SdkTracerProvider, InMemorySpanExporter,
400/// };
401/// use opentelemetry::trace::Tracer as _;
402/// use opentelemetry::trace::Span;
403/// use std::time::Duration;
404///
405/// // Step 1: Create an exporter (e.g., an In-Memory Exporter for demonstration).
406/// let exporter = InMemorySpanExporter::default();
407///
408/// // Step 2: Configure the BatchSpanProcessor.
409/// let batch_processor = BatchSpanProcessor::builder(exporter)
410///     .with_batch_config(
411///         BatchConfigBuilder::default()
412///             .with_max_queue_size(1024) // Buffer up to 1024 spans.
413///             .with_max_export_batch_size(256) // Export in batches of up to 256 spans.
414///             .with_scheduled_delay(Duration::from_secs(5)) // Export every 5 seconds.
415///             .build(),
416///     )
417///     .build();
418///
419/// // Step 3: Set up a TracerProvider with the configured processor.
420/// let provider = SdkTracerProvider::builder()
421///     .with_span_processor(batch_processor)
422///     .build();
423/// global::set_tracer_provider(provider.clone());
424///
425/// // Step 4: Create spans and record operations.
426/// let tracer = global::tracer("example-tracer");
427/// let mut span = tracer.start("example-span");
428/// span.end(); // Mark the span as completed.
429///
430/// // Step 5: Ensure all spans are flushed before exiting.
431/// provider.shutdown();
432/// # }
433/// ```
434use std::sync::mpsc::sync_channel;
435use std::sync::mpsc::Receiver;
436use std::sync::mpsc::RecvTimeoutError;
437use std::sync::mpsc::SyncSender;
438
439/// Messages exchanged between the main thread and the background thread.
440#[allow(clippy::large_enum_variant)]
441#[derive(Debug)]
442enum BatchMessage {
443    //ExportSpan(SpanData),
444    ExportSpan(Arc<AtomicBool>),
445    ForceFlush(SyncSender<OTelSdkResult>),
446    Shutdown(SyncSender<OTelSdkResult>),
447    SetResource(Arc<Resource>),
448}
449
450/// The `BatchSpanProcessor` collects finished spans in a buffer and exports them
451/// in batches to the configured `SpanExporter`. This processor is ideal for
452/// high-throughput environments, as it minimizes the overhead of exporting spans
453/// individually. It uses a **dedicated background thread** to manage and export spans
454/// asynchronously, ensuring that the application's main execution flow is not blocked.
455///
456/// This processor supports the following configurations:
457/// - **Queue size**: Maximum number of spans that can be buffered.
458/// - **Batch size**: Maximum number of spans to include in a single export.
459/// - **Scheduled delay**: Frequency at which the batch is exported.
460///
461/// When using this processor with the OTLP Exporter, the following exporter
462/// features are supported:
463/// - `grpc-tonic`: Requires `TracerProvider` to be created within a tokio runtime.
464/// - `reqwest-blocking-client`: Works with a regular `main` or `tokio::main`.
465///
466/// In other words, async HTTP clients like `reqwest-client` and `hyper-client`
467/// are not supported by this default processor. The OTLP HTTP exporter chooses
468/// its default HTTP client from enabled crate features and cannot tell which
469/// processor will drive it. If your dependency graph enables async HTTP client
470/// features, either pass an explicit blocking client for this processor or use
471/// the experimental async-runtime batch span processor.
472///
473/// `BatchSpanProcessor` buffers spans in memory and exports them in batches. An
474/// export is triggered when `max_export_batch_size` is reached or every
475/// `scheduled_delay` milliseconds. Users can explicitly trigger an export using
476/// the `force_flush` method. Shutdown also triggers an export of all buffered
477/// spans and is recommended to be called before the application exits to ensure
478/// all buffered spans are exported.
479///
480/// **Warning**: When using tokio's current-thread runtime, `shutdown()`, which
481/// is a blocking call ,should not be called from your main thread. This can
482/// cause deadlock. Instead, call `shutdown()` from a separate thread or use
483/// tokio's `spawn_blocking`.
484///
485#[derive(Debug)]
486pub struct BatchSpanProcessor {
487    span_sender: SyncSender<SpanData>, // Data channel to store spans
488    message_sender: SyncSender<BatchMessage>, // Control channel to store control messages.
489    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    // Self-diagnostics: otel.sdk.processor.span.processed counter, gated behind
498    // experimental_metrics_bound_instruments so the hot-path `add` is a single
499    // atomic increment with no per-call attribute resolution. The success count
500    // is recorded in the worker thread when a batch is submitted to the
501    // exporter; the drop counts are recorded here at enqueue time.
502    #[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    /// Creates a new instance of `BatchSpanProcessor`.
510    pub fn new<E>(
511        mut exporter: E,
512        config: BatchConfig,
513        //max_queue_size: usize,
514        //scheduled_delay: Duration,
515        //shutdown_timeout: Duration,
516    ) -> 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); // Is this a reasonable bound?
522        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        // Self-diagnostics: create the otel.sdk.processor.span.processed counter.
528        // Created before the worker thread is spawned so the success counter can
529        // be moved into the worker and incremented when a batch is exported.
530        #[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            // Attribute values follow the OTel semantic conventions for SDK metrics:
547            // https://github.com/open-telemetry/semantic-conventions/blob/main/docs/otel/sdk-metrics.md#metric-otelsdkprocessorspanprocessed
548            // https://github.com/open-telemetry/semantic-conventions/blob/main/docs/registry/attributes/otel.md#otel-component-attributes
549            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                // Counts spans for the otel.sdk.processor.span.processed metric;
586                // a no-op when the self-diagnostics feature is disabled.
587                #[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                                // Reset the export span message sent flag now it has has been processed.
603                                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                                    &current_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                                    &current_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                                    &current_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                                //
649                                // break out the loop and return from the current background thread.
650                                //
651                                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                                &current_batch_size,
668                                &config,
669                                &record_processed_success,
670                            );
671                        }
672                        Err(RecvTimeoutError::Disconnected) => {
673                            // Channel disconnected, only thing to do is break
674                            // out (i.e exit the thread)
675                            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"); //TODO: Handle thread spawn failure
688
689        Self {
690            span_sender,
691            message_sender,
692            handle: Mutex::new(Some(handle)),
693            forceflush_timeout: Duration::from_secs(5), // TODO: make this configurable
694            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    /// builder
707    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    // This method gets up to `max_export_batch_size` amount of spans from the channel and exports them.
718    // It returns the result of the export operation.
719    // It expects the spans vec to be empty when it's called.
720    #[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); // `target` is used to determine the stopping criteria for exporting spans.
735        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            // Get up to the remaining target batch size from the channel and push them to the spans vec
744            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(); // Count of spans that will be exported
752            if count_of_spans == 0 {
753                break;
754            }
755            total_exported_spans += count_of_spans;
756
757            // Count the batch as processed before invoking the exporter,
758            // regardless of the export outcome.
759            record_processed_success(count_of_spans as u64);
760
761            result = Self::export_batch_sync(exporter, spans, last_export_time); // This method clears the spans vec after exporting
762
763            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        // Splitting off batch clears the existing batch capacity, and is ready
784        // for re-use in the next export. The newly returned vec! from split_off
785        // is passed to the exporter.
786        // TODO: Compared to Logs, this requires new allocation for vec for
787        // every export. See if this can be optimized by
788        // *not* requiring ownership in the exporter.
789        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    /// Handles span start.
807    fn on_start(&self, _span: &mut Span, _cx: &Context) {
808        // Ignored
809    }
810
811    /// Handles span end.
812    fn on_end(&self, span: SpanData) {
813        // Count the span before enqueueing it so that a concurrent
814        // force_flush()/shutdown() drain never observes an
815        // enqueued-but-uncounted span and misses it (issue #3453). If the
816        // send fails, the increment is reverted in the error arms below.
817        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 for result and handle each separately
821        match result {
822            Ok(_) => {
823                // Successfully sent the span to the data channel.
824                // Check if the current batch size has reached the max export
825                // batch size.
826                if previous_batch_size + 1 >= self.max_export_batch_size {
827                    // Check if the a control message for exporting spans is
828                    // already sent to the worker thread. If not, send a control
829                    // message to export spans. `export_span_message_sent` is set
830                    // to false ONLY when the worker thread has processed the
831                    // control message.
832
833                    if !self.export_span_message_sent.load(Ordering::Relaxed) {
834                        // This is a cost-efficient check as atomic load
835                        // operations do not require exclusive access to cache
836                        // line. Perform atomic swap to
837                        // `export_span_message_sent` ONLY when the atomic load
838                        // operation above returns false. Atomic
839                        // swap/compare_exchange operations require exclusive
840                        // access to cache line on most processor architectures.
841                        // We could have used compare_exchange as well here, but
842                        // it's more verbose than swap.
843                        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                                    // Control message sent successfully.
849                                }
850                                Err(_err) => {
851                                    // TODO: Log error If the control message
852                                    // could not be sent, reset the
853                                    // `export_span_message_sent` flag.
854                                    self.export_span_message_sent
855                                        .store(false, Ordering::Relaxed);
856                                }
857                            }
858                        }
859                    }
860                }
861            }
862            Err(std::sync::mpsc::TrySendError::Full(_)) => {
863                // The span never entered the channel; revert the increment.
864                self.current_batch_size.fetch_sub(1, Ordering::AcqRel);
865                // Record queue-full drop in self-diagnostics.
866                #[cfg(feature = "experimental_metrics_bound_instruments")]
867                self.processed_queue_full.add(1);
868                // Increment dropped spans count. The first time we have to drop
869                // a span, emit a warning.
870                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                // The span never entered the channel; revert the increment.
877                self.current_batch_size.fetch_sub(1, Ordering::AcqRel);
878                // Record after-shutdown drop in self-diagnostics.
879                #[cfg(feature = "experimental_metrics_bound_instruments")]
880                self.processed_after_shutdown.add(1);
881                // Given background thread is the only receiver, and it's
882                // disconnected, it indicates the thread is shutdown
883                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    /// Flushes all pending spans.
892    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                // If the control message could not be sent, emit a warning.
909                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                // Given background thread is the only receiver, and it's
917                // disconnected, it indicates the thread is shutdown
918                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    /// Shuts down the processor.
929    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                        // join the background thread after receiving back the
948                        // shutdown signal
949                        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                // If the control message could not be sent, emit a warning.
973                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                // Given background thread is the only receiver, and it's
981                // disconnected, it indicates the thread is shutdown
982                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    /// Set the resource for the processor.
993    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/// Builder for `BatchSpanProcessorDedicatedThread`.
1002#[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    /// Set the BatchConfig for [BatchSpanProcessorBuilder]
1016    pub fn with_batch_config(self, config: BatchConfig) -> Self {
1017        BatchSpanProcessorBuilder { config, ..self }
1018    }
1019
1020    /// Build a new instance of `BatchSpanProcessor`.
1021    pub fn build(self) -> BatchSpanProcessor {
1022        BatchSpanProcessor::new(self.exporter, self.config)
1023    }
1024}
1025
1026/// Batch span processor configuration.
1027/// Use [`BatchConfigBuilder`] to configure your own instance of [`BatchConfig`].
1028#[derive(Debug)]
1029pub struct BatchConfig {
1030    /// The maximum queue size to buffer spans for delayed processing. If the
1031    /// queue gets full it drops the spans. The default value of is 2048.
1032    pub(crate) max_queue_size: usize,
1033
1034    /// The delay interval in milliseconds between two consecutive processing
1035    /// of batches. The default value is 5 seconds.
1036    pub(crate) scheduled_delay: Duration,
1037
1038    #[allow(dead_code)]
1039    /// The maximum number of spans to process in a single batch. If there are
1040    /// more than one batch worth of spans then it processes multiple batches
1041    /// of spans one batch after the other without any delay. The default value
1042    /// is 512.
1043    pub(crate) max_export_batch_size: usize,
1044
1045    #[allow(dead_code)]
1046    /// The maximum duration to export a batch of data.
1047    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/// A builder for creating [`BatchConfig`] instances.
1060#[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    /// Create a new [`BatchConfigBuilder`] initialized with default batch config values as per the specs.
1071    /// The values are overriden by environment variables if set.
1072    /// The supported environment variables are:
1073    /// * `OTEL_BSP_MAX_QUEUE_SIZE`
1074    /// * `OTEL_BSP_SCHEDULE_DELAY`
1075    /// * `OTEL_BSP_MAX_EXPORT_BATCH_SIZE`
1076    /// * `OTEL_BSP_EXPORT_TIMEOUT`
1077    /// * `OTEL_BSP_MAX_CONCURRENT_EXPORTS`
1078    ///
1079    /// Note: Programmatic configuration overrides any value set via the environment variable.
1080    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    /// Set max_queue_size for [`BatchConfigBuilder`].
1094    /// It's the maximum queue size to buffer spans for delayed processing.
1095    /// If the queue gets full it will drops the spans.
1096    /// The default value is 2048.
1097    ///
1098    /// Corresponding environment variable: `OTEL_BSP_MAX_QUEUE_SIZE`.
1099    ///
1100    /// Note: Programmatically setting this will override any value set via the environment variable.
1101    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    /// Set max_export_batch_size for [`BatchConfigBuilder`].
1107    /// It's the maximum number of spans to process in a single batch. If there are
1108    /// more than one batch worth of spans then it processes multiple batches
1109    /// of spans one batch after the other without any delay. The default value
1110    /// is 512.
1111    ///
1112    /// Corresponding environment variable: `OTEL_BSP_MAX_EXPORT_BATCH_SIZE`.
1113    ///
1114    /// Note: Programmatically setting this will override any value set via the environment variable.
1115    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    /// Set max_concurrent_exports for [`BatchConfigBuilder`].
1122    ///
1123    /// This value is honored by
1124    /// `span_processor_with_async_runtime::BatchSpanProcessor`, where it limits
1125    /// the number of concurrent export tasks.
1126    ///
1127    /// The thread-based `BatchSpanProcessor` exports serially and ignores this
1128    /// setting.
1129    ///
1130    /// Corresponding environment variable: `OTEL_BSP_MAX_CONCURRENT_EXPORTS`.
1131    ///
1132    /// For concurrent exports, enable
1133    /// `experimental_trace_batch_span_processor_with_async_runtime` and use the
1134    /// async-runtime processor.
1135    ///
1136    /// Note: Programmatically setting this will override any value set via the environment variable.
1137    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    /// Set scheduled_delay_duration for [`BatchConfigBuilder`].
1143    /// It's the delay interval in milliseconds between two consecutive processing of batches.
1144    /// The default value is 5000 milliseconds.
1145    ///
1146    /// Corresponding environment variable: `OTEL_BSP_SCHEDULE_DELAY`.
1147    ///
1148    /// Note: Programmatically setting this will override any value set via the environment variable.
1149    pub fn with_scheduled_delay(mut self, scheduled_delay: Duration) -> Self {
1150        self.scheduled_delay = scheduled_delay;
1151        self
1152    }
1153
1154    /// Set max_export_timeout for [`BatchConfigBuilder`].
1155    /// It's the maximum duration to export a batch of data.
1156    /// The The default value is 30000 milliseconds.
1157    ///
1158    /// Corresponding environment variable: `OTEL_BSP_EXPORT_TIMEOUT`.
1159    ///
1160    /// Note: Programmatically setting this will override any value set via the environment variable.
1161    #[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    /// Builds a `BatchConfig` enforcing the following invariants:
1168    /// * `max_export_batch_size` must be less than or equal to `max_queue_size`.
1169    pub fn build(self) -> BatchConfig {
1170        // max export batch size must be less or equal to max queue size.
1171        // we set max export batch size to max queue size if it's larger than max queue size.
1172        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        // max export batch size must be less or equal to max queue size.
1213        // we set max export batch size to max queue size if it's larger than max queue size.
1214        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    // cargo test trace::span_processor::tests:: --features=testing
1232    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        // Assume shutdown is called by ensuring spans are empty in the exporter
1292        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    // Helper function to create a default test span
1425    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    // Mock exporter to test functionality
1454    #[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        // Wait for flush interval to ensure the span is processed
1500        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(); // Shared access to verify exported spans
1511        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        // Create a test span and send it to the processor
1519        let test_span = create_test_span("force_flush_span");
1520        processor.on_end(test_span.clone());
1521
1522        // Call force_flush to immediately export the spans
1523        let flush_result = processor.force_flush();
1524        assert!(flush_result.is_ok(), "Force flush failed unexpectedly");
1525
1526        // Verify the exported spans in the mock exporter
1527        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            &current_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        // Since on_end() increments `current_batch_size` before enqueueing
1582        // (issue #3453), a concurrent drain can observe a counter that is
1583        // higher than the channel depth. The drain must export what is
1584        // available, keep the surplus count intact (no underflow), and pick
1585        // the late span up on a later cycle.
1586        let exporter = MockSpanExporter::new();
1587        let exported_spans = exporter.exported_spans.clone();
1588        let (sender, receiver) = std::sync::mpsc::sync_channel(4);
1589        // Two spans counted, but only one has landed in the channel so far.
1590        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            &current_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        // The late span lands; a later drain cycle must export it and settle
1623        // the counter back to zero.
1624        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            &current_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            // Block until the test releases the export.
1659            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        // Fill the queue to the export threshold; the worker drains all four
1683        // spans and blocks inside export().
1684        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        // While the worker is blocked, refill the queue and overflow it by
1692        // two spans, which must be dropped and their counts reverted.
1693        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        // 4 spans in-flight in the blocked export (not yet subtracted) plus
1702        // 4 spans queued. Without the queue-full revert this would read 10.
1703        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 the in-flight export and the one triggered by force_flush.
1710        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    /// A slow exporter that counts the number of spans received.
1729    /// Used for stress testing the BatchSpanProcessor.
1730    #[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            // Simulate slow export to cause queue buildup and drops
1739            std::thread::sleep(Duration::from_millis(20));
1740            Ok(())
1741        }
1742    }
1743
1744    /// Stress test that verifies all spans are accounted for.
1745    /// With multiple threads pushing spans faster than the exporter can
1746    /// handle, some spans will inevitably be dropped. This test validates:
1747    /// total_spans_sent == spans_received_by_exporter + spans_dropped
1748    #[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        // Shutdown the processor to ensure all buffered spans are flushed
1778        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        // The invariant: every span is either received or dropped
1784        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        // Setup exporter and processor - following the same pattern as test_batch_shutdown from logs
1794        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        // todo: expect to see errors here. How should we assert this?
1806        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        // This test verifies that BSP drops spans when the queue is full.
1815
1816        #[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                // Simulate slow export
1824                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        // Rapidly send many more spans than the queue can hold
1851        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        // Force flush any remaining spans - this waits for export to complete
1858        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        // Verify that dropped + exported = total (every span is accounted for)
1864        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        // With 100 spans sent rapidly and a slow exporter, we should have some drops
1874        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        // Create a span with attributes
1950        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        // Force flush to export the span
1958        let _ = processor.force_flush();
1959
1960        // Validate the exported attributes
1961        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        // Set a resource for the processor
1981        let resource = Resource::builder_empty()
1982            .with_attributes(vec![KeyValue::new("service.name", "test_service")])
1983            .build();
1984        processor.set_resource(&resource);
1985
1986        // Create a span and send it to the processor
1987        let test_span = create_test_span("resource_test");
1988        processor.on_end(test_span.clone());
1989
1990        // Force flush to ensure the span is exported
1991        let _ = processor.force_flush();
1992
1993        // Validate spans are exported
1994        let exported_spans = exporter_shared.lock().unwrap();
1995        assert_eq!(exported_spans.len(), 1);
1996
1997        // Validate the resource is correctly set in the exporter
1998        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        // Create the processor with the thread-safe exporter
2066        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        // Verify exported spans
2085        let exported_spans = exporter_shared.lock().unwrap();
2086        assert_eq!(exported_spans.len(), 10);
2087    }
2088
2089    /// Sums the values of `otel.sdk.processor.span.processed` data points whose
2090    /// `error.type` attribute equals `error_type` (or that have no `error.type`
2091    /// attribute when `error_type` is `None`).
2092    #[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        /// Verifies that `otel.sdk.processor.span.processed` counts spans (with no
2132        /// `error.type`) when the processor submits a batch to the exporter.
2133        ///
2134        /// `#[ignore]`d because it mutates process-wide state via
2135        /// `global::set_meter_provider()`. CI runs it in isolation via `test.sh`.
2136        #[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            // Flush so the batch is submitted to the exporter, which is when the
2161            // counter is incremented.
2162            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        /// Verifies that `otel.sdk.processor.span.processed` records queue-full drops
2176        /// with `error.type = queue_full` when spans overflow the queue while the
2177        /// worker is blocked exporting.
2178        ///
2179        /// `#[ignore]`d because it mutates process-wide state via
2180        /// `global::set_meter_provider()`. CI runs it in isolation via `test.sh`.
2181        #[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            // Fill the queue to the export threshold; the worker drains all four
2209            // spans and blocks inside export().
2210            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            // While the worker is blocked, refill the queue (4) and overflow it by
2218            // two spans, which must be dropped and counted as queue_full.
2219            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 the in-flight export and the one triggered by force_flush.
2227            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        /// Verifies that `otel.sdk.processor.span.processed` records post-shutdown
2243        /// emits with `error.type = already_shutdown`.
2244        ///
2245        /// `#[ignore]`d because it mutates process-wide state via
2246        /// `global::set_meter_provider()`. CI runs it in isolation via `test.sh`.
2247        #[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            // Shut the processor down so the worker thread (the only receiver)
2263            // disconnects; subsequent on_end calls hit the already_shutdown branch.
2264            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        /// Verifies that `otel.sdk.processor.span.processed` counts spans (with no
2282        /// `error.type`) when `SimpleSpanProcessor` submits them to the exporter.
2283        ///
2284        /// `#[ignore]`d because it mutates process-wide state via
2285        /// `global::set_meter_provider()`. CI runs it in isolation via `test.sh`.
2286        #[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        /// Verifies that `SimpleSpanProcessor` records post-shutdown spans with
2318        /// `error.type = already_shutdown` and does not count them as success.
2319        ///
2320        /// `#[ignore]`d because it mutates process-wide state via
2321        /// `global::set_meter_provider()`. CI runs it in isolation via `test.sh`.
2322        #[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            // Shut the processor down; subsequent on_end calls hit the
2338            // already_shutdown branch.
2339            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}