1pub 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
26pub const MAX_PAYLOAD_SIZE: usize = 25 * 1024 * 1024;
30const TOP_LEVEL_KEY: &str = "_top_level";
32const 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
39pub 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 span.service = get_v05_string(reader, dict, "service")?;
90 span.name = get_v05_string(reader, dict, "name")?;
92 span.resource = get_v05_string(reader, dict, "resource")?;
94
95 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 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 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 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 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 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 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 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 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#[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
274fn 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 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 if a.size + b.size < MAX_PAYLOAD_SIZE / 2 {
357 if b.tracer_payloads.append(&mut a.tracer_payloads) {
361 b.size += a.size;
362 return true;
363 }
364 }
365 }
366 false
367 });
368 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 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 !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
416pub 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 set_top_level_span(span)
437 }
438 }
439 None => {
440 set_top_level_span(span)
442 }
443 }
444 }
445}
446
447pub 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", };
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
591pub 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 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 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 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 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
699pub fn is_measured(span: &pb::Span) -> bool {
701 span.metrics.get(MEASURED_KEY).is_some_and(|v| *v == 1.0)
702}
703
704pub 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!(
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 }]
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, }],
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), 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(¬_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(¬_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 create_test_span(123, 1, 0, 1, false),
1104 create_test_span(123, 2, 1, 1, false),
1106 create_test_span(123, 4, 3, 1, false),
1109 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 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 },
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!(
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 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 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 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 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 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 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 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 let mut root = create_test_span(1, 1, 0, 1, true);
1527 root.meta.remove("version");
1528
1529 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}