Skip to main content

opentelemetry_sdk/logs/
batch_log_processor.rs

1//! # OpenTelemetry Batch Log Processor
2//! The `BatchLogProcessor` is one implementation of the `LogProcessor` interface.
3//!
4//! It buffers log records and sends them to the exporter
5//! in batches. This processor is designed for **production use** in high-throughput
6//! applications and reduces the overhead of frequent exports by using a background
7//! thread for batch processing.
8//!
9//! ## Diagram
10//!
11//! ```ascii
12//!   +-----+---------------+   +-----------------------+   +-------------------+
13//!   |     |               |   |                       |   |                   |
14//!   | SDK | Logger.emit() +---> (Batch)LogProcessor   +--->  (OTLPExporter)   |
15//!   +-----+---------------+   +-----------------------+   +-------------------+
16//! ```
17
18use crate::error::{OTelSdkError, OTelSdkResult};
19use crate::logs::log_processor::LogProcessor;
20use crate::{
21    logs::{LogBatch, LogExporter, SdkLogRecord},
22    Resource,
23};
24use std::sync::mpsc::{self, RecvTimeoutError, SyncSender};
25
26use opentelemetry::{otel_debug, otel_error, otel_warn, Context, InstrumentationScope};
27
28#[cfg(feature = "experimental_metrics_bound_instruments")]
29use opentelemetry::KeyValue;
30
31use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
32use std::{cmp::min, env, sync::Mutex};
33use std::{
34    fmt::{self, Debug, Formatter},
35    str::FromStr,
36    sync::Arc,
37    thread,
38    time::Duration,
39    time::Instant,
40};
41
42/// Environment variable for configuring the delay interval (in milliseconds)
43/// between two consecutive exports for the [`BatchLogProcessor`].
44pub const OTEL_BLRP_SCHEDULE_DELAY: &str = "OTEL_BLRP_SCHEDULE_DELAY";
45/// Default delay interval between two consecutive exports.
46pub const OTEL_BLRP_SCHEDULE_DELAY_DEFAULT: Duration = Duration::from_millis(1_000);
47/// Environment variable for configuring the maximum allowed time to export
48/// data.
49///
50/// This value is honored by
51/// `log_processor_with_async_runtime::BatchLogProcessor`. The thread-based
52/// [`BatchLogProcessor`] ignores this setting.
53pub const OTEL_BLRP_EXPORT_TIMEOUT: &str = "OTEL_BLRP_EXPORT_TIMEOUT";
54/// Default maximum allowed time to export data.
55///
56/// See [`OTEL_BLRP_EXPORT_TIMEOUT`] for which processors honor this value.
57pub const OTEL_BLRP_EXPORT_TIMEOUT_DEFAULT: Duration = Duration::from_millis(30_000);
58/// Environment variable for configuring the maximum queue size for the
59/// [`BatchLogProcessor`].
60pub const OTEL_BLRP_MAX_QUEUE_SIZE: &str = "OTEL_BLRP_MAX_QUEUE_SIZE";
61/// Default maximum queue size.
62pub const OTEL_BLRP_MAX_QUEUE_SIZE_DEFAULT: usize = 2_048;
63/// Environment variable for configuring the maximum batch size for the
64/// [`BatchLogProcessor`], must be less than or equal to
65/// `OTEL_BLRP_MAX_QUEUE_SIZE`.
66pub const OTEL_BLRP_MAX_EXPORT_BATCH_SIZE: &str = "OTEL_BLRP_MAX_EXPORT_BATCH_SIZE";
67/// Default maximum batch size.
68pub const OTEL_BLRP_MAX_EXPORT_BATCH_SIZE_DEFAULT: usize = 512;
69
70/// Messages sent between application thread and batch log processor's work thread.
71#[allow(clippy::large_enum_variant)]
72#[derive(Debug)]
73enum BatchMessage {
74    /// This is ONLY sent when the number of logs records in the data channel has reached `max_export_batch_size`.
75    ExportLog(Arc<AtomicBool>),
76    /// ForceFlush flushes the current buffer to the exporter.
77    ForceFlush(mpsc::SyncSender<OTelSdkResult>),
78    /// Shut down the worker thread, push all logs in buffer to the exporter.
79    Shutdown(mpsc::SyncSender<OTelSdkResult>),
80    /// Set the resource for the exporter.
81    SetResource(Arc<Resource>),
82}
83
84type LogsData = Box<(SdkLogRecord, InstrumentationScope)>;
85
86/// The `BatchLogProcessor` collects finished logs in a buffer and exports them
87/// in batches to the configured `LogExporter`. This processor is ideal for
88/// high-throughput environments, as it minimizes the overhead of exporting logs
89/// individually. It uses a **dedicated background thread** to manage and export logs
90/// asynchronously, ensuring that the application's main execution flow is not blocked.
91///
92/// This processor supports the following configurations:
93/// - **Queue size**: Maximum number of log records that can be buffered.
94/// - **Batch size**: Maximum number of log records to include in a single export.
95/// - **Scheduled delay**: Frequency at which the batch is exported.
96///
97/// When using this processor with the OTLP Exporter, the following exporter
98/// features are supported:
99/// - `grpc-tonic`: Requires `LoggerProvider` to be created within a tokio runtime.
100/// - `reqwest-blocking-client`: Works with a regular `main` or `tokio::main`.
101///
102/// In other words, async HTTP clients like `reqwest-client` and `hyper-client`
103/// are not supported by this default processor. The OTLP HTTP exporter chooses
104/// its default HTTP client from enabled crate features and cannot tell which
105/// processor will drive it. If your dependency graph enables async HTTP client
106/// features, either pass an explicit blocking client for this processor or use
107/// the experimental async-runtime batch log processor.
108///
109/// `BatchLogProcessor` buffers logs in memory and exports them in batches. An
110/// export is triggered when `max_export_batch_size` is reached or every
111/// `scheduled_delay` milliseconds. Users can explicitly trigger an export using
112/// the `force_flush` method. Shutdown also triggers an export of all buffered
113/// logs and is recommended to be called before the application exits to ensure
114/// all buffered logs are exported.
115///
116/// **Warning**: When using tokio's current-thread runtime, `shutdown()`, which
117/// is a blocking call ,should not be called from your main thread. This can
118/// cause deadlock. Instead, call `shutdown()` from a separate thread or use
119/// tokio's `spawn_blocking`.
120///
121///
122/// ### Using a BatchLogProcessor:
123///
124/// ```rust
125/// # #[cfg(feature = "testing")]
126/// # {
127/// use opentelemetry_sdk::logs::{BatchLogProcessor, BatchConfigBuilder, SdkLoggerProvider};
128/// use opentelemetry::global;
129/// use std::time::Duration;
130/// use opentelemetry_sdk::logs::InMemoryLogExporter;
131///
132/// let exporter = InMemoryLogExporter::default(); // Replace with an actual exporter
133/// let processor = BatchLogProcessor::builder(exporter)
134///     .with_batch_config(
135///         BatchConfigBuilder::default()
136///             .with_max_queue_size(2048)
137///             .with_max_export_batch_size(512)
138///             .with_scheduled_delay(Duration::from_secs(5))
139///             .build(),
140///     )
141///     .build();
142///
143/// let provider = SdkLoggerProvider::builder()
144///     .with_log_processor(processor)
145///     .build();
146/// # }
147///
148pub struct BatchLogProcessor {
149    logs_sender: SyncSender<LogsData>, // Data channel to store log records and instrumentation scopes
150    message_sender: SyncSender<BatchMessage>, // Control channel to store control messages for the worker thread
151    handle: Mutex<Option<thread::JoinHandle<()>>>,
152    forceflush_timeout: Duration,
153    export_log_message_sent: Arc<AtomicBool>,
154    current_batch_size: Arc<AtomicUsize>,
155    max_export_batch_size: usize,
156
157    // Track dropped logs - we'll log this at shutdown
158    dropped_logs_count: AtomicUsize,
159
160    // Track the maximum queue size that was configured for this processor
161    max_queue_size: usize,
162
163    // Self-diagnostics: otel.sdk.processor.log.processed counter.
164    // Gated behind experimental_metrics_bound_instruments so the hot-path
165    // `add` is a single atomic increment (~1.8 ns) with no per-call
166    // attribute resolution.
167    #[cfg(feature = "experimental_metrics_bound_instruments")]
168    processed_queue_full: opentelemetry::metrics::BoundCounter<u64>,
169    #[cfg(feature = "experimental_metrics_bound_instruments")]
170    processed_after_shutdown: opentelemetry::metrics::BoundCounter<u64>,
171}
172
173impl Debug for BatchLogProcessor {
174    fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
175        f.debug_struct("BatchLogProcessor")
176            .field("message_sender", &self.message_sender)
177            .finish()
178    }
179}
180
181impl LogProcessor for BatchLogProcessor {
182    fn emit(&self, record: &mut SdkLogRecord, instrumentation: &InstrumentationScope) {
183        // Count the log record before enqueueing it so that a concurrent
184        // force_flush()/shutdown() drain never observes an
185        // enqueued-but-uncounted record and misses it (issue #3453). If the
186        // send fails, the increment is reverted in the error arms below.
187        let previous_batch_size = self.current_batch_size.fetch_add(1, Ordering::AcqRel);
188        let result = self
189            .logs_sender
190            .try_send(Box::new((record.clone(), instrumentation.clone())));
191
192        // match for result and handle each separately
193        match result {
194            Ok(_) => {
195                // Successfully sent the log record to the data channel.
196                // `processed` success is counted when the batch is submitted to
197                // the exporter (in the worker thread), not here at enqueue.
198                //
199                // Check if the current batch size has reached the max export
200                // batch size.
201                if previous_batch_size + 1 >= self.max_export_batch_size {
202                    // Check if the a control message for exporting logs is
203                    // already sent to the worker thread. If not, send a control
204                    // message to export logs. `export_log_message_sent` is set
205                    // to false ONLY when the worker thread has processed the
206                    // control message.
207
208                    if !self.export_log_message_sent.load(Ordering::Relaxed) {
209                        // This is a cost-efficient check as atomic load
210                        // operations do not require exclusive access to cache
211                        // line. Perform atomic swap to
212                        // `export_log_message_sent` ONLY when the atomic load
213                        // operation above returns false. Atomic
214                        // swap/compare_exchange operations require exclusive
215                        // access to cache line on most processor architectures.
216                        // We could have used compare_exchange as well here, but
217                        // it's more verbose than swap.
218                        if !self.export_log_message_sent.swap(true, Ordering::Relaxed) {
219                            match self.message_sender.try_send(BatchMessage::ExportLog(
220                                self.export_log_message_sent.clone(),
221                            )) {
222                                Ok(_) => {
223                                    // Control message sent successfully.
224                                }
225                                Err(_err) => {
226                                    // TODO: Log error If the control message
227                                    // could not be sent, reset the
228                                    // `export_log_message_sent` flag.
229                                    self.export_log_message_sent.store(false, Ordering::Relaxed);
230                                }
231                            }
232                        }
233                    }
234                }
235            }
236            Err(mpsc::TrySendError::Full(_)) => {
237                // The record never entered the channel; revert the increment.
238                self.current_batch_size.fetch_sub(1, Ordering::AcqRel);
239                // Record queue-full drop in self-diagnostics
240                #[cfg(feature = "experimental_metrics_bound_instruments")]
241                self.processed_queue_full.add(1);
242
243                // Increment dropped logs count. The first time we have to drop
244                // a log, emit a warning.
245                if self.dropped_logs_count.fetch_add(1, Ordering::Relaxed) == 0 {
246                    otel_warn!(name: "BatchLogProcessor.LogDroppingStarted",
247                        message = "BatchLogProcessor dropped a LogRecord 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 logs dropped.");
248                }
249            }
250            Err(mpsc::TrySendError::Disconnected(_)) => {
251                // The record never entered the channel; revert the increment.
252                self.current_batch_size.fetch_sub(1, Ordering::AcqRel);
253                // Record after-shutdown drop in self-diagnostics
254                #[cfg(feature = "experimental_metrics_bound_instruments")]
255                self.processed_after_shutdown.add(1);
256
257                // The following `otel_warn!` may cause an infinite feedback loop of
258                // 'telemetry-induced-telemetry', potentially causing a stack overflow
259                let _guard = Context::enter_telemetry_suppressed_scope();
260
261                // Given background thread is the only receiver, and it's
262                // disconnected, it indicates the thread is shutdown
263                otel_warn!(
264                    name: "BatchLogProcessor.Emit.AfterShutdown",
265                    message = "Logs are being emitted even after Shutdown. This indicates incorrect lifecycle management of OTelLoggerProvider in application. Logs will not be exported."
266                );
267            }
268        }
269    }
270
271    fn force_flush(&self) -> OTelSdkResult {
272        let (sender, receiver) = mpsc::sync_channel(1);
273        match self
274            .message_sender
275            .try_send(BatchMessage::ForceFlush(sender))
276        {
277            Ok(_) => receiver
278                .recv_timeout(self.forceflush_timeout)
279                .map_err(|err| {
280                    if err == RecvTimeoutError::Timeout {
281                        OTelSdkError::Timeout(self.forceflush_timeout)
282                    } else {
283                        OTelSdkError::InternalFailure(format!("{err}"))
284                    }
285                })?,
286            Err(mpsc::TrySendError::Full(_)) => {
287                // If the control message could not be sent, emit a warning.
288                otel_debug!(
289                    name: "BatchLogProcessor.ForceFlush.ControlChannelFull",
290                    message = "Control message to flush the worker thread could not be sent as the control channel is full. This can occur if user repeatedily calls force_flush/shutdown without finishing the previous call."
291                );
292                Err(OTelSdkError::InternalFailure("ForceFlush cannot be performed as Control channel is full. This can occur if user repeatedily calls force_flush/shutdown without finishing the previous call.".into()))
293            }
294            Err(mpsc::TrySendError::Disconnected(_)) => {
295                // Given background thread is the only receiver, and it's
296                // disconnected, it indicates the thread is shutdown
297                otel_debug!(
298                    name: "BatchLogProcessor.ForceFlush.AlreadyShutdown",
299                    message = "ForceFlush invoked after Shutdown. This will not perform Flush and indicates a incorrect lifecycle management in Application."
300                );
301
302                Err(OTelSdkError::AlreadyShutdown)
303            }
304        }
305    }
306
307    fn shutdown_with_timeout(&self, timeout: Duration) -> OTelSdkResult {
308        let dropped_logs = self.dropped_logs_count.load(Ordering::Relaxed);
309        let max_queue_size = self.max_queue_size;
310        if dropped_logs > 0 {
311            otel_warn!(
312                name: "BatchLogProcessor.LogsDropped",
313                dropped_logs_count = dropped_logs,
314                max_queue_size = max_queue_size,
315                message = "Logs were dropped due to a queue being full. The count represents the total count of log records dropped in the lifetime of this BatchLogProcessor. Consider increasing the queue size and/or decrease delay between intervals."
316            );
317        }
318
319        let (sender, receiver) = mpsc::sync_channel(1);
320        match self.message_sender.try_send(BatchMessage::Shutdown(sender)) {
321            Ok(_) => {
322                receiver
323                    .recv_timeout(timeout)
324                    .map(|_| {
325                        // join the background thread after receiving back the
326                        // shutdown signal
327                        if let Some(handle) = self.handle.lock().unwrap().take() {
328                            handle.join().unwrap();
329                        }
330                        OTelSdkResult::Ok(())
331                    })
332                    .map_err(|err| match err {
333                        RecvTimeoutError::Timeout => {
334                            // TODO: When shutdown times out, log records still
335                            // in the queue or mid-export are silently lost. The
336                            // background thread is not joined and may continue
337                            // running. Consider: (1) recording the lost count
338                            // in the self-diagnostics counter, (2) joining the
339                            // thread with a best-effort wait, or (3) signalling
340                            // the thread to abort the current export.
341                            otel_error!(
342                                name: "BatchLogProcessor.Shutdown.Timeout",
343                                message = "BatchLogProcessor shutdown timing out."
344                            );
345                            OTelSdkError::Timeout(timeout)
346                        }
347                        _ => {
348                            otel_error!(
349                                name: "BatchLogProcessor.Shutdown.Error",
350                                error = format!("{}", err)
351                            );
352                            OTelSdkError::InternalFailure(format!("{err}"))
353                        }
354                    })?
355            }
356            Err(mpsc::TrySendError::Full(_)) => {
357                // If the control message could not be sent, emit a warning.
358                otel_debug!(
359                    name: "BatchLogProcessor.Shutdown.ControlChannelFull",
360                    message = "Control message to shutdown the worker thread could not be sent as the control channel is full. This can occur if user repeatedily calls force_flush/shutdown without finishing the previous call."
361                );
362                Err(OTelSdkError::InternalFailure("Shutdown cannot be performed as Control channel is full. This can occur if user repeatedily calls force_flush/shutdown without finishing the previous call.".into()))
363            }
364            Err(mpsc::TrySendError::Disconnected(_)) => {
365                // Given background thread is the only receiver, and it's
366                // disconnected, it indicates the thread is shutdown
367                otel_debug!(
368                    name: "BatchLogProcessor.Shutdown.AlreadyShutdown",
369                    message = "Shutdown is being invoked more than once. This is noop, but indicates a potential issue in the application's lifecycle management."
370                );
371
372                Err(OTelSdkError::AlreadyShutdown)
373            }
374        }
375    }
376
377    fn set_resource(&mut self, resource: &Resource) {
378        let resource = Arc::new(resource.clone());
379        let _ = self
380            .message_sender
381            .try_send(BatchMessage::SetResource(resource));
382    }
383}
384
385impl BatchLogProcessor {
386    pub(crate) fn new<E>(mut exporter: E, config: BatchConfig) -> Self
387    where
388        E: LogExporter + Send + Sync + 'static,
389    {
390        let (logs_sender, logs_receiver) = mpsc::sync_channel::<LogsData>(config.max_queue_size);
391        let (message_sender, message_receiver) = mpsc::sync_channel::<BatchMessage>(64); // Is this a reasonable bound?
392        let max_queue_size = config.max_queue_size;
393        let max_export_batch_size = config.max_export_batch_size;
394        let current_batch_size = Arc::new(AtomicUsize::new(0));
395        let current_batch_size_for_thread = current_batch_size.clone();
396
397        // Self-diagnostics: create the otel.sdk.processor.log.processed counter.
398        // Created before the worker thread is spawned so the success counter can
399        // be moved into the worker and incremented when a batch is exported.
400        #[cfg(feature = "experimental_metrics_bound_instruments")]
401        let (processed_success, processed_queue_full, processed_after_shutdown) = {
402            static INSTANCE_COUNTER: AtomicUsize = AtomicUsize::new(0);
403            let instance_id = INSTANCE_COUNTER.fetch_add(1, Ordering::Relaxed);
404            let component_name = format!("batching_log_processor/{instance_id}");
405
406            let meter = opentelemetry::global::meter("otel.sdk");
407            let counter = meter
408                .u64_counter("otel.sdk.processor.log.processed")
409                .with_description(
410                    "The number of log records for which the processing has finished, \
411                     either successful or failed.",
412                )
413                .with_unit("{log_record}")
414                .build();
415
416            // Self-diagnostics: otel.sdk.processor.log.queue.capacity. A weak
417            // reference ensures a dropped processor stops reporting, since
418            // observable callbacks live for the meter provider's lifetime and
419            // cannot be individually unregistered.
420            let capacity_attrs = [
421                KeyValue::new("otel.component.type", "batching_log_processor"),
422                KeyValue::new("otel.component.name", component_name.clone()),
423            ];
424            let capacity_state = Arc::downgrade(&current_batch_size);
425            let capacity_value = i64::try_from(max_queue_size).unwrap_or(i64::MAX);
426            let _ = meter
427                .i64_observable_up_down_counter("otel.sdk.processor.log.queue.capacity")
428                .with_description(
429                    "The maximum number of log records the queue of a given instance of \
430                     an SDK log processor can hold.",
431                )
432                .with_unit("{log_record}")
433                .with_callback(move |observer| {
434                    // The capacity value is constant; this is only a liveness
435                    // guard so a dropped processor stops emitting this
436                    // otherwise-unregisterable callback. `strong_count()` is
437                    // sufficient since the callback never reads the state.
438                    if capacity_state.strong_count() > 0 {
439                        observer.observe(capacity_value, &capacity_attrs);
440                    }
441                })
442                .build();
443
444            // Attribute values follow the OTel semantic conventions for SDK metrics:
445            // https://github.com/open-telemetry/semantic-conventions/blob/main/docs/otel/sdk-metrics.md#metric-otelsdkprocessorlogprocessed
446            // https://github.com/open-telemetry/semantic-conventions/blob/main/docs/registry/attributes/otel.md#otel-component-attributes
447            let success_attrs = [
448                KeyValue::new("otel.component.type", "batching_log_processor"),
449                KeyValue::new("otel.component.name", component_name.clone()),
450            ];
451            let queue_full_attrs = [
452                KeyValue::new("error.type", "queue_full"),
453                KeyValue::new("otel.component.type", "batching_log_processor"),
454                KeyValue::new("otel.component.name", component_name.clone()),
455            ];
456            let after_shutdown_attrs = [
457                KeyValue::new("error.type", "already_shutdown"),
458                KeyValue::new("otel.component.type", "batching_log_processor"),
459                KeyValue::new("otel.component.name", component_name),
460            ];
461
462            (
463                counter.bind(&success_attrs),
464                counter.bind(&queue_full_attrs),
465                counter.bind(&after_shutdown_attrs),
466            )
467        };
468
469        let handle = thread::Builder::new()
470            .name("OpenTelemetry.Logs.BatchProcessor".to_string())
471            .spawn(move || {
472                let _suppress_guard = Context::enter_telemetry_suppressed_scope();
473                otel_debug!(
474                    name: "BatchLogProcessor.ThreadStarted",
475                    interval_in_millisecs = config.scheduled_delay.as_millis(),
476                    max_export_batch_size = config.max_export_batch_size,
477                    max_queue_size = max_queue_size,
478                );
479                let mut last_export_time = Instant::now();
480                let mut logs = Vec::with_capacity(config.max_export_batch_size);
481                let current_batch_size = current_batch_size_for_thread;
482
483                // Counts records for the otel.sdk.processor.log.processed
484                // metric; a no-op when the self-diagnostics feature is disabled.
485                #[cfg(feature = "experimental_metrics_bound_instruments")]
486                let record_processed_success = move |count: u64| processed_success.add(count);
487                #[cfg(not(feature = "experimental_metrics_bound_instruments"))]
488                let record_processed_success = |_count: u64| {};
489
490                // This method gets up to `max_export_batch_size` amount of logs from the channel and exports them.
491                // It returns the result of the export operation.
492                // It expects the logs vec to be empty when it's called.
493                #[inline]
494                fn get_logs_and_export<E, F>(
495                    logs_receiver: &mpsc::Receiver<LogsData>,
496                    exporter: &E,
497                    logs: &mut Vec<LogsData>,
498                    last_export_time: &mut Instant,
499                    current_batch_size: &AtomicUsize,
500                    max_export_size: usize,
501                    record_processed_success: &F,
502                ) -> OTelSdkResult
503                where
504                    E: LogExporter + Send + Sync + 'static,
505                    F: Fn(u64),
506                {
507                    let target = current_batch_size.load(Ordering::Acquire); // `target` is used to determine the stopping criteria for exporting logs.
508                    let mut result = OTelSdkResult::Ok(());
509                    let mut total_exported_logs: usize = 0;
510
511                    while target > 0 && total_exported_logs < target {
512                        let batch_limit = max_export_size.min(target - total_exported_logs);
513
514                        // Get up to the remaining target batch size from the channel and push them to the logs vec
515                        while let Ok(log) = logs_receiver.try_recv() {
516                            logs.push(log);
517                            if logs.len() == batch_limit {
518                                break;
519                            }
520                        }
521
522                        let count_of_logs = logs.len(); // Count of logs that will be exported
523                        if count_of_logs == 0 {
524                            break;
525                        }
526                        total_exported_logs += count_of_logs;
527
528                        // Count the batch as processed before invoking the
529                        // exporter, regardless of the export outcome.
530                        record_processed_success(count_of_logs as u64);
531
532                        result = export_batch_sync(exporter, logs, last_export_time); // This method clears the logs vec after exporting
533
534                        current_batch_size.fetch_sub(count_of_logs, Ordering::AcqRel);
535                    }
536                    result
537                }
538
539                loop {
540                    let remaining_time = config
541                        .scheduled_delay
542                        .checked_sub(last_export_time.elapsed())
543                        .unwrap_or(config.scheduled_delay);
544
545                    match message_receiver.recv_timeout(remaining_time) {
546                        Ok(BatchMessage::ExportLog(export_log_message_sent)) => {
547                            // Reset the export log message sent flag now it has has been processed.
548                            export_log_message_sent.store(false, Ordering::Relaxed);
549
550                            otel_debug!(
551                                name: "BatchLogProcessor.ExportingDueToBatchSize",
552                            );
553
554                            let _ = get_logs_and_export(
555                                &logs_receiver,
556                                &exporter,
557                                &mut logs,
558                                &mut last_export_time,
559                                &current_batch_size,
560                                max_export_batch_size,
561                                &record_processed_success,
562                            );
563                        }
564                        Ok(BatchMessage::ForceFlush(sender)) => {
565                            otel_debug!(name: "BatchLogProcessor.ExportingDueToForceFlush");
566                            let result = get_logs_and_export(
567                                &logs_receiver,
568                                &exporter,
569                                &mut logs,
570                                &mut last_export_time,
571                                &current_batch_size,
572                                max_export_batch_size,
573                                &record_processed_success,
574                            );
575                            let _ = sender.send(result);
576                        }
577                        Ok(BatchMessage::Shutdown(sender)) => {
578                            otel_debug!(name: "BatchLogProcessor.ExportingDueToShutdown");
579                            let result = get_logs_and_export(
580                                &logs_receiver,
581                                &exporter,
582                                &mut logs,
583                                &mut last_export_time,
584                                &current_batch_size,
585                                max_export_batch_size,
586                                &record_processed_success,
587                            );
588                            let _ = exporter.shutdown();
589                            let _ = sender.send(result);
590
591                            otel_debug!(
592                                name: "BatchLogProcessor.ThreadExiting",
593                                reason = "ShutdownRequested"
594                            );
595                            //
596                            // break out the loop and return from the current background thread.
597                            //
598                            break;
599                        }
600                        Ok(BatchMessage::SetResource(resource)) => {
601                            exporter.set_resource(&resource);
602                        }
603                        Err(RecvTimeoutError::Timeout) => {
604                            otel_debug!(
605                                name: "BatchLogProcessor.ExportingDueToTimer",
606                            );
607
608                            let _ = get_logs_and_export(
609                                &logs_receiver,
610                                &exporter,
611                                &mut logs,
612                                &mut last_export_time,
613                                &current_batch_size,
614                                max_export_batch_size,
615                                &record_processed_success,
616                            );
617                        }
618                        Err(RecvTimeoutError::Disconnected) => {
619                            // Channel disconnected, only thing to do is break
620                            // out (i.e exit the thread)
621                            otel_debug!(
622                                name: "BatchLogProcessor.ThreadExiting",
623                                reason = "MessageSenderDisconnected"
624                            );
625                            break;
626                        }
627                    }
628                }
629                otel_debug!(
630                    name: "BatchLogProcessor.ThreadStopped"
631                );
632            })
633            .expect("Thread spawn failed."); //TODO: Handle thread spawn failure
634
635        // Return batch processor with link to worker
636        BatchLogProcessor {
637            logs_sender,
638            message_sender,
639            handle: Mutex::new(Some(handle)),
640            forceflush_timeout: Duration::from_secs(5), // TODO: make this configurable
641            dropped_logs_count: AtomicUsize::new(0),
642            max_queue_size,
643            export_log_message_sent: Arc::new(AtomicBool::new(false)),
644            current_batch_size,
645            max_export_batch_size,
646            #[cfg(feature = "experimental_metrics_bound_instruments")]
647            processed_queue_full,
648            #[cfg(feature = "experimental_metrics_bound_instruments")]
649            processed_after_shutdown,
650        }
651    }
652
653    /// Create a new batch processor builder
654    pub fn builder<E>(exporter: E) -> BatchLogProcessorBuilder<E>
655    where
656        E: LogExporter,
657    {
658        BatchLogProcessorBuilder {
659            exporter,
660            config: Default::default(),
661        }
662    }
663}
664
665#[allow(clippy::vec_box)]
666fn export_batch_sync<E>(
667    exporter: &E,
668    batch: &mut Vec<Box<(SdkLogRecord, InstrumentationScope)>>,
669    last_export_time: &mut Instant,
670) -> OTelSdkResult
671where
672    E: LogExporter + ?Sized,
673{
674    *last_export_time = Instant::now();
675
676    if batch.is_empty() {
677        return OTelSdkResult::Ok(());
678    }
679
680    let export = exporter.export(LogBatch::new_with_owned_data(batch.as_slice()));
681    let export_result = futures_executor::block_on(export);
682
683    // Clear the batch vec after exporting
684    batch.clear();
685
686    match export_result {
687        Ok(_) => OTelSdkResult::Ok(()),
688        Err(err) => {
689            otel_error!(
690                name: "BatchLogProcessor.ExportError",
691                error = format!("{}", err)
692            );
693            OTelSdkResult::Err(err)
694        }
695    }
696}
697
698///
699/// A builder for creating [`BatchLogProcessor`] instances.
700///
701#[derive(Debug)]
702pub struct BatchLogProcessorBuilder<E> {
703    exporter: E,
704    config: BatchConfig,
705}
706
707impl<E> BatchLogProcessorBuilder<E>
708where
709    E: LogExporter + 'static,
710{
711    /// Set the BatchConfig for [`BatchLogProcessorBuilder`]
712    pub fn with_batch_config(self, config: BatchConfig) -> Self {
713        BatchLogProcessorBuilder { config, ..self }
714    }
715
716    /// Build a batch processor
717    pub fn build(self) -> BatchLogProcessor {
718        BatchLogProcessor::new(self.exporter, self.config)
719    }
720}
721
722/// Batch log processor configuration.
723/// Use [`BatchConfigBuilder`] to configure your own instance of [`BatchConfig`].
724#[derive(Debug)]
725#[allow(dead_code)]
726pub struct BatchConfig {
727    /// The maximum queue size to buffer logs for delayed processing. If the
728    /// queue gets full it drops the logs. The default value of is 2048.
729    pub(crate) max_queue_size: usize,
730
731    /// The delay interval in milliseconds between two consecutive processing
732    /// of batches. The default value is 1 second.
733    pub(crate) scheduled_delay: Duration,
734
735    /// The maximum number of logs to process in a single batch. If there are
736    /// more than one batch worth of logs then it processes multiple batches
737    /// of logs one batch after the other without any delay. The default value
738    /// is 512.
739    pub(crate) max_export_batch_size: usize,
740
741    /// The maximum duration to export a batch of data.
742    #[cfg(feature = "experimental_logs_batch_log_processor_with_async_runtime")]
743    pub(crate) max_export_timeout: Duration,
744}
745
746impl Default for BatchConfig {
747    fn default() -> Self {
748        BatchConfigBuilder::default().build()
749    }
750}
751
752/// A builder for creating [`BatchConfig`] instances.
753#[derive(Debug)]
754pub struct BatchConfigBuilder {
755    max_queue_size: usize,
756    scheduled_delay: Duration,
757    max_export_batch_size: usize,
758    #[cfg(feature = "experimental_logs_batch_log_processor_with_async_runtime")]
759    max_export_timeout: Duration,
760}
761
762impl Default for BatchConfigBuilder {
763    /// Create a new [`BatchConfigBuilder`] initialized with default batch config values as per the specs.
764    /// The values are overridden by environment variables if set.
765    /// The supported environment variables are:
766    /// * `OTEL_BLRP_MAX_QUEUE_SIZE`
767    /// * `OTEL_BLRP_SCHEDULE_DELAY`
768    /// * `OTEL_BLRP_MAX_EXPORT_BATCH_SIZE`
769    /// * `OTEL_BLRP_EXPORT_TIMEOUT`
770    ///
771    /// Note: Programmatic configuration overrides any value set via the environment variable.
772    fn default() -> Self {
773        BatchConfigBuilder {
774            max_queue_size: OTEL_BLRP_MAX_QUEUE_SIZE_DEFAULT,
775            scheduled_delay: OTEL_BLRP_SCHEDULE_DELAY_DEFAULT,
776            max_export_batch_size: OTEL_BLRP_MAX_EXPORT_BATCH_SIZE_DEFAULT,
777            #[cfg(feature = "experimental_logs_batch_log_processor_with_async_runtime")]
778            max_export_timeout: OTEL_BLRP_EXPORT_TIMEOUT_DEFAULT,
779        }
780        .init_from_env_vars()
781    }
782}
783
784impl BatchConfigBuilder {
785    /// Set max_queue_size for [`BatchConfigBuilder`].
786    /// It's the maximum queue size to buffer logs for delayed processing.
787    /// If the queue gets full it will drop the logs.
788    /// The default value is 2048.
789    ///
790    /// Corresponding environment variable: `OTEL_BLRP_MAX_QUEUE_SIZE`.
791    ///
792    /// Note: Programmatically setting this will override any value set via the environment variable.
793    pub fn with_max_queue_size(mut self, max_queue_size: usize) -> Self {
794        self.max_queue_size = max_queue_size;
795        self
796    }
797
798    /// Set scheduled_delay for [`BatchConfigBuilder`].
799    /// It's the delay interval in milliseconds between two consecutive processing of batches.
800    /// The default value is 1000 milliseconds.
801    ///
802    /// Corresponding environment variable: `OTEL_BLRP_SCHEDULE_DELAY`.
803    ///
804    /// Note: Programmatically setting this will override any value set via the environment variable.
805    pub fn with_scheduled_delay(mut self, scheduled_delay: Duration) -> Self {
806        self.scheduled_delay = scheduled_delay;
807        self
808    }
809
810    /// Set max_export_timeout for [`BatchConfigBuilder`].
811    /// It's the maximum duration to export a batch of data.
812    /// The default value is 30000 milliseconds.
813    ///
814    /// Corresponding environment variable: `OTEL_BLRP_EXPORT_TIMEOUT`.
815    ///
816    /// Note: Programmatically setting this will override any value set via the environment variable.
817    #[cfg(feature = "experimental_logs_batch_log_processor_with_async_runtime")]
818    pub fn with_max_export_timeout(mut self, max_export_timeout: Duration) -> Self {
819        self.max_export_timeout = max_export_timeout;
820        self
821    }
822
823    /// Set max_export_batch_size for [`BatchConfigBuilder`].
824    /// It's the maximum number of logs to process in a single batch. If there are
825    /// more than one batch worth of logs then it processes multiple batches
826    /// of logs one batch after the other without any delay.
827    /// The default value is 512.
828    ///
829    /// Corresponding environment variable: `OTEL_BLRP_MAX_EXPORT_BATCH_SIZE`.
830    ///
831    /// Note: Programmatically setting this will override any value set via the environment variable.
832    pub fn with_max_export_batch_size(mut self, max_export_batch_size: usize) -> Self {
833        self.max_export_batch_size = max_export_batch_size;
834        self
835    }
836
837    /// Builds a `BatchConfig` enforcing the following invariants:
838    /// * `max_export_batch_size` must be less than or equal to `max_queue_size`.
839    pub fn build(self) -> BatchConfig {
840        // max export batch size must be less or equal to max queue size.
841        // we set max export batch size to max queue size if it's larger than max queue size.
842        let max_export_batch_size = min(self.max_export_batch_size, self.max_queue_size);
843
844        BatchConfig {
845            max_queue_size: self.max_queue_size,
846            scheduled_delay: self.scheduled_delay,
847            #[cfg(feature = "experimental_logs_batch_log_processor_with_async_runtime")]
848            max_export_timeout: self.max_export_timeout,
849            max_export_batch_size,
850        }
851    }
852
853    fn init_from_env_vars(mut self) -> Self {
854        if let Some(max_queue_size) = env::var(OTEL_BLRP_MAX_QUEUE_SIZE)
855            .ok()
856            .and_then(|queue_size| usize::from_str(&queue_size).ok())
857        {
858            self.max_queue_size = max_queue_size;
859        }
860
861        if let Some(max_export_batch_size) = env::var(OTEL_BLRP_MAX_EXPORT_BATCH_SIZE)
862            .ok()
863            .and_then(|batch_size| usize::from_str(&batch_size).ok())
864        {
865            self.max_export_batch_size = max_export_batch_size;
866        }
867
868        if let Some(scheduled_delay) = env::var(OTEL_BLRP_SCHEDULE_DELAY)
869            .ok()
870            .and_then(|delay| u64::from_str(&delay).ok())
871        {
872            self.scheduled_delay = Duration::from_millis(scheduled_delay);
873        }
874
875        #[cfg(feature = "experimental_logs_batch_log_processor_with_async_runtime")]
876        if let Some(max_export_timeout) = env::var(OTEL_BLRP_EXPORT_TIMEOUT)
877            .ok()
878            .and_then(|s| u64::from_str(&s).ok())
879        {
880            self.max_export_timeout = Duration::from_millis(max_export_timeout);
881        }
882
883        self
884    }
885}
886
887#[cfg(all(test, feature = "testing", feature = "logs"))]
888mod tests {
889    use super::{
890        BatchConfig, BatchConfigBuilder, BatchLogProcessor, OTEL_BLRP_MAX_EXPORT_BATCH_SIZE,
891        OTEL_BLRP_MAX_EXPORT_BATCH_SIZE_DEFAULT, OTEL_BLRP_MAX_QUEUE_SIZE,
892        OTEL_BLRP_MAX_QUEUE_SIZE_DEFAULT, OTEL_BLRP_SCHEDULE_DELAY,
893        OTEL_BLRP_SCHEDULE_DELAY_DEFAULT,
894    };
895    #[cfg(feature = "experimental_logs_batch_log_processor_with_async_runtime")]
896    use super::{OTEL_BLRP_EXPORT_TIMEOUT, OTEL_BLRP_EXPORT_TIMEOUT_DEFAULT};
897    use crate::error::OTelSdkResult;
898    use crate::logs::log_processor::tests::MockLogExporter;
899    use crate::logs::SdkLogRecord;
900    use crate::logs::{LogBatch, LogExporter};
901    use crate::{
902        logs::{InMemoryLogExporter, InMemoryLogExporterBuilder, LogProcessor, SdkLoggerProvider},
903        Resource,
904    };
905    use opentelemetry::logs::LogRecord;
906    use opentelemetry::InstrumentationScope;
907    use opentelemetry::KeyValue;
908    use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
909    use std::sync::mpsc;
910    use std::sync::{Arc, Mutex};
911    use std::time::Duration;
912
913    #[test]
914    fn test_default_const_values() {
915        assert_eq!(OTEL_BLRP_SCHEDULE_DELAY, "OTEL_BLRP_SCHEDULE_DELAY");
916        assert_eq!(OTEL_BLRP_SCHEDULE_DELAY_DEFAULT.as_millis(), 1_000);
917        #[cfg(feature = "experimental_logs_batch_log_processor_with_async_runtime")]
918        assert_eq!(OTEL_BLRP_EXPORT_TIMEOUT, "OTEL_BLRP_EXPORT_TIMEOUT");
919        #[cfg(feature = "experimental_logs_batch_log_processor_with_async_runtime")]
920        assert_eq!(OTEL_BLRP_EXPORT_TIMEOUT_DEFAULT.as_millis(), 30_000);
921        assert_eq!(OTEL_BLRP_MAX_QUEUE_SIZE, "OTEL_BLRP_MAX_QUEUE_SIZE");
922        assert_eq!(OTEL_BLRP_MAX_QUEUE_SIZE_DEFAULT, 2_048);
923        assert_eq!(
924            OTEL_BLRP_MAX_EXPORT_BATCH_SIZE,
925            "OTEL_BLRP_MAX_EXPORT_BATCH_SIZE"
926        );
927        assert_eq!(OTEL_BLRP_MAX_EXPORT_BATCH_SIZE_DEFAULT, 512);
928    }
929
930    #[test]
931    fn test_default_batch_config_adheres_to_specification() {
932        // The following environment variables are expected to be unset so that their default values are used.
933        let env_vars = vec![
934            OTEL_BLRP_SCHEDULE_DELAY,
935            #[cfg(feature = "experimental_logs_batch_log_processor_with_async_runtime")]
936            OTEL_BLRP_EXPORT_TIMEOUT,
937            OTEL_BLRP_MAX_QUEUE_SIZE,
938            OTEL_BLRP_MAX_EXPORT_BATCH_SIZE,
939        ];
940
941        let config = temp_env::with_vars_unset(env_vars, BatchConfig::default);
942
943        assert_eq!(config.scheduled_delay, OTEL_BLRP_SCHEDULE_DELAY_DEFAULT);
944        #[cfg(feature = "experimental_logs_batch_log_processor_with_async_runtime")]
945        assert_eq!(config.max_export_timeout, OTEL_BLRP_EXPORT_TIMEOUT_DEFAULT);
946        assert_eq!(config.max_queue_size, OTEL_BLRP_MAX_QUEUE_SIZE_DEFAULT);
947        assert_eq!(
948            config.max_export_batch_size,
949            OTEL_BLRP_MAX_EXPORT_BATCH_SIZE_DEFAULT
950        );
951    }
952
953    #[test]
954    fn test_code_based_config_overrides_env_vars() {
955        let env_vars = vec![
956            (OTEL_BLRP_SCHEDULE_DELAY, Some("2000")),
957            (OTEL_BLRP_MAX_QUEUE_SIZE, Some("4096")),
958            (OTEL_BLRP_MAX_EXPORT_BATCH_SIZE, Some("1024")),
959        ];
960
961        temp_env::with_vars(env_vars, || {
962            let config = BatchConfigBuilder::default()
963                .with_max_queue_size(2048)
964                .with_scheduled_delay(Duration::from_millis(1000))
965                .with_max_export_batch_size(512)
966                .build();
967
968            assert_eq!(config.scheduled_delay, Duration::from_millis(1000));
969            assert_eq!(config.max_queue_size, 2048);
970            assert_eq!(config.max_export_batch_size, 512);
971        });
972    }
973
974    #[test]
975    fn test_batch_config_configurable_by_env_vars() {
976        let env_vars = vec![
977            (OTEL_BLRP_SCHEDULE_DELAY, Some("2000")),
978            #[cfg(feature = "experimental_logs_batch_log_processor_with_async_runtime")]
979            (OTEL_BLRP_EXPORT_TIMEOUT, Some("60000")),
980            (OTEL_BLRP_MAX_QUEUE_SIZE, Some("4096")),
981            (OTEL_BLRP_MAX_EXPORT_BATCH_SIZE, Some("1024")),
982        ];
983
984        let config = temp_env::with_vars(env_vars, BatchConfig::default);
985
986        assert_eq!(config.scheduled_delay, Duration::from_millis(2000));
987        #[cfg(feature = "experimental_logs_batch_log_processor_with_async_runtime")]
988        assert_eq!(config.max_export_timeout, Duration::from_millis(60000));
989        assert_eq!(config.max_queue_size, 4096);
990        assert_eq!(config.max_export_batch_size, 1024);
991    }
992    #[test]
993    fn test_force_flush_being_called() {
994        #[derive(Debug, Clone)]
995        struct MockExporter {
996            export_called: Arc<AtomicBool>,
997        }
998        impl LogExporter for MockExporter {
999            async fn export(&self, _batch: LogBatch<'_>) -> OTelSdkResult {
1000                self.export_called.store(true, Ordering::SeqCst);
1001                Ok(())
1002            }
1003        }
1004        let exporter = MockExporter {
1005            export_called: Arc::new(AtomicBool::new(false)),
1006        };
1007        let processor = BatchLogProcessor::new(exporter.clone(), BatchConfig::default());
1008        let scope = opentelemetry::InstrumentationScope::builder("my-crate")
1009            .with_schema_url("https://opentelemetry.io/schemas/1.17.0")
1010            .build();
1011        processor.emit(&mut SdkLogRecord::new(), &scope);
1012        processor.force_flush().unwrap();
1013        assert!(exporter.export_called.load(Ordering::SeqCst));
1014    }
1015
1016    #[test]
1017    fn test_batch_config_max_export_batch_size_validation() {
1018        let env_vars = vec![
1019            (OTEL_BLRP_MAX_QUEUE_SIZE, Some("256")),
1020            (OTEL_BLRP_MAX_EXPORT_BATCH_SIZE, Some("1024")),
1021        ];
1022
1023        let config = temp_env::with_vars(env_vars, BatchConfig::default);
1024
1025        assert_eq!(config.max_queue_size, 256);
1026        assert_eq!(config.max_export_batch_size, 256);
1027        assert_eq!(config.scheduled_delay, OTEL_BLRP_SCHEDULE_DELAY_DEFAULT);
1028        #[cfg(feature = "experimental_logs_batch_log_processor_with_async_runtime")]
1029        assert_eq!(config.max_export_timeout, OTEL_BLRP_EXPORT_TIMEOUT_DEFAULT);
1030    }
1031
1032    #[test]
1033    fn test_batch_config_with_fields() {
1034        let batch_builder = BatchConfigBuilder::default()
1035            .with_max_export_batch_size(1)
1036            .with_scheduled_delay(Duration::from_millis(2))
1037            .with_max_queue_size(4);
1038
1039        #[cfg(feature = "experimental_logs_batch_log_processor_with_async_runtime")]
1040        let batch_builder = batch_builder.with_max_export_timeout(Duration::from_millis(3));
1041        let batch = batch_builder.build();
1042
1043        assert_eq!(batch.max_export_batch_size, 1);
1044        assert_eq!(batch.scheduled_delay, Duration::from_millis(2));
1045        #[cfg(feature = "experimental_logs_batch_log_processor_with_async_runtime")]
1046        assert_eq!(batch.max_export_timeout, Duration::from_millis(3));
1047        assert_eq!(batch.max_queue_size, 4);
1048    }
1049
1050    #[test]
1051    fn test_build_batch_log_processor_builder() {
1052        let mut env_vars = vec![
1053            (OTEL_BLRP_MAX_EXPORT_BATCH_SIZE, Some("500")),
1054            (OTEL_BLRP_SCHEDULE_DELAY, Some("I am not number")),
1055            #[cfg(feature = "experimental_logs_batch_log_processor_with_async_runtime")]
1056            (OTEL_BLRP_EXPORT_TIMEOUT, Some("2046")),
1057        ];
1058        temp_env::with_vars(env_vars.clone(), || {
1059            let builder = BatchLogProcessor::builder(InMemoryLogExporter::default());
1060
1061            assert_eq!(builder.config.max_export_batch_size, 500);
1062            assert_eq!(
1063                builder.config.scheduled_delay,
1064                OTEL_BLRP_SCHEDULE_DELAY_DEFAULT
1065            );
1066            assert_eq!(
1067                builder.config.max_queue_size,
1068                OTEL_BLRP_MAX_QUEUE_SIZE_DEFAULT
1069            );
1070
1071            #[cfg(feature = "experimental_logs_batch_log_processor_with_async_runtime")]
1072            assert_eq!(
1073                builder.config.max_export_timeout,
1074                Duration::from_millis(2046)
1075            );
1076        });
1077
1078        env_vars.push((OTEL_BLRP_MAX_QUEUE_SIZE, Some("120")));
1079
1080        temp_env::with_vars(env_vars, || {
1081            let builder = BatchLogProcessor::builder(InMemoryLogExporter::default());
1082            assert_eq!(builder.config.max_export_batch_size, 120);
1083            assert_eq!(builder.config.max_queue_size, 120);
1084        });
1085    }
1086
1087    #[test]
1088    fn test_build_batch_log_processor_builder_with_custom_config() {
1089        let expected = BatchConfigBuilder::default()
1090            .with_max_export_batch_size(1)
1091            .with_scheduled_delay(Duration::from_millis(2))
1092            .with_max_queue_size(4)
1093            .build();
1094
1095        let builder =
1096            BatchLogProcessor::builder(InMemoryLogExporter::default()).with_batch_config(expected);
1097
1098        let actual = &builder.config;
1099        assert_eq!(actual.max_export_batch_size, 1);
1100        assert_eq!(actual.scheduled_delay, Duration::from_millis(2));
1101        assert_eq!(actual.max_queue_size, 4);
1102    }
1103
1104    #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
1105    async fn test_set_resource_batch_processor() {
1106        let exporter = MockLogExporter {
1107            resource: Arc::new(Mutex::new(None)),
1108        };
1109        let processor = BatchLogProcessor::new(exporter.clone(), BatchConfig::default());
1110        let provider = SdkLoggerProvider::builder()
1111            .with_log_processor(processor)
1112            .with_resource(
1113                Resource::builder_empty()
1114                    .with_attributes([
1115                        KeyValue::new("k1", "v1"),
1116                        KeyValue::new("k2", "v3"),
1117                        KeyValue::new("k3", "v3"),
1118                        KeyValue::new("k4", "v4"),
1119                        KeyValue::new("k5", "v5"),
1120                    ])
1121                    .build(),
1122            )
1123            .build();
1124
1125        provider.force_flush().unwrap();
1126
1127        assert_eq!(exporter.get_resource().unwrap().into_iter().count(), 5);
1128        let _ = provider.shutdown();
1129    }
1130
1131    #[tokio::test(flavor = "multi_thread")]
1132    async fn test_batch_shutdown() {
1133        // assert we will receive an error
1134        // setup
1135        let exporter = InMemoryLogExporterBuilder::default()
1136            .keep_records_on_shutdown()
1137            .build();
1138        let processor = BatchLogProcessor::new(exporter.clone(), BatchConfig::default());
1139
1140        let mut record = SdkLogRecord::new();
1141        let instrumentation = InstrumentationScope::default();
1142
1143        processor.emit(&mut record, &instrumentation);
1144        processor.force_flush().unwrap();
1145        processor.shutdown().unwrap();
1146        // todo: expect to see errors here. How should we assert this?
1147        processor.emit(&mut record, &instrumentation);
1148        assert_eq!(1, exporter.get_emitted_logs().unwrap().len());
1149        assert!(exporter.is_shutdown_called());
1150    }
1151
1152    #[tokio::test(flavor = "current_thread")]
1153    async fn test_batch_log_processor_shutdown_under_async_runtime_current_flavor_multi_thread() {
1154        let exporter = InMemoryLogExporterBuilder::default().build();
1155        let processor = BatchLogProcessor::new(exporter.clone(), BatchConfig::default());
1156
1157        processor.shutdown().unwrap();
1158    }
1159
1160    #[tokio::test(flavor = "current_thread")]
1161    async fn test_batch_log_processor_shutdown_with_async_runtime_current_flavor_current_thread() {
1162        let exporter = InMemoryLogExporterBuilder::default().build();
1163        let processor = BatchLogProcessor::new(exporter.clone(), BatchConfig::default());
1164        processor.shutdown().unwrap();
1165    }
1166
1167    #[tokio::test(flavor = "multi_thread")]
1168    async fn test_batch_log_processor_shutdown_with_async_runtime_multi_flavor_multi_thread() {
1169        let exporter = InMemoryLogExporterBuilder::default().build();
1170        let processor = BatchLogProcessor::new(exporter.clone(), BatchConfig::default());
1171        processor.shutdown().unwrap();
1172    }
1173
1174    #[tokio::test(flavor = "multi_thread")]
1175    async fn test_batch_log_processor_shutdown_with_async_runtime_multi_flavor_current_thread() {
1176        let exporter = InMemoryLogExporterBuilder::default().build();
1177        let processor = BatchLogProcessor::new(exporter.clone(), BatchConfig::default());
1178        processor.shutdown().unwrap();
1179    }
1180
1181    #[derive(Debug)]
1182    struct BlockingExporter {
1183        exported_count: Arc<AtomicUsize>,
1184        export_started: mpsc::SyncSender<()>,
1185        release: Arc<Mutex<mpsc::Receiver<()>>>,
1186    }
1187
1188    impl LogExporter for BlockingExporter {
1189        async fn export(&self, batch: LogBatch<'_>) -> OTelSdkResult {
1190            let _ = self.export_started.try_send(());
1191            // Block until the test releases the export.
1192            let _ = self.release.lock().unwrap().recv();
1193            self.exported_count.fetch_add(batch.len(), Ordering::SeqCst);
1194            Ok(())
1195        }
1196    }
1197
1198    #[test]
1199    fn test_batch_log_processor_emit_reverts_count_when_queue_full() {
1200        let (started_sender, started_receiver) = mpsc::sync_channel(8);
1201        let (release_sender, release_receiver) = mpsc::sync_channel(8);
1202        let exported_count = Arc::new(AtomicUsize::new(0));
1203        let exporter = BlockingExporter {
1204            exported_count: exported_count.clone(),
1205            export_started: started_sender,
1206            release: Arc::new(Mutex::new(release_receiver)),
1207        };
1208        let config = BatchConfigBuilder::default()
1209            .with_max_queue_size(4)
1210            .with_max_export_batch_size(4)
1211            .with_scheduled_delay(Duration::from_secs(60))
1212            .build();
1213        let processor = BatchLogProcessor::new(exporter, config);
1214        let instrumentation = InstrumentationScope::default();
1215        let emit = || {
1216            let mut record = SdkLogRecord::new();
1217            record.set_body("test log".into());
1218            processor.emit(&mut record, &instrumentation);
1219        };
1220
1221        // Fill the queue to the export threshold; the worker drains all four
1222        // records and blocks inside export().
1223        for _ in 0..4 {
1224            emit();
1225        }
1226        started_receiver
1227            .recv_timeout(Duration::from_secs(5))
1228            .expect("worker should start exporting the first batch");
1229
1230        // While the worker is blocked, refill the queue and overflow it by
1231        // two records, which must be dropped and their counts reverted.
1232        for _ in 0..6 {
1233            emit();
1234        }
1235
1236        assert_eq!(processor.dropped_logs_count.load(Ordering::Relaxed), 2);
1237        // 4 records in-flight in the blocked export (not yet subtracted)
1238        // plus 4 records queued. Without the queue-full revert this would
1239        // read 10.
1240        assert_eq!(
1241            processor.current_batch_size.load(Ordering::Relaxed),
1242            8,
1243            "dropped logs must not remain counted as pending"
1244        );
1245
1246        // Release the in-flight export and the one triggered by force_flush.
1247        release_sender.send(()).unwrap();
1248        release_sender.send(()).unwrap();
1249        let flush_result = processor.force_flush();
1250        assert!(flush_result.is_ok(), "force flush failed unexpectedly");
1251
1252        assert_eq!(
1253            exported_count.load(Ordering::SeqCst),
1254            8,
1255            "all logs that entered the queue must be exported"
1256        );
1257        assert_eq!(
1258            processor.current_batch_size.load(Ordering::Relaxed),
1259            0,
1260            "counter should settle to zero; a leftover value indicates the \
1261             queue-full path did not revert its increment"
1262        );
1263    }
1264
1265    /// A slow exporter that counts the number of logs received.
1266    /// Used for stress testing the BatchLogProcessor.
1267    #[derive(Debug, Clone)]
1268    struct CountingExporter {
1269        count: Arc<AtomicUsize>,
1270    }
1271
1272    impl CountingExporter {
1273        fn new() -> Self {
1274            CountingExporter {
1275                count: Arc::new(AtomicUsize::new(0)),
1276            }
1277        }
1278    }
1279
1280    impl LogExporter for CountingExporter {
1281        async fn export(&self, batch: LogBatch<'_>) -> OTelSdkResult {
1282            self.count.fetch_add(batch.len(), Ordering::SeqCst);
1283            // Simulate slow export to cause queue buildup and drops
1284            std::thread::sleep(std::time::Duration::from_millis(20));
1285            Ok(())
1286        }
1287    }
1288
1289    /// Stress test that verifies all logs are accounted for.
1290    /// With multiple threads pushing logs faster than the exporter can handle,
1291    /// some logs will inevitably be dropped. This test validates:
1292    /// total_logs_sent == logs_received_by_exporter + logs_dropped
1293    #[test]
1294    fn test_batch_log_processor_all_logs_accounted_for() {
1295        let exporter = CountingExporter::new();
1296        let exporter_count = exporter.count.clone();
1297
1298        // Configure with small queue to force drops under pressure
1299        let config = BatchConfigBuilder::default()
1300            .with_max_queue_size(2048)
1301            .with_max_export_batch_size(512)
1302            .with_scheduled_delay(Duration::from_millis(5))
1303            .build();
1304
1305        let processor = BatchLogProcessor::new(exporter, config);
1306
1307        let total_logs_per_thread = 100_000;
1308        let num_threads = 4;
1309        let total_logs_to_emit = total_logs_per_thread * num_threads;
1310
1311        // Use scoped threads to safely share the processor reference
1312        std::thread::scope(|s| {
1313            for _ in 0..num_threads {
1314                s.spawn(|| {
1315                    for _ in 0..total_logs_per_thread {
1316                        let mut record = SdkLogRecord::new();
1317                        record.set_body("stress test log".into());
1318                        let instrumentation = InstrumentationScope::default();
1319                        processor.emit(&mut record, &instrumentation);
1320                    }
1321                });
1322            }
1323        });
1324
1325        // Shutdown the processor to ensure all buffered logs are flushed
1326        processor.shutdown().unwrap();
1327
1328        let logs_received = exporter_count.load(Ordering::SeqCst);
1329        let logs_dropped = processor.dropped_logs_count.load(Ordering::SeqCst);
1330
1331        // The invariant: every log is either received or dropped
1332        assert_eq!(
1333            logs_received + logs_dropped,
1334            total_logs_to_emit,
1335            "Logs unaccounted for! Received: {}, Dropped: {}, Total emitted: {}",
1336            logs_received,
1337            logs_dropped,
1338            total_logs_to_emit
1339        );
1340
1341        // Also verify that some logs were actually dropped (stress test validation)
1342        // If no logs were dropped, the test parameters may need adjustment
1343        assert!(
1344            logs_dropped > 0,
1345            "Expected some logs to be dropped under stress, but none were. \
1346             Consider reducing queue size or increasing thread count/log volume."
1347        );
1348    }
1349
1350    #[cfg(feature = "experimental_metrics_bound_instruments")]
1351    mod self_obs {
1352        use super::*;
1353
1354        /// Verifies that `otel.sdk.processor.log.processed` counter records
1355        /// successful log processing when `experimental_metrics_bound_instruments`
1356        /// is enabled and a real MeterProvider is set as global before creating
1357        /// the processor.
1358        ///
1359        /// This test is `#[ignore]`d because it calls
1360        /// `global::set_meter_provider()` which mutates process-wide state.
1361        /// CI runs it in isolation via `test.sh`.
1362        #[cfg(feature = "experimental_metrics_bound_instruments")]
1363        #[test]
1364        #[ignore]
1365        fn self_diagnostics_counter_records_success() {
1366            use crate::metrics::data::{AggregatedMetrics, MetricData};
1367            use crate::metrics::{InMemoryMetricExporter, SdkMeterProvider};
1368
1369            // Setup a real MeterProvider and set it as global BEFORE creating the
1370            // BatchLogProcessor, so the processor picks up a real meter.
1371            let metric_exporter = InMemoryMetricExporter::default();
1372            let meter_provider = SdkMeterProvider::builder()
1373                .with_periodic_exporter(metric_exporter.clone())
1374                .build();
1375            opentelemetry::global::set_meter_provider(meter_provider.clone());
1376
1377            let log_exporter = InMemoryLogExporter::default();
1378            let config = BatchConfigBuilder::default()
1379                .with_max_queue_size(256)
1380                .with_max_export_batch_size(64)
1381                .with_scheduled_delay(Duration::from_secs(60))
1382                .build();
1383            let processor = BatchLogProcessor::new(log_exporter, config);
1384
1385            // Emit 10 logs
1386            let instrumentation = InstrumentationScope::default();
1387            for _ in 0..10 {
1388                let mut record = SdkLogRecord::new();
1389                processor.emit(&mut record, &instrumentation);
1390            }
1391
1392            // Flush so the batch is submitted to the exporter, which is when the
1393            // counter is incremented.
1394            processor.force_flush().unwrap();
1395
1396            // Force a metrics collection
1397            meter_provider.force_flush().unwrap();
1398
1399            // Find the otel.sdk.processor.log.processed metric and sum all data points
1400            let metrics = metric_exporter.get_finished_metrics().unwrap();
1401            let mut found = false;
1402            let mut total_value: u64 = 0;
1403            for rm in &metrics {
1404                for sm in &rm.scope_metrics {
1405                    for metric in &sm.metrics {
1406                        if metric.name == "otel.sdk.processor.log.processed" {
1407                            found = true;
1408                            if let AggregatedMetrics::U64(MetricData::Sum(sum)) = &metric.data {
1409                                for dp in sum.data_points() {
1410                                    total_value += dp.value();
1411                                }
1412                            }
1413                        }
1414                    }
1415                }
1416            }
1417
1418            assert!(found, "otel.sdk.processor.log.processed metric not found");
1419            assert_eq!(
1420                total_value, 10,
1421                "Expected 10 processed logs, got {total_value}"
1422            );
1423
1424            processor.shutdown().unwrap();
1425            meter_provider.shutdown().unwrap();
1426        }
1427
1428        /// Verifies `otel.sdk.processor.log.queue.capacity` through a real
1429        /// `SdkLoggerProvider` + `BatchLogProcessor`. The metric reports the
1430        /// configured max queue size with the component identity attributes and
1431        /// stops reporting after the provider (and processor) is dropped.
1432        ///
1433        /// `#[ignore]`d because it mutates process-wide state via
1434        /// `global::set_meter_provider()`. CI runs it in isolation via `test.sh`.
1435        #[cfg(feature = "experimental_metrics_bound_instruments")]
1436        #[test]
1437        #[ignore]
1438        fn self_diagnostics_queue_capacity() {
1439            use crate::metrics::data::{AggregatedMetrics, MetricData};
1440            use crate::metrics::{InMemoryMetricExporter, SdkMeterProvider};
1441
1442            let metric_exporter = InMemoryMetricExporter::default();
1443            let meter_provider = SdkMeterProvider::builder()
1444                .with_periodic_exporter(metric_exporter.clone())
1445                .build();
1446            opentelemetry::global::set_meter_provider(meter_provider.clone());
1447
1448            let log_exporter = InMemoryLogExporter::default();
1449            let config = BatchConfigBuilder::default()
1450                .with_max_queue_size(256)
1451                .build();
1452            let processor = BatchLogProcessor::new(log_exporter, config);
1453            let provider = SdkLoggerProvider::builder()
1454                .with_log_processor(processor)
1455                .build();
1456
1457            // Force a metrics collection so the observable callbacks run. This does
1458            // NOT drain the log queue (that only happens on the provider/processor).
1459            meter_provider.force_flush().unwrap();
1460
1461            let read = |name: &str| -> Option<i64> {
1462                let metrics = metric_exporter.get_finished_metrics().unwrap();
1463                for rm in &metrics {
1464                    for sm in &rm.scope_metrics {
1465                        for metric in &sm.metrics {
1466                            if metric.name == name {
1467                                if let AggregatedMetrics::I64(MetricData::Sum(sum)) = &metric.data {
1468                                    for dp in sum.data_points() {
1469                                        let has_component = dp.attributes().any(|kv| {
1470                                            kv.key.as_str() == "otel.component.type"
1471                                                && kv.value.as_str() == "batching_log_processor"
1472                                        });
1473                                        if has_component {
1474                                            return Some(dp.value());
1475                                        }
1476                                    }
1477                                }
1478                            }
1479                        }
1480                    }
1481                }
1482                None
1483            };
1484
1485            assert_eq!(
1486                read("otel.sdk.processor.log.queue.capacity"),
1487                Some(256),
1488                "queue.capacity should equal the configured max_queue_size"
1489            );
1490
1491            // Dropping the provider shuts down and drops the processor, releasing the
1492            // Arc<AtomicUsize> the callback holds a Weak to. A subsequent collection
1493            // must therefore omit the metric because the Weak upgrade fails.
1494            metric_exporter.reset();
1495            provider.shutdown().unwrap();
1496            drop(provider);
1497            meter_provider.force_flush().unwrap();
1498
1499            assert_eq!(
1500                read("otel.sdk.processor.log.queue.capacity"),
1501                None,
1502                "queue.capacity must stop being reported after the processor is dropped"
1503            );
1504
1505            meter_provider.shutdown().unwrap();
1506        }
1507
1508        /// Sums the values of `otel.sdk.processor.log.processed` data points whose
1509        /// `error.type` attribute equals `error_type`.
1510        #[cfg(feature = "experimental_metrics_bound_instruments")]
1511        fn sum_processed_log_records_with_error_type(
1512            metric_exporter: &crate::metrics::InMemoryMetricExporter,
1513            error_type: &str,
1514        ) -> u64 {
1515            use crate::metrics::data::{AggregatedMetrics, MetricData};
1516
1517            let metrics = metric_exporter.get_finished_metrics().unwrap();
1518            let mut total: u64 = 0;
1519            for rm in &metrics {
1520                for sm in &rm.scope_metrics {
1521                    for metric in &sm.metrics {
1522                        if metric.name == "otel.sdk.processor.log.processed" {
1523                            if let AggregatedMetrics::U64(MetricData::Sum(sum)) = &metric.data {
1524                                for dp in sum.data_points() {
1525                                    let matches = dp.attributes().any(|kv| {
1526                                        kv.key.as_str() == "error.type"
1527                                            && kv.value.as_str() == error_type
1528                                    });
1529                                    if matches {
1530                                        total += dp.value();
1531                                    }
1532                                }
1533                            }
1534                        }
1535                    }
1536                }
1537            }
1538            total
1539        }
1540
1541        /// Verifies that `otel.sdk.processor.log.processed` records queue-full drops
1542        /// with `error.type = queue_full` when records overflow the queue while the
1543        /// worker is blocked exporting.
1544        ///
1545        /// `#[ignore]`d because it mutates process-wide state via
1546        /// `global::set_meter_provider()`. CI runs it in isolation via `test.sh`.
1547        #[cfg(feature = "experimental_metrics_bound_instruments")]
1548        #[test]
1549        #[ignore]
1550        fn self_diagnostics_counter_records_queue_full_drops() {
1551            use crate::metrics::{InMemoryMetricExporter, SdkMeterProvider};
1552
1553            let metric_exporter = InMemoryMetricExporter::default();
1554            let meter_provider = SdkMeterProvider::builder()
1555                .with_periodic_exporter(metric_exporter.clone())
1556                .build();
1557            opentelemetry::global::set_meter_provider(meter_provider.clone());
1558
1559            let (started_sender, started_receiver) = mpsc::sync_channel(8);
1560            let (release_sender, release_receiver) = mpsc::sync_channel(8);
1561            let exported_count = Arc::new(AtomicUsize::new(0));
1562            let exporter = BlockingExporter {
1563                exported_count: exported_count.clone(),
1564                export_started: started_sender,
1565                release: Arc::new(Mutex::new(release_receiver)),
1566            };
1567            let config = BatchConfigBuilder::default()
1568                .with_max_queue_size(4)
1569                .with_max_export_batch_size(4)
1570                .with_scheduled_delay(Duration::from_secs(60))
1571                .build();
1572            let processor = BatchLogProcessor::new(exporter, config);
1573            let instrumentation = InstrumentationScope::default();
1574            let emit = || {
1575                let mut record = SdkLogRecord::new();
1576                processor.emit(&mut record, &instrumentation);
1577            };
1578
1579            // Fill the queue to the export threshold; the worker drains all four
1580            // records and blocks inside export().
1581            for _ in 0..4 {
1582                emit();
1583            }
1584            started_receiver
1585                .recv_timeout(Duration::from_secs(5))
1586                .expect("worker should start exporting the first batch");
1587
1588            // While the worker is blocked, refill the queue (4) and overflow it by
1589            // two records, which must be dropped and counted as queue_full.
1590            for _ in 0..6 {
1591                emit();
1592            }
1593
1594            // Release the in-flight export and the one triggered by force_flush.
1595            release_sender.send(()).unwrap();
1596            release_sender.send(()).unwrap();
1597            processor.force_flush().unwrap();
1598
1599            meter_provider.force_flush().unwrap();
1600
1601            let queue_full =
1602                sum_processed_log_records_with_error_type(&metric_exporter, "queue_full");
1603            assert_eq!(
1604                queue_full, 2,
1605                "expected 2 queue_full drops, got {queue_full}"
1606            );
1607
1608            processor.shutdown().unwrap();
1609            meter_provider.shutdown().unwrap();
1610        }
1611
1612        /// Verifies that `otel.sdk.processor.log.processed` records post-shutdown
1613        /// emits with `error.type = already_shutdown`.
1614        ///
1615        /// `#[ignore]`d because it mutates process-wide state via
1616        /// `global::set_meter_provider()`. CI runs it in isolation via `test.sh`.
1617        #[cfg(feature = "experimental_metrics_bound_instruments")]
1618        #[test]
1619        #[ignore]
1620        fn self_diagnostics_counter_records_already_shutdown_drops() {
1621            use crate::metrics::{InMemoryMetricExporter, SdkMeterProvider};
1622
1623            let metric_exporter = InMemoryMetricExporter::default();
1624            let meter_provider = SdkMeterProvider::builder()
1625                .with_periodic_exporter(metric_exporter.clone())
1626                .build();
1627            opentelemetry::global::set_meter_provider(meter_provider.clone());
1628
1629            let log_exporter = InMemoryLogExporter::default();
1630            let processor = BatchLogProcessor::new(log_exporter, BatchConfig::default());
1631
1632            // Shut the processor down so the worker thread (the only receiver)
1633            // disconnects; subsequent emits hit the already_shutdown branch.
1634            processor.shutdown().unwrap();
1635
1636            let instrumentation = InstrumentationScope::default();
1637            for _ in 0..7 {
1638                let mut record = SdkLogRecord::new();
1639                processor.emit(&mut record, &instrumentation);
1640            }
1641
1642            meter_provider.force_flush().unwrap();
1643
1644            let already_shutdown =
1645                sum_processed_log_records_with_error_type(&metric_exporter, "already_shutdown");
1646            assert_eq!(
1647                already_shutdown, 7,
1648                "expected 7 already_shutdown drops, got {already_shutdown}"
1649            );
1650
1651            meter_provider.shutdown().unwrap();
1652        }
1653    }
1654}