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::TracerHeaderTags;
9use crate::tracer_payload::TracerPayloadCollection;
10use crate::tracer_payload::{self, TraceChunks, TraceEncoding};
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.
358                b.tracer_payloads.append(&mut a.tracer_payloads);
359                b.size += a.size;
360                return true;
361            }
362        }
363        false
364    });
365    // Merge chunks with common properties. Reduces requests for agentful mode.
366    // And reduces a little bit of data for agentless.
367    for send_data in data.iter_mut() {
368        send_data.tracer_payloads.merge();
369    }
370    data
371}
372
373pub fn get_root_span_index(trace: &[pb::Span]) -> anyhow::Result<usize> {
374    if trace.is_empty() {
375        anyhow::bail!("Cannot find root span index in an empty trace.");
376    }
377
378    // Do a first pass to find if we have an obvious root span (starting from the end) since some
379    // clients put the root span last.
380    for (i, span) in trace.iter().enumerate().rev() {
381        if span.parent_id == 0 {
382            return Ok(i);
383        }
384    }
385
386    let span_ids: HashSet<_> = trace.iter().map(|span| span.span_id).collect();
387
388    let mut root_span_id = None;
389    for (i, span) in trace.iter().enumerate() {
390        // If a span's parent is not in the trace, it is a root
391        if !span_ids.contains(&span.parent_id) {
392            if root_span_id.is_some() {
393                debug!(
394                    trace_id = &trace[0].trace_id,
395                    "trace has multiple root spans"
396                );
397            }
398            root_span_id = Some(i);
399        }
400    }
401    Ok(match root_span_id {
402        Some(i) => i,
403        None => {
404            debug!(
405                trace_id = &trace[0].trace_id,
406                "Could not find the root span for trace"
407            );
408            trace.len() - 1
409        }
410    })
411}
412
413/// Updates all the spans top-level attribute.
414/// A span is considered top-level if:
415///   - it's a root span
416///   - OR its parent is unknown (other part of the code, distributed trace)
417///   - OR its parent belongs to another service (in that case it's a "local root" being the highest
418///     ancestor of other spans belonging to this service and attached to it).
419pub fn compute_top_level_span(trace: &mut [pb::Span]) {
420    let mut span_id_to_service: HashMap<u64, String> = HashMap::new();
421    for span in trace.iter() {
422        span_id_to_service.insert(span.span_id, span.service.clone());
423    }
424    for span in trace.iter_mut() {
425        if span.parent_id == 0 {
426            set_top_level_span(span);
427            continue;
428        }
429        match span_id_to_service.get(&span.parent_id) {
430            Some(parent_span_service) => {
431                if !parent_span_service.eq(&span.service) {
432                    // parent is not in the same service
433                    set_top_level_span(span)
434                }
435            }
436            None => {
437                // span has no parent in chunk
438                set_top_level_span(span)
439            }
440        }
441    }
442}
443
444/// Return true if the span has a top level key set
445pub fn has_top_level(span: &pb::Span) -> bool {
446    span.metrics
447        .get(TRACER_TOP_LEVEL_KEY)
448        .is_some_and(|v| *v == 1.0)
449        || span.metrics.get(TOP_LEVEL_KEY).is_some_and(|v| *v == 1.0)
450}
451
452fn set_top_level_span(span: &mut pb::Span) {
453    span.metrics.insert(TOP_LEVEL_KEY.to_string(), 1.0);
454}
455
456pub fn set_serverless_root_span_tags(
457    span: &mut pb::Span,
458    app_name: Option<String>,
459    env_type: &EnvironmentType,
460) {
461    let origin_tag = match env_type {
462        EnvironmentType::CloudFunction => "cloudfunction",
463        EnvironmentType::AzureFunction => "azurefunction",
464        EnvironmentType::AzureSpringApp => "azurespringapp",
465        EnvironmentType::LambdaFunction => "lambda", // historical reasons
466    };
467    span.meta
468        .insert("_dd.origin".to_string(), origin_tag.to_string());
469    span.meta
470        .insert("origin".to_string(), origin_tag.to_string());
471
472    if let Some(function_name) = app_name {
473        match env_type {
474            EnvironmentType::CloudFunction
475            | EnvironmentType::AzureFunction
476            | EnvironmentType::LambdaFunction => {
477                span.meta.insert("functionname".to_string(), function_name);
478            }
479            _ => {}
480        }
481    }
482}
483
484fn update_tracer_top_level(span: &mut pb::Span) {
485    if span.metrics.contains_key(TRACER_TOP_LEVEL_KEY) {
486        span.metrics.insert(TOP_LEVEL_KEY.to_string(), 1.0);
487    }
488}
489
490#[derive(Clone, Debug, Eq, PartialEq)]
491pub enum EnvironmentType {
492    CloudFunction,
493    AzureFunction,
494    AzureSpringApp,
495    LambdaFunction,
496}
497
498#[derive(Clone, Debug, Eq, PartialEq)]
499pub struct MiniAgentMetadata {
500    pub azure_spring_app_hostname: Option<String>,
501    pub azure_spring_app_name: Option<String>,
502    pub gcp_project_id: Option<String>,
503    pub gcp_region: Option<String>,
504    pub version: Option<String>,
505}
506
507impl Default for MiniAgentMetadata {
508    fn default() -> Self {
509        MiniAgentMetadata {
510            azure_spring_app_hostname: Default::default(),
511            azure_spring_app_name: Default::default(),
512            gcp_project_id: Default::default(),
513            gcp_region: Default::default(),
514            version: env::var("DD_SERVERLESS_COMPAT_VERSION").ok(),
515        }
516    }
517}
518
519pub fn enrich_span_with_mini_agent_metadata(
520    span: &mut pb::Span,
521    mini_agent_metadata: &MiniAgentMetadata,
522) {
523    if let Some(azure_spring_app_hostname) = &mini_agent_metadata.azure_spring_app_hostname {
524        span.meta.insert(
525            "asa.hostname".to_string(),
526            azure_spring_app_hostname.to_string(),
527        );
528    }
529    if let Some(azure_spring_app_name) = &mini_agent_metadata.azure_spring_app_name {
530        span.meta
531            .insert("asa.name".to_string(), azure_spring_app_name.to_string());
532    }
533    if let Some(serverless_compat_version) = &mini_agent_metadata.version {
534        span.meta.insert(
535            "_dd.serverless_compat_version".to_string(),
536            serverless_compat_version.to_string(),
537        );
538    }
539}
540
541pub fn enrich_span_with_google_cloud_function_metadata(
542    span: &mut pb::Span,
543    mini_agent_metadata: &MiniAgentMetadata,
544    function: Option<String>,
545) {
546    #[allow(clippy::todo)]
547    let Some(region) = &mini_agent_metadata.gcp_region
548    else {
549        todo!()
550    };
551    #[allow(clippy::todo)]
552    let Some(project) = &mini_agent_metadata.gcp_project_id
553    else {
554        todo!()
555    };
556
557    if let Some(function) = function {
558        if !region.is_empty() && !project.is_empty() {
559            let resource_name = format!(
560                "projects/{}/locations/{}/functions/{}",
561                project, region, function
562            );
563
564            span.meta
565                .insert("gcrfx.location".to_string(), region.to_string());
566            span.meta
567                .insert("gcrfx.project_id".to_string(), project.to_string());
568            span.meta
569                .insert("gcrfx.resource_name".to_string(), resource_name.to_string());
570        }
571    }
572}
573
574pub fn enrich_span_with_azure_function_metadata(span: &mut pb::Span) {
575    if span.name == "azure.apim" {
576        return;
577    }
578
579    if let Some(aas_metadata) = &*azure_app_services::AAS_METADATA_FUNCTION {
580        span.meta.extend(
581            aas_metadata
582                .get_function_tags()
583                .map(|(name, value)| (name.to_string(), value.to_string())),
584        );
585    }
586}
587
588pub fn collect_trace_chunks<T: TraceData>(
589    traces: Vec<Vec<crate::span::v04::Span<T>>>,
590    format: TraceEncoding,
591) -> anyhow::Result<TraceChunks<T>> {
592    match format {
593        TraceEncoding::V05 => {
594            let mut shared_dict = SharedDict::default();
595            let mut v05_traces: Vec<Vec<v05::Span>> = Vec::with_capacity(traces.len());
596            for trace in traces {
597                let v05_trace = trace
598                    .into_iter()
599                    .map(|span| v05::from_v04_span(span, &mut shared_dict))
600                    .collect::<anyhow::Result<Vec<_>>>()?;
601                v05_traces.push(v05_trace);
602            }
603            Ok(TraceChunks::V05((shared_dict, v05_traces)))
604        }
605        TraceEncoding::V04 => Ok(TraceChunks::V04(traces)),
606    }
607}
608
609pub fn collect_pb_trace_chunks<T: tracer_payload::TraceChunkProcessor>(
610    mut traces: Vec<Vec<pb::Span>>,
611    tracer_header_tags: &TracerHeaderTags,
612    process_chunk: &mut T,
613    is_agentless: bool,
614) -> anyhow::Result<TracerPayloadCollection> {
615    let mut trace_chunks: Vec<pb::TraceChunk> = Vec::new();
616
617    // We'll skip setting the global metadata and rely on the agent to unpack these
618    let mut tracer_payload_tags = TracerPayloadTags::default();
619
620    for trace in traces.iter_mut() {
621        if is_agentless {
622            if let Err(e) = normalizer::normalize_trace(trace) {
623                error!("Error normalizing trace: {e}");
624            }
625        }
626
627        let mut chunk = construct_trace_chunk(trace.to_vec());
628
629        let root_span_index = match get_root_span_index(trace) {
630            Ok(res) => res,
631            Err(e) => {
632                error!("Error getting the root span index of a trace, skipping. {e}");
633                continue;
634            }
635        };
636
637        if let Err(e) = normalizer::normalize_chunk(&mut chunk, root_span_index) {
638            error!("Error normalizing trace chunk: {e}");
639        }
640
641        for span in chunk.spans.iter_mut() {
642            // TODO: obfuscate & truncate spans
643            if tracer_header_tags.client_computed_top_level {
644                update_tracer_top_level(span);
645            }
646        }
647
648        if !tracer_header_tags.client_computed_top_level {
649            compute_top_level_span(&mut chunk.spans);
650        }
651
652        process_chunk.process(&mut chunk, root_span_index);
653
654        trace_chunks.push(chunk);
655
656        if is_agentless {
657            // Check each field independently so that a later trace can fill in fields missing
658            // from an earlier trace.
659            let root = &trace[root_span_index];
660            if tracer_payload_tags.env.is_empty() {
661                if let Some(mut v) = search_trace_for_field(root, trace, "env") {
662                    // Normalize env tag in case the span it was pulled from was skipped during
663                    // normalization
664                    libdd_trace_normalization::normalize_utils::normalize_tag(&mut v);
665                    if !v.is_empty() {
666                        tracer_payload_tags.env = v;
667                    }
668                }
669            }
670            if tracer_payload_tags.app_version.is_empty() {
671                if let Some(v) = search_trace_for_field(root, trace, "version") {
672                    tracer_payload_tags.app_version = v;
673                }
674            }
675            if tracer_payload_tags.hostname.is_empty() {
676                if let Some(v) = search_trace_for_field(root, trace, "_dd.hostname") {
677                    tracer_payload_tags.hostname = v;
678                }
679            }
680            if tracer_payload_tags.runtime_id.is_empty() {
681                if let Some(v) = search_trace_for_field(root, trace, "runtime-id") {
682                    tracer_payload_tags.runtime_id = v;
683                }
684            }
685        }
686    }
687
688    Ok(TracerPayloadCollection::V07(vec![
689        construct_tracer_payload(trace_chunks, tracer_header_tags, tracer_payload_tags),
690    ]))
691}
692
693/// Returns true if a span should be measured (i.e., it should get trace metrics calculated).
694pub fn is_measured(span: &pb::Span) -> bool {
695    span.metrics.get(MEASURED_KEY).is_some_and(|v| *v == 1.0)
696}
697
698/// Returns true if the span is a partial snapshot.
699/// This kind of spans are partial images of long-running spans.
700/// When incomplete, a partial snapshot has a metric _dd.partial_version which is a positive
701/// integer. The metric usually increases each time a new version of the same span is sent by the
702/// tracer
703pub fn is_partial_snapshot(span: &pb::Span) -> bool {
704    span.metrics
705        .get(PARTIAL_VERSION_KEY)
706        .is_some_and(|v| *v >= 0.0)
707}
708
709#[cfg(test)]
710mod tests {
711    use super::*;
712    use crate::{
713        span::SharedDictBytes,
714        test_utils::{create_test_no_alloc_span, create_test_span},
715    };
716    use http::Request;
717    use libdd_common::{http_common, Endpoint};
718    use serde_json::json;
719
720    fn find_index_in_dict(dict: &SharedDictBytes, value: &str) -> Option<u32> {
721        let idx = dict.iter().position(|e| e.as_str() == value);
722        idx.map(|idx| idx.try_into().unwrap())
723    }
724
725    #[test]
726    fn test_coalescing_does_not_exceed_max_size() {
727        fn dummy() -> SendData {
728            SendData::new(
729                MAX_PAYLOAD_SIZE / 5 + 1,
730                TracerPayloadCollection::V07(vec![pb::TracerPayload {
731                    container_id: "".to_string(),
732                    language_name: "".to_string(),
733                    language_version: "".to_string(),
734                    tracer_version: "".to_string(),
735                    runtime_id: "".to_string(),
736                    chunks: vec![pb::TraceChunk {
737                        priority: 0,
738                        origin: "".to_string(),
739                        spans: vec![],
740                        tags: Default::default(),
741                        dropped_trace: false,
742                    }],
743                    tags: Default::default(),
744                    env: "".to_string(),
745                    hostname: "".to_string(),
746                    app_version: "".to_string(),
747                    container_debug: None,
748                }]),
749                TracerHeaderTags::default(),
750                &Endpoint::default(),
751            )
752        }
753        let coalesced = coalesce_send_data(vec![dummy(), dummy(), dummy(), dummy(), dummy()]);
754        assert_eq!(
755            5,
756            coalesced
757                .iter()
758                .map(|s| s.tracer_payloads.size())
759                .sum::<usize>()
760        );
761        // assert some chunks are actually coalesced
762        assert!(
763            coalesced
764                .iter()
765                .map(|s| {
766                    if let TracerPayloadCollection::V07(collection) = &s.tracer_payloads {
767                        collection.iter().map(|s| s.chunks.len()).max().unwrap()
768                    } else {
769                        0
770                    }
771                })
772                .max()
773                .unwrap()
774                > 1
775        );
776        assert!(coalesced.len() > 1 && coalesced.len() < 5);
777    }
778
779    #[tokio::test]
780    #[allow(clippy::type_complexity)]
781    #[cfg_attr(all(miri, target_os = "macos"), ignore)]
782    async fn test_get_v05_traces_from_request_body() {
783        let data: (
784            Vec<String>,
785            Vec<
786                Vec<(
787                    u8,
788                    u8,
789                    u8,
790                    u64,
791                    u64,
792                    u64,
793                    i64,
794                    i64,
795                    i32,
796                    HashMap<u8, u8>,
797                    HashMap<u8, f64>,
798                    u8,
799                )>,
800            >,
801        ) = (
802            vec![
803                "baggage".to_string(),
804                "item".to_string(),
805                "elasticsearch.version".to_string(),
806                "7.0".to_string(),
807                "my-name".to_string(),
808                "X".to_string(),
809                "my-service".to_string(),
810                "my-resource".to_string(),
811                "_dd.sampling_rate_whatever".to_string(),
812                "value whatever".to_string(),
813                "sql".to_string(),
814            ],
815            vec![vec![(
816                6,
817                4,
818                7,
819                1,
820                2,
821                3,
822                123,
823                456,
824                1,
825                HashMap::from([(8, 9), (0, 1), (2, 3)]),
826                HashMap::from([(5, 1.2)]),
827                10,
828            )]],
829        );
830        let bytes = rmp_serde::to_vec(&data).unwrap();
831        let res = get_v05_traces_from_request_body(http_common::Body::from(bytes)).await;
832        assert!(res.is_ok());
833        let (_, traces) = res.unwrap();
834        let span = traces[0][0].clone();
835        let test_span = pb::Span {
836            service: "my-service".to_string(),
837            name: "my-name".to_string(),
838            resource: "my-resource".to_string(),
839            trace_id: 1,
840            span_id: 2,
841            parent_id: 3,
842            start: 123,
843            duration: 456,
844            error: 1,
845            meta: HashMap::from([
846                ("baggage".to_string(), "item".to_string()),
847                ("elasticsearch.version".to_string(), "7.0".to_string()),
848                (
849                    "_dd.sampling_rate_whatever".to_string(),
850                    "value whatever".to_string(),
851                ),
852            ]),
853            metrics: HashMap::from([("X".to_string(), 1.2)]),
854            meta_struct: HashMap::default(),
855            r#type: "sql".to_string(),
856            span_links: vec![],
857            span_events: vec![],
858        };
859        assert_eq!(span, test_span);
860    }
861
862    #[tokio::test]
863    #[cfg_attr(miri, ignore)]
864    async fn test_get_traces_from_request_body() {
865        let pairs = vec![
866            (
867                json!([{
868                    "service": "test-service",
869                    "name": "test-service-name",
870                    "resource": "test-service-resource",
871                    "trace_id": 111,
872                    "span_id": 222,
873                    "parent_id": 333,
874                    "start": 1,
875                    "duration": 5,
876                    "error": 0,
877                    "meta": {},
878                    "metrics": {},
879                }]),
880                vec![vec![pb::Span {
881                    service: "test-service".to_string(),
882                    name: "test-service-name".to_string(),
883                    resource: "test-service-resource".to_string(),
884                    trace_id: 111,
885                    span_id: 222,
886                    parent_id: 333,
887                    start: 1,
888                    duration: 5,
889                    error: 0,
890                    meta: HashMap::new(),
891                    metrics: HashMap::new(),
892                    meta_struct: HashMap::new(),
893                    r#type: "".to_string(),
894                    span_links: vec![],
895                    span_events: vec![],
896                }]],
897            ),
898            (
899                json!([{
900                    "name": "test-service-name",
901                    "resource": "test-service-resource",
902                    "trace_id": 111,
903                    "span_id": 222,
904                    "start": 1,
905                    "duration": 5,
906                    "meta": {},
907                }]),
908                vec![vec![pb::Span {
909                    service: "".to_string(),
910                    name: "test-service-name".to_string(),
911                    resource: "test-service-resource".to_string(),
912                    trace_id: 111,
913                    span_id: 222,
914                    parent_id: 0,
915                    start: 1,
916                    duration: 5,
917                    error: 0,
918                    meta: HashMap::new(),
919                    metrics: HashMap::new(),
920                    meta_struct: HashMap::new(),
921                    r#type: "".to_string(),
922                    span_links: vec![],
923                    span_events: vec![],
924                }]],
925            ),
926        ];
927
928        for (trace_input, output) in pairs {
929            let bytes = rmp_serde::to_vec(&vec![&trace_input]).unwrap();
930            let request = Request::builder()
931                .body(http_common::Body::from(bytes))
932                .unwrap();
933            let res = get_traces_from_request_body(request.into_body()).await;
934            assert!(res.is_ok());
935            assert_eq!(res.unwrap().1, output);
936        }
937    }
938
939    #[tokio::test]
940    #[cfg_attr(miri, ignore)]
941    async fn test_get_traces_from_request_body_with_span_links() {
942        let trace_input = json!([[{
943            "service": "test-service",
944            "name": "test-name",
945            "resource": "test-resource",
946            "trace_id": 111,
947            "span_id": 222,
948            "parent_id": 333,
949            "start": 1,
950            "duration": 5,
951            "error": 0,
952            "meta": {},
953            "metrics": {},
954            "span_links": [{
955                "trace_id": 999,
956                "span_id": 888,
957                "trace_id_high": 777,
958                "attributes": {"key": "value"},
959                "tracestate": "vendor=value"
960                // flags field intentionally omitted
961            }]
962        }]]);
963
964        let expected_output = vec![vec![pb::Span {
965            service: "test-service".to_string(),
966            name: "test-name".to_string(),
967            resource: "test-resource".to_string(),
968            trace_id: 111,
969            span_id: 222,
970            parent_id: 333,
971            start: 1,
972            duration: 5,
973            error: 0,
974            meta: HashMap::new(),
975            metrics: HashMap::new(),
976            meta_struct: HashMap::new(),
977            r#type: String::new(),
978            span_links: vec![pb::SpanLink {
979                trace_id: 999,
980                span_id: 888,
981                trace_id_high: 777,
982                attributes: HashMap::from([("key".to_string(), "value".to_string())]),
983                tracestate: "vendor=value".to_string(),
984                flags: 0, // Should default to 0 when omitted
985            }],
986            span_events: vec![],
987        }]];
988
989        let bytes = rmp_serde::to_vec(&trace_input).unwrap();
990        let request = Request::builder()
991            .body(http_common::Body::from(bytes))
992            .unwrap();
993
994        let res = get_traces_from_request_body(request.into_body()).await;
995        assert!(res.is_ok(), "Failed to deserialize: {res:?}");
996        assert_eq!(res.unwrap().1, expected_output);
997    }
998
999    #[test]
1000    fn test_get_root_span_index_from_complete_trace() {
1001        let trace = vec![
1002            create_test_span(1234, 12341, 0, 1, false),
1003            create_test_span(1234, 12342, 12341, 1, false),
1004            create_test_span(1234, 12343, 12342, 1, false),
1005        ];
1006
1007        let root_span_index = get_root_span_index(&trace);
1008        assert!(root_span_index.is_ok());
1009        assert_eq!(root_span_index.unwrap(), 0);
1010    }
1011
1012    #[test]
1013    fn test_get_root_span_index_from_partial_trace() {
1014        let trace = vec![
1015            create_test_span(1234, 12342, 12341, 1, false),
1016            create_test_span(1234, 12341, 12340, 1, false), /* this is the root span, it's
1017                                                             * parent is not in the trace */
1018            create_test_span(1234, 12343, 12342, 1, false),
1019        ];
1020
1021        let root_span_index = get_root_span_index(&trace);
1022        assert!(root_span_index.is_ok());
1023        assert_eq!(root_span_index.unwrap(), 1);
1024    }
1025
1026    #[test]
1027    fn test_set_serverless_root_span_tags_azure_function() {
1028        let mut span = create_test_span(1234, 12342, 12341, 1, false);
1029        set_serverless_root_span_tags(
1030            &mut span,
1031            Some("test_function".to_string()),
1032            &EnvironmentType::AzureFunction,
1033        );
1034        assert_eq!(
1035            span.meta,
1036            HashMap::from([
1037                (
1038                    "runtime-id".to_string(),
1039                    "test-runtime-id-value".to_string()
1040                ),
1041                ("_dd.origin".to_string(), "azurefunction".to_string()),
1042                ("origin".to_string(), "azurefunction".to_string()),
1043                ("functionname".to_string(), "test_function".to_string()),
1044                ("env".to_string(), "test-env".to_string()),
1045                ("service".to_string(), "test-service".to_string())
1046            ]),
1047        );
1048    }
1049
1050    #[test]
1051    fn test_set_serverless_root_span_tags_cloud_function() {
1052        let mut span = create_test_span(1234, 12342, 12341, 1, false);
1053        set_serverless_root_span_tags(
1054            &mut span,
1055            Some("test_function".to_string()),
1056            &EnvironmentType::CloudFunction,
1057        );
1058        assert_eq!(
1059            span.meta,
1060            HashMap::from([
1061                (
1062                    "runtime-id".to_string(),
1063                    "test-runtime-id-value".to_string()
1064                ),
1065                ("_dd.origin".to_string(), "cloudfunction".to_string()),
1066                ("origin".to_string(), "cloudfunction".to_string()),
1067                ("functionname".to_string(), "test_function".to_string()),
1068                ("env".to_string(), "test-env".to_string()),
1069                ("service".to_string(), "test-service".to_string())
1070            ]),
1071        );
1072    }
1073
1074    #[test]
1075    fn test_has_top_level() {
1076        let top_level_span = create_test_span(123, 1234, 12, 1, true);
1077        let not_top_level_span = create_test_span(123, 1234, 12, 1, false);
1078        assert!(has_top_level(&top_level_span));
1079        assert!(!has_top_level(&not_top_level_span));
1080    }
1081
1082    #[test]
1083    fn test_is_measured() {
1084        let mut measured_span = create_test_span(123, 1234, 12, 1, true);
1085        measured_span.metrics.insert(MEASURED_KEY.into(), 1.0);
1086        let not_measured_span = create_test_span(123, 1234, 12, 1, true);
1087        assert!(is_measured(&measured_span));
1088        assert!(!is_measured(&not_measured_span));
1089    }
1090
1091    #[test]
1092    fn test_compute_top_level() {
1093        let mut span_with_different_service = create_test_span(123, 5, 2, 1, false);
1094        span_with_different_service.service = "another_service".into();
1095        let mut trace = vec![
1096            // Root span, should be marked as top-level
1097            create_test_span(123, 1, 0, 1, false),
1098            // Should not be marked as top-level
1099            create_test_span(123, 2, 1, 1, false),
1100            // No parent in local trace, should be marked as
1101            // top-level
1102            create_test_span(123, 4, 3, 1, false),
1103            // Parent belongs to another service, should be marked
1104            // as top-level
1105            span_with_different_service,
1106        ];
1107
1108        compute_top_level_span(trace.as_mut_slice());
1109
1110        let spans_marked_as_top_level: Vec<u64> = trace
1111            .iter()
1112            .filter_map(|span| {
1113                if has_top_level(span) {
1114                    Some(span.span_id)
1115                } else {
1116                    None
1117                }
1118            })
1119            .collect();
1120        assert_eq!(spans_marked_as_top_level, [1, 4, 5])
1121    }
1122
1123    #[test]
1124    fn test_collect_trace_chunks_v05() {
1125        let chunk = vec![create_test_no_alloc_span(123, 456, 789, 1, true)];
1126
1127        let collection = collect_trace_chunks(vec![chunk], TraceEncoding::V05).unwrap();
1128
1129        let (dict, traces) = match collection {
1130            TraceChunks::V05(payload) => payload,
1131            _ => panic!("Unexpected type"),
1132        };
1133
1134        assert_eq!(dict.len(), 16);
1135
1136        let span = &traces[0][0];
1137        assert_eq!(span.service, 1);
1138        assert_eq!(span.name, 2);
1139        assert_eq!(span.resource, 3);
1140        assert_eq!(span.trace_id, 123);
1141        assert_eq!(span.span_id, 456);
1142        assert_eq!(span.parent_id, 789);
1143        assert_eq!(span.start, 1);
1144        assert_eq!(span.error, 0);
1145        assert_eq!(span.error, 0);
1146        assert_eq!(span.r#type, 15);
1147        assert_eq!(
1148            *span
1149                .meta
1150                .get(&find_index_in_dict(&dict, "service").unwrap())
1151                .unwrap(),
1152            find_index_in_dict(&dict, "test-service").unwrap()
1153        );
1154        assert_eq!(
1155            *span
1156                .meta
1157                .get(&find_index_in_dict(&dict, "env").unwrap())
1158                .unwrap(),
1159            find_index_in_dict(&dict, "test-env").unwrap()
1160        );
1161        assert_eq!(
1162            *span
1163                .meta
1164                .get(&find_index_in_dict(&dict, "runtime-id").unwrap())
1165                .unwrap(),
1166            find_index_in_dict(&dict, "test-runtime-id-value").unwrap()
1167        );
1168        assert_eq!(
1169            *span
1170                .meta
1171                .get(&find_index_in_dict(&dict, "_dd.origin").unwrap())
1172                .unwrap(),
1173            find_index_in_dict(&dict, "cloudfunction").unwrap()
1174        );
1175        assert_eq!(
1176            *span
1177                .meta
1178                .get(&find_index_in_dict(&dict, "origin").unwrap())
1179                .unwrap(),
1180            find_index_in_dict(&dict, "cloudfunction").unwrap()
1181        );
1182        assert_eq!(
1183            *span
1184                .meta
1185                .get(&find_index_in_dict(&dict, "functionname").unwrap())
1186                .unwrap(),
1187            find_index_in_dict(&dict, "dummy_function_name").unwrap()
1188        );
1189        assert_eq!(
1190            *span
1191                .metrics
1192                .get(&find_index_in_dict(&dict, "_top_level").unwrap())
1193                .unwrap(),
1194            1.0
1195        );
1196    }
1197
1198    #[test]
1199    fn test_collect_trace_chunks_v04() {
1200        let chunk = vec![create_test_no_alloc_span(123, 456, 789, 1, true)];
1201
1202        let collection = collect_trace_chunks(vec![chunk], TraceEncoding::V04).unwrap();
1203
1204        let traces = match collection {
1205            TraceChunks::V04(traces) => traces,
1206            _ => panic!("Unexpected type"),
1207        };
1208
1209        assert_eq!(traces.len(), 1);
1210        assert_eq!(traces[0].len(), 1);
1211        let span = &traces[0][0];
1212        assert_eq!(span.trace_id, 123);
1213        assert_eq!(span.span_id, 456);
1214        assert_eq!(span.parent_id, 789);
1215        assert_eq!(span.start, 1);
1216        assert_eq!(span.error, 0);
1217    }
1218
1219    #[test]
1220    fn test_rmp_serde_deserialize_meta_with_null_values() {
1221        // Create a JSON representation with null value in meta
1222        let span_json = json!({
1223            "service": "test-service",
1224            "name": "test_name",
1225            "resource": "test-resource",
1226            "trace_id": 1_u64,
1227            "span_id": 2_u64,
1228            "parent_id": 0_u64,
1229            "start": 0_i64,
1230            "duration": 5_i64,
1231            "error": 0_i32,
1232            "meta": {
1233                "service": "test-service",
1234                "env": "test-env",
1235                "runtime-id": "test-runtime-id-value",
1236                "problematic_key": null  // Ensure this null value does not cause an error
1237            },
1238            "metrics": {},
1239            "type": "",
1240            "meta_struct": {},
1241            "span_links": [],
1242            "span_events": []
1243        });
1244
1245        let traces_json = vec![vec![span_json]];
1246        let encoded_data = rmp_serde::to_vec(&traces_json).unwrap();
1247        let traces: Vec<Vec<pb::Span>> = rmp_serde::from_read(&encoded_data[..])
1248            .expect("Failed to deserialize traces with null values in meta");
1249
1250        assert_eq!(1, traces.len());
1251        assert_eq!(1, traces[0].len());
1252        let decoded_span = &traces[0][0];
1253
1254        assert_eq!("test-service", decoded_span.service);
1255        assert_eq!("test_name", decoded_span.name);
1256        assert_eq!("test-resource", decoded_span.resource);
1257        assert_eq!("test-service", decoded_span.meta.get("service").unwrap());
1258        assert_eq!("test-env", decoded_span.meta.get("env").unwrap());
1259        assert_eq!(
1260            "test-runtime-id-value",
1261            decoded_span.meta.get("runtime-id").unwrap()
1262        );
1263        // Assert that the null value was filtered out (key not present in map)
1264        assert!(
1265            !decoded_span.meta.contains_key("problematic_key"),
1266            "Null value should be skipped, but key was present"
1267        );
1268    }
1269
1270    #[test]
1271    fn test_enrich_span_with_azure_function_metadata_adds_tags_for_non_apim() {
1272        let mut span = create_test_span(1234, 12342, 12341, 1, false);
1273        span.name = "azure.function".to_string();
1274
1275        enrich_span_with_azure_function_metadata(&mut span);
1276
1277        // If AAS_METADATA_FUNCTION is available, verify aas.* tags were added
1278        // If not available (most test environments), this is a no-op
1279        // This test primarily ensures the function doesn't skip non-apim spans
1280        if azure_app_services::AAS_METADATA_FUNCTION.is_some() {
1281            assert!(span.meta.contains_key("aas.resource.id"));
1282            assert!(span.meta.contains_key("aas.environment.instance_id"));
1283            assert!(span.meta.contains_key("aas.environment.instance_name"));
1284            assert!(span.meta.contains_key("aas.subscription.id"));
1285            assert!(span.meta.contains_key("aas.environment.os"));
1286            assert!(span.meta.contains_key("aas.environment.runtime"));
1287            assert!(span.meta.contains_key("aas.environment.runtime_version"));
1288            assert!(span.meta.contains_key("aas.environment.function_runtime"));
1289            assert!(span.meta.contains_key("aas.resource.group"));
1290            assert!(span.meta.contains_key("aas.site.name"));
1291            assert!(span.meta.contains_key("aas.site.kind"));
1292            assert!(span.meta.contains_key("aas.site.type"));
1293        }
1294    }
1295
1296    #[test]
1297    fn test_enrich_span_with_azure_function_metadata_skips_azure_apim() {
1298        let mut span = create_test_span(1234, 12342, 12341, 1, false);
1299        span.name = "azure.apim".to_string();
1300
1301        enrich_span_with_azure_function_metadata(&mut span);
1302
1303        // Verify no aas.* tags were added
1304        assert!(!span.meta.contains_key("aas.resource.id"));
1305        assert!(!span.meta.contains_key("aas.environment.instance_id"));
1306        assert!(!span.meta.contains_key("aas.environment.instance_name"));
1307        assert!(!span.meta.contains_key("aas.subscription.id"));
1308        assert!(!span.meta.contains_key("aas.environment.os"));
1309        assert!(!span.meta.contains_key("aas.environment.runtime"));
1310        assert!(!span.meta.contains_key("aas.environment.runtime_version"));
1311        assert!(!span.meta.contains_key("aas.environment.function_runtime"));
1312        assert!(!span.meta.contains_key("aas.resource.group"));
1313        assert!(!span.meta.contains_key("aas.site.name"));
1314        assert!(!span.meta.contains_key("aas.site.kind"));
1315        assert!(!span.meta.contains_key("aas.site.type"));
1316    }
1317
1318    #[test]
1319    fn test_collect_pb_trace_chunks_searches_multiple_root_spans_for_fields() {
1320        // First trace root span has no fields. Second trace root span has all fields.
1321        // The second root span should populate all fields.
1322        let mut first_root_span = create_test_span(1, 1, 0, 1, true);
1323        first_root_span.meta.remove("env");
1324        first_root_span.meta.remove("runtime-id");
1325
1326        let mut second_root_span = create_test_span(2, 3, 0, 1, true);
1327        second_root_span
1328            .meta
1329            .insert("version".to_string(), "1.2.3".to_string());
1330        second_root_span
1331            .meta
1332            .insert("env".to_string(), "prod".to_string());
1333        second_root_span
1334            .meta
1335            .insert("_dd.hostname".to_string(), "my-host".to_string());
1336        second_root_span
1337            .meta
1338            .insert("runtime-id".to_string(), "123".to_string());
1339
1340        let result = collect_pb_trace_chunks(
1341            vec![vec![first_root_span], vec![second_root_span]],
1342            &TracerHeaderTags::default(),
1343            &mut tracer_payload::DefaultTraceChunkProcessor,
1344            true,
1345        )
1346        .unwrap();
1347
1348        let TracerPayloadCollection::V07(payloads) = result else {
1349            panic!("expected TracerPayloadCollection::V07");
1350        };
1351        assert_eq!(payloads[0].app_version, "1.2.3");
1352        assert_eq!(payloads[0].env, "prod");
1353        assert_eq!(payloads[0].hostname, "my-host");
1354        assert_eq!(payloads[0].runtime_id, "123");
1355    }
1356
1357    #[test]
1358    fn test_collect_pb_trace_chunks_searches_non_root_spans_for_fields() {
1359        // Root span has no fields. Child span has all fields.
1360        // The child span should populate all fields.
1361        let mut root_span = create_test_span(1, 1, 0, 1, true);
1362        root_span.meta.remove("env");
1363        root_span.meta.remove("runtime-id");
1364        let mut child_span = create_test_span(1, 2, 1, 1, false);
1365        child_span
1366            .meta
1367            .insert("version".to_string(), "1.2.3".to_string());
1368        child_span
1369            .meta
1370            .insert("env".to_string(), "prod".to_string());
1371        child_span
1372            .meta
1373            .insert("_dd.hostname".to_string(), "my-host".to_string());
1374        child_span
1375            .meta
1376            .insert("runtime-id".to_string(), "123".to_string());
1377
1378        let result = collect_pb_trace_chunks(
1379            vec![vec![root_span, child_span]],
1380            &TracerHeaderTags::default(),
1381            &mut tracer_payload::DefaultTraceChunkProcessor,
1382            true,
1383        )
1384        .unwrap();
1385
1386        let TracerPayloadCollection::V07(payloads) = result else {
1387            panic!("expected TracerPayloadCollection::V07");
1388        };
1389        assert_eq!(payloads[0].app_version, "1.2.3");
1390        assert_eq!(payloads[0].env, "prod");
1391        assert_eq!(payloads[0].hostname, "my-host");
1392        assert_eq!(payloads[0].runtime_id, "123");
1393    }
1394
1395    #[test]
1396    fn test_collect_pb_trace_chunks_root_span_takes_priority_over_child() {
1397        // Root span has all fields. Child has different values for all fields.
1398        // The root span should populate all fields.
1399        let mut root_span = create_test_span(1, 1, 0, 1, true);
1400        root_span
1401            .meta
1402            .insert("version".to_string(), "root-version".to_string());
1403        root_span
1404            .meta
1405            .insert("env".to_string(), "root-env".to_string());
1406        root_span
1407            .meta
1408            .insert("_dd.hostname".to_string(), "root-host".to_string());
1409        root_span
1410            .meta
1411            .insert("runtime-id".to_string(), "root-runtime-id".to_string());
1412
1413        let mut child_span = create_test_span(1, 2, 1, 1, false);
1414        child_span
1415            .meta
1416            .insert("version".to_string(), "child-version".to_string());
1417        child_span
1418            .meta
1419            .insert("env".to_string(), "child-env".to_string());
1420        child_span
1421            .meta
1422            .insert("_dd.hostname".to_string(), "child-host".to_string());
1423        child_span
1424            .meta
1425            .insert("runtime-id".to_string(), "child-runtime-id".to_string());
1426
1427        let result = collect_pb_trace_chunks(
1428            vec![vec![root_span, child_span]],
1429            &TracerHeaderTags::default(),
1430            &mut tracer_payload::DefaultTraceChunkProcessor,
1431            true,
1432        )
1433        .unwrap();
1434
1435        let TracerPayloadCollection::V07(payloads) = result else {
1436            panic!("expected TracerPayloadCollection::V07");
1437        };
1438        assert_eq!(payloads[0].app_version, "root-version");
1439        assert_eq!(payloads[0].env, "root-env");
1440        assert_eq!(payloads[0].hostname, "root-host");
1441        assert_eq!(payloads[0].runtime_id, "root-runtime-id");
1442    }
1443
1444    #[test]
1445    fn test_collect_pb_trace_chunks_skips_empty_root_span_value() {
1446        // Root span has empty values for all fields. Child span has non-empty values.
1447        // The child span should populate all fields.
1448        let mut root_span = create_test_span(1, 1, 0, 1, true);
1449        root_span.meta.insert("version".to_string(), "".to_string());
1450        root_span.meta.insert("env".to_string(), "".to_string());
1451        root_span
1452            .meta
1453            .insert("_dd.hostname".to_string(), "".to_string());
1454        root_span
1455            .meta
1456            .insert("runtime-id".to_string(), "".to_string());
1457
1458        let mut child_span = create_test_span(1, 2, 1, 1, false);
1459        child_span
1460            .meta
1461            .insert("version".to_string(), "1.2.3".to_string());
1462        child_span
1463            .meta
1464            .insert("env".to_string(), "prod".to_string());
1465        child_span
1466            .meta
1467            .insert("_dd.hostname".to_string(), "my-host".to_string());
1468        child_span
1469            .meta
1470            .insert("runtime-id".to_string(), "123".to_string());
1471
1472        let result = collect_pb_trace_chunks(
1473            vec![vec![root_span, child_span]],
1474            &TracerHeaderTags::default(),
1475            &mut tracer_payload::DefaultTraceChunkProcessor,
1476            true,
1477        )
1478        .unwrap();
1479
1480        let TracerPayloadCollection::V07(payloads) = result else {
1481            panic!("expected TracerPayloadCollection::V07");
1482        };
1483        assert_eq!(payloads[0].app_version, "1.2.3");
1484        assert_eq!(payloads[0].env, "prod");
1485        assert_eq!(payloads[0].hostname, "my-host");
1486        assert_eq!(payloads[0].runtime_id, "123");
1487    }
1488
1489    #[test]
1490    fn test_collect_pb_trace_chunks_normalizes_env() {
1491        let mut root = create_test_span(1, 1, 0, 1, true);
1492        root.meta
1493            .insert("env".to_string(), "PRODUCTION".to_string());
1494
1495        let result = collect_pb_trace_chunks(
1496            vec![vec![root]],
1497            &TracerHeaderTags::default(),
1498            &mut tracer_payload::DefaultTraceChunkProcessor,
1499            true,
1500        )
1501        .unwrap();
1502
1503        let TracerPayloadCollection::V07(payloads) = result else {
1504            panic!("expected TracerPayloadCollection::V07");
1505        };
1506        assert_eq!(payloads[0].env, "production");
1507    }
1508
1509    #[test]
1510    fn test_collect_pb_trace_chunks_skips_env_empty_after_normalization() {
1511        // First root span has an env that normalizes to empty (all invalid characters).
1512        // Second root span has an env should populate env fields.
1513        let mut first_root_span = create_test_span(1, 1, 0, 1, true);
1514        first_root_span
1515            .meta
1516            .insert("env".to_string(), "!!!".to_string());
1517
1518        let mut second_root_span = create_test_span(2, 3, 0, 1, true);
1519        second_root_span
1520            .meta
1521            .insert("env".to_string(), "prod".to_string());
1522
1523        let result = collect_pb_trace_chunks(
1524            vec![vec![first_root_span], vec![second_root_span]],
1525            &TracerHeaderTags::default(),
1526            &mut tracer_payload::DefaultTraceChunkProcessor,
1527            true,
1528        )
1529        .unwrap();
1530
1531        let TracerPayloadCollection::V07(payloads) = result else {
1532            panic!("expected TracerPayloadCollection::V07");
1533        };
1534        assert_eq!(payloads[0].env, "prod");
1535    }
1536
1537    #[test]
1538    fn test_search_trace_for_field_skips_span_with_same_id_as_root() {
1539        // A span with the same span_id as root is treated as the root and skipped
1540        // in the child span search. Only the root spans own meta is checked for it.
1541        let mut root = create_test_span(1, 1, 0, 1, true);
1542        root.meta.remove("version");
1543
1544        // This span shares the same span_id as the root span, it should be skipped.
1545        let mut duplicate = create_test_span(1, 1, 0, 1, false);
1546        duplicate
1547            .meta
1548            .insert("version".to_string(), "should-not-appear".to_string());
1549
1550        let trace = vec![root.clone(), duplicate];
1551        assert_eq!(search_trace_for_field(&root, &trace, "version"), None);
1552    }
1553}