1use opentelemetry::{
13 Key,
14 logs::{AnyValue, LogRecord, Logger, LoggerProvider, Severity},
15};
16use std::collections::HashMap;
17use tracing::Subscriber;
18use tracing_subscriber::{
19 Layer, Registry,
20 layer::Context,
21 registry::{LookupSpan, SpanRef},
22};
23
24#[derive(Default)]
26struct TrackedSpanFields(HashMap<&'static str, String>);
27
28struct FieldCollector<'a> {
30 fields: &'a mut HashMap<&'static str, String>,
31 keys: &'static [&'static str],
32}
33
34impl tracing::field::Visit for FieldCollector<'_> {
35 fn record_str(&mut self, field: &tracing::field::Field, value: &str) {
36 if self.keys.contains(&field.name()) {
37 self.fields.insert(field.name(), value.to_owned());
38 }
39 }
40
41 fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
42 if self.keys.contains(&field.name()) {
43 self.fields.insert(
44 field.name(),
45 format!("{value:?}").trim_matches('"').to_owned(),
46 );
47 }
48 }
49}
50
51struct LogRecordVisitor<'a, LR: LogRecord>(&'a mut LR);
53
54impl<LR: LogRecord> tracing::field::Visit for LogRecordVisitor<'_, LR> {
55 fn record_str(&mut self, field: &tracing::field::Field, value: &str) {
56 if field.name() == "message" {
57 self.0.set_body(AnyValue::String(value.to_owned().into()));
58 } else {
59 self.0.add_attribute(
60 Key::new(field.name()),
61 AnyValue::String(value.to_owned().into()),
62 );
63 }
64 }
65
66 fn record_bool(&mut self, field: &tracing::field::Field, value: bool) {
67 self.0
68 .add_attribute(Key::new(field.name()), AnyValue::Boolean(value));
69 }
70
71 fn record_i64(&mut self, field: &tracing::field::Field, value: i64) {
72 self.0
73 .add_attribute(Key::new(field.name()), AnyValue::Int(value));
74 }
75
76 fn record_f64(&mut self, field: &tracing::field::Field, value: f64) {
77 self.0
78 .add_attribute(Key::new(field.name()), AnyValue::Double(value));
79 }
80
81 fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
82 if field.name() == "message" {
83 self.0
84 .set_body(AnyValue::String(format!("{value:?}").into()));
85 } else {
86 self.0.add_attribute(
87 Key::new(field.name()),
88 AnyValue::String(format!("{value:?}").into()),
89 );
90 }
91 }
92}
93
94fn severity_of_level(level: &tracing::Level) -> Severity {
95 match *level {
96 tracing::Level::TRACE => Severity::Trace,
97 tracing::Level::DEBUG => Severity::Debug,
98 tracing::Level::INFO => Severity::Info,
99 tracing::Level::WARN => Severity::Warn,
100 tracing::Level::ERROR => Severity::Error,
101 }
102}
103
104#[derive(Default)]
114pub struct SpanLogAttrs(pub(crate) Vec<(Key, AnyValue)>);
115
116pub fn record_span_log_attr(key: Key, value: AnyValue) {
122 record_span_log_attr_on(&tracing::Span::current(), key, value);
123}
124
125pub fn record_span_log_attr_on(span: &tracing::Span, key: Key, value: AnyValue) {
130 span.with_subscriber(|(id, dispatch)| {
131 if let Some(registry) = dispatch.downcast_ref::<Registry>() {
132 if let Some(span_ref) = registry.span(id) {
133 let mut ext = span_ref.extensions_mut();
134 if ext.get_mut::<SpanLogAttrs>().is_none() {
135 ext.insert(SpanLogAttrs::default());
136 }
137 let attrs = ext.get_mut::<SpanLogAttrs>().unwrap();
138 if let Some(existing) = attrs.0.iter_mut().find(|(k, _)| k == &key) {
139 existing.1 = value;
140 } else {
141 attrs.0.push((key, value));
142 }
143 }
144 }
145 });
146}
147
148pub const PROPAGATED_SPAN_FIELDS: &[&str] = &[
154 "request.id",
155 "enduser.id",
156 "enduser.org_id",
157 "enduser.org_path",
158 "enduser.principal_kind",
159 "http.request.method",
160 "http.response.status_code",
161 "http.route",
162];
163
164pub struct SpanAwareLogBridge<P: LoggerProvider> {
166 logger: P::Logger,
167 span_fields: &'static [&'static str],
168}
169
170impl<P: LoggerProvider + Send + Sync> SpanAwareLogBridge<P> {
171 pub fn new(provider: &P, span_fields: &'static [&'static str]) -> Self {
179 Self {
180 logger: provider.logger("otel-bootstrap"),
181 span_fields,
182 }
183 }
184}
185
186impl<S, P> Layer<S> for SpanAwareLogBridge<P>
187where
188 S: Subscriber + for<'a> LookupSpan<'a>,
189 P: LoggerProvider + Send + Sync + 'static,
190 P::Logger: Logger + Send + Sync,
191{
192 fn on_new_span(
193 &self,
194 attrs: &tracing::span::Attributes<'_>,
195 id: &tracing::span::Id,
196 ctx: Context<'_, S>,
197 ) {
198 let mut tracked = TrackedSpanFields::default();
199 attrs.record(&mut FieldCollector {
200 fields: &mut tracked.0,
201 keys: self.span_fields,
202 });
203 if !tracked.0.is_empty() {
204 if let Some(span) = ctx.span(id) {
205 span.extensions_mut().insert(tracked);
206 }
207 }
208 }
209
210 fn on_event(&self, event: &tracing::Event<'_>, ctx: Context<'_, S>) {
211 let meta = event.metadata();
212 let mut log_record = self.logger.create_log_record();
213
214 log_record.set_severity_number(severity_of_level(meta.level()));
215 log_record.set_severity_text(meta.level().as_str());
216 log_record.set_target(meta.target());
217 log_record.set_event_name(meta.name());
218
219 event.record(&mut LogRecordVisitor(&mut log_record));
220
221 if let Some(span) = ctx.event_span(event) {
222 inject_span_context(&span, &mut log_record);
223 }
224
225 self.logger.emit(log_record);
226 }
227}
228
229fn inject_span_context<S, LR>(span: &SpanRef<'_, S>, log_record: &mut LR)
230where
231 S: Subscriber + for<'a> LookupSpan<'a>,
232 LR: LogRecord,
233{
234 for ancestor in span.scope() {
235 if let Some(tracked) = ancestor.extensions().get::<TrackedSpanFields>() {
237 for (k, v) in &tracked.0 {
238 log_record.add_attribute(Key::new(*k), AnyValue::String(v.clone().into()));
239 }
240 }
241
242 if let Some(log_attrs) = ancestor.extensions().get::<SpanLogAttrs>() {
245 for (k, v) in &log_attrs.0 {
246 log_record.add_attribute(k.clone(), v.clone());
247 }
248 }
249 }
250 }
253
254#[cfg(test)]
255mod tests {
256 use super::*;
257 use opentelemetry::{InstrumentationScope, logs::Logger, logs::LoggerProvider};
258 use std::sync::{Arc, Mutex};
259 use tracing_subscriber::layer::SubscriberExt;
260
261 #[derive(Default, Clone)]
262 struct CapturedRecord {
263 body: Option<AnyValue>,
264 attributes: Vec<(Key, AnyValue)>,
265 severity: Option<Severity>,
266 }
267
268 #[derive(Default, Clone)]
269 struct CapturingLogRecord(Arc<Mutex<CapturedRecord>>);
270
271 impl opentelemetry::logs::LogRecord for CapturingLogRecord {
272 fn set_event_name(&mut self, _name: &'static str) {}
273 fn set_target<T: Into<std::borrow::Cow<'static, str>>>(&mut self, _target: T) {}
274 fn set_timestamp(&mut self, _ts: std::time::SystemTime) {}
275 fn set_observed_timestamp(&mut self, _ts: std::time::SystemTime) {}
276 fn set_severity_text(&mut self, _text: &'static str) {}
277 fn set_severity_number(&mut self, sev: opentelemetry::logs::Severity) {
278 self.0.lock().unwrap().severity = Some(sev);
279 }
280 fn set_body(&mut self, body: AnyValue) {
281 self.0.lock().unwrap().body = Some(body);
282 }
283 fn add_attributes<I, K, V>(&mut self, attributes: I)
284 where
285 I: IntoIterator<Item = (K, V)>,
286 K: Into<Key>,
287 V: Into<AnyValue>,
288 {
289 let mut guard = self.0.lock().unwrap();
290 for (k, v) in attributes {
291 guard.attributes.push((k.into(), v.into()));
292 }
293 }
294 fn add_attribute<K, V>(&mut self, key: K, value: V)
295 where
296 K: Into<Key>,
297 V: Into<AnyValue>,
298 {
299 self.0
300 .lock()
301 .unwrap()
302 .attributes
303 .push((key.into(), value.into()));
304 }
305 }
306
307 #[derive(Clone, Default)]
308 struct CapturingLogger {
309 records: Arc<Mutex<Vec<CapturedRecord>>>,
310 }
311
312 impl Logger for CapturingLogger {
313 type LogRecord = CapturingLogRecord;
314
315 fn create_log_record(&self) -> Self::LogRecord {
316 CapturingLogRecord(Arc::new(Mutex::new(CapturedRecord::default())))
317 }
318
319 fn emit(&self, record: Self::LogRecord) {
320 let captured = record.0.lock().unwrap().clone();
321 self.records.lock().unwrap().push(captured);
322 }
323
324 fn event_enabled(
325 &self,
326 _level: opentelemetry::logs::Severity,
327 _target: &str,
328 _name: Option<&str>,
329 ) -> bool {
330 true
331 }
332 }
333
334 #[derive(Clone, Default)]
335 struct CapturingLoggerProvider {
336 logger: CapturingLogger,
337 }
338
339 impl LoggerProvider for CapturingLoggerProvider {
340 type Logger = CapturingLogger;
341
342 fn logger_with_scope(&self, _scope: InstrumentationScope) -> Self::Logger {
343 self.logger.clone()
344 }
345 }
346
347 fn make_subscriber(
348 fields: &'static [&'static str],
349 ) -> (impl tracing::Subscriber, Arc<Mutex<Vec<CapturedRecord>>>) {
350 let provider = CapturingLoggerProvider::default();
351 let records = provider.logger.records.clone();
352 let bridge = SpanAwareLogBridge::new(&provider, fields);
353 (tracing_subscriber::registry().with(bridge), records)
354 }
355
356 fn attr_str<'a>(record: &'a CapturedRecord, key: &str) -> Option<&'a str> {
357 record.attributes.iter().find_map(|(k, v)| {
358 if k.as_str() == key {
359 if let AnyValue::String(s) = v {
360 Some(s.as_str())
361 } else {
362 None
363 }
364 } else {
365 None
366 }
367 })
368 }
369
370 #[test]
371 fn request_id_tracing_field_propagates_to_log_record() {
372 let (sub, records) = make_subscriber(&["request.id"]);
373 let _guard = tracing::subscriber::set_default(sub);
374
375 let span = tracing::info_span!("request", "request.id" = "test-uuid-1234");
376 let _enter = span.enter();
377 tracing::info!("hello from inside the span");
378
379 let recs = records.lock().unwrap();
380 assert!(!recs.is_empty());
381 assert_eq!(attr_str(&recs[0], "request.id"), Some("test-uuid-1234"));
382 }
383
384 #[test]
385 fn span_log_attrs_propagate_to_log_record() {
386 let (sub, records) = make_subscriber(&[]);
387 let _guard = tracing::subscriber::set_default(sub);
388
389 let span = tracing::info_span!("req");
390 let _enter = span.enter();
391 record_span_log_attr_on(&span, Key::new("x-custom"), AnyValue::String("val".into()));
392 tracing::info!("inside");
393
394 let recs = records.lock().unwrap();
395 assert!(!recs.is_empty());
396 assert_eq!(attr_str(&recs[0], "x-custom"), Some("val"));
397 }
398
399 #[test]
400 fn span_log_attrs_update_existing_key() {
401 let (sub, records) = make_subscriber(&[]);
402 let _guard = tracing::subscriber::set_default(sub);
403
404 let span = tracing::info_span!("req");
405 let _enter = span.enter();
406 record_span_log_attr_on(&span, Key::new("k"), AnyValue::String("first".into()));
407 record_span_log_attr_on(&span, Key::new("k"), AnyValue::String("second".into()));
408 tracing::info!("inside");
409
410 let recs = records.lock().unwrap();
411 assert!(!recs.is_empty());
412 let vals: Vec<_> = recs[0]
413 .attributes
414 .iter()
415 .filter(|(k, _)| k.as_str() == "k")
416 .collect();
417 assert_eq!(vals.len(), 1, "duplicate key must be deduplicated");
418 assert!(matches!(&vals[0].1, AnyValue::String(s) if s.as_str() == "second"));
419 }
420
421 #[test]
422 fn record_span_log_attr_current_span_no_op_outside_span() {
423 record_span_log_attr(Key::new("k"), AnyValue::String("v".into()));
425 }
426
427 #[test]
428 fn on_new_span_non_matching_fields_no_extension_inserted() {
429 let (sub, records) = make_subscriber(&["request.id"]);
430 let _guard = tracing::subscriber::set_default(sub);
431
432 let span = tracing::info_span!("plain");
434 let _enter = span.enter();
435 tracing::info!("msg");
436
437 let recs = records.lock().unwrap();
438 assert!(!recs.is_empty());
439 assert!(attr_str(&recs[0], "request.id").is_none());
440 }
441
442 #[test]
443 fn log_outside_span_no_panic() {
444 let (sub, records) = make_subscriber(&["request.id"]);
445 let _guard = tracing::subscriber::set_default(sub);
446 tracing::info!("no span active");
447 assert!(!records.lock().unwrap().is_empty());
448 }
449
450 #[test]
451 fn nested_span_outer_field_appears_in_inner_log() {
452 let (sub, records) = make_subscriber(&["request.id"]);
453 let _guard = tracing::subscriber::set_default(sub);
454
455 let outer = tracing::info_span!("outer", "request.id" = "outer-id");
456 let _e1 = outer.enter();
457 let inner = tracing::info_span!("inner");
458 let _e2 = inner.enter();
459 tracing::info!("deep log");
460
461 let recs = records.lock().unwrap();
462 assert!(!recs.is_empty());
463 assert_eq!(attr_str(&recs[0], "request.id"), Some("outer-id"));
464 }
465
466 #[test]
467 fn log_record_visitor_bool_field() {
468 let (sub, records) = make_subscriber(&[]);
469 let _guard = tracing::subscriber::set_default(sub);
470 tracing::info!(flag = true, "msg");
471 let recs = records.lock().unwrap();
472 assert!(!recs.is_empty());
473 let has_flag = recs[0]
474 .attributes
475 .iter()
476 .any(|(k, v)| k.as_str() == "flag" && matches!(v, AnyValue::Boolean(true)));
477 assert!(has_flag);
478 }
479
480 #[test]
481 fn log_record_visitor_i64_field() {
482 let (sub, records) = make_subscriber(&[]);
483 let _guard = tracing::subscriber::set_default(sub);
484 tracing::info!(count = 42i64, "msg");
485 let recs = records.lock().unwrap();
486 assert!(!recs.is_empty());
487 let has_count = recs[0]
488 .attributes
489 .iter()
490 .any(|(k, v)| k.as_str() == "count" && matches!(v, AnyValue::Int(42)));
491 assert!(has_count);
492 }
493
494 #[test]
495 fn log_record_visitor_f64_field() {
496 let (sub, records) = make_subscriber(&[]);
497 let _guard = tracing::subscriber::set_default(sub);
498 tracing::info!(ratio = 0.5f64, "msg");
499 let recs = records.lock().unwrap();
500 assert!(!recs.is_empty());
501 let has_ratio = recs[0].attributes.iter().any(|(k, v)| {
502 k.as_str() == "ratio" && matches!(v, AnyValue::Double(f) if (f - 0.5).abs() < 1e-9)
503 });
504 assert!(has_ratio);
505 }
506
507 #[test]
508 fn log_record_visitor_message_becomes_body() {
509 let (sub, records) = make_subscriber(&[]);
510 let _guard = tracing::subscriber::set_default(sub);
511 tracing::info!("the body text");
512 let recs = records.lock().unwrap();
513 assert!(!recs.is_empty());
514 assert!(
515 matches!(&recs[0].body, Some(AnyValue::String(s)) if s.as_str() == "the body text")
516 );
517 }
518
519 #[test]
520 fn severity_warn() {
521 let (sub, records) = make_subscriber(&[]);
522 let _guard = tracing::subscriber::set_default(sub);
523 tracing::warn!("warn msg");
524 let recs = records.lock().unwrap();
525 assert!(!recs.is_empty());
526 assert_eq!(recs[0].severity, Some(Severity::Warn));
527 }
528
529 #[test]
530 fn severity_error() {
531 let (sub, records) = make_subscriber(&[]);
532 let _guard = tracing::subscriber::set_default(sub);
533 tracing::error!("err msg");
534 let recs = records.lock().unwrap();
535 assert!(!recs.is_empty());
536 assert_eq!(recs[0].severity, Some(Severity::Error));
537 }
538
539 #[test]
540 fn severity_debug() {
541 let (sub, records) = make_subscriber(&[]);
542 let _guard = tracing::subscriber::set_default(sub);
543 tracing::debug!("dbg msg");
544 let recs = records.lock().unwrap();
545 assert!(!recs.is_empty());
546 assert_eq!(recs[0].severity, Some(Severity::Debug));
547 }
548
549 #[test]
550 fn severity_trace() {
551 let (sub, records) = make_subscriber(&[]);
552 let _guard = tracing::subscriber::set_default(sub);
553 tracing::trace!("trc msg");
554 let recs = records.lock().unwrap();
555 assert!(!recs.is_empty());
556 assert_eq!(recs[0].severity, Some(Severity::Trace));
557 }
558}