Skip to main content

libdd_trace_utils/agentless_encoder/
mod.rs

1// Copyright 2024-Present Datadog, Inc. https://www.datadoghq.com/
2// SPDX-License-Identifier: Apache-2.0
3
4//! Agentless APM JSON encoder.
5//!
6//! Encodes Datadog v04 trace chunks to the JSON body
7//! accepted by the Datadog HTTP trace intake (`POST /v1/input`).
8//!
9//! ## Differences from the regular agent (msgpack v04) encoding
10//!
11//! - **Wire format**: JSON, wrapped as `{"traces": [ {hostname, env, ..., spans: [...] }, ... ]}`.
12//!   Per-trace metadata (hostname, env, language*, tracerVersion, runtimeID, containerID) is
13//!   inlined on each trace instead of being carried in request headers. Hostname is always emitted
14//! - **IDs**: `trace_id`, `span_id`, `parent_id` are lowercase hex strings (16 chars; 32 for
15//!   span-link trace IDs)
16//! - **128-bit trace IDs**: only the low 64 bits go into `trace_id`; the `_dd.p.tid` meta tag
17//!   carries upper 64 bits
18//! - **Span links / events**: not top-level fields. They are JSON-stringified into
19//!   `meta["_dd.span_links"]` and `meta["events"]`, each truncated to 25_000 chars. No top-level
20//!   `links` field is emitted. If meta attributes for "_dd.span_links" of "events" are already
21//!   attached to the span we will keep existing fields and the span level field will be dropped
22//! - **Stats / top-level flags**: the intake has no trace-agent to compute them, so the encoder
23//!   injects `meta["_dd.compute_stats"]="1"` on the first span of each chunk and
24//!   `metrics["_trace_root"]=1` where applicable.
25//! - **Non-finite metrics** (NaN/Inf) are dropped (JSON can't represent them).
26//!
27//! TODO: span normalization (service/name/resource/type truncation + defaults)
28
29use crate::span::v04::{AttributeAnyValue, AttributeArrayValue, Span, SpanEvent, SpanLink};
30use crate::span::{TraceData, SPAN_LINK_FLAGS_SET_SENTINEL};
31use crate::tracer_metadata::TracerMetadata;
32use serde::{
33    ser::{SerializeMap, SerializeSeq},
34    Serializer,
35};
36use std::borrow::Borrow;
37
38/// Maximum allowed size of a `meta` value before truncation.
39const MAX_META_VALUE_LEN: usize = 25_000;
40/// Suffix appended when a `meta` value is truncated.
41const TRUNCATION_SUFFIX: &str = "...";
42
43/// # Why are we doing this?
44///
45/// The JSON agentless format is different from the in-memory model of v04 spans (and to v03/v04/v05
46/// on-the-wire schemas) In order to not have to copy to intermediary structs, we have to write a
47/// manual encoder. For JSON there is no widely available JSON emitter in rust other than serde
48/// JSON. But serde does not let us drive serialization other than through the serde::Serialize
49/// trait.
50///
51/// Defining structs implementing serde::Serialize for every nested object in the span is heavy,
52///
53/// This macro captures parameters from the environment and creates a local struct implementing
54/// serde::Serialize, with a custom implementation.
55///
56/// # Usage
57///
58/// The shape of the input of the macro is made to look like a closure
59/// Contrary to a closure, the names of the types have to be named in full
60/// ```ignore
61///        Optional generic| serializer|
62///            parameter   |.  |               Captured variables from env       
63///         --------------  --- ---------------------------------------------------------
64/// ser_fn!(<T: TraceData> |ser, traces: &'a [Vec<Span<T>>], metadata: &'a TracerMetadata| {
65///     // Body of the closure
66/// }
67/// ```
68macro_rules! ser_fn {
69    ($(<$generic:ident $(: $bound:ident )?>)? |$serializer:ident , $($captured:ident : $ty:ty),+ $(,)?| { $($body:tt)* }) => {
70        {
71            struct SerializeClosure<'a, $($generic $(: $bound + 'a)? ,)?>(($(&'a $ty ,)*));
72
73            impl <'a, $($generic $(: $bound + 'a)?,)?> serde::Serialize for SerializeClosure<'a, $($generic,)?> {
74                #[inline]
75                fn serialize<S: serde::Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
76                    let captured = self.0;
77                    (|$serializer: S , ($(& $captured, )*) : ($(&'a $ty ,)*)| {
78                        $($body)*
79                    })(serializer, captured)
80                }
81            }
82
83            SerializeClosure(($(& $captured ,)*))
84        }
85    }
86}
87
88/// Encode the given `traces` to the agentless JSON payload (`/v1/input` body).
89///
90/// Returns the serialized JSON bytes on success.
91pub fn encode_payload<T: TraceData>(
92    traces: &[Vec<Span<T>>],
93    metadata: &TracerMetadata,
94) -> Result<Vec<u8>, serde_json::Error> {
95    let mut bytes = Vec::new();
96    let mut serializer = serde_json::Serializer::new(&mut bytes);
97
98    let mut map_ser = serializer.serialize_map(Some(1))?;
99    map_ser.serialize_entry(
100        "traces",
101        &ser_fn!(<T: TraceData> |ser, traces: &'a [Vec<Span<T>>], metadata: &'a TracerMetadata| {
102            let mut traces_serializer = ser.serialize_seq(Some(traces.len()))?;
103            for chunk in traces {
104                traces_serializer.serialize_element(&ser_fn!(<T: TraceData> |ser, chunk: &'a Vec<Span<T>>, metadata: &'a TracerMetadata| {
105                    encode_trace(ser, chunk, metadata)
106                }))?;
107            }
108            traces_serializer.end()
109        }),
110    )?;
111    SerializeMap::end(map_ser)?;
112    Ok(bytes)
113}
114
115fn encode_trace<T: TraceData, S: Serializer>(
116    ser: S,
117    chunk: &[Span<T>],
118    metadata: &TracerMetadata,
119) -> Result<S::Ok, S::Error> {
120    let mut map = ser.serialize_map(None)?;
121
122    // Per-trace metadata. Always include hostname; other fields when set.
123    map.serialize_entry("hostname", &metadata.hostname)?;
124    if !metadata.env.is_empty() {
125        map.serialize_entry("env", &metadata.env)?;
126    }
127    if !metadata.language.is_empty() {
128        map.serialize_entry("languageName", &metadata.language)?;
129    }
130    if !metadata.language_version.is_empty() {
131        map.serialize_entry("languageVersion", &metadata.language_version)?;
132    }
133    if !metadata.tracer_version.is_empty() {
134        map.serialize_entry("tracerVersion", &metadata.tracer_version)?;
135    }
136    if !metadata.runtime_id.is_empty() {
137        map.serialize_entry("runtimeID", &metadata.runtime_id)?;
138    }
139    if let Some(container_id) = libdd_common::entity_id::get_container_id() {
140        map.serialize_entry("containerID", container_id)?;
141    }
142
143    map.serialize_entry(
144        "spans",
145        &ser_fn!(<T: TraceData> |ser, chunk: &'a [Span<T>]| {
146            let mut seq = ser.serialize_seq(Some(chunk.len()))?;
147            for (i, span) in chunk.iter().enumerate() {
148                let is_first = i == 0;
149                seq.serialize_element(&ser_fn!(<T: TraceData> |ser, span: &'a Span<T>, is_first: bool| {
150                    encode_span(ser, span, is_first)
151                }))?;
152            }
153            seq.end()
154        }),
155    )?;
156
157    map.end()
158}
159
160fn encode_span<T: TraceData, S: Serializer>(
161    ser: S,
162    span: &Span<T>,
163    is_first_in_trace: bool,
164) -> Result<S::Ok, S::Error> {
165    let mut map = ser.serialize_map(None)?;
166
167    let trace_id = span.trace_id;
168    map.serialize_entry(
169        "trace_id",
170        &ser_fn!(|ser, trace_id: u128| {
171            ser.collect_str(&format_args!("{:016x}", trace_id as u64))
172        }),
173    )?;
174    let span_id = span.span_id;
175    map.serialize_entry(
176        "span_id",
177        &ser_fn!(|ser, span_id: u64| { ser.collect_str(&format_args!("{:016x}", span_id as u64)) }),
178    )?;
179    let parent_id = span.parent_id;
180    map.serialize_entry(
181        "parent_id",
182        &ser_fn!(|ser, parent_id: u64| {
183            ser.collect_str(&format_args!("{:016x}", parent_id as u64))
184        }),
185    )?;
186
187    // Resource defaults to name when empty.
188    let name_str: &str = span.name.borrow();
189    let resource_str: &str = span.resource.borrow();
190    let service_str: &str = span.service.borrow();
191    map.serialize_entry("name", name_str)?;
192    map.serialize_entry(
193        "resource",
194        if resource_str.is_empty() {
195            name_str
196        } else {
197            resource_str
198        },
199    )?;
200    map.serialize_entry("service", service_str)?;
201    map.serialize_entry("error", &span.error)?;
202    map.serialize_entry("start", &span.start.max(0))?;
203    map.serialize_entry("duration", &span.duration)?;
204
205    let type_str: &str = span.r#type.borrow();
206    if !type_str.is_empty() {
207        map.serialize_entry("type", type_str)?;
208    }
209
210    map.serialize_entry(
211        "meta",
212        &ser_fn!(<T: TraceData> |ser, span: &'a Span<T>, is_first_in_trace: bool| {
213            let upper_bits = (span.trace_id >> 64) as u64;
214            let mut p_tid_seen = false;
215            let mut span_links_seen = false;
216            let mut events_seen = false;
217            let mut compute_stats_seen = false;
218
219            let mut meta = ser.serialize_map(None)?;
220            for (k, v) in span.meta.iter() {
221                let key: &str = k.borrow();
222                match key {
223                    "_dd.p.tid" => p_tid_seen = true,
224                    "_dd.span_links" => span_links_seen = true,
225                    "events" => events_seen = true,
226                    "_dd.compute_stats"=> compute_stats_seen = true,
227                    _ => {}
228                };
229                let val: &str = v.borrow();
230                meta.serialize_entry(key, val)?;
231            }
232            if !p_tid_seen && upper_bits != 0 {
233                meta.serialize_entry(
234                    "_dd.p.tid",
235                    &ser_fn!(|ser, upper_bits: u64| {
236                        ser.collect_str(&format_args!("{:016x}", upper_bits as u64))
237                    }),
238                )?;
239            }
240            if !span_links_seen && !span.span_links.is_empty() {
241                if let Some(s) = serialize_span_links(&span.span_links) {
242                    meta.serialize_entry("_dd.span_links", &s)?;
243                }
244            }
245            if !events_seen && !span.span_events.is_empty() {
246                if let Some(s) = serialize_span_events(&span.span_events) {
247                    meta.serialize_entry("events", &s)?;
248                }
249            }
250            if !compute_stats_seen && is_first_in_trace {
251                meta.serialize_entry("_dd.compute_stats", "1")?;
252            }
253            meta.end()
254        }),
255    )?;
256
257    map.serialize_entry(
258        "metrics",
259        &ser_fn!(<T: TraceData> |ser, span: &'a Span<T>| {
260            let mut metrics = ser.serialize_map(None)?;
261            let mut trace_root_seen = false;
262            for (k, v) in span.metrics.iter() {
263                let key: &str = k.borrow();
264                // serde_json refuses to serialize NaN/Inf; drop them silently.
265                if v.is_finite() {
266                    match key {
267                        "_trace_root" => trace_root_seen = true,
268                        "_top_level" => {
269                            metrics.serialize_entry(key, &(*v as u32))?;
270                            continue
271                        },
272                        _ => {},
273                    }
274                    metrics.serialize_entry(key, v)?
275                }
276            }
277            if !trace_root_seen && span.parent_id == 0 {
278                metrics.serialize_entry("_trace_root", &1u32)?;
279            }
280            metrics.end()
281        }),
282    )?;
283
284    if !span.meta_struct.is_empty() {
285        map.serialize_entry(
286            "meta_struct",
287            &ser_fn!(<T: TraceData> |ser, span: &'a Span<T>| {
288                let mut ms = ser.serialize_map(None)?;
289                for (k, v) in span.meta_struct.iter() {
290                    let key: &str = k.borrow();
291                    let bytes: &[u8] = v.borrow();
292
293                    // abort whole payload on malformed entry
294                    ms.serialize_entry(key, &MsgpackAsJson(bytes))?;
295                }
296                ms.end()
297            }),
298        )?;
299    }
300    map.end()
301}
302
303/// Serialize span links to a JSON string suitable for `meta['_dd.span_links']`.
304///
305/// Returns `None` if serialization fails. The result is truncated to
306/// [`MAX_META_VALUE_LEN`] characters with a trailing `"..."` if it would
307/// otherwise exceed that limit.
308fn serialize_span_links<T: TraceData>(links: &[SpanLink<T>]) -> Option<String> {
309    let s = serde_json::to_string(&ser_fn!(<T: TraceData> |ser, links: &'a [SpanLink<T>]| {
310        let mut seq = ser.serialize_seq(Some(links.len()))?;
311        for link in links {
312            seq.serialize_element(&ser_fn!(<T: TraceData> |ser, link: &'a SpanLink<T>| {
313                encode_span_link(ser, link)
314            }))?;
315        }
316        seq.end()
317    }))
318    .ok()?;
319    Some(truncate_with_ellipsis(s, MAX_META_VALUE_LEN))
320}
321
322fn encode_span_link<T: TraceData, S: Serializer>(
323    ser: S,
324    link: &SpanLink<T>,
325) -> Result<S::Ok, S::Error> {
326    let mut map = ser.serialize_map(None)?;
327    let trace_id_128: u128 = ((link.trace_id_high as u128) << 64) | (link.trace_id as u128);
328    map.serialize_entry("trace_id", &format!("{:032x}", trace_id_128))?;
329    map.serialize_entry("span_id", &format!("{:016x}", link.span_id))?;
330    if !link.attributes.is_empty() {
331        map.serialize_entry(
332            "attributes",
333            &ser_fn!(<T: TraceData> |ser, link: &'a SpanLink<T>| {
334                let mut attrs = ser.serialize_map(Some(link.attributes.len()))?;
335                for (k, v) in link.attributes.iter() {
336                    let key: &str = k.borrow();
337                    let val: &str = v.borrow();
338                    attrs.serialize_entry(key, val)?;
339                }
340                attrs.end()
341            }),
342        )?;
343    }
344    // When `flags` is 0, no sampling decision exists, so omit the field. Before emission,
345    // mask off the internal "explicitly set" sentinel (bit 31), because this JSON field uses
346    // the same `_dd.span_links` key that the v0.5 encoder produces and must match its output.
347    if link.flags != 0 {
348        map.serialize_entry(
349            "flags",
350            &((link.flags & !SPAN_LINK_FLAGS_SET_SENTINEL) as u64),
351        )?;
352    }
353    let tracestate: &str = link.tracestate.borrow();
354    if !tracestate.is_empty() {
355        map.serialize_entry("tracestate", tracestate)?;
356    }
357    map.end()
358}
359
360/// Serialize span events to a JSON string suitable for `meta['events']`.
361fn serialize_span_events<T: TraceData>(events: &[SpanEvent<T>]) -> Option<String> {
362    let s = serde_json::to_string(&ser_fn!(<T: TraceData> |ser, events: &'a [SpanEvent<T>]| {
363        let mut seq = ser.serialize_seq(Some(events.len()))?;
364        for event in events {
365            seq.serialize_element(&ser_fn!(<T: TraceData> |ser, event: &'a SpanEvent<T>| {
366                encode_span_event(ser, event)
367            }))?;
368        }
369        seq.end()
370    }))
371    .ok()?;
372    Some(truncate_with_ellipsis(s, MAX_META_VALUE_LEN))
373}
374
375fn encode_span_event<T: TraceData, S: Serializer>(
376    ser: S,
377    event: &SpanEvent<T>,
378) -> Result<S::Ok, S::Error> {
379    let mut map = ser.serialize_map(None)?;
380    let name: &str = event.name.borrow();
381    map.serialize_entry("name", name)?;
382    map.serialize_entry("time_unix_nano", &event.time_unix_nano)?;
383    if !event.attributes.is_empty() {
384        map.serialize_entry(
385            "attributes",
386            &ser_fn!(<T: TraceData> |ser, event: &'a SpanEvent<T>| {
387                let mut attrs = ser.serialize_map(Some(event.attributes.len()))?;
388                for (k, v) in event.attributes.iter() {
389                    let key: &str = k.borrow();
390                    attrs.serialize_entry(key, &ser_fn!(<T: TraceData> |ser, v: &'a AttributeAnyValue<T> | {
391                        match v {
392                            AttributeAnyValue::SingleValue(v) => serialize_scalar(ser, v),
393                            AttributeAnyValue::Array(values) => {
394                                let mut seq = ser.serialize_seq(Some(values.len()))?;
395                                for v in values {
396                                    seq.serialize_element(&ser_fn!(<T: TraceData> |ser, v: &'a AttributeArrayValue<T>| {
397                                        serialize_scalar(ser, v)
398                                    }))?;
399                                }
400                                seq.end()
401                            }
402                        }
403                    }))?;
404                }
405                attrs.end()
406            }),
407        )?;
408    }
409    map.end()
410}
411
412fn serialize_scalar<S: serde::Serializer, T: TraceData>(
413    ser: S,
414    s: &AttributeArrayValue<T>,
415) -> Result<S::Ok, S::Error> {
416    match s {
417        AttributeArrayValue::String(s) => {
418            let s: &str = s.borrow();
419            ser.serialize_str(s)
420        }
421        AttributeArrayValue::Boolean(b) => ser.serialize_bool(*b),
422        AttributeArrayValue::Integer(i) => ser.serialize_i64(*i),
423        AttributeArrayValue::Double(d) => {
424            if d.is_finite() {
425                ser.serialize_f64(*d)
426            } else {
427                // NaN/Inf become JSON null.
428                ser.serialize_unit()
429            }
430        }
431    }
432}
433
434/// `serde::Serialize` adapter that interprets `bytes` as a self-describing
435/// msgpack value and transcodes it into the destination serializer.
436///
437/// Used to inline `meta_struct` values (which are stored as msgpack-encoded
438/// bytes) into the agentless JSON payload as real JSON objects.
439struct MsgpackAsJson<'a>(&'a [u8]);
440
441impl serde::Serialize for MsgpackAsJson<'_> {
442    fn serialize<S: serde::Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
443        let mut de = rmp_serde::Deserializer::from_read_ref(self.0);
444        serde_transcode::transcode(&mut de, serializer)
445    }
446}
447
448/// Truncate `s` to at most `max_len` bytes, appending `"..."` when truncation occurs.
449fn truncate_with_ellipsis(mut s: String, max_len: usize) -> String {
450    if s.len() <= max_len {
451        return s;
452    }
453    let suffix_len = TRUNCATION_SUFFIX.len();
454    let cut = max_len.saturating_sub(suffix_len);
455    // Find the previous char boundary so we don't slice in the middle of a UTF-8 sequence.
456    let mut end = cut;
457    while end > 0 && !s.is_char_boundary(end) {
458        end -= 1;
459    }
460    s.truncate(end);
461    s.push_str(TRUNCATION_SUFFIX);
462    s
463}
464
465#[cfg(test)]
466mod tests;