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::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
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 b.tracer_payloads.append(&mut a.tracer_payloads);
359 b.size += a.size;
360 return true;
361 }
362 }
363 false
364 });
365 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 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 !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
413pub 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 set_top_level_span(span)
434 }
435 }
436 None => {
437 set_top_level_span(span)
439 }
440 }
441 }
442}
443
444pub 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", };
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 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 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 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 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
693pub fn is_measured(span: &pb::Span) -> bool {
695 span.metrics.get(MEASURED_KEY).is_some_and(|v| *v == 1.0)
696}
697
698pub 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!(
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 }]
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, }],
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), 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(¬_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(¬_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 create_test_span(123, 1, 0, 1, false),
1098 create_test_span(123, 2, 1, 1, false),
1100 create_test_span(123, 4, 3, 1, false),
1103 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 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 },
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!(
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 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 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 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 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 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 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 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 let mut root = create_test_span(1, 1, 0, 1, true);
1542 root.meta.remove("version");
1543
1544 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}