Skip to main content

datadog_formatting_layer/
layer.rs

1use crate::{
2    datadog_ids,
3    event_sink::{EventSink, StdoutSink},
4    fields::{self, FieldPair, FieldStore},
5    formatting::DatadogLog,
6};
7use chrono::Utc;
8use std::sync::{Arc, OnceLock};
9use tracing::{dispatcher::WeakDispatch, span::Attributes, Dispatch, Event, Id, Subscriber};
10use tracing_subscriber::{layer::Context, registry::LookupSpan, Layer};
11
12/// The layer responsible for formatting tracing events in a way datadog can parse them
13#[non_exhaustive]
14#[derive(Debug, Clone)]
15pub struct DatadogFormattingLayer<Sink: EventSink + 'static> {
16    event_sink: Sink,
17    // Captured during `on_register_dispatch` so that `on_event` can ask the OTel
18    // layer for the active OpenTelemetry context via `get_otel_context`.
19    dispatch: Arc<OnceLock<WeakDispatch>>,
20}
21
22impl<S: EventSink + 'static> DatadogFormattingLayer<S> {
23    /// Create a new `DatadogFormattingLayer` with the provided event sink
24    ///
25    /// # Example
26    /// ```
27    /// use datadog_formatting_layer::{DatadogFormattingLayer, EventSink, StdoutSink};
28    ///
29    /// let layer: DatadogFormattingLayer<StdoutSink> =
30    ///     DatadogFormattingLayer::with_sink(StdoutSink::default());
31    /// ```
32    pub fn with_sink(sink: S) -> Self {
33        Self {
34            event_sink: sink,
35            dispatch: Arc::new(OnceLock::new()),
36        }
37    }
38}
39
40impl Default for DatadogFormattingLayer<StdoutSink> {
41    fn default() -> Self {
42        Self::with_sink(StdoutSink::default())
43    }
44}
45
46impl<S: Subscriber + for<'a> LookupSpan<'a>, Sink: EventSink + 'static> Layer<S>
47    for DatadogFormattingLayer<Sink>
48{
49    fn on_register_dispatch(&self, subscriber: &Dispatch) {
50        // `on_register_dispatch` can fire more than once for a given subscriber if
51        // the layer is reused; subsequent `set` calls return Err and are ignored.
52        #[allow(clippy::let_underscore_must_use, clippy::let_underscore_untyped)]
53        let _ = self.dispatch.set(subscriber.downgrade());
54    }
55
56    fn on_new_span(&self, span_attrs: &Attributes<'_>, id: &Id, ctx: Context<'_, S>) {
57        #[allow(clippy::expect_used)]
58        let span = ctx.span(id).expect("Span not found, this is a bug");
59
60        let mut extensions = span.extensions_mut();
61
62        let fields = fields::from_attributes(span_attrs);
63
64        // insert fields from new span e.g #[instrument(fields(hello = "world"))]
65        if extensions.get_mut::<FieldStore>().is_none() {
66            extensions.insert(FieldStore { fields });
67        }
68    }
69
70    // IDEA: maybe an on record implementation is required here
71
72    fn on_event(&self, event: &Event<'_>, ctx: Context<'_, S>) {
73        let event_fields = fields::from_event(event);
74
75        // find message if present in event fields
76        let message = event_fields
77            .iter()
78            .find(|pair| pair.name == "message")
79            .map(|pair| pair.value.clone())
80            .unwrap_or_default();
81
82        let all_fields: Vec<FieldPair> = Vec::default()
83            .into_iter()
84            .chain(fields::from_spans(&ctx, event))
85            .chain(event_fields)
86            .collect();
87
88        // look for datadog trace- and span-id
89        let datadog_ids = datadog_ids::read_from_context(&ctx, self.dispatch.get());
90
91        let log = DatadogLog {
92            timestamp: Utc::now(),
93            level: event.metadata().level().to_owned(),
94            message,
95            fields: all_fields,
96            target: event.metadata().target().to_string(),
97            datadog_ids,
98        };
99
100        let serialized_event = log.format();
101
102        self.event_sink.write(serialized_event);
103    }
104}