Skip to main content

otel_bootstrap/
log_bridge.rs

1/// Span-aware OTLP log bridge.
2///
3/// Replaces `opentelemetry_appender_tracing::layer::OpenTelemetryTracingBridge` with a layer
4/// that also propagates selected span fields (e.g. `"request.id"`) into every emitted log
5/// record. Trace/span context is attached automatically by the SDK logger via
6/// `opentelemetry::Context::current()` at emit time.
7///
8/// Background: the upstream bridge only visits log-event fields. Span fields set via
9/// `info_span!("request", "request.id" = ...)` are invisible to it. This bridge captures
10/// those fields in `on_new_span` and replays them onto each log record that fires within
11/// the span's scope.
12use 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/// Span fields captured at span-creation time and stored in span extensions.
25#[derive(Default)]
26struct TrackedSpanFields(HashMap<&'static str, String>);
27
28/// Tracing field visitor that captures a fixed set of field names.
29struct 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
51/// Tracing field visitor that sets event fields on an OTLP log record.
52struct 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/// Span-level log attributes written via [`record_span_log_attr_on`].
105///
106/// Stored in span extensions alongside `TrackedSpanFields`. Populated by
107/// `span_enrichment::emit_enduser_fields` and `emit_request_id` — fields that
108/// arrive via `OpenTelemetrySpanExt::set_attribute` (post-creation, not tracing
109/// fields) and therefore cannot be captured by `FieldCollector` in `on_new_span`.
110///
111/// Uses the same `with_subscriber` + `Registry` downcast pattern as
112/// `tracing_opentelemetry::OtelData` — safe to call from any non-Layer context.
113#[derive(Default)]
114pub struct SpanLogAttrs(pub(crate) Vec<(Key, AnyValue)>);
115
116/// Write a key-value pair into the current span's [`SpanLogAttrs`] extension.
117///
118/// Callable from anywhere (middleware, enrichers) — not limited to Layer
119/// `on_*` hooks. No-ops when no span is active or the subscriber is not
120/// `tracing_subscriber::Registry`-backed.
121pub fn record_span_log_attr(key: Key, value: AnyValue) {
122    record_span_log_attr_on(&tracing::Span::current(), key, value);
123}
124
125/// Write a key-value pair into a specific span's [`SpanLogAttrs`] extension.
126///
127/// Prefer [`record_span_log_attr`] for the current span; use this variant
128/// in `make_span_with` closures that thread an explicit span reference.
129pub 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
148/// Span fields propagated automatically from ancestor spans into every log record.
149///
150/// These must be tracing span fields (set at `info_span!("name", key = val)` creation
151/// time) — not attributes added later via `OpenTelemetrySpanExt::set_attribute`.
152/// The bridge captures them in `on_new_span` via `FieldCollector`.
153pub 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
164/// OTLP log bridge that propagates span-level fields and trace context into log records.
165pub struct SpanAwareLogBridge<P: LoggerProvider> {
166    logger: P::Logger,
167    span_fields: &'static [&'static str],
168}
169
170impl<P: LoggerProvider + Send + Sync> SpanAwareLogBridge<P> {
171    /// Construct the bridge.
172    ///
173    /// `span_fields` is the set of tracing span field names whose values are
174    /// captured in `on_new_span` and replayed onto every log record emitted
175    /// within the span. Pass [`PROPAGATED_SPAN_FIELDS`] for the platform
176    /// default set; callers may extend it via
177    /// [`TelemetryBuilder::with_propagated_span_fields`].
178    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        // Path 1: tracing-native fields captured at span creation (on_new_span).
236        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        // Path 2: dynamic attributes written via record_span_log_attr_on
243        // (e.g. request.id / enduser.* set after span creation).
244        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    // Trace context (trace_id / span_id) is attached automatically by the SDK
251    // logger via opentelemetry::Context::current() at emit time.
252}
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        // Must not panic when called without an active span.
424        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        // Span has no fields that match our keys list.
433        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}