1use 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
42pub const OTEL_BLRP_SCHEDULE_DELAY: &str = "OTEL_BLRP_SCHEDULE_DELAY";
45pub const OTEL_BLRP_SCHEDULE_DELAY_DEFAULT: Duration = Duration::from_millis(1_000);
47pub const OTEL_BLRP_EXPORT_TIMEOUT: &str = "OTEL_BLRP_EXPORT_TIMEOUT";
54pub const OTEL_BLRP_EXPORT_TIMEOUT_DEFAULT: Duration = Duration::from_millis(30_000);
58pub const OTEL_BLRP_MAX_QUEUE_SIZE: &str = "OTEL_BLRP_MAX_QUEUE_SIZE";
61pub const OTEL_BLRP_MAX_QUEUE_SIZE_DEFAULT: usize = 2_048;
63pub const OTEL_BLRP_MAX_EXPORT_BATCH_SIZE: &str = "OTEL_BLRP_MAX_EXPORT_BATCH_SIZE";
67pub const OTEL_BLRP_MAX_EXPORT_BATCH_SIZE_DEFAULT: usize = 512;
69
70#[allow(clippy::large_enum_variant)]
72#[derive(Debug)]
73enum BatchMessage {
74 ExportLog(Arc<AtomicBool>),
76 ForceFlush(mpsc::SyncSender<OTelSdkResult>),
78 Shutdown(mpsc::SyncSender<OTelSdkResult>),
80 SetResource(Arc<Resource>),
82}
83
84type LogsData = Box<(SdkLogRecord, InstrumentationScope)>;
85
86pub struct BatchLogProcessor {
149 logs_sender: SyncSender<LogsData>, message_sender: SyncSender<BatchMessage>, 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 dropped_logs_count: AtomicUsize,
159
160 max_queue_size: usize,
162
163 #[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 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 result {
194 Ok(_) => {
195 if previous_batch_size + 1 >= self.max_export_batch_size {
202 if !self.export_log_message_sent.load(Ordering::Relaxed) {
209 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 }
225 Err(_err) => {
226 self.export_log_message_sent.store(false, Ordering::Relaxed);
230 }
231 }
232 }
233 }
234 }
235 }
236 Err(mpsc::TrySendError::Full(_)) => {
237 self.current_batch_size.fetch_sub(1, Ordering::AcqRel);
239 #[cfg(feature = "experimental_metrics_bound_instruments")]
241 self.processed_queue_full.add(1);
242
243 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 self.current_batch_size.fetch_sub(1, Ordering::AcqRel);
253 #[cfg(feature = "experimental_metrics_bound_instruments")]
255 self.processed_after_shutdown.add(1);
256
257 let _guard = Context::enter_telemetry_suppressed_scope();
260
261 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 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 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 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 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 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 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); 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 #[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 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(¤t_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 if capacity_state.strong_count() > 0 {
439 observer.observe(capacity_value, &capacity_attrs);
440 }
441 })
442 .build();
443
444 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 #[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 #[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); 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 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(); if count_of_logs == 0 {
524 break;
525 }
526 total_exported_logs += count_of_logs;
527
528 record_processed_success(count_of_logs as u64);
531
532 result = export_batch_sync(exporter, logs, last_export_time); 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 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 ¤t_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 ¤t_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 ¤t_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 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 ¤t_batch_size,
614 max_export_batch_size,
615 &record_processed_success,
616 );
617 }
618 Err(RecvTimeoutError::Disconnected) => {
619 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."); BatchLogProcessor {
637 logs_sender,
638 message_sender,
639 handle: Mutex::new(Some(handle)),
640 forceflush_timeout: Duration::from_secs(5), 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 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 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#[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 pub fn with_batch_config(self, config: BatchConfig) -> Self {
713 BatchLogProcessorBuilder { config, ..self }
714 }
715
716 pub fn build(self) -> BatchLogProcessor {
718 BatchLogProcessor::new(self.exporter, self.config)
719 }
720}
721
722#[derive(Debug)]
725#[allow(dead_code)]
726pub struct BatchConfig {
727 pub(crate) max_queue_size: usize,
730
731 pub(crate) scheduled_delay: Duration,
734
735 pub(crate) max_export_batch_size: usize,
740
741 #[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#[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 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 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 pub fn with_scheduled_delay(mut self, scheduled_delay: Duration) -> Self {
806 self.scheduled_delay = scheduled_delay;
807 self
808 }
809
810 #[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 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 pub fn build(self) -> BatchConfig {
840 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 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 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 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 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 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 for _ in 0..6 {
1233 emit();
1234 }
1235
1236 assert_eq!(processor.dropped_logs_count.load(Ordering::Relaxed), 2);
1237 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_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 #[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 std::thread::sleep(std::time::Duration::from_millis(20));
1285 Ok(())
1286 }
1287 }
1288
1289 #[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 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 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 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 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 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 #[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 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 let instrumentation = InstrumentationScope::default();
1387 for _ in 0..10 {
1388 let mut record = SdkLogRecord::new();
1389 processor.emit(&mut record, &instrumentation);
1390 }
1391
1392 processor.force_flush().unwrap();
1395
1396 meter_provider.force_flush().unwrap();
1398
1399 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 #[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 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 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 #[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 #[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 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 for _ in 0..6 {
1591 emit();
1592 }
1593
1594 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 #[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 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}