1use crate::error::{OTelSdkError, OTelSdkResult};
20use crate::logs::log_processor::LogProcessor;
21use crate::{
22 logs::{LogBatch, LogExporter, SdkLogRecord},
23 Resource,
24};
25
26#[cfg(feature = "experimental_metrics_bound_instruments")]
27use opentelemetry::KeyValue;
28use opentelemetry::{otel_debug, otel_error, otel_warn, Context, InstrumentationScope};
29
30use std::fmt::Debug;
31use std::sync::atomic::AtomicBool;
32#[cfg(feature = "experimental_metrics_bound_instruments")]
33use std::sync::atomic::AtomicUsize;
34use std::sync::Mutex;
35use std::time::Duration;
36
37#[derive(Debug)]
69pub struct SimpleLogProcessor<T: LogExporter> {
70 exporter: Mutex<T>,
71 is_shutdown: AtomicBool,
72
73 #[cfg(feature = "experimental_metrics_bound_instruments")]
80 processed_success: opentelemetry::metrics::BoundCounter<u64>,
81 #[cfg(feature = "experimental_metrics_bound_instruments")]
82 processed_after_shutdown: opentelemetry::metrics::BoundCounter<u64>,
83}
84
85impl<T: LogExporter> SimpleLogProcessor<T> {
86 pub fn new(exporter: T) -> Self {
88 #[cfg(feature = "experimental_metrics_bound_instruments")]
89 let (processed_success, processed_after_shutdown) = {
90 static INSTANCE_COUNTER: AtomicUsize = AtomicUsize::new(0);
91 let instance_id = INSTANCE_COUNTER.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
92 let component_name = format!("simple_log_processor/{instance_id}");
93
94 let meter = opentelemetry::global::meter("otel.sdk");
95 let counter = meter
96 .u64_counter("otel.sdk.processor.log.processed")
97 .with_description(
98 "The number of log records for which the processing has finished, \
99 either successful or failed.",
100 )
101 .with_unit("{log_record}")
102 .build();
103
104 let success_attrs = [
108 KeyValue::new("otel.component.type", "simple_log_processor"),
109 KeyValue::new("otel.component.name", component_name.clone()),
110 ];
111 let after_shutdown_attrs = [
112 KeyValue::new("error.type", "already_shutdown"),
113 KeyValue::new("otel.component.type", "simple_log_processor"),
114 KeyValue::new("otel.component.name", component_name),
115 ];
116
117 (
118 counter.bind(&success_attrs),
119 counter.bind(&after_shutdown_attrs),
120 )
121 };
122
123 SimpleLogProcessor {
124 exporter: Mutex::new(exporter),
125 is_shutdown: AtomicBool::new(false),
126 #[cfg(feature = "experimental_metrics_bound_instruments")]
127 processed_success,
128 #[cfg(feature = "experimental_metrics_bound_instruments")]
129 processed_after_shutdown,
130 }
131 }
132}
133
134impl<T: LogExporter> LogProcessor for SimpleLogProcessor<T> {
135 fn emit(&self, record: &mut SdkLogRecord, instrumentation: &InstrumentationScope) {
136 let _suppress_guard = Context::enter_telemetry_suppressed_scope();
137 if self.is_shutdown.load(std::sync::atomic::Ordering::Relaxed) {
139 #[cfg(feature = "experimental_metrics_bound_instruments")]
141 self.processed_after_shutdown.add(1);
142
143 otel_warn!(
145 name: "SimpleLogProcessor.Emit.ProcessorShutdown",
146 );
147 return;
148 }
149
150 let result = match self.exporter.lock() {
151 Ok(exporter) => {
152 let log_tuple = &[(record as &SdkLogRecord, instrumentation)];
153 #[cfg(feature = "experimental_metrics_bound_instruments")]
157 self.processed_success.add(1);
158 futures_executor::block_on(exporter.export(LogBatch::new(log_tuple)))
159 }
160 Err(_) => Err(OTelSdkError::InternalFailure(
161 "SimpleLogProcessor mutex poison".into(),
162 )),
163 };
164 match result {
166 Err(OTelSdkError::InternalFailure(_)) => {
167 otel_debug!(
169 name: "SimpleLogProcessor.Emit.MutexPoisoning",
170 );
171 }
172 Err(err) => {
173 otel_error!(
174 name: "SimpleLogProcessor.Emit.ExportError",
175 error = format!("{}",err)
176 );
177 }
178 _ => {}
179 }
180 }
181
182 fn force_flush(&self) -> OTelSdkResult {
183 Ok(())
184 }
185
186 fn shutdown_with_timeout(&self, timeout: Duration) -> OTelSdkResult {
187 self.is_shutdown
188 .store(true, std::sync::atomic::Ordering::Relaxed);
189 if let Ok(exporter) = self.exporter.lock() {
190 exporter.shutdown_with_timeout(timeout)
191 } else {
192 Err(OTelSdkError::InternalFailure(
193 "SimpleLogProcessor mutex poison at shutdown".into(),
194 ))
195 }
196 }
197
198 fn set_resource(&mut self, resource: &Resource) {
199 if let Ok(mut exporter) = self.exporter.lock() {
200 exporter.set_resource(resource);
201 }
202 }
203
204 #[inline]
205 fn event_enabled(
206 &self,
207 level: opentelemetry::logs::Severity,
208 target: &str,
209 name: Option<&str>,
210 ) -> bool {
211 if let Ok(exporter) = self.exporter.lock() {
212 exporter.event_enabled(level, target, name)
213 } else {
214 true
215 }
216 }
217}
218
219#[cfg(all(test, feature = "testing", feature = "logs"))]
220mod tests {
221 use crate::logs::log_processor::tests::MockLogExporter;
222 use crate::logs::{LogBatch, LogExporter, SdkLogRecord, SdkLogger};
223 use crate::{
224 error::OTelSdkResult,
225 logs::{InMemoryLogExporterBuilder, LogProcessor, SdkLoggerProvider, SimpleLogProcessor},
226 Resource,
227 };
228 use opentelemetry::logs::{LogRecord, Logger, LoggerProvider};
229 use opentelemetry::InstrumentationScope;
230 use opentelemetry::KeyValue;
231 use std::sync::atomic::{AtomicUsize, Ordering};
232 use std::sync::{Arc, Mutex};
233 use std::time;
234 use std::time::Duration;
235
236 #[derive(Debug, Clone)]
237 struct LogExporterThatRequiresTokio {
238 export_count: Arc<AtomicUsize>,
239 }
240
241 impl LogExporterThatRequiresTokio {
242 fn new() -> Self {
244 LogExporterThatRequiresTokio {
245 export_count: Arc::new(AtomicUsize::new(0)),
246 }
247 }
248
249 fn len(&self) -> usize {
251 self.export_count.load(Ordering::Acquire)
252 }
253 }
254
255 impl LogExporter for LogExporterThatRequiresTokio {
256 async fn export(&self, batch: LogBatch<'_>) -> OTelSdkResult {
257 tokio::time::sleep(Duration::from_millis(50)).await;
259
260 for _ in batch.iter() {
261 self.export_count.fetch_add(1, Ordering::Acquire);
262 }
263 Ok(())
264 }
265 fn shutdown_with_timeout(&self, _timeout: time::Duration) -> OTelSdkResult {
266 Ok(())
267 }
268 }
269
270 #[test]
271 fn test_set_resource_simple_processor() {
272 let exporter = MockLogExporter {
273 resource: Arc::new(Mutex::new(None)),
274 };
275 let processor = SimpleLogProcessor::new(exporter.clone());
276 let _ = SdkLoggerProvider::builder()
277 .with_log_processor(processor)
278 .with_resource(
279 Resource::builder_empty()
280 .with_attributes([
281 KeyValue::new("k1", "v1"),
282 KeyValue::new("k2", "v3"),
283 KeyValue::new("k3", "v3"),
284 KeyValue::new("k4", "v4"),
285 KeyValue::new("k5", "v5"),
286 ])
287 .build(),
288 )
289 .build();
290 assert_eq!(exporter.get_resource().unwrap().into_iter().count(), 5);
291 }
292
293 #[test]
294 fn test_simple_shutdown() {
295 let exporter = InMemoryLogExporterBuilder::default()
296 .keep_records_on_shutdown()
297 .build();
298 let processor = SimpleLogProcessor::new(exporter.clone());
299
300 let mut record: SdkLogRecord = SdkLogRecord::new();
301 let instrumentation: InstrumentationScope = Default::default();
302
303 processor.emit(&mut record, &instrumentation);
304
305 processor.shutdown().unwrap();
306
307 let is_shutdown = processor
308 .is_shutdown
309 .load(std::sync::atomic::Ordering::Relaxed);
310 assert!(is_shutdown);
311
312 processor.emit(&mut record, &instrumentation);
313
314 assert_eq!(1, exporter.get_emitted_logs().unwrap().len());
315 assert!(exporter.is_shutdown_called());
316 }
317
318 #[test]
319 fn test_simple_processor_sync_exporter_without_runtime() {
320 let exporter = InMemoryLogExporterBuilder::default().build();
321 let processor = SimpleLogProcessor::new(exporter.clone());
322
323 let mut record: SdkLogRecord = SdkLogRecord::new();
324 let instrumentation: InstrumentationScope = Default::default();
325
326 processor.emit(&mut record, &instrumentation);
327
328 assert_eq!(exporter.get_emitted_logs().unwrap().len(), 1);
329 }
330
331 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
332 async fn test_simple_processor_sync_exporter_with_runtime() {
333 let exporter = InMemoryLogExporterBuilder::default().build();
334 let processor = SimpleLogProcessor::new(exporter.clone());
335
336 let mut record: SdkLogRecord = SdkLogRecord::new();
337 let instrumentation: InstrumentationScope = Default::default();
338
339 processor.emit(&mut record, &instrumentation);
340
341 assert_eq!(exporter.get_emitted_logs().unwrap().len(), 1);
342 }
343
344 #[tokio::test(flavor = "multi_thread")]
345 async fn test_simple_processor_sync_exporter_with_multi_thread_runtime() {
346 let exporter = InMemoryLogExporterBuilder::default().build();
347 let processor = Arc::new(SimpleLogProcessor::new(exporter.clone()));
348
349 let mut handles = vec![];
350 for _ in 0..10 {
351 let processor_clone = Arc::clone(&processor);
352 let handle = tokio::spawn(async move {
353 let mut record: SdkLogRecord = SdkLogRecord::new();
354 let instrumentation: InstrumentationScope = Default::default();
355 processor_clone.emit(&mut record, &instrumentation);
356 });
357 handles.push(handle);
358 }
359
360 for handle in handles {
361 handle.await.unwrap();
362 }
363
364 assert_eq!(exporter.get_emitted_logs().unwrap().len(), 10);
365 }
366
367 #[tokio::test(flavor = "current_thread")]
368 async fn test_simple_processor_sync_exporter_with_current_thread_runtime() {
369 let exporter = InMemoryLogExporterBuilder::default().build();
370 let processor = SimpleLogProcessor::new(exporter.clone());
371
372 let mut record: SdkLogRecord = SdkLogRecord::new();
373 let instrumentation: InstrumentationScope = Default::default();
374
375 processor.emit(&mut record, &instrumentation);
376
377 assert_eq!(exporter.get_emitted_logs().unwrap().len(), 1);
378 }
379
380 #[test]
381 fn test_simple_processor_async_exporter_without_runtime() {
382 let result = std::panic::catch_unwind(|| {
384 let exporter = LogExporterThatRequiresTokio::new();
385 let processor = SimpleLogProcessor::new(exporter.clone());
386
387 let mut record: SdkLogRecord = SdkLogRecord::new();
388 let instrumentation: InstrumentationScope = Default::default();
389
390 processor.emit(&mut record, &instrumentation);
392 });
393
394 assert!(
396 result.is_err(),
397 "The test should fail due to missing Tokio runtime, but it did not."
398 );
399 let panic_payload = result.unwrap_err();
400 let panic_message = panic_payload
401 .downcast_ref::<String>()
402 .map(|s| s.as_str())
403 .or_else(|| panic_payload.downcast_ref::<&str>().copied())
404 .unwrap_or("No panic message");
405
406 assert!(
407 panic_message.contains("no reactor running")
408 || panic_message.contains("must be called from the context of a Tokio 1.x runtime"),
409 "Expected panic message about missing Tokio runtime, but got: {panic_message}"
410 );
411 }
412
413 #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
414 #[ignore]
415 async fn test_simple_processor_async_exporter_with_all_runtime_worker_threads_blocked() {
432 let exporter = LogExporterThatRequiresTokio::new();
433 let processor = Arc::new(SimpleLogProcessor::new(exporter.clone()));
434
435 let concurrent_emit = 4; let mut handles = vec![];
438 for _ in 0..concurrent_emit {
440 let processor_clone = Arc::clone(&processor);
441 let handle = tokio::spawn(async move {
442 let mut record: SdkLogRecord = SdkLogRecord::new();
443 let instrumentation: InstrumentationScope = Default::default();
444 processor_clone.emit(&mut record, &instrumentation);
445 });
446 handles.push(handle);
447 }
448
449 for handle in handles {
451 handle.await.unwrap();
452 }
453 assert_eq!(exporter.len(), concurrent_emit);
454 }
455
456 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
457 async fn test_simple_processor_async_exporter_with_runtime() {
463 let exporter = LogExporterThatRequiresTokio::new();
464 let processor = SimpleLogProcessor::new(exporter.clone());
465
466 let mut record: SdkLogRecord = SdkLogRecord::new();
467 let instrumentation: InstrumentationScope = Default::default();
468
469 processor.emit(&mut record, &instrumentation);
470
471 assert_eq!(exporter.len(), 1);
472 }
473
474 #[tokio::test(flavor = "multi_thread")]
475 async fn test_simple_processor_async_exporter_with_multi_thread_runtime() {
481 let exporter = LogExporterThatRequiresTokio::new();
482
483 let processor = SimpleLogProcessor::new(exporter.clone());
484
485 let mut record: SdkLogRecord = SdkLogRecord::new();
486 let instrumentation: InstrumentationScope = Default::default();
487
488 processor.emit(&mut record, &instrumentation);
489
490 assert_eq!(exporter.len(), 1);
491 }
492
493 #[tokio::test(flavor = "current_thread")]
494 #[ignore]
495 async fn test_simple_processor_async_exporter_with_current_thread_runtime() {
501 let exporter = LogExporterThatRequiresTokio::new();
502
503 let processor = SimpleLogProcessor::new(exporter.clone());
504
505 let mut record: SdkLogRecord = SdkLogRecord::new();
506 let instrumentation: InstrumentationScope = Default::default();
507
508 processor.emit(&mut record, &instrumentation);
509
510 assert_eq!(exporter.len(), 1);
511 }
512
513 #[derive(Debug, Clone)]
514 struct ReentrantLogExporter {
515 logger: Arc<Mutex<Option<SdkLogger>>>,
516 }
517
518 impl ReentrantLogExporter {
519 fn new() -> Self {
520 Self {
521 logger: Arc::new(Mutex::new(None)),
522 }
523 }
524
525 fn set_logger(&self, logger: SdkLogger) {
526 let mut guard = self.logger.lock().unwrap();
527 *guard = Some(logger);
528 }
529 }
530
531 impl LogExporter for ReentrantLogExporter {
532 async fn export(&self, _batch: LogBatch<'_>) -> OTelSdkResult {
533 let logger = self.logger.lock().unwrap();
534 if let Some(logger) = logger.as_ref() {
535 let mut log_record = logger.create_log_record();
536 log_record.set_severity_number(opentelemetry::logs::Severity::Error);
537 logger.emit(log_record);
538 }
539
540 Ok(())
541 }
542 }
543
544 #[test]
545 fn exporter_internal_log_does_not_deadlock_with_simple_processor() {
546 let exporter: ReentrantLogExporter = ReentrantLogExporter::new();
550 let logger_provider = SdkLoggerProvider::builder()
551 .with_simple_exporter(exporter.clone())
552 .build();
553 exporter.set_logger(logger_provider.logger("processor-logger"));
554
555 let logger = logger_provider.logger("test-logger");
556 let mut log_record = logger.create_log_record();
557 log_record.set_severity_number(opentelemetry::logs::Severity::Error);
558 logger.emit(log_record);
559 }
560
561 #[cfg(feature = "experimental_metrics_bound_instruments")]
562 mod self_obs {
563 use super::*;
564
565 #[cfg(feature = "experimental_metrics_bound_instruments")]
569 fn sum_processed_log_records(
570 metric_exporter: &crate::metrics::InMemoryMetricExporter,
571 error_type: Option<&str>,
572 ) -> u64 {
573 use crate::metrics::data::{AggregatedMetrics, MetricData};
574
575 let metrics = metric_exporter.get_finished_metrics().unwrap();
576 let mut total: u64 = 0;
577 for rm in &metrics {
578 for sm in &rm.scope_metrics {
579 for metric in &sm.metrics {
580 if metric.name == "otel.sdk.processor.log.processed" {
581 if let AggregatedMetrics::U64(MetricData::Sum(sum)) = &metric.data {
582 for dp in sum.data_points() {
583 let dp_error_type = dp
584 .attributes()
585 .find(|kv| kv.key.as_str() == "error.type")
586 .map(|kv| kv.value.as_str().to_string());
587 let matches = match error_type {
588 Some(expected) => {
589 dp_error_type.as_deref() == Some(expected)
590 }
591 None => dp_error_type.is_none(),
592 };
593 if matches {
594 total += dp.value();
595 }
596 }
597 }
598 }
599 }
600 }
601 }
602 total
603 }
604
605 #[cfg(feature = "experimental_metrics_bound_instruments")]
612 #[test]
613 #[ignore]
614 fn self_diagnostics_counter_records_success() {
615 use crate::logs::InMemoryLogExporterBuilder;
616 use crate::metrics::{InMemoryMetricExporter, SdkMeterProvider};
617
618 let metric_exporter = InMemoryMetricExporter::default();
619 let meter_provider = SdkMeterProvider::builder()
620 .with_periodic_exporter(metric_exporter.clone())
621 .build();
622 opentelemetry::global::set_meter_provider(meter_provider.clone());
623
624 let log_exporter = InMemoryLogExporterBuilder::default().build();
625 let processor = SimpleLogProcessor::new(log_exporter);
626
627 let instrumentation = InstrumentationScope::default();
628 for _ in 0..10 {
629 let mut record = SdkLogRecord::new();
630 processor.emit(&mut record, &instrumentation);
631 }
632
633 meter_provider.force_flush().unwrap();
634
635 let processed = sum_processed_log_records(&metric_exporter, None);
636 assert_eq!(processed, 10, "expected 10 processed logs, got {processed}");
637
638 meter_provider.shutdown().unwrap();
639 }
640
641 #[cfg(feature = "experimental_metrics_bound_instruments")]
647 #[test]
648 #[ignore]
649 fn self_diagnostics_counter_records_already_shutdown_drops() {
650 use crate::logs::InMemoryLogExporterBuilder;
651 use crate::metrics::{InMemoryMetricExporter, SdkMeterProvider};
652
653 let metric_exporter = InMemoryMetricExporter::default();
654 let meter_provider = SdkMeterProvider::builder()
655 .with_periodic_exporter(metric_exporter.clone())
656 .build();
657 opentelemetry::global::set_meter_provider(meter_provider.clone());
658
659 let log_exporter = InMemoryLogExporterBuilder::default().build();
660 let processor = SimpleLogProcessor::new(log_exporter);
661
662 processor.shutdown().unwrap();
664
665 let instrumentation = InstrumentationScope::default();
666 for _ in 0..7 {
667 let mut record = SdkLogRecord::new();
668 processor.emit(&mut record, &instrumentation);
669 }
670
671 meter_provider.force_flush().unwrap();
672
673 let already_shutdown =
674 sum_processed_log_records(&metric_exporter, Some("already_shutdown"));
675 assert_eq!(
676 already_shutdown, 7,
677 "expected 7 already_shutdown drops, got {already_shutdown}"
678 );
679 let success = sum_processed_log_records(&metric_exporter, None);
680 assert_eq!(
681 success, 0,
682 "post-shutdown emits must not be counted as success, got {success}"
683 );
684
685 meter_provider.shutdown().unwrap();
686 }
687 }
688}