Skip to main content

libdd_trace_utils/
trace_utils.rs

1// Copyright 2023-Present Datadog, Inc. https://www.datadoghq.com/
2// SPDX-License-Identifier: Apache-2.0
3
4pub use crate::send_data::send_data_result::SendDataResult;
5pub use crate::send_data::SendData;
6use crate::span::v05::dict::SharedDict;
7use crate::span::{v05, TraceData};
8pub use crate::tracer_header_tags::{TracerGenericTags, TracerHeaderTags};
9use crate::tracer_payload::TracerPayloadCollection;
10use crate::tracer_payload::{self, TraceChunks};
11use anyhow::anyhow;
12use bytes::buf::Reader;
13use bytes::Buf;
14use http_body_util::BodyExt;
15use libdd_common::azure_app_services;
16use libdd_trace_normalization::normalizer;
17use libdd_trace_protobuf::pb;
18use rmp::decode::read_array_len;
19use rmpv::decode::read_value;
20use rmpv::{Integer, Value};
21use std::cmp::Ordering;
22use std::collections::{HashMap, HashSet};
23use std::env;
24use tracing::{debug, error};
25
26/// The maximum payload size for a single request that can be sent to the trace agent. Payloads
27/// larger than this size will be dropped and the agent will return a 413 error if
28/// `datadog-send-real-http-status` is set.
29pub const MAX_PAYLOAD_SIZE: usize = 25 * 1024 * 1024;
30/// Span metric the mini agent must set for the backend to recognize top level span
31const TOP_LEVEL_KEY: &str = "_top_level";
32/// Span metric the tracer sets to denote a top level span
33const TRACER_TOP_LEVEL_KEY: &str = "_dd.top_level";
34const MEASURED_KEY: &str = "_dd.measured";
35const PARTIAL_VERSION_KEY: &str = "_dd.partial_version";
36const MAX_STRING_DICT_SIZE: u32 = 25_000_000;
37const SPAN_ELEMENT_COUNT: usize = 12;
38
39/// First value of returned tuple is the payload size
40pub async fn get_traces_from_request_body<B>(body: B) -> anyhow::Result<(usize, Vec<Vec<pb::Span>>)>
41where
42    B: http_body::Body,
43    B::Error: std::error::Error + Send + Sync + 'static,
44{
45    let buffer = body.collect().await?.aggregate();
46    let size = buffer.remaining();
47
48    let traces: Vec<Vec<pb::Span>> = match rmp_serde::from_read(buffer.reader()) {
49        Ok(res) => res,
50        Err(err) => {
51            anyhow::bail!("Error deserializing trace from request body: {err}")
52        }
53    };
54
55    Ok((size, traces))
56}
57
58#[inline]
59fn get_v05_strings_dict(reader: &mut Reader<impl Buf>) -> anyhow::Result<Vec<String>> {
60    let dict_size =
61        read_array_len(reader).map_err(|err| anyhow!("Error reading dict size: {err}"))?;
62    if dict_size > MAX_STRING_DICT_SIZE {
63        anyhow::bail!(
64            "Error deserializing strings dictionary. Dict size is too large: {dict_size}"
65        );
66    }
67    let mut dict: Vec<String> = Vec::with_capacity(dict_size.try_into()?);
68    for _ in 0..dict_size {
69        match read_value(reader)? {
70            Value::String(s) => {
71                let parsed_string = s.into_str().ok_or_else(|| anyhow!("Error reading string dict"))?;
72                dict.push(parsed_string);
73            }
74            val => anyhow::bail!("Error deserializing strings dictionary. Value in string dict is not a string: {val}")
75        }
76    }
77    Ok(dict)
78}
79
80#[inline]
81fn get_v05_span(reader: &mut Reader<impl Buf>, dict: &[String]) -> anyhow::Result<pb::Span> {
82    let mut span: pb::Span = Default::default();
83    let span_size = rmp::decode::read_array_len(reader)
84        .map_err(|err| anyhow!("Error reading span size: {err}"))? as usize;
85    if span_size != SPAN_ELEMENT_COUNT {
86        anyhow::bail!("Expected an array of exactly 12 elements in a span, got {span_size}");
87    }
88    // 0 - service
89    span.service = get_v05_string(reader, dict, "service")?;
90    // 1 - name
91    span.name = get_v05_string(reader, dict, "name")?;
92    // 2 - resource
93    span.resource = get_v05_string(reader, dict, "resource")?;
94
95    // 3 - trace_id
96    match read_value(reader)? {
97        Value::Integer(i) => {
98            span.trace_id = i.as_u64().ok_or_else(|| {
99                anyhow!("Error reading span trace_id, value is not an integer: {i}")
100            })?;
101        }
102        val => anyhow::bail!("Error reading span trace_id, value is not an integer: {val}"),
103    };
104    // 4 - span_id
105    match read_value(reader)? {
106        Value::Integer(i) => {
107            span.span_id = i.as_u64().ok_or_else(|| {
108                anyhow!("Error reading span span_id, value is not an integer: {i}")
109            })?;
110        }
111        val => anyhow::bail!("Error reading span span_id, value is not an integer: {val}"),
112    };
113    // 5 - parent_id
114    match read_value(reader)? {
115        Value::Integer(i) => {
116            span.parent_id = i.as_u64().ok_or_else(|| {
117                anyhow!("Error reading span parent_id, value is not an integer: {i}")
118            })?;
119        }
120        val => anyhow::bail!("Error reading span parent_id, value is not an integer: {val}"),
121    };
122    // 6 - start
123    match read_value(reader)? {
124        Value::Integer(i) => {
125            span.start = i
126                .as_i64()
127                .ok_or_else(|| anyhow!("Error reading span start, value is not an integer: {i}"))?;
128        }
129        val => anyhow::bail!("Error reading span start, value is not an integer: {val}"),
130    };
131    // 7 - duration
132    match read_value(reader)? {
133        Value::Integer(i) => {
134            span.duration = i.as_i64().ok_or_else(|| {
135                anyhow!("Error reading span duration, value is not an integer: {i}")
136            })?;
137        }
138        val => anyhow::bail!("Error reading span duration, value is not an integer: {val}"),
139    };
140    // 8 - error
141    match read_value(reader)? {
142        Value::Integer(i) => {
143            span.error = i
144                .as_i64()
145                .ok_or_else(|| anyhow!("Error reading span error, value is not an integer: {i}"))?
146                as i32;
147        }
148        val => anyhow::bail!("Error reading span error, value is not an integer: {val}"),
149    }
150    // 9 - meta
151    match read_value(reader)? {
152        Value::Map(meta) => {
153            for (k, v) in meta.iter() {
154                match k {
155                    Value::Integer(k) => {
156                        match v {
157                            Value::Integer(v) => {
158                                let key = str_from_dict(dict, *k)?;
159                                let val = str_from_dict(dict, *v)?;
160                                span.meta.insert(key, val);
161                            }
162                            _ => anyhow::bail!("Error reading span meta, value is not an integer and can't be looked up in dict: {v}")
163                        }
164                    }
165                    _ => anyhow::bail!("Error reading span meta, key is not an integer and can't be looked up in dict: {k}")
166                }
167            }
168        }
169        val => anyhow::bail!("Error reading span meta, value is not a map: {val}"),
170    }
171    // 10 - metrics
172    match read_value(reader)? {
173        Value::Map(metrics) => {
174            for (k, v) in metrics.iter() {
175                match k {
176                    Value::Integer(k) => {
177                        match v {
178                            Value::Integer(v) => {
179                                let key = str_from_dict(dict, *k)?;
180                                span.metrics.insert(key, v.as_f64().ok_or_else(||anyhow!("Error reading span metrics, value is not an integer: {v}"))?);
181                            }
182                            Value::F64(v) => {
183                                let key = str_from_dict(dict, *k)?;
184                                span.metrics.insert(key, *v);
185                            }
186                            _ => anyhow::bail!(
187                                "Error reading span metrics, value is not a float or integer: {v}"
188                            ),
189                        }
190                    }
191                    _ => anyhow::bail!("Error reading span metrics, key is not an integer: {k}"),
192                }
193            }
194        }
195        val => anyhow::bail!("Error reading span metrics, value is not a map: {val}"),
196    }
197
198    // 11 - type
199    match read_value(reader)? {
200        Value::Integer(s) => span.r#type = str_from_dict(dict, s)?,
201        val => anyhow::bail!("Error reading span type, value is not an integer: {val}"),
202    }
203    Ok(span)
204}
205
206#[inline]
207fn str_from_dict(dict: &[String], id: Integer) -> anyhow::Result<String> {
208    let id = id
209        .as_i64()
210        .ok_or_else(|| anyhow!("Error reading string from dict, id is not an integer: {id}"))?
211        as usize;
212    if id >= dict.len() {
213        anyhow::bail!("Error reading string from dict, id out of bounds: {id}");
214    }
215    Ok(dict[id].to_string())
216}
217
218#[inline]
219fn get_v05_string(
220    reader: &mut Reader<impl Buf>,
221    dict: &[String],
222    field_name: &str,
223) -> anyhow::Result<String> {
224    match read_value(reader)? {
225        Value::Integer(s) => {
226            str_from_dict(dict, s)
227        },
228        val => anyhow::bail!("Error reading {field_name}, value is not an integer and can't be looked up in dict: {val}")
229    }
230}
231
232pub async fn get_v05_traces_from_request_body<B>(
233    body: B,
234) -> anyhow::Result<(usize, Vec<Vec<pb::Span>>)>
235where
236    B: http_body::Body,
237    B::Error: std::error::Error + Send + Sync + 'static,
238{
239    let buffer = body.collect().await?.aggregate();
240    let body_size = buffer.remaining();
241    let mut reader = buffer.reader();
242    let wrapper_size = read_array_len(&mut reader)?;
243    if wrapper_size != 2 {
244        anyhow::bail!("Expected an arrary of exactly 2 elements, got {wrapper_size}");
245    }
246
247    let dict = get_v05_strings_dict(&mut reader)?;
248
249    let traces_size = rmp::decode::read_array_len(&mut reader)?;
250    let mut traces: Vec<Vec<pb::Span>> = Default::default();
251
252    for _ in 0..traces_size {
253        let spans_size = rmp::decode::read_array_len(&mut reader)?;
254        let mut trace: Vec<pb::Span> = Default::default();
255
256        for _ in 0..spans_size {
257            let span = get_v05_span(&mut reader, &dict)?;
258            trace.push(span);
259        }
260        traces.push(trace);
261    }
262    Ok((body_size, traces))
263}
264
265/// Tags extracted from a tracer payload's traces, used to populate top level tracer payload fields.
266#[derive(Default)]
267pub struct TracerPayloadTags {
268    pub env: String,
269    pub app_version: String,
270    pub hostname: String,
271    pub runtime_id: String,
272}
273
274/// Returns the first non-empty value of `field` found in `trace`, searching the root span first
275/// then all other spans.
276fn search_trace_for_field(root: &pb::Span, trace: &[pb::Span], field: &str) -> Option<String> {
277    if let Some(v) = root.meta.get(field) {
278        if !v.is_empty() {
279            return Some(v.clone());
280        }
281    }
282    for span in trace {
283        if span.span_id == root.span_id {
284            continue;
285        }
286        if let Some(v) = span.meta.get(field) {
287            if !v.is_empty() {
288                return Some(v.clone());
289            }
290        }
291    }
292    None
293}
294
295pub(crate) fn construct_trace_chunk(trace: Vec<pb::Span>) -> pb::TraceChunk {
296    pb::TraceChunk {
297        priority: normalizer::SamplerPriority::None as i32,
298        origin: "".to_string(),
299        spans: trace,
300        tags: HashMap::new(),
301        dropped_trace: false,
302    }
303}
304
305pub(crate) fn construct_tracer_payload(
306    chunks: Vec<pb::TraceChunk>,
307    tracer_tags: &TracerHeaderTags,
308    tracer_payload_tags: TracerPayloadTags,
309) -> pb::TracerPayload {
310    pb::TracerPayload {
311        app_version: tracer_payload_tags.app_version,
312        language_name: tracer_tags.lang.to_string(),
313        container_id: tracer_tags.container_id.to_string(),
314        env: tracer_payload_tags.env,
315        runtime_id: tracer_payload_tags.runtime_id,
316        chunks,
317        hostname: tracer_payload_tags.hostname,
318        language_version: tracer_tags.lang_version.to_string(),
319        tags: HashMap::new(),
320        tracer_version: tracer_tags.tracer_version.to_string(),
321        container_debug: None,
322    }
323}
324
325pub(crate) fn cmp_send_data_payloads(a: &pb::TracerPayload, b: &pb::TracerPayload) -> Ordering {
326    a.tracer_version
327        .cmp(&b.tracer_version)
328        .then(a.language_version.cmp(&b.language_version))
329        .then(a.language_name.cmp(&b.language_name))
330        .then(a.hostname.cmp(&b.hostname))
331        .then(a.container_id.cmp(&b.container_id))
332        .then(a.runtime_id.cmp(&b.runtime_id))
333        .then(a.env.cmp(&b.env))
334        .then(a.app_version.cmp(&b.app_version))
335        .then(a.container_debug.cmp(&b.container_debug))
336}
337
338pub fn coalesce_send_data(mut data: Vec<SendData>) -> Vec<SendData> {
339    // TODO trace payloads with identical data except for chunk could be merged?
340
341    data.sort_unstable_by(|a, b| {
342        a.get_target()
343            .url
344            .to_string()
345            .cmp(&b.get_target().url.to_string())
346            .then(a.get_target().test_token.cmp(&b.get_target().test_token))
347    });
348    data.dedup_by(|a, b| {
349        if a.get_target().url == b.get_target().url
350            && a.get_target().test_token == b.get_target().test_token
351        {
352            // Size is only an approximation. In practice it won't vary much, but be safe here.
353            // We also don't care about the exact maximum size, like two 25 MB or one 50 MB request
354            // has similar results. The primary goal here is avoiding many small requests.
355            // TODO: maybe make the MAX_PAYLOAD_SIZE configurable?
356            if a.size + b.size < MAX_PAYLOAD_SIZE / 2 {
357                // Note: dedup_by drops a, and retains b. Only drop a if the append actually
358                // merged its data into b; otherwise keep both entries (e.g. diverging V1
359                // tracer metadata) so a's traces aren't silently lost.
360                if b.tracer_payloads.append(&mut a.tracer_payloads) {
361                    b.size += a.size;
362                    return true;
363                }
364            }
365        }
366        false
367    });
368    // Merge chunks with common properties. Reduces requests for agentful mode.
369    // And reduces a little bit of data for agentless.
370    for send_data in data.iter_mut() {
371        send_data.tracer_payloads.merge();
372    }
373    data
374}
375
376pub fn get_root_span_index(trace: &[pb::Span]) -> anyhow::Result<usize> {
377    if trace.is_empty() {
378        anyhow::bail!("Cannot find root span index in an empty trace.");
379    }
380
381    // Do a first pass to find if we have an obvious root span (starting from the end) since some
382    // clients put the root span last.
383    for (i, span) in trace.iter().enumerate().rev() {
384        if span.parent_id == 0 {
385            return Ok(i);
386        }
387    }
388
389    let span_ids: HashSet<_> = trace.iter().map(|span| span.span_id).collect();
390
391    let mut root_span_id = None;
392    for (i, span) in trace.iter().enumerate() {
393        // If a span's parent is not in the trace, it is a root
394        if !span_ids.contains(&span.parent_id) {
395            if root_span_id.is_some() {
396                debug!(
397                    trace_id = &trace[0].trace_id,
398                    "trace has multiple root spans"
399                );
400            }
401            root_span_id = Some(i);
402        }
403    }
404    Ok(match root_span_id {
405        Some(i) => i,
406        None => {
407            debug!(
408                trace_id = &trace[0].trace_id,
409                "Could not find the root span for trace"
410            );
411            trace.len() - 1
412        }
413    })
414}
415
416/// Updates all the spans top-level attribute.
417/// A span is considered top-level if:
418///   - it's a root span
419///   - OR its parent is unknown (other part of the code, distributed trace)
420///   - OR its parent belongs to another service (in that case it's a "local root" being the highest
421///     ancestor of other spans belonging to this service and attached to it).
422pub fn compute_top_level_span(trace: &mut [pb::Span]) {
423    let mut span_id_to_service: HashMap<u64, String> = HashMap::new();
424    for span in trace.iter() {
425        span_id_to_service.insert(span.span_id, span.service.clone());
426    }
427    for span in trace.iter_mut() {
428        if span.parent_id == 0 {
429            set_top_level_span(span);
430            continue;
431        }
432        match span_id_to_service.get(&span.parent_id) {
433            Some(parent_span_service) => {
434                if !parent_span_service.eq(&span.service) {
435                    // parent is not in the same service
436                    set_top_level_span(span)
437                }
438            }
439            None => {
440                // span has no parent in chunk
441                set_top_level_span(span)
442            }
443        }
444    }
445}
446
447/// Return true if the span has a top level key set
448pub fn has_top_level(span: &pb::Span) -> bool {
449    span.metrics
450        .get(TRACER_TOP_LEVEL_KEY)
451        .is_some_and(|v| *v == 1.0)
452        || span.metrics.get(TOP_LEVEL_KEY).is_some_and(|v| *v == 1.0)
453}
454
455fn set_top_level_span(span: &mut pb::Span) {
456    span.metrics.insert(TOP_LEVEL_KEY.to_string(), 1.0);
457}
458
459pub fn set_serverless_root_span_tags(
460    span: &mut pb::Span,
461    app_name: Option<String>,
462    env_type: &EnvironmentType,
463) {
464    let origin_tag = match env_type {
465        EnvironmentType::CloudFunction => "cloudfunction",
466        EnvironmentType::AzureFunction => "azurefunction",
467        EnvironmentType::AzureSpringApp => "azurespringapp",
468        EnvironmentType::LambdaFunction => "lambda", // historical reasons
469    };
470    span.meta
471        .insert("_dd.origin".to_string(), origin_tag.to_string());
472    span.meta
473        .insert("origin".to_string(), origin_tag.to_string());
474
475    if let Some(function_name) = app_name {
476        match env_type {
477            EnvironmentType::CloudFunction
478            | EnvironmentType::AzureFunction
479            | EnvironmentType::LambdaFunction => {
480                span.meta.insert("functionname".to_string(), function_name);
481            }
482            _ => {}
483        }
484    }
485}
486
487fn update_tracer_top_level(span: &mut pb::Span) {
488    if span.metrics.contains_key(TRACER_TOP_LEVEL_KEY) {
489        span.metrics.insert(TOP_LEVEL_KEY.to_string(), 1.0);
490    }
491}
492
493#[derive(Clone, Debug, Eq, PartialEq)]
494pub enum EnvironmentType {
495    CloudFunction,
496    AzureFunction,
497    AzureSpringApp,
498    LambdaFunction,
499}
500
501#[derive(Clone, Debug, Eq, PartialEq)]
502pub struct MiniAgentMetadata {
503    pub azure_spring_app_hostname: Option<String>,
504    pub azure_spring_app_name: Option<String>,
505    pub gcp_project_id: Option<String>,
506    pub gcp_region: Option<String>,
507    pub version: Option<String>,
508}
509
510impl Default for MiniAgentMetadata {
511    fn default() -> Self {
512        MiniAgentMetadata {
513            azure_spring_app_hostname: Default::default(),
514            azure_spring_app_name: Default::default(),
515            gcp_project_id: Default::default(),
516            gcp_region: Default::default(),
517            version: env::var("DD_SERVERLESS_COMPAT_VERSION").ok(),
518        }
519    }
520}
521
522pub fn enrich_span_with_mini_agent_metadata(
523    span: &mut pb::Span,
524    mini_agent_metadata: &MiniAgentMetadata,
525) {
526    if let Some(azure_spring_app_hostname) = &mini_agent_metadata.azure_spring_app_hostname {
527        span.meta.insert(
528            "asa.hostname".to_string(),
529            azure_spring_app_hostname.to_string(),
530        );
531    }
532    if let Some(azure_spring_app_name) = &mini_agent_metadata.azure_spring_app_name {
533        span.meta
534            .insert("asa.name".to_string(), azure_spring_app_name.to_string());
535    }
536    if let Some(serverless_compat_version) = &mini_agent_metadata.version {
537        span.meta.insert(
538            "_dd.serverless_compat_version".to_string(),
539            serverless_compat_version.to_string(),
540        );
541    }
542}
543
544pub fn enrich_span_with_google_cloud_function_metadata(
545    span: &mut pb::Span,
546    mini_agent_metadata: &MiniAgentMetadata,
547    function: Option<String>,
548) {
549    #[allow(clippy::todo)]
550    let Some(region) = &mini_agent_metadata.gcp_region
551    else {
552        todo!()
553    };
554    #[allow(clippy::todo)]
555    let Some(project) = &mini_agent_metadata.gcp_project_id
556    else {
557        todo!()
558    };
559
560    if let Some(function) = function {
561        if !region.is_empty() && !project.is_empty() {
562            let resource_name = format!(
563                "projects/{}/locations/{}/functions/{}",
564                project, region, function
565            );
566
567            span.meta
568                .insert("gcrfx.location".to_string(), region.to_string());
569            span.meta
570                .insert("gcrfx.project_id".to_string(), project.to_string());
571            span.meta
572                .insert("gcrfx.resource_name".to_string(), resource_name.to_string());
573        }
574    }
575}
576
577pub fn enrich_span_with_azure_function_metadata(span: &mut pb::Span) {
578    if span.name == "azure.apim" {
579        return;
580    }
581
582    if let Some(aas_metadata) = &*azure_app_services::AAS_METADATA_FUNCTION {
583        span.meta.extend(
584            aas_metadata
585                .get_function_tags()
586                .map(|(name, value)| (name.to_string(), value.to_string())),
587        );
588    }
589}
590
591/// Converts v0.4-shaped span chunks into the v0.5 wire representation.
592///
593/// v0.5 deduplicates every string field across the whole payload through a shared dictionary
594/// and replaces them with `u32` indices. This walks each span via [`v05::from_v04_span`],
595/// interning strings into the [`SharedDict`] as it goes, and returns the resulting
596/// `(dict, traces)` pair wrapped in [`TraceChunks::V05`].
597///
598/// Returns `Err` if any span fails to convert (e.g. unsupported field value); the partial
599/// dictionary built so far is discarded.
600pub fn convert_trace_chunks_v04_to_v05<T: TraceData>(
601    traces: Vec<Vec<crate::span::v04::Span<T>>>,
602) -> anyhow::Result<TraceChunks<T>> {
603    let mut shared_dict = SharedDict::default();
604    let mut v05_traces: Vec<Vec<v05::Span>> = Vec::with_capacity(traces.len());
605    for trace in traces {
606        let v05_trace = trace
607            .into_iter()
608            .map(|span| v05::from_v04_span(span, &mut shared_dict))
609            .collect::<anyhow::Result<Vec<_>>>()?;
610        v05_traces.push(v05_trace);
611    }
612    Ok(TraceChunks::V05((shared_dict, v05_traces)))
613}
614
615pub fn collect_pb_trace_chunks<T: tracer_payload::TraceChunkProcessor>(
616    mut traces: Vec<Vec<pb::Span>>,
617    tracer_header_tags: &TracerHeaderTags,
618    process_chunk: &mut T,
619    is_agentless: bool,
620) -> anyhow::Result<TracerPayloadCollection> {
621    let mut trace_chunks: Vec<pb::TraceChunk> = Vec::new();
622
623    // We'll skip setting the global metadata and rely on the agent to unpack these
624    let mut tracer_payload_tags = TracerPayloadTags::default();
625
626    for trace in traces.iter_mut() {
627        if is_agentless {
628            if let Err(e) = normalizer::normalize_trace(trace) {
629                error!("Error normalizing trace: {e}");
630            }
631        }
632
633        let mut chunk = construct_trace_chunk(trace.to_vec());
634
635        let root_span_index = match get_root_span_index(trace) {
636            Ok(res) => res,
637            Err(e) => {
638                error!("Error getting the root span index of a trace, skipping. {e}");
639                continue;
640            }
641        };
642
643        if let Err(e) = normalizer::normalize_chunk(&mut chunk, root_span_index) {
644            error!("Error normalizing trace chunk: {e}");
645        }
646
647        for span in chunk.spans.iter_mut() {
648            // TODO: obfuscate & truncate spans
649            if tracer_header_tags.generic.client_computed_top_level {
650                update_tracer_top_level(span);
651            }
652        }
653
654        if !tracer_header_tags.generic.client_computed_top_level {
655            compute_top_level_span(&mut chunk.spans);
656        }
657
658        process_chunk.process(&mut chunk, root_span_index);
659
660        trace_chunks.push(chunk);
661
662        if is_agentless {
663            // Check each field independently so that a later trace can fill in fields missing
664            // from an earlier trace.
665            let root = &trace[root_span_index];
666            if tracer_payload_tags.env.is_empty() {
667                if let Some(mut v) = search_trace_for_field(root, trace, "env") {
668                    // Normalize env tag in case the span it was pulled from was skipped during
669                    // normalization
670                    libdd_trace_normalization::normalize_utils::normalize_tag(&mut v);
671                    if !v.is_empty() {
672                        tracer_payload_tags.env = v;
673                    }
674                }
675            }
676            if tracer_payload_tags.app_version.is_empty() {
677                if let Some(v) = search_trace_for_field(root, trace, "version") {
678                    tracer_payload_tags.app_version = v;
679                }
680            }
681            if tracer_payload_tags.hostname.is_empty() {
682                if let Some(v) = search_trace_for_field(root, trace, "_dd.hostname") {
683                    tracer_payload_tags.hostname = v;
684                }
685            }
686            if tracer_payload_tags.runtime_id.is_empty() {
687                if let Some(v) = search_trace_for_field(root, trace, "runtime-id") {
688                    tracer_payload_tags.runtime_id = v;
689                }
690            }
691        }
692    }
693
694    Ok(TracerPayloadCollection::V07(vec![
695        construct_tracer_payload(trace_chunks, tracer_header_tags, tracer_payload_tags),
696    ]))
697}
698
699/// Returns true if a span should be measured (i.e., it should get trace metrics calculated).
700pub fn is_measured(span: &pb::Span) -> bool {
701    span.metrics.get(MEASURED_KEY).is_some_and(|v| *v == 1.0)
702}
703
704/// Returns true if the span is a partial snapshot.
705/// This kind of spans are partial images of long-running spans.
706/// When incomplete, a partial snapshot has a metric _dd.partial_version which is a positive
707/// integer. The metric usually increases each time a new version of the same span is sent by the
708/// tracer
709pub fn is_partial_snapshot(span: &pb::Span) -> bool {
710    span.metrics
711        .get(PARTIAL_VERSION_KEY)
712        .is_some_and(|v| *v >= 0.0)
713}
714
715#[cfg(test)]
716mod tests {
717    use super::*;
718    use crate::{
719        span::SharedDictBytes,
720        test_utils::{create_test_no_alloc_span, create_test_span},
721    };
722    use http::Request;
723    use libdd_common::{http_common, Endpoint};
724    use serde_json::json;
725
726    fn find_index_in_dict(dict: &SharedDictBytes, value: &str) -> Option<u32> {
727        let idx = dict.iter().position(|e| e.as_str() == value);
728        idx.map(|idx| idx.try_into().unwrap())
729    }
730
731    #[test]
732    fn test_coalescing_does_not_exceed_max_size() {
733        fn dummy() -> SendData {
734            SendData::new(
735                MAX_PAYLOAD_SIZE / 5 + 1,
736                TracerPayloadCollection::V07(vec![pb::TracerPayload {
737                    container_id: "".to_string(),
738                    language_name: "".to_string(),
739                    language_version: "".to_string(),
740                    tracer_version: "".to_string(),
741                    runtime_id: "".to_string(),
742                    chunks: vec![pb::TraceChunk {
743                        priority: 0,
744                        origin: "".to_string(),
745                        spans: vec![],
746                        tags: Default::default(),
747                        dropped_trace: false,
748                    }],
749                    tags: Default::default(),
750                    env: "".to_string(),
751                    hostname: "".to_string(),
752                    app_version: "".to_string(),
753                    container_debug: None,
754                }]),
755                TracerHeaderTags::default(),
756                &Endpoint::default(),
757            )
758        }
759        let coalesced = coalesce_send_data(vec![dummy(), dummy(), dummy(), dummy(), dummy()]);
760        assert_eq!(
761            5,
762            coalesced
763                .iter()
764                .map(|s| s.tracer_payloads.size())
765                .sum::<usize>()
766        );
767        // assert some chunks are actually coalesced
768        assert!(
769            coalesced
770                .iter()
771                .map(|s| {
772                    if let TracerPayloadCollection::V07(collection) = &s.tracer_payloads {
773                        collection.iter().map(|s| s.chunks.len()).max().unwrap()
774                    } else {
775                        0
776                    }
777                })
778                .max()
779                .unwrap()
780                > 1
781        );
782        assert!(coalesced.len() > 1 && coalesced.len() < 5);
783    }
784
785    #[tokio::test]
786    #[allow(clippy::type_complexity)]
787    #[cfg_attr(all(miri, target_os = "macos"), ignore)]
788    async fn test_get_v05_traces_from_request_body() {
789        let data: (
790            Vec<String>,
791            Vec<
792                Vec<(
793                    u8,
794                    u8,
795                    u8,
796                    u64,
797                    u64,
798                    u64,
799                    i64,
800                    i64,
801                    i32,
802                    HashMap<u8, u8>,
803                    HashMap<u8, f64>,
804                    u8,
805                )>,
806            >,
807        ) = (
808            vec![
809                "baggage".to_string(),
810                "item".to_string(),
811                "elasticsearch.version".to_string(),
812                "7.0".to_string(),
813                "my-name".to_string(),
814                "X".to_string(),
815                "my-service".to_string(),
816                "my-resource".to_string(),
817                "_dd.sampling_rate_whatever".to_string(),
818                "value whatever".to_string(),
819                "sql".to_string(),
820            ],
821            vec![vec![(
822                6,
823                4,
824                7,
825                1,
826                2,
827                3,
828                123,
829                456,
830                1,
831                HashMap::from([(8, 9), (0, 1), (2, 3)]),
832                HashMap::from([(5, 1.2)]),
833                10,
834            )]],
835        );
836        let bytes = rmp_serde::to_vec(&data).unwrap();
837        let res = get_v05_traces_from_request_body(http_common::Body::from(bytes)).await;
838        assert!(res.is_ok());
839        let (_, traces) = res.unwrap();
840        let span = traces[0][0].clone();
841        let test_span = pb::Span {
842            service: "my-service".to_string(),
843            name: "my-name".to_string(),
844            resource: "my-resource".to_string(),
845            trace_id: 1,
846            span_id: 2,
847            parent_id: 3,
848            start: 123,
849            duration: 456,
850            error: 1,
851            meta: HashMap::from([
852                ("baggage".to_string(), "item".to_string()),
853                ("elasticsearch.version".to_string(), "7.0".to_string()),
854                (
855                    "_dd.sampling_rate_whatever".to_string(),
856                    "value whatever".to_string(),
857                ),
858            ]),
859            metrics: HashMap::from([("X".to_string(), 1.2)]),
860            meta_struct: HashMap::default(),
861            r#type: "sql".to_string(),
862            span_links: vec![],
863            span_events: vec![],
864        };
865        assert_eq!(span, test_span);
866    }
867
868    #[tokio::test]
869    #[cfg_attr(miri, ignore)]
870    async fn test_get_traces_from_request_body() {
871        let pairs = vec![
872            (
873                json!([{
874                    "service": "test-service",
875                    "name": "test-service-name",
876                    "resource": "test-service-resource",
877                    "trace_id": 111,
878                    "span_id": 222,
879                    "parent_id": 333,
880                    "start": 1,
881                    "duration": 5,
882                    "error": 0,
883                    "meta": {},
884                    "metrics": {},
885                }]),
886                vec![vec![pb::Span {
887                    service: "test-service".to_string(),
888                    name: "test-service-name".to_string(),
889                    resource: "test-service-resource".to_string(),
890                    trace_id: 111,
891                    span_id: 222,
892                    parent_id: 333,
893                    start: 1,
894                    duration: 5,
895                    error: 0,
896                    meta: HashMap::new(),
897                    metrics: HashMap::new(),
898                    meta_struct: HashMap::new(),
899                    r#type: "".to_string(),
900                    span_links: vec![],
901                    span_events: vec![],
902                }]],
903            ),
904            (
905                json!([{
906                    "name": "test-service-name",
907                    "resource": "test-service-resource",
908                    "trace_id": 111,
909                    "span_id": 222,
910                    "start": 1,
911                    "duration": 5,
912                    "meta": {},
913                }]),
914                vec![vec![pb::Span {
915                    service: "".to_string(),
916                    name: "test-service-name".to_string(),
917                    resource: "test-service-resource".to_string(),
918                    trace_id: 111,
919                    span_id: 222,
920                    parent_id: 0,
921                    start: 1,
922                    duration: 5,
923                    error: 0,
924                    meta: HashMap::new(),
925                    metrics: HashMap::new(),
926                    meta_struct: HashMap::new(),
927                    r#type: "".to_string(),
928                    span_links: vec![],
929                    span_events: vec![],
930                }]],
931            ),
932        ];
933
934        for (trace_input, output) in pairs {
935            let bytes = rmp_serde::to_vec(&vec![&trace_input]).unwrap();
936            let request = Request::builder()
937                .body(http_common::Body::from(bytes))
938                .unwrap();
939            let res = get_traces_from_request_body(request.into_body()).await;
940            assert!(res.is_ok());
941            assert_eq!(res.unwrap().1, output);
942        }
943    }
944
945    #[tokio::test]
946    #[cfg_attr(miri, ignore)]
947    async fn test_get_traces_from_request_body_with_span_links() {
948        let trace_input = json!([[{
949            "service": "test-service",
950            "name": "test-name",
951            "resource": "test-resource",
952            "trace_id": 111,
953            "span_id": 222,
954            "parent_id": 333,
955            "start": 1,
956            "duration": 5,
957            "error": 0,
958            "meta": {},
959            "metrics": {},
960            "span_links": [{
961                "trace_id": 999,
962                "span_id": 888,
963                "trace_id_high": 777,
964                "attributes": {"key": "value"},
965                "tracestate": "vendor=value"
966                // flags field intentionally omitted
967            }]
968        }]]);
969
970        let expected_output = vec![vec![pb::Span {
971            service: "test-service".to_string(),
972            name: "test-name".to_string(),
973            resource: "test-resource".to_string(),
974            trace_id: 111,
975            span_id: 222,
976            parent_id: 333,
977            start: 1,
978            duration: 5,
979            error: 0,
980            meta: HashMap::new(),
981            metrics: HashMap::new(),
982            meta_struct: HashMap::new(),
983            r#type: String::new(),
984            span_links: vec![pb::SpanLink {
985                trace_id: 999,
986                span_id: 888,
987                trace_id_high: 777,
988                attributes: HashMap::from([("key".to_string(), "value".to_string())]),
989                tracestate: "vendor=value".to_string(),
990                flags: 0, // Should default to 0 when omitted
991            }],
992            span_events: vec![],
993        }]];
994
995        let bytes = rmp_serde::to_vec(&trace_input).unwrap();
996        let request = Request::builder()
997            .body(http_common::Body::from(bytes))
998            .unwrap();
999
1000        let res = get_traces_from_request_body(request.into_body()).await;
1001        assert!(res.is_ok(), "Failed to deserialize: {res:?}");
1002        assert_eq!(res.unwrap().1, expected_output);
1003    }
1004
1005    #[test]
1006    fn test_get_root_span_index_from_complete_trace() {
1007        let trace = vec![
1008            create_test_span(1234, 12341, 0, 1, false),
1009            create_test_span(1234, 12342, 12341, 1, false),
1010            create_test_span(1234, 12343, 12342, 1, false),
1011        ];
1012
1013        let root_span_index = get_root_span_index(&trace);
1014        assert!(root_span_index.is_ok());
1015        assert_eq!(root_span_index.unwrap(), 0);
1016    }
1017
1018    #[test]
1019    fn test_get_root_span_index_from_partial_trace() {
1020        let trace = vec![
1021            create_test_span(1234, 12342, 12341, 1, false),
1022            create_test_span(1234, 12341, 12340, 1, false), /* this is the root span, it's
1023                                                             * parent is not in the trace */
1024            create_test_span(1234, 12343, 12342, 1, false),
1025        ];
1026
1027        let root_span_index = get_root_span_index(&trace);
1028        assert!(root_span_index.is_ok());
1029        assert_eq!(root_span_index.unwrap(), 1);
1030    }
1031
1032    #[test]
1033    fn test_set_serverless_root_span_tags_azure_function() {
1034        let mut span = create_test_span(1234, 12342, 12341, 1, false);
1035        set_serverless_root_span_tags(
1036            &mut span,
1037            Some("test_function".to_string()),
1038            &EnvironmentType::AzureFunction,
1039        );
1040        assert_eq!(
1041            span.meta,
1042            HashMap::from([
1043                (
1044                    "runtime-id".to_string(),
1045                    "test-runtime-id-value".to_string()
1046                ),
1047                ("_dd.origin".to_string(), "azurefunction".to_string()),
1048                ("origin".to_string(), "azurefunction".to_string()),
1049                ("functionname".to_string(), "test_function".to_string()),
1050                ("env".to_string(), "test-env".to_string()),
1051                ("service".to_string(), "test-service".to_string())
1052            ]),
1053        );
1054    }
1055
1056    #[test]
1057    fn test_set_serverless_root_span_tags_cloud_function() {
1058        let mut span = create_test_span(1234, 12342, 12341, 1, false);
1059        set_serverless_root_span_tags(
1060            &mut span,
1061            Some("test_function".to_string()),
1062            &EnvironmentType::CloudFunction,
1063        );
1064        assert_eq!(
1065            span.meta,
1066            HashMap::from([
1067                (
1068                    "runtime-id".to_string(),
1069                    "test-runtime-id-value".to_string()
1070                ),
1071                ("_dd.origin".to_string(), "cloudfunction".to_string()),
1072                ("origin".to_string(), "cloudfunction".to_string()),
1073                ("functionname".to_string(), "test_function".to_string()),
1074                ("env".to_string(), "test-env".to_string()),
1075                ("service".to_string(), "test-service".to_string())
1076            ]),
1077        );
1078    }
1079
1080    #[test]
1081    fn test_has_top_level() {
1082        let top_level_span = create_test_span(123, 1234, 12, 1, true);
1083        let not_top_level_span = create_test_span(123, 1234, 12, 1, false);
1084        assert!(has_top_level(&top_level_span));
1085        assert!(!has_top_level(&not_top_level_span));
1086    }
1087
1088    #[test]
1089    fn test_is_measured() {
1090        let mut measured_span = create_test_span(123, 1234, 12, 1, true);
1091        measured_span.metrics.insert(MEASURED_KEY.into(), 1.0);
1092        let not_measured_span = create_test_span(123, 1234, 12, 1, true);
1093        assert!(is_measured(&measured_span));
1094        assert!(!is_measured(&not_measured_span));
1095    }
1096
1097    #[test]
1098    fn test_compute_top_level() {
1099        let mut span_with_different_service = create_test_span(123, 5, 2, 1, false);
1100        span_with_different_service.service = "another_service".into();
1101        let mut trace = vec![
1102            // Root span, should be marked as top-level
1103            create_test_span(123, 1, 0, 1, false),
1104            // Should not be marked as top-level
1105            create_test_span(123, 2, 1, 1, false),
1106            // No parent in local trace, should be marked as
1107            // top-level
1108            create_test_span(123, 4, 3, 1, false),
1109            // Parent belongs to another service, should be marked
1110            // as top-level
1111            span_with_different_service,
1112        ];
1113
1114        compute_top_level_span(trace.as_mut_slice());
1115
1116        let spans_marked_as_top_level: Vec<u64> = trace
1117            .iter()
1118            .filter_map(|span| {
1119                if has_top_level(span) {
1120                    Some(span.span_id)
1121                } else {
1122                    None
1123                }
1124            })
1125            .collect();
1126        assert_eq!(spans_marked_as_top_level, [1, 4, 5])
1127    }
1128
1129    #[test]
1130    fn test_convert_trace_chunks_v04_to_v05() {
1131        let chunk = vec![create_test_no_alloc_span(123, 456, 789, 1, true)];
1132
1133        let collection = convert_trace_chunks_v04_to_v05(vec![chunk]).unwrap();
1134
1135        let (dict, traces) = match collection {
1136            TraceChunks::V05(payload) => payload,
1137            _ => panic!("Unexpected type"),
1138        };
1139
1140        assert_eq!(dict.len(), 16);
1141
1142        let span = &traces[0][0];
1143        assert_eq!(span.service, 1);
1144        assert_eq!(span.name, 2);
1145        assert_eq!(span.resource, 3);
1146        assert_eq!(span.trace_id, 123);
1147        assert_eq!(span.span_id, 456);
1148        assert_eq!(span.parent_id, 789);
1149        assert_eq!(span.start, 1);
1150        assert_eq!(span.error, 0);
1151        assert_eq!(span.error, 0);
1152        assert_eq!(span.r#type, 15);
1153        assert_eq!(
1154            *span
1155                .meta
1156                .get(&find_index_in_dict(&dict, "service").unwrap())
1157                .unwrap(),
1158            find_index_in_dict(&dict, "test-service").unwrap()
1159        );
1160        assert_eq!(
1161            *span
1162                .meta
1163                .get(&find_index_in_dict(&dict, "env").unwrap())
1164                .unwrap(),
1165            find_index_in_dict(&dict, "test-env").unwrap()
1166        );
1167        assert_eq!(
1168            *span
1169                .meta
1170                .get(&find_index_in_dict(&dict, "runtime-id").unwrap())
1171                .unwrap(),
1172            find_index_in_dict(&dict, "test-runtime-id-value").unwrap()
1173        );
1174        assert_eq!(
1175            *span
1176                .meta
1177                .get(&find_index_in_dict(&dict, "_dd.origin").unwrap())
1178                .unwrap(),
1179            find_index_in_dict(&dict, "cloudfunction").unwrap()
1180        );
1181        assert_eq!(
1182            *span
1183                .meta
1184                .get(&find_index_in_dict(&dict, "origin").unwrap())
1185                .unwrap(),
1186            find_index_in_dict(&dict, "cloudfunction").unwrap()
1187        );
1188        assert_eq!(
1189            *span
1190                .meta
1191                .get(&find_index_in_dict(&dict, "functionname").unwrap())
1192                .unwrap(),
1193            find_index_in_dict(&dict, "dummy_function_name").unwrap()
1194        );
1195        assert_eq!(
1196            *span
1197                .metrics
1198                .get(&find_index_in_dict(&dict, "_top_level").unwrap())
1199                .unwrap(),
1200            1.0
1201        );
1202    }
1203
1204    #[test]
1205    fn test_rmp_serde_deserialize_meta_with_null_values() {
1206        // Create a JSON representation with null value in meta
1207        let span_json = json!({
1208            "service": "test-service",
1209            "name": "test_name",
1210            "resource": "test-resource",
1211            "trace_id": 1_u64,
1212            "span_id": 2_u64,
1213            "parent_id": 0_u64,
1214            "start": 0_i64,
1215            "duration": 5_i64,
1216            "error": 0_i32,
1217            "meta": {
1218                "service": "test-service",
1219                "env": "test-env",
1220                "runtime-id": "test-runtime-id-value",
1221                "problematic_key": null  // Ensure this null value does not cause an error
1222            },
1223            "metrics": {},
1224            "type": "",
1225            "meta_struct": {},
1226            "span_links": [],
1227            "span_events": []
1228        });
1229
1230        let traces_json = vec![vec![span_json]];
1231        let encoded_data = rmp_serde::to_vec(&traces_json).unwrap();
1232        let traces: Vec<Vec<pb::Span>> = rmp_serde::from_read(&encoded_data[..])
1233            .expect("Failed to deserialize traces with null values in meta");
1234
1235        assert_eq!(1, traces.len());
1236        assert_eq!(1, traces[0].len());
1237        let decoded_span = &traces[0][0];
1238
1239        assert_eq!("test-service", decoded_span.service);
1240        assert_eq!("test_name", decoded_span.name);
1241        assert_eq!("test-resource", decoded_span.resource);
1242        assert_eq!("test-service", decoded_span.meta.get("service").unwrap());
1243        assert_eq!("test-env", decoded_span.meta.get("env").unwrap());
1244        assert_eq!(
1245            "test-runtime-id-value",
1246            decoded_span.meta.get("runtime-id").unwrap()
1247        );
1248        // Assert that the null value was filtered out (key not present in map)
1249        assert!(
1250            !decoded_span.meta.contains_key("problematic_key"),
1251            "Null value should be skipped, but key was present"
1252        );
1253    }
1254
1255    #[test]
1256    fn test_enrich_span_with_azure_function_metadata_adds_tags_for_non_apim() {
1257        let mut span = create_test_span(1234, 12342, 12341, 1, false);
1258        span.name = "azure.function".to_string();
1259
1260        enrich_span_with_azure_function_metadata(&mut span);
1261
1262        // If AAS_METADATA_FUNCTION is available, verify aas.* tags were added
1263        // If not available (most test environments), this is a no-op
1264        // This test primarily ensures the function doesn't skip non-apim spans
1265        if azure_app_services::AAS_METADATA_FUNCTION.is_some() {
1266            assert!(span.meta.contains_key("aas.resource.id"));
1267            assert!(span.meta.contains_key("aas.environment.instance_id"));
1268            assert!(span.meta.contains_key("aas.environment.instance_name"));
1269            assert!(span.meta.contains_key("aas.subscription.id"));
1270            assert!(span.meta.contains_key("aas.environment.os"));
1271            assert!(span.meta.contains_key("aas.environment.runtime"));
1272            assert!(span.meta.contains_key("aas.environment.runtime_version"));
1273            assert!(span.meta.contains_key("aas.environment.function_runtime"));
1274            assert!(span.meta.contains_key("aas.resource.group"));
1275            assert!(span.meta.contains_key("aas.site.name"));
1276            assert!(span.meta.contains_key("aas.site.kind"));
1277            assert!(span.meta.contains_key("aas.site.type"));
1278        }
1279    }
1280
1281    #[test]
1282    fn test_enrich_span_with_azure_function_metadata_skips_azure_apim() {
1283        let mut span = create_test_span(1234, 12342, 12341, 1, false);
1284        span.name = "azure.apim".to_string();
1285
1286        enrich_span_with_azure_function_metadata(&mut span);
1287
1288        // Verify no aas.* tags were added
1289        assert!(!span.meta.contains_key("aas.resource.id"));
1290        assert!(!span.meta.contains_key("aas.environment.instance_id"));
1291        assert!(!span.meta.contains_key("aas.environment.instance_name"));
1292        assert!(!span.meta.contains_key("aas.subscription.id"));
1293        assert!(!span.meta.contains_key("aas.environment.os"));
1294        assert!(!span.meta.contains_key("aas.environment.runtime"));
1295        assert!(!span.meta.contains_key("aas.environment.runtime_version"));
1296        assert!(!span.meta.contains_key("aas.environment.function_runtime"));
1297        assert!(!span.meta.contains_key("aas.resource.group"));
1298        assert!(!span.meta.contains_key("aas.site.name"));
1299        assert!(!span.meta.contains_key("aas.site.kind"));
1300        assert!(!span.meta.contains_key("aas.site.type"));
1301    }
1302
1303    #[test]
1304    fn test_collect_pb_trace_chunks_searches_multiple_root_spans_for_fields() {
1305        // First trace root span has no fields. Second trace root span has all fields.
1306        // The second root span should populate all fields.
1307        let mut first_root_span = create_test_span(1, 1, 0, 1, true);
1308        first_root_span.meta.remove("env");
1309        first_root_span.meta.remove("runtime-id");
1310
1311        let mut second_root_span = create_test_span(2, 3, 0, 1, true);
1312        second_root_span
1313            .meta
1314            .insert("version".to_string(), "1.2.3".to_string());
1315        second_root_span
1316            .meta
1317            .insert("env".to_string(), "prod".to_string());
1318        second_root_span
1319            .meta
1320            .insert("_dd.hostname".to_string(), "my-host".to_string());
1321        second_root_span
1322            .meta
1323            .insert("runtime-id".to_string(), "123".to_string());
1324
1325        let result = collect_pb_trace_chunks(
1326            vec![vec![first_root_span], vec![second_root_span]],
1327            &TracerHeaderTags::default(),
1328            &mut tracer_payload::DefaultTraceChunkProcessor,
1329            true,
1330        )
1331        .unwrap();
1332
1333        let TracerPayloadCollection::V07(payloads) = result else {
1334            panic!("expected TracerPayloadCollection::V07");
1335        };
1336        assert_eq!(payloads[0].app_version, "1.2.3");
1337        assert_eq!(payloads[0].env, "prod");
1338        assert_eq!(payloads[0].hostname, "my-host");
1339        assert_eq!(payloads[0].runtime_id, "123");
1340    }
1341
1342    #[test]
1343    fn test_collect_pb_trace_chunks_searches_non_root_spans_for_fields() {
1344        // Root span has no fields. Child span has all fields.
1345        // The child span should populate all fields.
1346        let mut root_span = create_test_span(1, 1, 0, 1, true);
1347        root_span.meta.remove("env");
1348        root_span.meta.remove("runtime-id");
1349        let mut child_span = create_test_span(1, 2, 1, 1, false);
1350        child_span
1351            .meta
1352            .insert("version".to_string(), "1.2.3".to_string());
1353        child_span
1354            .meta
1355            .insert("env".to_string(), "prod".to_string());
1356        child_span
1357            .meta
1358            .insert("_dd.hostname".to_string(), "my-host".to_string());
1359        child_span
1360            .meta
1361            .insert("runtime-id".to_string(), "123".to_string());
1362
1363        let result = collect_pb_trace_chunks(
1364            vec![vec![root_span, child_span]],
1365            &TracerHeaderTags::default(),
1366            &mut tracer_payload::DefaultTraceChunkProcessor,
1367            true,
1368        )
1369        .unwrap();
1370
1371        let TracerPayloadCollection::V07(payloads) = result else {
1372            panic!("expected TracerPayloadCollection::V07");
1373        };
1374        assert_eq!(payloads[0].app_version, "1.2.3");
1375        assert_eq!(payloads[0].env, "prod");
1376        assert_eq!(payloads[0].hostname, "my-host");
1377        assert_eq!(payloads[0].runtime_id, "123");
1378    }
1379
1380    #[test]
1381    fn test_collect_pb_trace_chunks_root_span_takes_priority_over_child() {
1382        // Root span has all fields. Child has different values for all fields.
1383        // The root span should populate all fields.
1384        let mut root_span = create_test_span(1, 1, 0, 1, true);
1385        root_span
1386            .meta
1387            .insert("version".to_string(), "root-version".to_string());
1388        root_span
1389            .meta
1390            .insert("env".to_string(), "root-env".to_string());
1391        root_span
1392            .meta
1393            .insert("_dd.hostname".to_string(), "root-host".to_string());
1394        root_span
1395            .meta
1396            .insert("runtime-id".to_string(), "root-runtime-id".to_string());
1397
1398        let mut child_span = create_test_span(1, 2, 1, 1, false);
1399        child_span
1400            .meta
1401            .insert("version".to_string(), "child-version".to_string());
1402        child_span
1403            .meta
1404            .insert("env".to_string(), "child-env".to_string());
1405        child_span
1406            .meta
1407            .insert("_dd.hostname".to_string(), "child-host".to_string());
1408        child_span
1409            .meta
1410            .insert("runtime-id".to_string(), "child-runtime-id".to_string());
1411
1412        let result = collect_pb_trace_chunks(
1413            vec![vec![root_span, child_span]],
1414            &TracerHeaderTags::default(),
1415            &mut tracer_payload::DefaultTraceChunkProcessor,
1416            true,
1417        )
1418        .unwrap();
1419
1420        let TracerPayloadCollection::V07(payloads) = result else {
1421            panic!("expected TracerPayloadCollection::V07");
1422        };
1423        assert_eq!(payloads[0].app_version, "root-version");
1424        assert_eq!(payloads[0].env, "root-env");
1425        assert_eq!(payloads[0].hostname, "root-host");
1426        assert_eq!(payloads[0].runtime_id, "root-runtime-id");
1427    }
1428
1429    #[test]
1430    fn test_collect_pb_trace_chunks_skips_empty_root_span_value() {
1431        // Root span has empty values for all fields. Child span has non-empty values.
1432        // The child span should populate all fields.
1433        let mut root_span = create_test_span(1, 1, 0, 1, true);
1434        root_span.meta.insert("version".to_string(), "".to_string());
1435        root_span.meta.insert("env".to_string(), "".to_string());
1436        root_span
1437            .meta
1438            .insert("_dd.hostname".to_string(), "".to_string());
1439        root_span
1440            .meta
1441            .insert("runtime-id".to_string(), "".to_string());
1442
1443        let mut child_span = create_test_span(1, 2, 1, 1, false);
1444        child_span
1445            .meta
1446            .insert("version".to_string(), "1.2.3".to_string());
1447        child_span
1448            .meta
1449            .insert("env".to_string(), "prod".to_string());
1450        child_span
1451            .meta
1452            .insert("_dd.hostname".to_string(), "my-host".to_string());
1453        child_span
1454            .meta
1455            .insert("runtime-id".to_string(), "123".to_string());
1456
1457        let result = collect_pb_trace_chunks(
1458            vec![vec![root_span, child_span]],
1459            &TracerHeaderTags::default(),
1460            &mut tracer_payload::DefaultTraceChunkProcessor,
1461            true,
1462        )
1463        .unwrap();
1464
1465        let TracerPayloadCollection::V07(payloads) = result else {
1466            panic!("expected TracerPayloadCollection::V07");
1467        };
1468        assert_eq!(payloads[0].app_version, "1.2.3");
1469        assert_eq!(payloads[0].env, "prod");
1470        assert_eq!(payloads[0].hostname, "my-host");
1471        assert_eq!(payloads[0].runtime_id, "123");
1472    }
1473
1474    #[test]
1475    fn test_collect_pb_trace_chunks_normalizes_env() {
1476        let mut root = create_test_span(1, 1, 0, 1, true);
1477        root.meta
1478            .insert("env".to_string(), "PRODUCTION".to_string());
1479
1480        let result = collect_pb_trace_chunks(
1481            vec![vec![root]],
1482            &TracerHeaderTags::default(),
1483            &mut tracer_payload::DefaultTraceChunkProcessor,
1484            true,
1485        )
1486        .unwrap();
1487
1488        let TracerPayloadCollection::V07(payloads) = result else {
1489            panic!("expected TracerPayloadCollection::V07");
1490        };
1491        assert_eq!(payloads[0].env, "production");
1492    }
1493
1494    #[test]
1495    fn test_collect_pb_trace_chunks_skips_env_empty_after_normalization() {
1496        // First root span has an env that normalizes to empty (all invalid characters).
1497        // Second root span has an env should populate env fields.
1498        let mut first_root_span = create_test_span(1, 1, 0, 1, true);
1499        first_root_span
1500            .meta
1501            .insert("env".to_string(), "!!!".to_string());
1502
1503        let mut second_root_span = create_test_span(2, 3, 0, 1, true);
1504        second_root_span
1505            .meta
1506            .insert("env".to_string(), "prod".to_string());
1507
1508        let result = collect_pb_trace_chunks(
1509            vec![vec![first_root_span], vec![second_root_span]],
1510            &TracerHeaderTags::default(),
1511            &mut tracer_payload::DefaultTraceChunkProcessor,
1512            true,
1513        )
1514        .unwrap();
1515
1516        let TracerPayloadCollection::V07(payloads) = result else {
1517            panic!("expected TracerPayloadCollection::V07");
1518        };
1519        assert_eq!(payloads[0].env, "prod");
1520    }
1521
1522    #[test]
1523    fn test_search_trace_for_field_skips_span_with_same_id_as_root() {
1524        // A span with the same span_id as root is treated as the root and skipped
1525        // in the child span search. Only the root spans own meta is checked for it.
1526        let mut root = create_test_span(1, 1, 0, 1, true);
1527        root.meta.remove("version");
1528
1529        // This span shares the same span_id as the root span, it should be skipped.
1530        let mut duplicate = create_test_span(1, 1, 0, 1, false);
1531        duplicate
1532            .meta
1533            .insert("version".to_string(), "should-not-appear".to_string());
1534
1535        let trace = vec![root.clone(), duplicate];
1536        assert_eq!(search_trace_for_field(&root, &trace, "version"), None);
1537    }
1538}