Skip to main content

opentelemetry_sdk/logs/
mod.rs

1//! # OpenTelemetry Log SDK
2mod batch_log_processor;
3mod export;
4mod log_processor;
5mod logger;
6mod logger_provider;
7pub(crate) mod record;
8mod simple_log_processor;
9
10/// In-Memory log exporter for testing purpose.
11#[cfg(any(feature = "testing", test))]
12#[cfg_attr(docsrs, doc(cfg(any(feature = "testing", test))))]
13pub mod in_memory_exporter;
14#[cfg(any(feature = "testing", test))]
15#[cfg_attr(docsrs, doc(cfg(any(feature = "testing", test))))]
16pub use in_memory_exporter::{InMemoryLogExporter, InMemoryLogExporterBuilder};
17
18pub use batch_log_processor::{
19    BatchConfig, BatchConfigBuilder, BatchLogProcessor, BatchLogProcessorBuilder,
20    OTEL_BLRP_EXPORT_TIMEOUT, OTEL_BLRP_EXPORT_TIMEOUT_DEFAULT, OTEL_BLRP_MAX_EXPORT_BATCH_SIZE,
21    OTEL_BLRP_MAX_EXPORT_BATCH_SIZE_DEFAULT, OTEL_BLRP_MAX_QUEUE_SIZE,
22    OTEL_BLRP_MAX_QUEUE_SIZE_DEFAULT, OTEL_BLRP_SCHEDULE_DELAY, OTEL_BLRP_SCHEDULE_DELAY_DEFAULT,
23};
24pub use export::{LogBatch, LogExporter};
25pub use log_processor::LogProcessor;
26pub use logger::SdkLogger;
27pub use logger_provider::{LoggerProviderBuilder, SdkLoggerProvider};
28pub use record::{SdkLogRecord, TraceContext};
29pub use simple_log_processor::SimpleLogProcessor;
30
31#[cfg(feature = "experimental_logs_batch_log_processor_with_async_runtime")]
32/// Module for BatchLogProcessor with async runtime.
33pub mod log_processor_with_async_runtime;
34
35#[cfg(all(test, feature = "testing"))]
36mod tests {
37    use super::*;
38    use crate::error::OTelSdkResult;
39    use crate::Resource;
40    use opentelemetry::baggage::BaggageExt;
41    use opentelemetry::logs::LogRecord;
42    use opentelemetry::logs::{Logger, LoggerProvider, Severity};
43    use opentelemetry::{logs::AnyValue, Key, KeyValue};
44    use opentelemetry::{Context, InstrumentationScope};
45    use std::borrow::Borrow;
46    use std::collections::HashMap;
47    use std::sync::{Arc, Mutex};
48
49    #[test]
50    fn logging_sdk_test() {
51        // Arrange
52        let resource = Resource::builder_empty()
53            .with_attributes([
54                KeyValue::new("k1", "v1"),
55                KeyValue::new("k2", "v2"),
56                KeyValue::new("k3", "v3"),
57                KeyValue::new("k4", "v4"),
58            ])
59            .build();
60        let exporter: InMemoryLogExporter = InMemoryLogExporter::default();
61        let logger_provider = SdkLoggerProvider::builder()
62            .with_resource(resource.clone())
63            .with_log_processor(SimpleLogProcessor::new(exporter.clone()))
64            .build();
65
66        // Act
67        let logger = logger_provider.logger("test-logger");
68        let mut log_record = logger.create_log_record();
69        log_record.set_severity_number(Severity::Error);
70        log_record.set_severity_text("Error");
71
72        // Adding attributes using a vector with explicitly constructed Key and AnyValue objects.
73        log_record.add_attributes(vec![
74            (Key::new("key1"), AnyValue::from("value1")),
75            (Key::new("key2"), AnyValue::from("value2")),
76        ]);
77
78        // Adding attributes using an array with explicitly constructed Key and AnyValue objects.
79        log_record.add_attributes([
80            (Key::new("key3"), AnyValue::from("value3")),
81            (Key::new("key4"), AnyValue::from("value4")),
82        ]);
83
84        // Adding attributes using a vector with tuple auto-conversion to Key and AnyValue.
85        log_record.add_attributes(vec![("key5", "value5"), ("key6", "value6")]);
86
87        // Adding attributes using an array with tuple auto-conversion to Key and AnyValue.
88        log_record.add_attributes([("key7", "value7"), ("key8", "value8")]);
89
90        // Adding Attributes from a HashMap
91        let mut attributes_map = HashMap::new();
92        attributes_map.insert("key9", "value9");
93        attributes_map.insert("key10", "value10");
94
95        log_record.add_attributes(attributes_map);
96
97        logger.emit(log_record);
98
99        // Assert
100        let exported_logs = exporter
101            .get_emitted_logs()
102            .expect("Logs are expected to be exported.");
103        assert_eq!(exported_logs.len(), 1);
104        let log = exported_logs
105            .first()
106            .expect("Atleast one log is expected to be present.");
107        assert_eq!(log.instrumentation.name(), "test-logger");
108        assert_eq!(log.record.severity_number, Some(Severity::Error));
109        assert_eq!(log.record.attributes_len(), 10);
110        for i in 1..=10 {
111            assert!(log.record.attributes_contains(
112                &Key::new(format!("key{i}")),
113                &AnyValue::String(format!("value{i}").into())
114            ));
115        }
116
117        // validate Resource
118        assert_eq!(&resource, log.resource.borrow());
119    }
120
121    #[test]
122    fn logger_attributes() {
123        let exporter: InMemoryLogExporter = InMemoryLogExporter::default();
124        let provider = SdkLoggerProvider::builder()
125            .with_log_processor(SimpleLogProcessor::new(exporter.clone()))
126            .build();
127
128        let scope = InstrumentationScope::builder("test_logger")
129            .with_schema_url("https://opentelemetry.io/schemas/1.0.0")
130            .with_attributes(vec![(KeyValue::new("test_k", "test_v"))])
131            .build();
132
133        let logger = provider.logger_with_scope(scope);
134
135        let mut log_record = logger.create_log_record();
136        log_record.set_severity_number(Severity::Error);
137
138        logger.emit(log_record);
139
140        let mut exported_logs = exporter
141            .get_emitted_logs()
142            .expect("Logs are expected to be exported.");
143        assert_eq!(exported_logs.len(), 1);
144        let log = exported_logs.remove(0);
145        assert_eq!(log.record.severity_number, Some(Severity::Error));
146
147        let instrumentation_scope = log.instrumentation;
148        assert_eq!(instrumentation_scope.name(), "test_logger");
149        assert_eq!(
150            instrumentation_scope.schema_url(),
151            Some("https://opentelemetry.io/schemas/1.0.0")
152        );
153        assert!(instrumentation_scope
154            .attributes()
155            .eq(&[KeyValue::new("test_k", "test_v")]));
156    }
157
158    #[derive(Debug)]
159    struct EnrichWithBaggageProcessor;
160    impl LogProcessor for EnrichWithBaggageProcessor {
161        fn emit(&self, data: &mut SdkLogRecord, _instrumentation: &InstrumentationScope) {
162            Context::map_current(|cx| {
163                for (kk, vv) in cx.baggage().iter() {
164                    data.add_attribute(kk.clone(), vv.0.clone());
165                }
166            });
167        }
168
169        fn force_flush(&self) -> crate::error::OTelSdkResult {
170            Ok(())
171        }
172
173        fn shutdown_with_timeout(&self, _timeout: std::time::Duration) -> OTelSdkResult {
174            Ok(())
175        }
176    }
177    #[test]
178    fn log_and_baggage() {
179        // Arrange
180        let exporter: InMemoryLogExporter = InMemoryLogExporter::default();
181        let logger_provider = SdkLoggerProvider::builder()
182            .with_log_processor(EnrichWithBaggageProcessor)
183            .with_log_processor(SimpleLogProcessor::new(exporter.clone()))
184            .build();
185
186        // Act
187        let logger = logger_provider.logger("test-logger");
188        let context_with_baggage =
189            Context::current_with_baggage(vec![KeyValue::new("key-from-bag", "value-from-bag")]);
190        let _cx_guard = context_with_baggage.attach();
191        let mut log_record = logger.create_log_record();
192        log_record.add_attribute("key", "value");
193        logger.emit(log_record);
194
195        // Assert
196        let exported_logs = exporter
197            .get_emitted_logs()
198            .expect("Logs are expected to be exported.");
199        assert_eq!(exported_logs.len(), 1);
200        let log = exported_logs
201            .first()
202            .expect("Atleast one log is expected to be present.");
203        assert_eq!(log.instrumentation.name(), "test-logger");
204        assert_eq!(log.record.attributes_len(), 2);
205
206        // Assert that the log record contains the baggage attribute
207        // and the attribute added to the log record.
208        assert!(log
209            .record
210            .attributes_contains(&Key::new("key"), &AnyValue::String("value".into())));
211        assert!(log.record.attributes_contains(
212            &Key::new("key-from-bag"),
213            &AnyValue::String("value-from-bag".into())
214        ));
215    }
216
217    #[test]
218    fn log_suppression() {
219        // Arrange
220        let exporter: InMemoryLogExporter = InMemoryLogExporter::default();
221        let logger_provider = SdkLoggerProvider::builder()
222            .with_simple_exporter(exporter.clone())
223            .build();
224
225        // Act
226        let logger = logger_provider.logger("test-logger");
227        let log_record = logger.create_log_record();
228        {
229            let _suppressed_context = Context::enter_telemetry_suppressed_scope();
230            // This log emission should be suppressed and not exported.
231            logger.emit(log_record);
232        }
233
234        // Assert
235        let exported_logs = exporter.get_emitted_logs().expect("this should not fail.");
236        assert_eq!(
237            exported_logs.len(),
238            0,
239            "There should be a no logs as log emission is done inside a suppressed context"
240        );
241    }
242
243    #[derive(Debug, Clone)]
244    struct ReentrantLogProcessor {
245        logger: Arc<Mutex<Option<SdkLogger>>>,
246    }
247
248    impl ReentrantLogProcessor {
249        fn new() -> Self {
250            Self {
251                logger: Arc::new(Mutex::new(None)),
252            }
253        }
254
255        fn set_logger(&self, logger: SdkLogger) {
256            let mut guard = self.logger.lock().unwrap();
257            *guard = Some(logger);
258        }
259    }
260
261    impl LogProcessor for ReentrantLogProcessor {
262        fn emit(&self, _data: &mut SdkLogRecord, _instrumentation: &InstrumentationScope) {
263            let _suppress = Context::enter_telemetry_suppressed_scope();
264            // Without the suppression above, the logger.emit(log_record) below will cause a deadlock,
265            // as it emits another log, which will attempt to acquire the same lock that is
266            // already held by itself!
267            let logger = self.logger.lock().unwrap();
268            if let Some(logger) = logger.as_ref() {
269                let mut log_record = logger.create_log_record();
270                log_record.set_severity_number(Severity::Error);
271                logger.emit(log_record);
272            }
273        }
274
275        fn force_flush(&self) -> OTelSdkResult {
276            Ok(())
277        }
278
279        fn shutdown_with_timeout(&self, _timeout: std::time::Duration) -> OTelSdkResult {
280            Ok(())
281        }
282    }
283
284    #[test]
285    fn processor_internal_log_does_not_deadlock_with_suppression_enabled() {
286        let processor: ReentrantLogProcessor = ReentrantLogProcessor::new();
287        let logger_provider = SdkLoggerProvider::builder()
288            .with_log_processor(processor.clone())
289            .build();
290        processor.set_logger(logger_provider.logger("processor-logger"));
291
292        let logger = logger_provider.logger("test-logger");
293        let mut log_record = logger.create_log_record();
294        log_record.set_severity_number(Severity::Error);
295        logger.emit(log_record);
296    }
297}