1mod batch_log_processor;
3mod export;
4mod log_processor;
5mod logger;
6mod logger_provider;
7pub(crate) mod record;
8mod simple_log_processor;
9
10#[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")]
32pub 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 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 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 log_record.add_attributes(vec![
74 (Key::new("key1"), AnyValue::from("value1")),
75 (Key::new("key2"), AnyValue::from("value2")),
76 ]);
77
78 log_record.add_attributes([
80 (Key::new("key3"), AnyValue::from("value3")),
81 (Key::new("key4"), AnyValue::from("value4")),
82 ]);
83
84 log_record.add_attributes(vec![("key5", "value5"), ("key6", "value6")]);
86
87 log_record.add_attributes([("key7", "value7"), ("key8", "value8")]);
89
90 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 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 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 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 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 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!(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 let exporter: InMemoryLogExporter = InMemoryLogExporter::default();
221 let logger_provider = SdkLoggerProvider::builder()
222 .with_simple_exporter(exporter.clone())
223 .build();
224
225 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 logger.emit(log_record);
232 }
233
234 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 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}