1use std::{
6 sync::Arc,
7 time::{Duration, SystemTime, UNIX_EPOCH},
8};
9
10use agent_data_plane_config::domains;
11use async_trait::async_trait;
12use saluki_core::{
13 accounting::{MemoryBounds, MemoryBoundsBuilder},
14 components::{transforms::*, BuildContext},
15 data_model::{
16 event::{
17 trace::{AttributeValue, Trace},
18 trace_stats::{ClientStatsPayload, TraceStats},
19 Event, EventType,
20 },
21 origin::OriginTagCardinality,
22 tags::TagSet,
23 },
24 topology::OutputDefinition,
25};
26use saluki_env::{
27 host::providers::BoxedHostProvider, workload::EntityId, EnvironmentProvider, HostProvider, WorkloadProvider,
28};
29use saluki_error::{ErrorContext as _, GenericError};
30use stringtheory::MetaString;
31use tokio::{select, time::interval};
32use tracing::{debug, error};
33
34use crate::common::otlp::util::extract_container_tags_from_attributes_map;
35
36mod aggregation;
37pub(crate) use self::aggregation::{process_tags_hash, PayloadAggregationKey};
38
39mod peer_ip_quantize;
40
41mod peer_tags;
42
43mod span_concentrator;
44pub(crate) use self::span_concentrator::{InfraTags, SpanConcentrator};
45
46mod statsraw;
47
48mod weight;
49use self::weight::weight;
50
51const DEFAULT_FLUSH_INTERVAL: Duration = Duration::from_secs(10);
53
54const TAG_PROCESS_TAGS: &str = "_dd.tags.process";
56
57const MAX_STATS_GROUPS_PER_EVENT: usize = 4000;
59
60pub struct ApmStatsTransformConfiguration {
65 compute_stats_by_span_kind: bool,
66 peer_tags_aggregation: bool,
67 peer_tags: Vec<MetaString>,
68 default_env: MetaString,
69 default_hostname: Option<String>,
70 workload_provider: Option<Arc<dyn WorkloadProvider + Send + Sync>>,
71}
72
73impl ApmStatsTransformConfiguration {
74 pub fn from_configuration(config: &domains::traces::Domain) -> Self {
76 Self {
77 compute_stats_by_span_kind: config.compute_stats_by_span_kind,
78 peer_tags_aggregation: config.peer_tags_aggregation,
79 peer_tags: config.peer_tags.iter().cloned().map(MetaString::from).collect(),
80 default_env: MetaString::from(config.default_env.clone()),
81 default_hostname: None,
82 workload_provider: None,
83 }
84 }
85
86 pub async fn with_environment_provider<E>(mut self, env_provider: E) -> Result<Self, GenericError>
88 where
89 E: EnvironmentProvider<Host = BoxedHostProvider>,
90 {
91 let hostname = env_provider.host().get_hostname().await?;
92 self.default_hostname = Some(hostname);
93 Ok(self)
94 }
95
96 pub fn with_workload_provider<W>(mut self, workload_provider: W) -> Self
100 where
101 W: WorkloadProvider + Send + Sync + 'static,
102 {
103 self.workload_provider = Some(Arc::new(workload_provider));
104 self
105 }
106}
107
108#[async_trait]
109impl TransformBuilder for ApmStatsTransformConfiguration {
110 async fn build(&self, _context: BuildContext) -> Result<Box<dyn Transform + Send>, GenericError> {
111 let agent_hostname = MetaString::from(self.default_hostname.clone().unwrap_or_default());
112 let concentrator = SpanConcentrator::new(
113 self.compute_stats_by_span_kind,
114 self.peer_tags_aggregation,
115 &self.peer_tags,
116 now_nanos(),
117 );
118
119 Ok(Box::new(ApmStats {
120 concentrator,
121 flush_interval: DEFAULT_FLUSH_INTERVAL,
122 agent_env: self.default_env.clone(),
123 agent_hostname,
124 workload_provider: self.workload_provider.clone(),
125 }))
126 }
127
128 fn input_event_type(&self) -> EventType {
129 EventType::Trace
130 }
131
132 fn outputs(&self) -> &[OutputDefinition<EventType>] {
133 static OUTPUTS: &[OutputDefinition<EventType>] = &[OutputDefinition::default_output(EventType::TraceStats)];
134 OUTPUTS
135 }
136}
137
138impl MemoryBounds for ApmStatsTransformConfiguration {
139 fn specify_bounds(&self, builder: &mut MemoryBoundsBuilder) {
140 builder.minimum().with_single_value::<ApmStats>("component struct");
141 }
143}
144
145struct ApmStats {
146 concentrator: SpanConcentrator,
147 flush_interval: Duration,
148 agent_env: MetaString,
149 agent_hostname: MetaString,
150 workload_provider: Option<Arc<dyn WorkloadProvider + Send + Sync>>,
151}
152
153impl ApmStats {
154 fn process_trace(&mut self, trace: &Trace) {
155 let root_span = trace
156 .spans()
157 .iter()
158 .find(|s| s.parent_id() == 0)
159 .or_else(|| trace.spans().first());
160
161 let trace_weight = root_span.map(weight).unwrap_or(1.0);
162
163 let process_tags = extract_process_tags(trace);
164
165 let payload_key = self.build_payload_key(trace, &process_tags);
166 let infra_tags = self.build_infra_tags(trace, &process_tags);
167
168 let origin = trace
169 .spans()
170 .first()
171 .and_then(|s| s.attributes.get("_dd.origin").and_then(AttributeValue::as_string))
172 .map(|s| s.as_ref())
173 .unwrap_or("");
174
175 for span in trace.spans() {
176 if let Some(stat_span) = self.concentrator.new_stat_span_from_span(span) {
177 self.concentrator
178 .add_span(&stat_span, trace_weight, &payload_key, &infra_tags, origin);
179 }
180 }
181 }
182
183 fn build_infra_tags(&self, trace: &Trace, process_tags: &str) -> InfraTags {
184 let container_id = trace.payload.container_id.clone();
185 let mut container_tags = if container_id.is_empty() {
186 TagSet::default()
187 } else {
188 let mut tags = TagSet::default();
189 extract_container_tags_from_attributes_map(&trace.attributes, &mut tags);
190 tags
191 };
192
193 if !container_id.is_empty() {
194 if let Some(workload_provider) = &self.workload_provider {
195 let entity_id = EntityId::Container(container_id.clone());
196 if let Some(tags) = workload_provider.get_tags_for_entity(&entity_id, OriginTagCardinality::Low) {
197 container_tags.merge_shared(&tags);
198 }
199 }
200 }
201
202 InfraTags::new(container_id, container_tags, process_tags)
203 }
204
205 fn build_payload_key(&self, trace: &Trace, process_tags: &str) -> PayloadAggregationKey {
206 let root_span = trace
207 .spans()
208 .iter()
209 .find(|s| s.parent_id() == 0)
210 .or_else(|| trace.spans().first());
211
212 let env = root_span
215 .and_then(|s| {
216 s.attributes
217 .get("env")
218 .and_then(AttributeValue::as_string)
219 .filter(|s| !s.is_empty())
220 })
221 .cloned()
222 .or_else(|| {
223 trace.spans().iter().find_map(|s| {
224 s.attributes
225 .get("env")
226 .and_then(AttributeValue::as_string)
227 .filter(|s| !s.is_empty())
228 .cloned()
229 })
230 })
231 .unwrap_or_else(|| {
232 if !trace.payload.env.is_empty() {
233 trace.payload.env.clone()
234 } else {
235 self.agent_env.clone()
236 }
237 });
238
239 let hostname = root_span
240 .and_then(|s| {
241 s.attributes
242 .get("_dd.hostname")
243 .and_then(AttributeValue::as_string)
244 .filter(|s| !s.is_empty())
245 })
246 .cloned()
247 .unwrap_or_else(|| {
248 if !trace.payload.hostname.is_empty() {
249 trace.payload.hostname.clone()
250 } else {
251 self.agent_hostname.clone()
252 }
253 });
254
255 let version = root_span
259 .and_then(|s| {
260 s.attributes
261 .get("version")
262 .and_then(AttributeValue::as_string)
263 .filter(|s| !s.is_empty())
264 })
265 .cloned()
266 .unwrap_or_else(|| trace.payload.app_version.clone());
267
268 let container_id = if !trace.payload.container_id.is_empty() {
269 trace.payload.container_id.clone()
270 } else {
271 root_span
272 .and_then(|s| s.attributes.get("_dd.container_id").and_then(AttributeValue::as_string))
273 .cloned()
274 .unwrap_or_default()
275 };
276
277 let git_commit_sha = root_span
278 .and_then(|s| {
279 s.attributes
280 .get("_dd.git.commit.sha")
281 .and_then(AttributeValue::as_string)
282 .filter(|s| !s.is_empty())
283 })
284 .cloned()
285 .unwrap_or_default();
286
287 let image_tag = root_span
288 .and_then(|s| {
289 s.attributes
290 .get("_dd.image_tag")
291 .and_then(AttributeValue::as_string)
292 .filter(|s| !s.is_empty())
293 })
294 .cloned()
295 .unwrap_or_default();
296
297 let lang = if !trace.payload.language_name.is_empty() {
298 trace.payload.language_name.clone()
299 } else {
300 root_span
301 .and_then(|s| s.attributes.get("language").and_then(AttributeValue::as_string))
302 .cloned()
303 .unwrap_or_default()
304 };
305
306 PayloadAggregationKey {
307 env,
308 hostname,
309 version,
310 container_id,
311 git_commit_sha,
312 image_tag,
313 lang,
314 process_tags_hash: process_tags_hash(process_tags),
315 }
316 }
317}
318
319fn split_into_trace_stats(client_payloads: Vec<ClientStatsPayload>, max_entries_per_event: usize) -> Vec<TraceStats> {
325 if client_payloads.is_empty() {
326 return Vec::new();
327 }
328
329 let total_grouped_entries = client_payloads
332 .iter()
333 .map(|p| p.stats().iter().map(|b| b.stats().len()).sum::<usize>())
334 .sum::<usize>();
335 if total_grouped_entries <= max_entries_per_event {
336 return vec![TraceStats::new(client_payloads)];
337 }
338
339 let mut events = Vec::new();
340 let mut current_client_payloads = Vec::new();
341 let mut current_event_len = 0;
342
343 for mut client_payload in client_payloads {
344 let client_payload_len = client_payload.stats().iter().map(|b| b.stats().len()).sum::<usize>();
346 if current_event_len + client_payload_len <= max_entries_per_event {
347 current_client_payloads.push(client_payload);
348 current_event_len += client_payload_len;
349 continue;
350 }
351
352 let mut current_client_stats_buckets = Vec::new();
359 for mut client_stats_bucket in client_payload.take_stats() {
360 let bucket_len = client_stats_bucket.stats().len();
361 if current_event_len + bucket_len <= max_entries_per_event {
363 current_client_stats_buckets.push(client_stats_bucket);
364 current_event_len += bucket_len;
365 continue;
366 }
367
368 let mut bucket_entries = client_stats_bucket.take_stats();
371 while current_event_len + bucket_entries.len() > max_entries_per_event {
372 let split_amount = max_entries_per_event - current_event_len;
376 let split_point = bucket_entries.len() - split_amount;
377 let split_entries = bucket_entries.split_off(split_point);
378
379 let split_bucket = client_stats_bucket.clone().with_stats(split_entries);
383 current_client_stats_buckets.push(split_bucket);
384
385 let split_client_payload = client_payload
386 .clone()
387 .with_stats(std::mem::take(&mut current_client_stats_buckets));
388 current_client_payloads.push(split_client_payload);
389
390 events.push(TraceStats::new(std::mem::take(&mut current_client_payloads)));
391 current_event_len = 0;
392 }
393
394 if !bucket_entries.is_empty() {
397 current_event_len += bucket_entries.len();
398 current_client_stats_buckets.push(client_stats_bucket.with_stats(bucket_entries));
399 }
400 }
401
402 if !current_client_stats_buckets.is_empty() {
404 current_client_payloads.push(client_payload.with_stats(current_client_stats_buckets));
405 }
406 }
407
408 if !current_client_payloads.is_empty() {
410 events.push(TraceStats::new(current_client_payloads));
411 }
412
413 events
414}
415
416#[async_trait]
417impl Transform for ApmStats {
418 async fn run(mut self: Box<Self>, mut context: TransformContext) -> Result<(), GenericError> {
419 let mut health = context.take_health_handle();
420
421 let mut flush_ticker = interval(self.flush_interval);
422 flush_ticker.tick().await;
423
424 let mut final_flush = false;
425
426 health.mark_ready();
427 debug!("APM Stats transform started.");
428
429 loop {
430 select! {
431 _ = health.live() => continue,
432
433 _ = flush_ticker.tick() => {
434 let stats_payloads = self.concentrator.flush(now_nanos(), final_flush);
435 if !stats_payloads.is_empty() {
436 debug!(stats_payloads = stats_payloads.len(), "Flushing APM stats.");
437
438 let events = split_into_trace_stats(stats_payloads, MAX_STATS_GROUPS_PER_EVENT);
439 let dispatcher = context.dispatcher().buffered()
440 .error_context("Default output should be available.")?;
441
442 if let Err(e) = dispatcher.send_all(events.into_iter().map(Event::TraceStats)).await {
443 error!(error = %e, "Failed to dispatch events.");
444 }
445 }
446
447 if final_flush {
448 debug!("Final APM stats flush complete.");
449 break;
450 }
451 },
452
453 maybe_events = context.events().next(), if !final_flush => {
454 match maybe_events {
455 Some(events) => {
456 for event in events {
457 if let Event::Trace(trace) = event {
458 self.process_trace(&trace);
459 }
460 }
461 },
462 None => {
463 final_flush = true;
466 flush_ticker.reset_immediately();
467 debug!("APM Stats transform stopping, triggering final flush...");
468 }
469 }
470 },
471 }
472 }
473
474 debug!("APM Stats transform stopped.");
475 Ok(())
476 }
477}
478
479fn now_nanos() -> u64 {
481 SystemTime::now()
482 .duration_since(UNIX_EPOCH)
483 .unwrap_or_default()
484 .as_nanos() as u64
485}
486
487fn extract_process_tags(trace: &Trace) -> MetaString {
489 let root_span = trace
490 .spans()
491 .iter()
492 .find(|s| s.parent_id() == 0)
493 .or_else(|| trace.spans().first());
494 if let Some(span) = root_span {
495 if let Some(tags) = span
496 .attributes
497 .get(TAG_PROCESS_TAGS)
498 .and_then(AttributeValue::as_string)
499 .filter(|s| !s.is_empty())
500 {
501 return tags.clone();
502 }
503 }
504 if let Some(AttributeValue::String(tags)) = trace.attributes.get(TAG_PROCESS_TAGS) {
505 if !tags.is_empty() {
506 return tags.clone();
507 }
508 }
509 MetaString::empty()
510}
511
512#[cfg(test)]
513mod tests {
514 use proptest::prelude::*;
515 use saluki_common::collections::FastHashMap;
516 use saluki_core::data_model::event::trace::{AttributeValue, Span};
517 use saluki_core::data_model::event::trace_stats::ClientGroupedStats;
518 use saluki_core::data_model::event::trace_stats::ClientStatsBucket;
519
520 use super::aggregation::BUCKET_DURATION_NS;
521 use super::span_concentrator::METRIC_PARTIAL_VERSION;
522 use super::*;
523
524 fn align_ts(ts: u64, bsize: u64) -> u64 {
526 ts - ts % bsize
527 }
528
529 #[allow(clippy::too_many_arguments)]
531 fn test_span(
532 aligned_now: u64, span_id: u64, parent_id: u64, duration: u64, bucket_offset: u64, service: &str,
533 resource: &str, error: i32, meta: Option<FastHashMap<MetaString, MetaString>>,
534 metrics: Option<FastHashMap<MetaString, f64>>,
535 ) -> Span {
536 let bucket_start = aligned_now - bucket_offset * BUCKET_DURATION_NS;
537 let start = bucket_start - duration;
538
539 let mut attrs: FastHashMap<MetaString, AttributeValue> = FastHashMap::default();
540 if let Some(m) = meta {
541 attrs.extend(m.into_iter().map(|(k, v)| (k, AttributeValue::String(v))));
542 }
543 if let Some(m) = metrics {
544 attrs.extend(m.into_iter().map(|(k, v)| (k, AttributeValue::Float(v))));
545 }
546 Span::new(
547 service, "query", resource, "db", span_id, parent_id, start, duration, error,
548 )
549 .with_attributes(attrs)
550 }
551
552 fn make_test_span(service: &str, name: &str, resource: &str) -> Span {
554 let mut attrs = FastHashMap::default();
555 attrs.insert(MetaString::from("_dd.measured"), AttributeValue::Float(1.0));
556 Span::new(service, name, resource, "web", 1, 0, 1000000000, 100000000, 0).with_attributes(attrs)
557 }
558
559 fn make_top_level_span(
561 aligned_now: u64, span_id: u64, duration: u64, bucket_offset: u64, service: &str, resource: &str, error: i32,
562 meta: Option<FastHashMap<MetaString, MetaString>>,
563 ) -> Span {
564 let mut metrics = FastHashMap::default();
565 metrics.insert(MetaString::from("_top_level"), 1.0);
566 test_span(
567 aligned_now,
568 span_id,
569 0,
570 duration,
571 bucket_offset,
572 service,
573 resource,
574 error,
575 meta,
576 Some(metrics),
577 )
578 }
579
580 #[test]
581 fn test_process_trace_creates_stats() {
582 let now = now_nanos();
583
584 let concentrator = SpanConcentrator::new(true, true, &[], now);
585 let mut transform = ApmStats {
586 concentrator,
587 flush_interval: DEFAULT_FLUSH_INTERVAL,
588 agent_env: MetaString::from("none"),
589 agent_hostname: MetaString::default(),
590 workload_provider: None,
591 };
592
593 let span = make_test_span("test-service", "test-operation", "test-resource");
594 let trace = Trace::new(vec![span]);
595
596 transform.process_trace(&trace);
597
598 let stats = transform.concentrator.flush(now + BUCKET_DURATION_NS * 2, true);
600 assert!(!stats.is_empty(), "Expected stats to be produced");
601 }
602
603 #[test]
604 fn test_weight_applied_to_stats() {
605 let now = now_nanos();
606
607 let concentrator = SpanConcentrator::new(true, true, &[], now);
608 let mut transform = ApmStats {
609 concentrator,
610 flush_interval: DEFAULT_FLUSH_INTERVAL,
611 agent_env: MetaString::from("none"),
612 agent_hostname: MetaString::default(),
613 workload_provider: None,
614 };
615
616 let mut attrs = FastHashMap::default();
618 attrs.insert(MetaString::from("_dd.measured"), AttributeValue::Float(1.0));
619 attrs.insert(MetaString::from("_sample_rate"), AttributeValue::Float(0.5));
620
621 let span = Span::new(
622 "test-service",
623 "test-op",
624 "test-resource",
625 "web",
626 1,
627 0,
628 now,
629 100000000,
630 0,
631 )
632 .with_attributes(attrs);
633
634 let trace = Trace::new(vec![span]);
635 transform.process_trace(&trace);
636
637 let stats = transform.concentrator.flush(now + BUCKET_DURATION_NS * 2, true);
638 assert!(!stats.is_empty());
639
640 let bucket = &stats[0].stats()[0];
642 let grouped = &bucket.stats()[0];
643 assert!(grouped.hits() >= 1, "Expected weighted hits");
645 }
646
647 #[test]
648 fn test_force_flush() {
649 let now = now_nanos();
650 let aligned_now = align_ts(now, BUCKET_DURATION_NS);
651
652 let mut concentrator = SpanConcentrator::new(true, true, &[], now);
653
654 let span = make_top_level_span(aligned_now, 1, 50, 5, "A1", "resource1", 0, None);
656 let trace = Trace::new(vec![span]);
657
658 let payload_key = PayloadAggregationKey {
659 env: MetaString::from("test"),
660 hostname: MetaString::from("host"),
661 ..Default::default()
662 };
663 let infra_tags = InfraTags::default();
664
665 for span in trace.spans() {
666 if let Some(stat_span) = concentrator.new_stat_span_from_span(span) {
667 concentrator.add_span(&stat_span, 1.0, &payload_key, &infra_tags, "");
668 }
669 }
670
671 let ts: u64 = 0;
673
674 let stats = concentrator.flush(ts, false);
676 assert!(stats.is_empty(), "Non-force flush should return empty");
677
678 let stats = concentrator.flush(ts, true);
680 assert!(!stats.is_empty(), "Force flush should return stats");
681 assert_eq!(stats[0].stats().len(), 1, "Should have 1 bucket");
682 }
683
684 #[test]
685 fn test_ignores_partial_spans() {
686 let now = now_nanos();
687 let aligned_now = align_ts(now, BUCKET_DURATION_NS);
688
689 let mut concentrator = SpanConcentrator::new(true, true, &[], now);
690
691 let mut metrics = FastHashMap::default();
693 metrics.insert(MetaString::from("_top_level"), 1.0);
694 metrics.insert(MetaString::from(METRIC_PARTIAL_VERSION), 830604.0);
695
696 let span = test_span(aligned_now, 1, 0, 50, 5, "A1", "resource1", 0, None, Some(metrics));
697 let trace = Trace::new(vec![span]);
698
699 let payload_key = PayloadAggregationKey {
700 env: MetaString::from("test"),
701 hostname: MetaString::from("tracer-hostname"),
702 ..Default::default()
703 };
704 let infra_tags = InfraTags::default();
705
706 for span in trace.spans() {
707 if let Some(stat_span) = concentrator.new_stat_span_from_span(span) {
708 concentrator.add_span(&stat_span, 1.0, &payload_key, &infra_tags, "");
709 }
710 }
711
712 let stats = concentrator.flush(now + BUCKET_DURATION_NS * 3, true);
714 assert!(stats.is_empty(), "Partial spans should be ignored");
715 }
716
717 #[test]
718 fn test_concentrator_stats_totals() {
719 let now = now_nanos();
720 let aligned_now = align_ts(now, BUCKET_DURATION_NS);
721
722 let oldest_ts = aligned_now - 2 * BUCKET_DURATION_NS;
724 let mut concentrator = SpanConcentrator::new(true, true, &[], oldest_ts);
725
726 let spans = vec![
728 make_top_level_span(aligned_now, 1, 50, 5, "A1", "resource1", 0, None),
729 make_top_level_span(aligned_now, 2, 40, 4, "A1", "resource1", 0, None),
730 make_top_level_span(aligned_now, 3, 30, 3, "A1", "resource1", 0, None),
731 make_top_level_span(aligned_now, 4, 20, 2, "A1", "resource1", 0, None),
732 make_top_level_span(aligned_now, 5, 10, 1, "A1", "resource1", 0, None),
733 make_top_level_span(aligned_now, 6, 1, 0, "A1", "resource1", 0, None),
734 ];
735
736 let payload_key = PayloadAggregationKey {
737 env: MetaString::from("none"),
738 ..Default::default()
739 };
740 let infra_tags = InfraTags::default();
741
742 for span in &spans {
743 if let Some(stat_span) = concentrator.new_stat_span_from_span(span) {
744 concentrator.add_span(&stat_span, 1.0, &payload_key, &infra_tags, "");
745 }
746 }
747
748 let all_stats = concentrator.flush(now + BUCKET_DURATION_NS * 10, true);
750
751 let mut total_duration: u64 = 0;
752 let mut total_hits: u64 = 0;
753 let mut total_errors: u64 = 0;
754 let mut total_top_level_hits: u64 = 0;
755
756 for payload in &all_stats {
757 for bucket in payload.stats() {
758 for grouped in bucket.stats() {
759 total_duration += grouped.duration();
760 total_hits += grouped.hits();
761 total_errors += grouped.errors();
762 total_top_level_hits += grouped.top_level_hits();
763 }
764 }
765 }
766
767 assert_eq!(total_duration, 50 + 40 + 30 + 20 + 10 + 1, "Wrong total duration");
768 assert_eq!(total_hits, 6, "Wrong total hits");
769 assert_eq!(total_top_level_hits, 6, "Wrong total top level hits");
770 assert_eq!(total_errors, 0, "Wrong total errors");
771 }
772
773 #[test]
774 fn test_root_tag() {
775 let now = now_nanos();
776 let aligned_now = align_ts(now, BUCKET_DURATION_NS);
777
778 let mut concentrator = SpanConcentrator::new(true, true, &[], now);
779
780 let mut root_metrics = FastHashMap::default();
782 root_metrics.insert(MetaString::from("_top_level"), 1.0);
783 let root_span = test_span(
784 aligned_now,
785 1,
786 0,
787 40,
788 10,
789 "A1",
790 "resource1",
791 0,
792 None,
793 Some(root_metrics),
794 );
795
796 let mut top_level_metrics = FastHashMap::default();
798 top_level_metrics.insert(MetaString::from("_top_level"), 1.0);
799 let top_level_span = test_span(
800 aligned_now,
801 4,
802 1000,
803 10,
804 10,
805 "A1",
806 "resource1",
807 0,
808 None,
809 Some(top_level_metrics),
810 );
811
812 let mut client_meta = FastHashMap::default();
814 client_meta.insert(MetaString::from("span.kind"), MetaString::from("client"));
815 let client_span = test_span(aligned_now, 3, 2, 20, 10, "A1", "resource1", 0, Some(client_meta), None);
816
817 let spans = vec![root_span, top_level_span, client_span];
818
819 let payload_key = PayloadAggregationKey {
820 env: MetaString::from("none"),
821 ..Default::default()
822 };
823 let infra_tags = InfraTags::default();
824
825 for span in &spans {
826 if let Some(stat_span) = concentrator.new_stat_span_from_span(span) {
827 concentrator.add_span(&stat_span, 1.0, &payload_key, &infra_tags, "");
828 }
829 }
830
831 let stats = concentrator.flush(now + BUCKET_DURATION_NS * 20, true);
832 assert!(!stats.is_empty(), "Should have stats");
833
834 let mut total_grouped = 0;
836 let mut root_count = 0;
837 let mut non_root_count = 0;
838
839 for payload in &stats {
840 for bucket in payload.stats() {
841 for grouped in bucket.stats() {
842 total_grouped += 1;
843 match grouped.is_trace_root() {
844 Some(true) => root_count += 1,
845 Some(false) => non_root_count += 1,
846 None => {}
847 }
848 }
849 }
850 }
851
852 assert_eq!(total_grouped, 3, "Expected 3 grouped stats");
857 assert_eq!(root_count, 1, "Expected 1 root span");
858 assert_eq!(non_root_count, 2, "Expected 2 non-root spans");
859 }
860
861 #[test]
862 fn test_compute_stats_through_span_kind_check() {
863 let now = now_nanos();
864
865 {
867 let mut concentrator = SpanConcentrator::new(false, true, &[], now);
868
869 let mut attrs = FastHashMap::default();
870 attrs.insert(MetaString::from("_top_level"), AttributeValue::Float(1.0));
871 let span = Span::new("myservice", "query", "GET /users", "web", 1, 0, now, 500, 0).with_attributes(attrs);
872
873 let payload_key = PayloadAggregationKey {
874 env: MetaString::from("test"),
875 ..Default::default()
876 };
877 let infra_tags = InfraTags::default();
878
879 if let Some(stat_span) = concentrator.new_stat_span_from_span(&span) {
880 concentrator.add_span(&stat_span, 1.0, &payload_key, &infra_tags, "");
881 }
882
883 let mut client_attrs = FastHashMap::default();
886 client_attrs.insert(
887 MetaString::from("span.kind"),
888 AttributeValue::String(MetaString::from("client")),
889 );
890 let client_span = Span::new("myservice", "postgres.query", "SELECT ...", "db", 2, 1, now, 75, 0)
891 .with_attributes(client_attrs);
892
893 if let Some(stat_span) = concentrator.new_stat_span_from_span(&client_span) {
894 concentrator.add_span(&stat_span, 1.0, &payload_key, &infra_tags, "");
895 }
896
897 let stats = concentrator.flush(now + BUCKET_DURATION_NS * 3, true);
898
899 let mut count = 0;
900 for payload in &stats {
901 for bucket in payload.stats() {
902 count += bucket.stats().len();
903 }
904 }
905
906 assert_eq!(count, 1, "Expected 1 stat when span kind check disabled");
908 }
909
910 {
912 let mut concentrator = SpanConcentrator::new(true, true, &[], now);
913
914 let mut attrs = FastHashMap::default();
915 attrs.insert(MetaString::from("_top_level"), AttributeValue::Float(1.0));
916 let span = Span::new("myservice", "query", "GET /users", "web", 1, 0, now, 500, 0).with_attributes(attrs);
917
918 let payload_key = PayloadAggregationKey {
919 env: MetaString::from("test"),
920 ..Default::default()
921 };
922 let infra_tags = InfraTags::default();
923
924 if let Some(stat_span) = concentrator.new_stat_span_from_span(&span) {
925 concentrator.add_span(&stat_span, 1.0, &payload_key, &infra_tags, "");
926 }
927
928 let mut client_attrs = FastHashMap::default();
931 client_attrs.insert(
932 MetaString::from("span.kind"),
933 AttributeValue::String(MetaString::from("client")),
934 );
935 let client_span = Span::new("myservice", "postgres.query", "SELECT ...", "db", 2, 1, now, 75, 0)
936 .with_attributes(client_attrs);
937
938 if let Some(stat_span) = concentrator.new_stat_span_from_span(&client_span) {
939 concentrator.add_span(&stat_span, 1.0, &payload_key, &infra_tags, "");
940 }
941
942 let stats = concentrator.flush(now + BUCKET_DURATION_NS * 3, true);
943
944 let mut count = 0;
945 for payload in &stats {
946 for bucket in payload.stats() {
947 count += bucket.stats().len();
948 }
949 }
950
951 assert_eq!(count, 2, "Expected 2 stats when span kind check enabled");
953 }
954 }
955
956 #[test]
957 fn test_peer_tags() {
958 let now = now_nanos();
959
960 {
962 let mut concentrator = SpanConcentrator::new(true, false, &[], now);
963
964 let mut attrs = FastHashMap::default();
965 attrs.insert(
966 MetaString::from("span.kind"),
967 AttributeValue::String(MetaString::from("client")),
968 );
969 attrs.insert(
970 MetaString::from("db.instance"),
971 AttributeValue::String(MetaString::from("i-1234")),
972 );
973 attrs.insert(
974 MetaString::from("db.system"),
975 AttributeValue::String(MetaString::from("postgres")),
976 );
977 attrs.insert(MetaString::from("_dd.measured"), AttributeValue::Float(1.0));
978 let client_span =
979 Span::new("myservice", "postgres.query", "SELECT ...", "db", 2, 1, now, 75, 0).with_attributes(attrs);
980
981 let payload_key = PayloadAggregationKey {
982 env: MetaString::from("test"),
983 ..Default::default()
984 };
985 let infra_tags = InfraTags::default();
986
987 if let Some(stat_span) = concentrator.new_stat_span_from_span(&client_span) {
988 concentrator.add_span(&stat_span, 1.0, &payload_key, &infra_tags, "");
989 }
990
991 let stats = concentrator.flush(now + BUCKET_DURATION_NS * 3, true);
992
993 for payload in &stats {
995 for bucket in payload.stats() {
996 for grouped in bucket.stats() {
997 assert!(
998 grouped.peer_tags().is_empty(),
999 "Peer tags should be empty when peer_tags_aggregation is false"
1000 );
1001 }
1002 }
1003 }
1004 }
1005
1006 {
1008 let mut concentrator = SpanConcentrator::new(true, true, &[], now);
1010
1011 let mut attrs = FastHashMap::default();
1012 attrs.insert(
1013 MetaString::from("span.kind"),
1014 AttributeValue::String(MetaString::from("client")),
1015 );
1016 attrs.insert(
1017 MetaString::from("db.instance"),
1018 AttributeValue::String(MetaString::from("i-1234")),
1019 );
1020 attrs.insert(
1021 MetaString::from("db.system"),
1022 AttributeValue::String(MetaString::from("postgres")),
1023 );
1024 attrs.insert(MetaString::from("_dd.measured"), AttributeValue::Float(1.0));
1025 let client_span =
1026 Span::new("myservice", "postgres.query", "SELECT ...", "db", 2, 1, now, 75, 0).with_attributes(attrs);
1027
1028 let payload_key = PayloadAggregationKey {
1029 env: MetaString::from("test"),
1030 ..Default::default()
1031 };
1032 let infra_tags = InfraTags::default();
1033
1034 if let Some(stat_span) = concentrator.new_stat_span_from_span(&client_span) {
1035 concentrator.add_span(&stat_span, 1.0, &payload_key, &infra_tags, "");
1036 }
1037
1038 let stats = concentrator.flush(now + BUCKET_DURATION_NS * 3, true);
1039
1040 let mut found_client_with_peer_tags = false;
1042 for payload in &stats {
1043 for bucket in payload.stats() {
1044 for grouped in bucket.stats() {
1045 if grouped.resource() == "SELECT ..." {
1046 assert!(!grouped.peer_tags().is_empty(), "Client span should have peer tags");
1047 let peer_tags: Vec<&str> = grouped.peer_tags().iter().map(|s| s.as_ref()).collect();
1049 assert!(
1050 peer_tags.iter().any(|t| t.starts_with("db.instance:")),
1051 "Should have db.instance peer tag"
1052 );
1053 assert!(
1054 peer_tags.iter().any(|t| t.starts_with("db.system:")),
1055 "Should have db.system peer tag"
1056 );
1057 found_client_with_peer_tags = true;
1058 }
1059 }
1060 }
1061 }
1062 assert!(
1063 found_client_with_peer_tags,
1064 "Should have found client span with peer tags"
1065 );
1066 }
1067 }
1068
1069 #[test]
1070 fn test_concentrator_oldest_ts() {
1071 let now = now_nanos();
1072 let aligned_now = align_ts(now, BUCKET_DURATION_NS);
1073
1074 {
1076 let mut concentrator = SpanConcentrator::new(true, true, &[], now);
1078
1079 let spans = vec![
1081 make_top_level_span(aligned_now, 1, 50, 5, "A1", "resource1", 0, None),
1082 make_top_level_span(aligned_now, 2, 40, 4, "A1", "resource1", 0, None),
1083 make_top_level_span(aligned_now, 3, 30, 3, "A1", "resource1", 0, None),
1084 make_top_level_span(aligned_now, 4, 20, 2, "A1", "resource1", 0, None),
1085 make_top_level_span(aligned_now, 5, 10, 1, "A1", "resource1", 0, None),
1086 make_top_level_span(aligned_now, 6, 1, 0, "A1", "resource1", 0, None),
1087 ];
1088
1089 let payload_key = PayloadAggregationKey {
1090 env: MetaString::from("none"),
1091 ..Default::default()
1092 };
1093 let infra_tags = InfraTags::default();
1094
1095 for span in &spans {
1096 if let Some(stat_span) = concentrator.new_stat_span_from_span(span) {
1097 concentrator.add_span(&stat_span, 1.0, &payload_key, &infra_tags, "");
1098 }
1099 }
1100
1101 let mut flush_time = now;
1103 let buffer_len = 2; for _ in 0..buffer_len {
1106 let stats = concentrator.flush(flush_time, false);
1107 assert!(stats.is_empty(), "Should not flush before buffer fills");
1108 flush_time += BUCKET_DURATION_NS;
1109 }
1110
1111 let stats = concentrator.flush(flush_time, false);
1113 assert!(!stats.is_empty(), "Should flush after buffer fills");
1114
1115 let mut total_hits: u64 = 0;
1117 let mut total_duration: u64 = 0;
1118 for payload in &stats {
1119 for bucket in payload.stats() {
1120 for grouped in bucket.stats() {
1121 total_hits += grouped.hits();
1122 total_duration += grouped.duration();
1123 }
1124 }
1125 }
1126
1127 assert_eq!(total_hits, 6, "All 6 spans should be counted");
1128 assert_eq!(
1129 total_duration,
1130 50 + 40 + 30 + 20 + 10 + 1,
1131 "Total duration should match"
1132 );
1133 }
1134 }
1135
1136 #[test]
1137 fn test_compute_stats_for_span_kind() {
1138 use super::span_concentrator::compute_stats_for_span_kind;
1139
1140 assert!(compute_stats_for_span_kind("server"));
1142 assert!(compute_stats_for_span_kind("consumer"));
1143 assert!(compute_stats_for_span_kind("client"));
1144 assert!(compute_stats_for_span_kind("producer"));
1145
1146 assert!(compute_stats_for_span_kind("SERVER"));
1148 assert!(compute_stats_for_span_kind("CONSUMER"));
1149 assert!(compute_stats_for_span_kind("CLIENT"));
1150 assert!(compute_stats_for_span_kind("PRODUCER"));
1151
1152 assert!(compute_stats_for_span_kind("SErVER"));
1154 assert!(compute_stats_for_span_kind("COnSUMER"));
1155 assert!(compute_stats_for_span_kind("CLiENT"));
1156 assert!(compute_stats_for_span_kind("PRoDUCER"));
1157
1158 assert!(!compute_stats_for_span_kind("internal"));
1160 assert!(!compute_stats_for_span_kind("INTERNAL"));
1161 assert!(!compute_stats_for_span_kind("INtERNAL"));
1162 assert!(!compute_stats_for_span_kind(""));
1163 }
1164
1165 #[test]
1166 fn test_extract_process_tags() {
1167 {
1169 let span = Span::default();
1170 let trace = Trace::new(vec![span]);
1171 let process_tags = extract_process_tags(&trace);
1172 assert!(process_tags.is_empty(), "Should be empty when no _dd.tags.process");
1173 }
1174
1175 {
1177 let mut attrs = FastHashMap::default();
1178 attrs.insert(
1179 MetaString::from(TAG_PROCESS_TAGS),
1180 AttributeValue::String(MetaString::from("a:1,b:2,c:3")),
1181 );
1182 let span = Span::default().with_attributes(attrs);
1183 let trace = Trace::new(vec![span]);
1184 let process_tags = extract_process_tags(&trace);
1185 assert_eq!(process_tags, "a:1,b:2,c:3");
1186 }
1187
1188 {
1190 let mut attrs = FastHashMap::default();
1191 attrs.insert(
1192 MetaString::from(TAG_PROCESS_TAGS),
1193 AttributeValue::String(MetaString::from("")),
1194 );
1195 let span = Span::default().with_attributes(attrs);
1196 let trace = Trace::new(vec![span]);
1197 let process_tags = extract_process_tags(&trace);
1198 assert!(
1199 process_tags.is_empty(),
1200 "Should be empty when _dd.tags.process is empty string"
1201 );
1202 }
1203
1204 {
1206 let trace = Trace::new(vec![]);
1207 let process_tags = extract_process_tags(&trace);
1208 assert!(process_tags.is_empty(), "Should be empty when trace has no spans");
1209 }
1210 }
1211
1212 #[test]
1213 fn test_process_tags_hash_computation() {
1214 use super::aggregation::process_tags_hash;
1215
1216 assert_eq!(process_tags_hash(""), 0);
1218
1219 let hash1 = process_tags_hash("a:1,b:2,c:3");
1221 let hash2 = process_tags_hash("a:1,b:2,c:3");
1222 assert_eq!(hash1, hash2);
1223
1224 let hash3 = process_tags_hash("a:1,b:2");
1226 assert_ne!(hash1, hash3);
1227 }
1228
1229 fn make_grouped_stats(service: &str, resource: &str) -> ClientGroupedStats {
1231 ClientGroupedStats::new(service, "operation", resource)
1232 .with_hits(1)
1233 .with_duration(100)
1234 }
1235
1236 fn make_bucket_with_stats(n: usize) -> ClientStatsBucket {
1238 let stats: Vec<ClientGroupedStats> = (0..n)
1239 .map(|i| make_grouped_stats("service", &format!("resource-{}", i)))
1240 .collect();
1241 ClientStatsBucket::new(1000, 10_000_000_000, stats)
1242 }
1243
1244 fn make_payload_with_buckets(hostname: &str, buckets: Vec<ClientStatsBucket>) -> ClientStatsPayload {
1246 ClientStatsPayload::new(hostname, "test-env", "1.0.0")
1247 .with_stats(buckets)
1248 .with_container_id("container-123")
1249 .with_lang("rust")
1250 }
1251
1252 fn count_grouped_stats(trace_stats: &TraceStats) -> usize {
1254 trace_stats
1255 .stats()
1256 .iter()
1257 .flat_map(|p| p.stats())
1258 .map(|b| b.stats().len())
1259 .sum()
1260 }
1261
1262 #[test]
1263 fn test_split_into_trace_stats_empty_input() {
1264 let result = split_into_trace_stats(vec![], 100);
1265 assert!(result.is_empty());
1266 }
1267
1268 #[test]
1269 fn test_split_into_trace_stats_no_split_needed() {
1270 let bucket = make_bucket_with_stats(50);
1272 let payload = make_payload_with_buckets("host1", vec![bucket]);
1273
1274 let result = split_into_trace_stats(vec![payload], 100);
1275
1276 assert_eq!(result.len(), 1);
1277 assert_eq!(count_grouped_stats(&result[0]), 50);
1278 }
1279
1280 #[test]
1281 fn test_split_into_trace_stats_exact_threshold() {
1282 let bucket = make_bucket_with_stats(100);
1284 let payload = make_payload_with_buckets("host1", vec![bucket]);
1285
1286 let result = split_into_trace_stats(vec![payload], 100);
1287
1288 assert_eq!(result.len(), 1);
1289 assert_eq!(count_grouped_stats(&result[0]), 100);
1290 }
1291
1292 #[test]
1293 fn test_split_into_trace_stats_splits_single_bucket() {
1294 let bucket = make_bucket_with_stats(250);
1296 let payload = make_payload_with_buckets("host1", vec![bucket]);
1297
1298 let result = split_into_trace_stats(vec![payload], 100);
1299
1300 assert_eq!(result.len(), 3);
1301 assert_eq!(count_grouped_stats(&result[0]), 100);
1302 assert_eq!(count_grouped_stats(&result[1]), 100);
1303 assert_eq!(count_grouped_stats(&result[2]), 50);
1304
1305 let total: usize = result.iter().map(count_grouped_stats).sum();
1307 assert_eq!(total, 250);
1308 }
1309
1310 #[test]
1311 fn test_split_into_trace_stats_splits_across_payloads() {
1312 let payload1 = make_payload_with_buckets("host1", vec![make_bucket_with_stats(60)]);
1314 let payload2 = make_payload_with_buckets("host2", vec![make_bucket_with_stats(60)]);
1315
1316 let result = split_into_trace_stats(vec![payload1, payload2], 100);
1317
1318 assert_eq!(result.len(), 2);
1319 assert_eq!(count_grouped_stats(&result[0]), 100);
1321 assert_eq!(count_grouped_stats(&result[1]), 20);
1323 }
1324
1325 #[test]
1326 fn test_split_into_trace_stats_splits_single_payload_multiple_buckets() {
1327 let bucket1 = make_bucket_with_stats(70);
1329 let bucket2 = make_bucket_with_stats(80);
1330 let payload = make_payload_with_buckets("host1", vec![bucket1, bucket2]);
1331
1332 let result = split_into_trace_stats(vec![payload], 100);
1333
1334 assert_eq!(result.len(), 2);
1335 let total: usize = result.iter().map(count_grouped_stats).sum();
1336 assert_eq!(total, 150);
1337 }
1338
1339 #[test]
1340 fn test_split_into_trace_stats_preserves_metadata() {
1341 let bucket = make_bucket_with_stats(250);
1342 let payload = ClientStatsPayload::new("test-host", "prod", "2.0.0")
1343 .with_stats(vec![bucket])
1344 .with_container_id("container-abc")
1345 .with_lang("go")
1346 .with_git_commit_sha("abc123")
1347 .with_image_tag("v1.2.3")
1348 .with_process_tags_hash(12345)
1349 .with_process_tags("tag1,tag2");
1350
1351 let result = split_into_trace_stats(vec![payload], 100);
1352
1353 for trace_stats in &result {
1355 for p in trace_stats.stats() {
1356 assert_eq!(p.hostname(), "test-host");
1357 assert_eq!(p.env(), "prod");
1358 assert_eq!(p.version(), "2.0.0");
1359 assert_eq!(p.container_id(), "container-abc");
1360 assert_eq!(p.lang(), "go");
1361 assert_eq!(p.git_commit_sha(), "abc123");
1362 assert_eq!(p.image_tag(), "v1.2.3");
1363 assert_eq!(p.process_tags_hash(), 12345);
1364 assert_eq!(p.process_tags(), "tag1,tag2");
1365 }
1366 }
1367 }
1368
1369 #[test]
1370 fn test_split_into_trace_stats_preserves_bucket_metadata() {
1371 let stats: Vec<ClientGroupedStats> = (0..150)
1373 .map(|i| make_grouped_stats("svc", &format!("res-{}", i)))
1374 .collect();
1375 let bucket = ClientStatsBucket::new(999_000_000, 10_000_000_000, stats).with_agent_time_shift(42);
1376 let payload = make_payload_with_buckets("host1", vec![bucket]);
1377
1378 let result = split_into_trace_stats(vec![payload], 100);
1379
1380 for trace_stats in &result {
1382 for p in trace_stats.stats() {
1383 for b in p.stats() {
1384 assert_eq!(b.start(), 999_000_000);
1385 assert_eq!(b.duration(), 10_000_000_000);
1386 assert_eq!(b.agent_time_shift(), 42);
1387 }
1388 }
1389 }
1390 }
1391
1392 #[test]
1393 fn test_split_into_trace_stats_handles_empty_bucket() {
1394 let empty_bucket = ClientStatsBucket::new(1000, 10_000_000_000, vec![]);
1395 let payload = make_payload_with_buckets("host1", vec![empty_bucket]);
1396
1397 let result = split_into_trace_stats(vec![payload], 100);
1398
1399 assert_eq!(result.len(), 1);
1400 assert_eq!(count_grouped_stats(&result[0]), 0);
1401 }
1402
1403 #[test]
1404 fn test_split_into_trace_stats_large_split() {
1405 let bucket = make_bucket_with_stats(10_000);
1407 let payload = make_payload_with_buckets("host1", vec![bucket]);
1408
1409 let result = split_into_trace_stats(vec![payload], 4000);
1410
1411 assert_eq!(result.len(), 3);
1412 assert_eq!(count_grouped_stats(&result[0]), 4000);
1413 assert_eq!(count_grouped_stats(&result[1]), 4000);
1414 assert_eq!(count_grouped_stats(&result[2]), 2000);
1415 }
1416
1417 fn arb_grouped_stats() -> impl Strategy<Value = ClientGroupedStats> {
1421 (0..100u64, 0..1000u64).prop_map(|(hits, duration)| {
1422 ClientGroupedStats::new("service", "operation", "resource")
1423 .with_hits(hits)
1424 .with_duration(duration)
1425 })
1426 }
1427
1428 fn arb_bucket(max_stats_per_bucket: usize) -> impl Strategy<Value = ClientStatsBucket> {
1430 proptest::collection::vec(arb_grouped_stats(), 0..=max_stats_per_bucket)
1431 .prop_map(|stats| ClientStatsBucket::new(1000, 10_000_000_000, stats))
1432 }
1433
1434 fn arb_payload(max_buckets: usize, max_stats_per_bucket: usize) -> impl Strategy<Value = ClientStatsPayload> {
1436 proptest::collection::vec(arb_bucket(max_stats_per_bucket), 1..=max_buckets)
1437 .prop_map(|buckets| ClientStatsPayload::new("host", "env", "1.0.0").with_stats(buckets))
1438 }
1439
1440 fn arb_split_inputs() -> impl Strategy<Value = (Vec<ClientStatsPayload>, usize)> {
1448 let payloads_strategy = proptest::collection::vec(arb_payload(5, 500), 1..=10);
1449 let max_entries_strategy = 1..=1000usize;
1450
1451 (payloads_strategy, max_entries_strategy)
1452 }
1453
1454 #[test]
1459 fn test_version_span_beats_resource_for_otlp() {
1460 let now = now_nanos();
1461 let concentrator = SpanConcentrator::new(true, true, &[], now);
1462 let transform = ApmStats {
1463 concentrator,
1464 flush_interval: DEFAULT_FLUSH_INTERVAL,
1465 agent_env: MetaString::default(),
1466 agent_hostname: MetaString::default(),
1467 workload_provider: None,
1468 };
1469
1470 let mut attrs = FastHashMap::default();
1473 attrs.insert(
1474 MetaString::from("version"),
1475 AttributeValue::String(MetaString::from("span-v2")),
1476 );
1477 let root_span = Span::new("svc", "op", "res", "web", 1, 0, now, 1_000_000, 0).with_attributes(attrs);
1478
1479 let mut trace = Trace::new(vec![root_span]);
1480 trace.payload.app_version = MetaString::from("resource-v1");
1482
1483 let key = transform.build_payload_key(&trace, "");
1484
1485 assert_eq!(
1488 key.version.as_ref(),
1489 "span-v2",
1490 "span-level version must take precedence over resource-level payload.app_version for OTLP traces"
1491 );
1492 }
1493
1494 proptest! {
1495 #[test]
1496 #[cfg_attr(miri, ignore)]
1497 fn property_test_split_respects_max_entries((payloads, max_entries_per_event) in arb_split_inputs()) {
1498 let input_total: usize = payloads.iter().flat_map(|p| p.stats()).map(|b| b.stats().len()).sum();
1499
1500 let result = split_into_trace_stats(payloads, max_entries_per_event);
1501
1502 for trace_stats in &result {
1504 let count = count_grouped_stats(trace_stats);
1505 prop_assert!(
1506 count <= max_entries_per_event,
1507 "TraceStats has {} grouped stats, exceeds max of {}",
1508 count,
1509 max_entries_per_event
1510 );
1511 }
1512
1513 let output_total: usize = result.iter().map(count_grouped_stats).sum();
1515 prop_assert_eq!(input_total, output_total, "Total stats count should be preserved");
1516 }
1517 }
1518}