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::v1;
31use crate::span::{TraceData, SPAN_LINK_FLAGS_SET_SENTINEL};
32use crate::tracer_metadata::TracerMetadata;
33use serde::{
34    ser::{SerializeMap, SerializeSeq},
35    Serializer,
36};
37use std::borrow::{Borrow, Cow};
38use std::collections::HashSet;
39use std::fmt::Write as _;
40
41/// Maximum allowed size of a `meta` value before truncation.
42const MAX_META_VALUE_LEN: usize = 25_000;
43/// Suffix appended when a `meta` value is truncated.
44const TRUNCATION_SUFFIX: &str = "...";
45
46/// # Why are we doing this?
47///
48/// The JSON agentless format is different from the in-memory model of v04 spans (and to v03/v04/v05
49/// on-the-wire schemas) In order to not have to copy to intermediary structs, we have to write a
50/// manual encoder. For JSON there is no widely available JSON emitter in rust other than serde
51/// JSON. But serde does not let us drive serialization other than through the serde::Serialize
52/// trait.
53///
54/// Defining structs implementing serde::Serialize for every nested object in the span is heavy,
55///
56/// This macro captures parameters from the environment and creates a local struct implementing
57/// serde::Serialize, with a custom implementation.
58///
59/// # Usage
60///
61/// The shape of the input of the macro is made to look like a closure
62/// Contrary to a closure, the names of the types have to be named in full
63/// ```ignore
64///        Optional generic| serializer|
65///            parameter   |.  |               Captured variables from env
66///         --------------  --- ---------------------------------------------------------
67/// ser_fn!(<T: TraceData> |ser, traces: &'a [Vec<Span<T>>], metadata: &'a TracerMetadata| {
68///     // Body of the closure
69/// }
70/// ```
71macro_rules! ser_fn {
72    ($(<$generic:ident $(: $bound:ident )?>)? |$serializer:ident , $($captured:ident : $ty:ty),+ $(,)?| { $($body:tt)* }) => {
73        {
74            struct SerializeClosure<'a, $($generic $(: $bound + 'a)? ,)?>(($(&'a $ty ,)*));
75
76            impl <'a, $($generic $(: $bound + 'a)?,)?> serde::Serialize for SerializeClosure<'a, $($generic,)?> {
77                #[inline]
78                fn serialize<S: serde::Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
79                    let captured = self.0;
80                    (|$serializer: S , ($(& $captured, )*) : ($(&'a $ty ,)*)| {
81                        $($body)*
82                    })(serializer, captured)
83                }
84            }
85
86            SerializeClosure(($(& $captured ,)*))
87        }
88    }
89}
90
91/// Encode the given `traces` to the agentless JSON payload (`/v1/input` body).
92///
93/// When `client_side_stats` is `true`, the encoder will **not** inject
94/// `meta["_dd.compute_stats"]="1"` on the first span of each chunk.  Set this
95/// when the caller is already computing and exporting stats locally so that the intake does not
96/// double- count the same traces.
97///
98/// Returns the serialized JSON bytes on success.
99pub fn encode_payload<T: TraceData>(
100    traces: &[Vec<Span<T>>],
101    metadata: &TracerMetadata,
102    client_side_stats: bool,
103) -> Result<Vec<u8>, serde_json::Error> {
104    let mut bytes = Vec::new();
105    let mut serializer = serde_json::Serializer::new(&mut bytes);
106
107    let mut map_ser = serializer.serialize_map(Some(1))?;
108    map_ser.serialize_entry(
109        "traces",
110        &ser_fn!(<T: TraceData> |ser, traces: &'a [Vec<Span<T>>], metadata: &'a TracerMetadata, client_side_stats: bool| {
111            let mut traces_serializer = ser.serialize_seq(Some(traces.len()))?;
112            for chunk in traces {
113                traces_serializer.serialize_element(&ser_fn!(<T: TraceData> |ser, chunk: &'a Vec<Span<T>>, metadata: &'a TracerMetadata, client_side_stats: bool| {
114                    encode_trace(ser, chunk, metadata, client_side_stats)
115                }))?;
116            }
117            traces_serializer.end()
118        }),
119    )?;
120    SerializeMap::end(map_ser)?;
121    Ok(bytes)
122}
123
124fn encode_trace<T: TraceData, S: Serializer>(
125    ser: S,
126    chunk: &[Span<T>],
127    metadata: &TracerMetadata,
128    client_side_stats: bool,
129) -> Result<S::Ok, S::Error> {
130    let mut map = ser.serialize_map(None)?;
131
132    // Per-trace metadata. Always include hostname; other fields when set.
133    map.serialize_entry("hostname", &metadata.hostname)?;
134    if !metadata.env.is_empty() {
135        map.serialize_entry("env", &metadata.env)?;
136    }
137    if !metadata.language.is_empty() {
138        map.serialize_entry("languageName", &metadata.language)?;
139    }
140    if !metadata.language_version.is_empty() {
141        map.serialize_entry("languageVersion", &metadata.language_version)?;
142    }
143    if !metadata.tracer_version.is_empty() {
144        map.serialize_entry("tracerVersion", &metadata.tracer_version)?;
145    }
146    if !metadata.runtime_id.is_empty() {
147        map.serialize_entry("runtimeID", &metadata.runtime_id)?;
148    }
149    if let Some(container_id) = libdd_common::entity_id::get_container_id() {
150        map.serialize_entry("containerID", container_id)?;
151    }
152
153    map.serialize_entry(
154        "spans",
155        &ser_fn!(<T: TraceData> |ser, chunk: &'a [Span<T>], client_side_stats: bool| {
156            let mut seq = ser.serialize_seq(Some(chunk.len()))?;
157            for (i, span) in chunk.iter().enumerate() {
158                let is_first = i == 0;
159                seq.serialize_element(&ser_fn!(<T: TraceData> |ser, span: &'a Span<T>, is_first: bool, client_side_stats: bool| {
160                    encode_span(ser, span, is_first, client_side_stats)
161                }))?;
162            }
163            seq.end()
164        }),
165    )?;
166
167    map.end()
168}
169
170fn encode_span<T: TraceData, S: Serializer>(
171    ser: S,
172    span: &Span<T>,
173    is_first_in_trace: bool,
174    client_side_stats: bool,
175) -> Result<S::Ok, S::Error> {
176    let mut map = ser.serialize_map(None)?;
177
178    let trace_id = span.trace_id;
179    map.serialize_entry(
180        "trace_id",
181        &ser_fn!(|ser, trace_id: u128| {
182            ser.collect_str(&format_args!("{:016x}", trace_id as u64))
183        }),
184    )?;
185    let span_id = span.span_id;
186    map.serialize_entry(
187        "span_id",
188        &ser_fn!(|ser, span_id: u64| { ser.collect_str(&format_args!("{:016x}", span_id as u64)) }),
189    )?;
190    let parent_id = span.parent_id;
191    map.serialize_entry(
192        "parent_id",
193        &ser_fn!(|ser, parent_id: u64| {
194            ser.collect_str(&format_args!("{:016x}", parent_id as u64))
195        }),
196    )?;
197
198    // Resource defaults to name when empty.
199    let name_str: &str = span.name.borrow();
200    let resource_str: &str = span.resource.borrow();
201    let service_str: &str = span.service.borrow();
202    map.serialize_entry("name", name_str)?;
203    map.serialize_entry(
204        "resource",
205        if resource_str.is_empty() {
206            name_str
207        } else {
208            resource_str
209        },
210    )?;
211    map.serialize_entry("service", service_str)?;
212    map.serialize_entry("error", &span.error)?;
213    map.serialize_entry("start", &span.start.max(0))?;
214    map.serialize_entry("duration", &span.duration)?;
215
216    let type_str: &str = span.r#type.borrow();
217    if !type_str.is_empty() {
218        map.serialize_entry("type", type_str)?;
219    }
220
221    map.serialize_entry(
222        "meta",
223        &ser_fn!(<T: TraceData> |ser, span: &'a Span<T>, is_first_in_trace: bool, client_side_stats: bool| {
224            let upper_bits = (span.trace_id >> 64) as u64;
225            let mut p_tid_seen = false;
226            let mut span_links_seen = false;
227            let mut events_seen = false;
228            let mut compute_stats_seen = false;
229
230            let mut meta = ser.serialize_map(None)?;
231            for (k, v) in span.meta.defensive_dedup().iter() {
232                let key: &str = k.borrow();
233                match key {
234                    "_dd.p.tid" => p_tid_seen = true,
235                    "_dd.span_links" => span_links_seen = true,
236                    "events" => events_seen = true,
237                    "_dd.compute_stats"=> compute_stats_seen = true,
238                    _ => {}
239                };
240                let val: &str = v.borrow();
241                meta.serialize_entry(key, val)?;
242            }
243            if !p_tid_seen && upper_bits != 0 {
244                meta.serialize_entry(
245                    "_dd.p.tid",
246                    &ser_fn!(|ser, upper_bits: u64| {
247                        ser.collect_str(&format_args!("{:016x}", upper_bits as u64))
248                    }),
249                )?;
250            }
251            if !span_links_seen && !span.span_links.is_empty() {
252                if let Some(s) = serialize_span_links(&span.span_links) {
253                    meta.serialize_entry("_dd.span_links", &s)?;
254                }
255            }
256            if !events_seen && !span.span_events.is_empty() {
257                if let Some(s) = serialize_span_events(&span.span_events) {
258                    meta.serialize_entry("events", &s)?;
259                }
260            }
261            if !compute_stats_seen && is_first_in_trace && !client_side_stats {
262                meta.serialize_entry("_dd.compute_stats", "1")?;
263            }
264            meta.end()
265        }),
266    )?;
267
268    map.serialize_entry(
269        "metrics",
270        &ser_fn!(<T: TraceData> |ser, span: &'a Span<T>| {
271            let mut metrics = ser.serialize_map(None)?;
272            let mut trace_root_seen = false;
273            for (k, v) in span.metrics.defensive_dedup().iter() {
274                let key: &str = k.borrow();
275                // serde_json refuses to serialize NaN/Inf; drop them silently.
276                if v.is_finite() {
277                    match key {
278                        "_trace_root" => trace_root_seen = true,
279                        "_top_level" => {
280                            metrics.serialize_entry(key, &(*v as u32))?;
281                            continue
282                        },
283                        _ => {},
284                    }
285                    metrics.serialize_entry(key, v)?
286                }
287            }
288            if !trace_root_seen && span.parent_id == 0 {
289                metrics.serialize_entry("_trace_root", &1u32)?;
290            }
291            metrics.end()
292        }),
293    )?;
294
295    if !span.meta_struct.is_empty() {
296        map.serialize_entry(
297            "meta_struct",
298            &ser_fn!(<T: TraceData> |ser, span: &'a Span<T>| {
299                let mut ms = ser.serialize_map(None)?;
300                for (k, v) in span.meta_struct.iter() {
301                    let key: &str = k.borrow();
302                    let bytes: &[u8] = v.borrow();
303
304                    // abort whole payload on malformed entry
305                    ms.serialize_entry(key, &MsgpackAsJson(bytes))?;
306                }
307                ms.end()
308            }),
309        )?;
310    }
311    map.end()
312}
313
314/// Serialize span links to a JSON string suitable for `meta['_dd.span_links']`.
315///
316/// Returns `None` if serialization fails. The result is truncated to
317/// [`MAX_META_VALUE_LEN`] characters with a trailing `"..."` if it would
318/// otherwise exceed that limit.
319fn serialize_span_links<T: TraceData>(links: &[SpanLink<T>]) -> Option<String> {
320    let s = serde_json::to_string(&ser_fn!(<T: TraceData> |ser, links: &'a [SpanLink<T>]| {
321        let mut seq = ser.serialize_seq(Some(links.len()))?;
322        for link in links {
323            seq.serialize_element(&ser_fn!(<T: TraceData> |ser, link: &'a SpanLink<T>| {
324                encode_span_link(ser, link)
325            }))?;
326        }
327        seq.end()
328    }))
329    .ok()?;
330    Some(truncate_with_ellipsis(s, MAX_META_VALUE_LEN))
331}
332
333fn encode_span_link<T: TraceData, S: Serializer>(
334    ser: S,
335    link: &SpanLink<T>,
336) -> Result<S::Ok, S::Error> {
337    let mut map = ser.serialize_map(None)?;
338    let trace_id_128: u128 = ((link.trace_id_high as u128) << 64) | (link.trace_id as u128);
339    map.serialize_entry("trace_id", &format!("{:032x}", trace_id_128))?;
340    map.serialize_entry("span_id", &format!("{:016x}", link.span_id))?;
341    if !link.attributes.is_empty() {
342        map.serialize_entry(
343            "attributes",
344            &ser_fn!(<T: TraceData> |ser, link: &'a SpanLink<T>| {
345                let mut attrs = ser.serialize_map(Some(link.attributes.len()))?;
346                for (k, v) in link.attributes.iter() {
347                    let key: &str = k.borrow();
348                    let val: &str = v.borrow();
349                    attrs.serialize_entry(key, val)?;
350                }
351                attrs.end()
352            }),
353        )?;
354    }
355    // When `flags` is 0, no sampling decision exists, so omit the field. Before emission,
356    // mask off the internal "explicitly set" sentinel (bit 31), because this JSON field uses
357    // the same `_dd.span_links` key that the v0.5 encoder produces and must match its output.
358    if link.flags != 0 {
359        map.serialize_entry(
360            "flags",
361            &((link.flags & !SPAN_LINK_FLAGS_SET_SENTINEL) as u64),
362        )?;
363    }
364    let tracestate: &str = link.tracestate.borrow();
365    if !tracestate.is_empty() {
366        map.serialize_entry("tracestate", tracestate)?;
367    }
368    map.end()
369}
370
371/// Serialize span events to a JSON string suitable for `meta['events']`.
372fn serialize_span_events<T: TraceData>(events: &[SpanEvent<T>]) -> Option<String> {
373    let s = serde_json::to_string(&ser_fn!(<T: TraceData> |ser, events: &'a [SpanEvent<T>]| {
374        let mut seq = ser.serialize_seq(Some(events.len()))?;
375        for event in events {
376            seq.serialize_element(&ser_fn!(<T: TraceData> |ser, event: &'a SpanEvent<T>| {
377                encode_span_event(ser, event)
378            }))?;
379        }
380        seq.end()
381    }))
382    .ok()?;
383    Some(truncate_with_ellipsis(s, MAX_META_VALUE_LEN))
384}
385
386fn encode_span_event<T: TraceData, S: Serializer>(
387    ser: S,
388    event: &SpanEvent<T>,
389) -> Result<S::Ok, S::Error> {
390    let mut map = ser.serialize_map(None)?;
391    let name: &str = event.name.borrow();
392    map.serialize_entry("name", name)?;
393    map.serialize_entry("time_unix_nano", &event.time_unix_nano)?;
394    if !event.attributes.is_empty() {
395        map.serialize_entry(
396            "attributes",
397            &ser_fn!(<T: TraceData> |ser, event: &'a SpanEvent<T>| {
398                let mut attrs = ser.serialize_map(Some(event.attributes.len()))?;
399                for (k, v) in event.attributes.iter() {
400                    let key: &str = k.borrow();
401                    attrs.serialize_entry(key, &ser_fn!(<T: TraceData> |ser, v: &'a AttributeAnyValue<T> | {
402                        match v {
403                            AttributeAnyValue::SingleValue(v) => serialize_scalar(ser, v),
404                            AttributeAnyValue::Array(values) => {
405                                let mut seq = ser.serialize_seq(Some(values.len()))?;
406                                for v in values {
407                                    seq.serialize_element(&ser_fn!(<T: TraceData> |ser, v: &'a AttributeArrayValue<T>| {
408                                        serialize_scalar(ser, v)
409                                    }))?;
410                                }
411                                seq.end()
412                            }
413                        }
414                    }))?;
415                }
416                attrs.end()
417            }),
418        )?;
419    }
420    map.end()
421}
422
423fn serialize_scalar<S: serde::Serializer, T: TraceData>(
424    ser: S,
425    s: &AttributeArrayValue<T>,
426) -> Result<S::Ok, S::Error> {
427    match s {
428        AttributeArrayValue::String(s) => {
429            let s: &str = s.borrow();
430            ser.serialize_str(s)
431        }
432        AttributeArrayValue::Boolean(b) => ser.serialize_bool(*b),
433        AttributeArrayValue::Integer(i) => ser.serialize_i64(*i),
434        AttributeArrayValue::Double(d) => {
435            if d.is_finite() {
436                ser.serialize_f64(*d)
437            } else {
438                // NaN/Inf become JSON null.
439                ser.serialize_unit()
440            }
441        }
442    }
443}
444
445/// Reserved v0.4 `meta`/`metrics` key names written from dedicated typed fields (`env`, chunk
446/// `origin`, ...) rather than from the v1 attribute map — the dedicated field always wins and
447/// a colliding attribute is dropped. See
448/// [`crate::msgpack_encoder::v04::span_v1::PROMOTED_ATTR_KEYS`] for the sibling list on the
449/// msgpack side (this one omits `_dd.p.tid`, which is handled below with "seen" tracking
450/// instead, matching this module's own `_dd.p.tid`/`_dd.span_links`/`events` convention).
451const PROMOTED_ATTR_KEYS_V1: &[&str] = &[
452    "env",
453    "version",
454    "component",
455    "span.kind",
456    "_dd.origin",
457    "_dd.p.dm",
458    "_sampling_priority_v1",
459];
460
461/// Maps a `SpanKind` to its v0.4 `span.kind` meta string. Returns `None` for `Internal` so
462/// callers can skip emitting the default value.
463fn span_kind_to_meta_v1(kind: v1::SpanKind) -> Option<&'static str> {
464    match kind {
465        v1::SpanKind::Internal => None,
466        v1::SpanKind::Server => Some("server"),
467        v1::SpanKind::Client => Some("client"),
468        v1::SpanKind::Producer => Some("producer"),
469        v1::SpanKind::Consumer => Some("consumer"),
470    }
471}
472
473/// Drops entries whose key was already seen, keeping the first occurrence: two distinct
474/// attributes can flatten to the same dotted key.
475fn dedup_first_wins_v1<'a, V>(leaves: &mut Vec<(Cow<'a, str>, V)>) {
476    let keep: Vec<bool> = {
477        let mut seen: HashSet<&str> = HashSet::with_capacity(leaves.len());
478        leaves
479            .iter()
480            .map(|(k, _)| seen.insert(k.as_ref()))
481            .collect()
482    };
483    let mut keep = keep.into_iter();
484    leaves.retain(|_| keep.next().unwrap_or(false));
485}
486
487/// Recursively flattens a `List`/`KeyValue` attribute into dotted-key leaf entries for `meta`,
488/// `metrics`, and `meta_struct`.
489///
490/// Only called for nested attributes, so the key is always freshly built and owned — unlike the
491/// top-level scalar fast path in [`collect_attrs_v1`], which borrows straight from the source.
492///
493/// `key` is a reused buffer: pushed to on the way down, truncated back on the way up, so it's
494/// unchanged once the call returns.
495///
496/// ## Examples
497/// ```text
498/// key="a",    v=KeyValue{b: "v"}           =>  meta_out    += ("a.b", "v")
499///
500/// key="list", v=List[String("x"), Int(2)]  =>  meta_out    += ("list.0", "x")
501///                                              metrics_out += ("list.1", 2.0)
502/// ```
503fn flatten_attr_into_v1<'a, T: TraceData>(
504    key: &mut String,
505    v: &'a v1::AttributeValue<T>,
506    meta_out: &mut Vec<(Cow<'a, str>, Cow<'a, str>)>,
507    metrics_out: &mut Vec<(Cow<'a, str>, f64)>,
508    bytes_out: &mut Vec<(Cow<'a, str>, T::Bytes)>,
509) {
510    match v {
511        v1::AttributeValue::String(s) => {
512            meta_out.push((Cow::Owned(key.clone()), Cow::Borrowed(s.borrow())))
513        }
514        v1::AttributeValue::Bool(b) => meta_out.push((
515            Cow::Owned(key.clone()),
516            Cow::Borrowed(if *b { "true" } else { "false" }),
517        )),
518        v1::AttributeValue::Int(i) => metrics_out.push((Cow::Owned(key.clone()), *i as f64)),
519        v1::AttributeValue::Float(f) => metrics_out.push((Cow::Owned(key.clone()), *f)),
520        v1::AttributeValue::Bytes(b) => bytes_out.push((Cow::Owned(key.clone()), b.clone())),
521        v1::AttributeValue::List(items) => {
522            let base_len = key.len();
523            for (i, item) in items.iter().enumerate() {
524                key.push('.');
525                let _ = write!(key, "{i}");
526                flatten_attr_into_v1(key, item, meta_out, metrics_out, bytes_out);
527                key.truncate(base_len);
528            }
529        }
530        v1::AttributeValue::KeyValue(map) => {
531            let base_len = key.len();
532            for (k, v) in map.defensive_dedup().iter() {
533                key.push('.');
534                key.push_str(k.borrow());
535                flatten_attr_into_v1(key, v, meta_out, metrics_out, bytes_out);
536                key.truncate(base_len);
537            }
538        }
539    }
540}
541
542/// Leaves collected by [`collect_attrs_v1`]: `meta`, `metrics`, and `meta_struct` (`Bytes`)
543/// entries. Keys/values borrow from the source span/chunk attributes wherever possible (the
544/// common top-level scalar case); only entries produced by flattening a nested `List`/`KeyValue`
545/// need an owned, freshly built dotted key.
546struct CollectedAttrsV1<'a, T: TraceData> {
547    meta: Vec<(Cow<'a, str>, Cow<'a, str>)>,
548    metrics: Vec<(Cow<'a, str>, f64)>,
549    meta_struct: Vec<(Cow<'a, str>, T::Bytes)>,
550}
551
552/// Merges a span's attributes with its chunk's (span overrides chunk on key collision),
553/// drops attributes colliding with a [`PROMOTED_ATTR_KEYS_V1`] name, and splits the rest into
554/// `meta` leaves, `metrics` leaves, and `meta_struct` (`Bytes`) leaves.
555fn collect_attrs_v1<'a, T: TraceData>(
556    span: &'a v1::Span<T>,
557    chunk: &'a v1::TraceChunk<T>,
558) -> CollectedAttrsV1<'a, T> {
559    let span_attrs_dd = span.attributes.defensive_dedup();
560    let chunk_attrs_dd = chunk.attributes.defensive_dedup();
561    let merged_attrs = span_attrs_dd
562        .iter()
563        .filter(|(k, _)| !PROMOTED_ATTR_KEYS_V1.contains(&(*k).borrow()))
564        .chain(chunk_attrs_dd.iter().filter(|(k, _)| {
565            !PROMOTED_ATTR_KEYS_V1.contains(&(*k).borrow())
566                && !span_attrs_dd.iter().any(|(k2, _)| k2 == *k)
567        }));
568
569    let mut meta_leaves: Vec<(Cow<'a, str>, Cow<'a, str>)> = Vec::new();
570    let mut metrics_leaves: Vec<(Cow<'a, str>, f64)> = Vec::new();
571    let mut bytes_leaves: Vec<(Cow<'a, str>, T::Bytes)> = Vec::new();
572    let mut key_buf = String::new();
573    for (k, v) in merged_attrs {
574        match v {
575            // Common case: a top-level scalar attribute maps 1:1 onto a leaf, so its key/value
576            // can be borrowed straight from the source attribute — no allocation.
577            v1::AttributeValue::String(s) => {
578                meta_leaves.push((Cow::Borrowed(k.borrow()), Cow::Borrowed(s.borrow())))
579            }
580            v1::AttributeValue::Bool(b) => meta_leaves.push((
581                Cow::Borrowed(k.borrow()),
582                Cow::Borrowed(if *b { "true" } else { "false" }),
583            )),
584            v1::AttributeValue::Int(i) => {
585                metrics_leaves.push((Cow::Borrowed(k.borrow()), *i as f64))
586            }
587            v1::AttributeValue::Float(f) => metrics_leaves.push((Cow::Borrowed(k.borrow()), *f)),
588            v1::AttributeValue::Bytes(b) => {
589                bytes_leaves.push((Cow::Borrowed(k.borrow()), b.clone()))
590            }
591            // Nested case: the leaf key has to be built (`key.0`, `key.a.b`, ...), so it can no
592            // longer borrow the original attribute name alone.
593            v1::AttributeValue::List(_) | v1::AttributeValue::KeyValue(_) => {
594                key_buf.clear();
595                key_buf.push_str(k.borrow());
596                flatten_attr_into_v1(
597                    &mut key_buf,
598                    v,
599                    &mut meta_leaves,
600                    &mut metrics_leaves,
601                    &mut bytes_leaves,
602                );
603            }
604        }
605    }
606    // A nested attribute can flatten into a promoted name (e.g. `span = {kind: "client"}` ->
607    // `span.kind`) even though its unflattened top-level key wasn't caught by the filter above
608    // — drop those too, so the dedicated field always wins as documented.
609    meta_leaves.retain(|(k, _)| !PROMOTED_ATTR_KEYS_V1.contains(&k.as_ref()));
610    metrics_leaves.retain(|(k, _)| !PROMOTED_ATTR_KEYS_V1.contains(&k.as_ref()));
611
612    dedup_first_wins_v1(&mut meta_leaves);
613    dedup_first_wins_v1(&mut metrics_leaves);
614    dedup_first_wins_v1(&mut bytes_leaves);
615
616    CollectedAttrsV1 {
617        meta: meta_leaves,
618        metrics: metrics_leaves,
619        meta_struct: bytes_leaves,
620    }
621}
622
623/// V1-native analog of [`encode_payload`]. Downgrades v1's unified attribute model back to the
624/// same `meta`/`metrics`/`meta_struct`-shaped wire fields, so the emitted JSON body is
625/// equivalent to what a v0.4 tracer would produce for the same trace — see
626/// [`crate::msgpack_encoder::v04::span_v1`] for the mapping table this mirrors. Chunk-level
627/// context (`trace_id`, `origin`, `priority`, `sampling_mechanism`, `dropped_trace`, chunk
628/// attributes) is propagated into every span, matching the [`v1::TraceChunk`]-level granularity v1
629/// operates at.
630pub fn encode_payload_from_v1<T: TraceData>(
631    chunks: &[v1::TraceChunk<T>],
632    metadata: &TracerMetadata,
633) -> Result<Vec<u8>, serde_json::Error> {
634    let mut bytes = Vec::new();
635    let mut serializer = serde_json::Serializer::new(&mut bytes);
636
637    let mut map_ser = serializer.serialize_map(Some(1))?;
638    map_ser.serialize_entry(
639        "traces",
640        &ser_fn!(<T: TraceData> |ser, chunks: &'a [v1::TraceChunk<T>], metadata: &'a TracerMetadata| {
641            let mut traces_serializer = ser.serialize_seq(Some(chunks.len()))?;
642            for chunk in chunks {
643                traces_serializer.serialize_element(&ser_fn!(<T: TraceData> |ser, chunk: &'a v1::TraceChunk<T>, metadata: &'a TracerMetadata| {
644                    encode_trace_v1(ser, chunk, metadata)
645                }))?;
646            }
647            traces_serializer.end()
648        }),
649    )?;
650    SerializeMap::end(map_ser)?;
651    Ok(bytes)
652}
653
654fn encode_trace_v1<T: TraceData, S: Serializer>(
655    ser: S,
656    chunk: &v1::TraceChunk<T>,
657    metadata: &TracerMetadata,
658) -> Result<S::Ok, S::Error> {
659    let container_id = libdd_common::entity_id::get_container_id();
660    let len = 2 // hostname + spans
661        + usize::from(!metadata.env.is_empty())
662        + usize::from(!metadata.language.is_empty())
663        + usize::from(!metadata.language_version.is_empty())
664        + usize::from(!metadata.tracer_version.is_empty())
665        + usize::from(!metadata.runtime_id.is_empty())
666        + usize::from(container_id.is_some());
667    let mut map = ser.serialize_map(Some(len))?;
668
669    map.serialize_entry("hostname", &metadata.hostname)?;
670    if !metadata.env.is_empty() {
671        map.serialize_entry("env", &metadata.env)?;
672    }
673    if !metadata.language.is_empty() {
674        map.serialize_entry("languageName", &metadata.language)?;
675    }
676    if !metadata.language_version.is_empty() {
677        map.serialize_entry("languageVersion", &metadata.language_version)?;
678    }
679    if !metadata.tracer_version.is_empty() {
680        map.serialize_entry("tracerVersion", &metadata.tracer_version)?;
681    }
682    if !metadata.runtime_id.is_empty() {
683        map.serialize_entry("runtimeID", &metadata.runtime_id)?;
684    }
685    if let Some(container_id) = container_id {
686        map.serialize_entry("containerID", container_id)?;
687    }
688
689    map.serialize_entry(
690        "spans",
691        &ser_fn!(<T: TraceData> |ser, chunk: &'a v1::TraceChunk<T>| {
692            let mut seq = ser.serialize_seq(Some(chunk.spans.len()))?;
693            for (i, span) in chunk.spans.iter().enumerate() {
694                let is_first = i == 0;
695                seq.serialize_element(&ser_fn!(<T: TraceData> |ser, chunk: &'a v1::TraceChunk<T>, span: &'a v1::Span<T>, is_first: bool| {
696                    encode_span_v1(ser, chunk, span, is_first)
697                }))?;
698            }
699            seq.end()
700        }),
701    )?;
702
703    map.end()
704}
705
706fn encode_span_v1<'a, T: TraceData, S: Serializer>(
707    ser: S,
708    chunk: &'a v1::TraceChunk<T>,
709    span: &'a v1::Span<T>,
710    is_first_in_trace: bool,
711) -> Result<S::Ok, S::Error> {
712    let mut map = ser.serialize_map(None)?;
713
714    let mut trace_id_low_bytes = [0u8; 8];
715    let mut trace_id_high_bytes = [0u8; 8];
716    trace_id_low_bytes.copy_from_slice(&chunk.trace_id[8..16]);
717    trace_id_high_bytes.copy_from_slice(&chunk.trace_id[0..8]);
718    let trace_id_low = u64::from_be_bytes(trace_id_low_bytes);
719    let trace_id_high = u64::from_be_bytes(trace_id_high_bytes);
720    map.serialize_entry(
721        "trace_id",
722        &ser_fn!(|ser, trace_id_low: u64| {
723            ser.collect_str(&format_args!("{trace_id_low:016x}"))
724        }),
725    )?;
726    let span_id = span.span_id;
727    map.serialize_entry(
728        "span_id",
729        &ser_fn!(|ser, span_id: u64| { ser.collect_str(&format_args!("{span_id:016x}")) }),
730    )?;
731    let parent_id = span.parent_id;
732    map.serialize_entry(
733        "parent_id",
734        &ser_fn!(|ser, parent_id: u64| { ser.collect_str(&format_args!("{parent_id:016x}")) }),
735    )?;
736
737    // Resource defaults to name when empty.
738    let name_str: &str = span.name.borrow();
739    let resource_str: &str = span.resource.borrow();
740    let service_str: &str = span.service.borrow();
741    map.serialize_entry("name", name_str)?;
742    map.serialize_entry(
743        "resource",
744        if resource_str.is_empty() {
745            name_str
746        } else {
747            resource_str
748        },
749    )?;
750    map.serialize_entry("service", service_str)?;
751    // v0.4's `error` is emitted as an integer (0/1) on the wire, unlike v1's own `bool` field.
752    map.serialize_entry("error", &(span.error as i32))?;
753    map.serialize_entry("start", &span.start.max(0))?;
754    map.serialize_entry("duration", &span.duration)?;
755
756    let type_str: &str = span.r#type.borrow();
757    if !type_str.is_empty() {
758        map.serialize_entry("type", type_str)?;
759    }
760
761    let collected = collect_attrs_v1(span, chunk);
762    let meta_leaves = &collected.meta;
763    let metrics_leaves = &collected.metrics;
764    let bytes_attrs = &collected.meta_struct;
765    let priority = if chunk.dropped_trace {
766        // v0.4 has no wire-level equivalent of `dropped_trace`; force `USER_REJECT` (-1)
767        // unless the chunk already carries a negative (reject-like) priority — same
768        // convention as the msgpack downgrade encoder.
769        Some(chunk.priority.filter(|&p| p < 0).unwrap_or(-1))
770    } else {
771        chunk.priority
772    };
773
774    map.serialize_entry(
775        "meta",
776        &ser_fn!(<T: TraceData> |ser, span: &'a v1::Span<T>, chunk: &'a v1::TraceChunk<T>, meta_leaves: &'a Vec<(Cow<'a, str>, Cow<'a, str>)>, is_first_in_trace: bool, trace_id_high: u64| {
777            let mut meta = ser.serialize_map(None)?;
778
779            let env: &str = span.env.borrow();
780            if !env.is_empty() {
781                meta.serialize_entry("env", env)?;
782            }
783            let version: &str = span.version.borrow();
784            if !version.is_empty() {
785                meta.serialize_entry("version", version)?;
786            }
787            let component: &str = span.component.borrow();
788            if !component.is_empty() {
789                meta.serialize_entry("component", component)?;
790            }
791            if let Some(kind) = span_kind_to_meta_v1(span.span_kind) {
792                meta.serialize_entry("span.kind", kind)?;
793            }
794            let origin: &str = chunk.origin.borrow();
795            if !origin.is_empty() {
796                meta.serialize_entry("_dd.origin", origin)?;
797            }
798            if let Some(mechanism) = chunk.sampling_mechanism {
799                let mut buf = itoa::Buffer::new();
800                meta.serialize_entry("_dd.p.dm", buf.format(-(mechanism as i64)))?;
801            }
802
803            let mut p_tid_seen = false;
804            let mut span_links_seen = false;
805            let mut events_seen = false;
806            let mut compute_stats_seen = false;
807            for (key, val) in meta_leaves.iter() {
808                match key.as_ref() {
809                    "_dd.p.tid" => p_tid_seen = true,
810                    "_dd.span_links" => span_links_seen = true,
811                    "events" => events_seen = true,
812                    "_dd.compute_stats" => compute_stats_seen = true,
813                    _ => {}
814                };
815                meta.serialize_entry(key, val)?;
816            }
817            if !p_tid_seen && trace_id_high != 0 {
818                meta.serialize_entry(
819                    "_dd.p.tid",
820                    &ser_fn!(|ser, trace_id_high: u64| {
821                        ser.collect_str(&format_args!("{trace_id_high:016x}"))
822                    }),
823                )?;
824            }
825            if !span_links_seen && !span.span_links.is_empty() {
826                if let Some(s) = serialize_span_links_v1(&span.span_links) {
827                    meta.serialize_entry("_dd.span_links", &s)?;
828                }
829            }
830            if !events_seen && !span.span_events.is_empty() {
831                if let Some(s) = serialize_span_events_v1(&span.span_events) {
832                    meta.serialize_entry("events", &s)?;
833                }
834            }
835            if !compute_stats_seen && is_first_in_trace {
836                meta.serialize_entry("_dd.compute_stats", "1")?;
837            }
838            meta.end()
839        }),
840    )?;
841
842    map.serialize_entry(
843        "metrics",
844        &ser_fn!(<T: TraceData> |ser, span: &'a v1::Span<T>, metrics_leaves: &'a Vec<(Cow<'a, str>, f64)>, priority: Option<i32>| {
845            let mut metrics = ser.serialize_map(None)?;
846            let mut trace_root_seen = false;
847            for (key, val) in metrics_leaves.iter() {
848                // serde_json refuses to serialize NaN/Inf; drop them silently.
849                if !val.is_finite() {
850                    continue;
851                }
852                match key.as_ref() {
853                    "_trace_root" => trace_root_seen = true,
854                    "_top_level" => {
855                        metrics.serialize_entry(key, &(*val as u32))?;
856                        continue;
857                    }
858                    _ => {}
859                }
860                metrics.serialize_entry(key, val)?;
861            }
862            if let Some(p) = priority {
863                metrics.serialize_entry("_sampling_priority_v1", &(p as f64))?;
864            }
865            if !trace_root_seen && span.parent_id == 0 {
866                metrics.serialize_entry("_trace_root", &1u32)?;
867            }
868            metrics.end()
869        }),
870    )?;
871
872    if !bytes_attrs.is_empty() {
873        map.serialize_entry(
874            "meta_struct",
875            &ser_fn!(<T: TraceData> |ser, span: &'a v1::Span<T>, bytes_attrs: &'a Vec<(Cow<'a, str>, T::Bytes)>| {
876                let _ = span;
877                let mut ms = ser.serialize_map(None)?;
878                for (k, v) in bytes_attrs.iter() {
879                    let raw: &[u8] = v.borrow();
880                    ms.serialize_entry(k, &MsgpackAsJson(raw))?;
881                }
882                ms.end()
883            }),
884        )?;
885    }
886    map.end()
887}
888
889/// Serialize v1 span links to a JSON string suitable for `meta['_dd.span_links']`. Same
890/// truncation convention as [`serialize_span_links`].
891fn serialize_span_links_v1<T: TraceData>(links: &[v1::SpanLink<T>]) -> Option<String> {
892    let s = serde_json::to_string(
893        &ser_fn!(<T: TraceData> |ser, links: &'a [v1::SpanLink<T>]| {
894            let mut seq = ser.serialize_seq(Some(links.len()))?;
895            for link in links {
896                seq.serialize_element(&ser_fn!(<T: TraceData> |ser, link: &'a v1::SpanLink<T>| {
897                    encode_span_link_v1(ser, link)
898                }))?;
899            }
900            seq.end()
901        }),
902    )
903    .ok()?;
904    Some(truncate_with_ellipsis(s, MAX_META_VALUE_LEN))
905}
906
907fn encode_span_link_v1<T: TraceData, S: Serializer>(
908    ser: S,
909    link: &v1::SpanLink<T>,
910) -> Result<S::Ok, S::Error> {
911    let mut map = ser.serialize_map(None)?;
912    let trace_id_128 = u128::from_be_bytes(link.trace_id);
913    map.serialize_entry("trace_id", &format!("{trace_id_128:032x}"))?;
914    map.serialize_entry("span_id", &format!("{:016x}", link.span_id))?;
915    let attrs_dd = link.attributes.defensive_dedup();
916    let attrs_dd = &attrs_dd;
917    let has_attributes = attrs_dd.iter().any(|(_, v)| {
918        matches!(
919            v,
920            v1::AttributeValue::String(_) | v1::AttributeValue::Bool(_)
921        )
922    });
923    if has_attributes {
924        map.serialize_entry(
925            "attributes",
926            &ser_fn!(<T: TraceData> |ser, attrs_dd: &'a crate::span::vec_map::DedupedVecMap<'a, T::Text, v1::AttributeValue<T>>| {
927                let mut attrs = ser.serialize_map(None)?;
928                for (k, v) in attrs_dd.iter() {
929                    let key: &str = k.borrow();
930                    match v {
931                        v1::AttributeValue::String(s) => attrs.serialize_entry(key, s.borrow() as &str)?,
932                        v1::AttributeValue::Bool(b) => {
933                            attrs.serialize_entry(key, if *b { "true" } else { "false" })?
934                        }
935                        _ => {}
936                    }
937                }
938                attrs.end()
939            }),
940        )?;
941    }
942    // When `flags` is 0, no sampling decision exists, so omit the field. Mask off the internal
943    // "explicitly set" sentinel (bit 31) before emission, same as `encode_span_link`.
944    if link.flags != 0 {
945        map.serialize_entry(
946            "flags",
947            &((link.flags & !SPAN_LINK_FLAGS_SET_SENTINEL) as u64),
948        )?;
949    }
950    let tracestate: &str = link.tracestate.borrow();
951    if !tracestate.is_empty() {
952        map.serialize_entry("tracestate", tracestate)?;
953    }
954    map.end()
955}
956
957/// Serialize v1 span events to a JSON string suitable for `meta['events']`. Same truncation
958/// convention as [`serialize_span_events`].
959fn serialize_span_events_v1<T: TraceData>(events: &[v1::SpanEvent<T>]) -> Option<String> {
960    let s = serde_json::to_string(
961        &ser_fn!(<T: TraceData> |ser, events: &'a [v1::SpanEvent<T>]| {
962            let mut seq = ser.serialize_seq(Some(events.len()))?;
963            for event in events {
964                seq.serialize_element(&ser_fn!(<T: TraceData> |ser, event: &'a v1::SpanEvent<T>| {
965                    encode_span_event_v1(ser, event)
966                }))?;
967            }
968            seq.end()
969        }),
970    )
971    .ok()?;
972    Some(truncate_with_ellipsis(s, MAX_META_VALUE_LEN))
973}
974
975/// Returns `true` when `v` can be downgraded to a v0.4 event-attribute (scalar or scalar list).
976fn is_supported_event_attr_v1<T: TraceData>(v: &v1::AttributeValue<T>) -> bool {
977    matches!(
978        v,
979        v1::AttributeValue::String(_)
980            | v1::AttributeValue::Bool(_)
981            | v1::AttributeValue::Int(_)
982            | v1::AttributeValue::Float(_)
983            | v1::AttributeValue::List(_)
984    )
985}
986
987/// Returns `true` when `v` is a scalar that fits in a v0.4 array element (no nesting).
988fn is_scalar_array_elem_v1<T: TraceData>(v: &v1::AttributeValue<T>) -> bool {
989    matches!(
990        v,
991        v1::AttributeValue::String(_)
992            | v1::AttributeValue::Bool(_)
993            | v1::AttributeValue::Int(_)
994            | v1::AttributeValue::Float(_)
995    )
996}
997
998fn encode_span_event_v1<T: TraceData, S: Serializer>(
999    ser: S,
1000    event: &v1::SpanEvent<T>,
1001) -> Result<S::Ok, S::Error> {
1002    let mut map = ser.serialize_map(None)?;
1003    let name: &str = event.name.borrow();
1004    map.serialize_entry("name", name)?;
1005    map.serialize_entry("time_unix_nano", &event.time_unix_nano)?;
1006    let attrs_dd = event.attributes.defensive_dedup();
1007    let attrs_dd = &attrs_dd;
1008    let has_attributes = attrs_dd.iter().any(|(_, v)| is_supported_event_attr_v1(v));
1009    if has_attributes {
1010        map.serialize_entry(
1011            "attributes",
1012            &ser_fn!(<T: TraceData> |ser, attrs_dd: &'a crate::span::vec_map::DedupedVecMap<'a, T::Text, v1::AttributeValue<T>>| {
1013                let mut attrs = ser.serialize_map(None)?;
1014                for (k, v) in attrs_dd.iter().filter(|(_, v)| is_supported_event_attr_v1(v)) {
1015                    let key: &str = k.borrow();
1016                    attrs.serialize_entry(key, &ser_fn!(<T: TraceData> |ser, v: &'a v1::AttributeValue<T>| {
1017                        encode_event_attr_value_v1(ser, v)
1018                    }))?;
1019                }
1020                attrs.end()
1021            }),
1022        )?;
1023    }
1024    map.end()
1025}
1026
1027/// Serializes a v1 event attribute value as a plain JSON value — same shape as
1028/// [`serialize_scalar`] on the v04-native path, which is what the agentless intake expects.
1029/// `List` produces a plain JSON array, filtering out non-scalar entries (no equivalent for them).
1030fn encode_event_attr_value_v1<T: TraceData, S: Serializer>(
1031    ser: S,
1032    v: &v1::AttributeValue<T>,
1033) -> Result<S::Ok, S::Error> {
1034    match v {
1035        v1::AttributeValue::List(items) => {
1036            let scalars: Vec<_> = items
1037                .iter()
1038                .filter(|e| is_scalar_array_elem_v1(e))
1039                .collect();
1040            let mut seq = ser.serialize_seq(Some(scalars.len()))?;
1041            for elem in scalars {
1042                seq.serialize_element(
1043                    &ser_fn!(<T: TraceData> |ser, elem: &'a v1::AttributeValue<T>| {
1044                        encode_event_scalar_v1(ser, elem)
1045                    }),
1046                )?;
1047            }
1048            seq.end()
1049        }
1050        other => encode_event_scalar_v1(ser, other),
1051    }
1052}
1053
1054fn encode_event_scalar_v1<T: TraceData, S: Serializer>(
1055    ser: S,
1056    v: &v1::AttributeValue<T>,
1057) -> Result<S::Ok, S::Error> {
1058    match v {
1059        v1::AttributeValue::String(s) => {
1060            let s: &str = s.borrow();
1061            ser.serialize_str(s)
1062        }
1063        v1::AttributeValue::Bool(b) => ser.serialize_bool(*b),
1064        v1::AttributeValue::Int(i) => ser.serialize_i64(*i),
1065        v1::AttributeValue::Float(f) => {
1066            if f.is_finite() {
1067                ser.serialize_f64(*f)
1068            } else {
1069                // NaN/Inf become JSON null, matching `serialize_scalar` on the v04-native path.
1070                ser.serialize_unit()
1071            }
1072        }
1073        _ => unreachable!("filtered by is_scalar_array_elem_v1"),
1074    }
1075}
1076
1077/// `serde::Serialize` adapter that interprets `bytes` as a self-describing
1078/// msgpack value and transcodes it into the destination serializer.
1079///
1080/// Used to inline `meta_struct` values (which are stored as msgpack-encoded
1081/// bytes) into the agentless JSON payload as real JSON objects.
1082struct MsgpackAsJson<'a>(&'a [u8]);
1083
1084impl serde::Serialize for MsgpackAsJson<'_> {
1085    fn serialize<S: serde::Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
1086        let mut de = rmp_serde::Deserializer::from_read_ref(self.0);
1087        serde_transcode::transcode(&mut de, serializer)
1088    }
1089}
1090
1091/// Truncate `s` to at most `max_len` bytes, appending `"..."` when truncation occurs.
1092fn truncate_with_ellipsis(mut s: String, max_len: usize) -> String {
1093    if s.len() <= max_len {
1094        return s;
1095    }
1096    let suffix_len = TRUNCATION_SUFFIX.len();
1097    let cut = max_len.saturating_sub(suffix_len);
1098    // Find the previous char boundary so we don't slice in the middle of a UTF-8 sequence.
1099    let mut end = cut;
1100    while end > 0 && !s.is_char_boundary(end) {
1101        end -= 1;
1102    }
1103    s.truncate(end);
1104    s.push_str(TRUNCATION_SUFFIX);
1105    s
1106}
1107
1108#[cfg(test)]
1109mod tests;
1110#[cfg(test)]
1111mod tests_v1;