saluki_components/transforms/apm_stats/
mod.rs

1//! APM Stats transform.
2//!
3//! Aggregates traces into time-bucketed statistics, producing `TraceStats` events.
4
5use 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
51/// Default flush interval for the APM stats transform.
52const DEFAULT_FLUSH_INTERVAL: Duration = Duration::from_secs(10);
53
54/// Tag key for process tags in span meta.
55const TAG_PROCESS_TAGS: &str = "_dd.tags.process";
56
57/// Maximum number of `ClientGroupedStats` entries per `TraceStats` event.
58const MAX_STATS_GROUPS_PER_EVENT: usize = 4000;
59
60/// APM Stats transform configuration.
61///
62/// Aggregates incoming `Trace` events into time-bucketed statistics, emitting
63/// `TraceStats` events.
64pub 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    /// Creates a new `ApmStatsTransformConfiguration` from the resolved traces configuration.
75    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    /// Sets the default hostname using the environment provider.
87    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    /// Sets the workload provider.
97    ///
98    /// Defaults to unset.
99    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        // TODO: Think about everything we need to account for here.
142    }
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        // env fallback order mirrors get_trace_env (Go agent pkg/trace/traceutil/trace.go:GetEnv):
213        // root span attrs → all span attrs → trace.payload.env (ADP extension) → agent_env.
214        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        // Version resolution mirrors Go agent pkg/trace/version/version.go:GetAppVersionFromTrace:
256        // root span meta["version"] (set by otel_span_to_dd_span with span-over-resource precedence)
257        // takes priority, falling back to the resource-level payload.app_version.
258        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
319/// Splits stats payloads into multiple `TraceStats` events, each containing at most `max_entries_per_event` grouped
320/// entries.
321///
322/// This function won't attempt to collapse/pack payloads to optimize the number of events generated: there will always
323/// be at least as many output client payloads as there are input client payloads.
324fn 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    // If the total number of grouped stats entries is below the specified threshold, then we don't do any splitting or
330    // collapsing or anything.
331    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        // If the current payload can fit entirely in the current event, then collect it as-is.
345        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        // Consume all the stats buckets from the current client payload, and set ourselves up to go through each
353        // bucket, splitting them as necessary.
354        //
355        // We basically iterate until we find a bucket where adding its entries would cause us to exceed our threshold,
356        // and then we split that bucket, creating a new client payload in the process, and keep repeating that until
357        // we've exhausted all buckets for our starting client payload.
358        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 this bucket can fit in the current event, then collect it as-is.
362            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            // We have to split this bucket. We take the grouped entries it has, and we subdivide those into a new
369            // client bucket, or buckets in order to ensure we don't exceed our threshold.
370            let mut bucket_entries = client_stats_bucket.take_stats();
371            while current_event_len + bucket_entries.len() > max_entries_per_event {
372                // We calculate the split point this way because the returned vector from `split_off` is referenced
373                // against the end of the vector, so if we have 100 items, and we only want 20, we need to split at
374                // index 80 (100 - 20 == 80).
375                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                // Create a new "split" bucket based on the current bucket (`client_stats_bucket`) but containing only
380                // the split-off entries. This will feed into the current client payload (`client_payload`) which we'll
381                // finalize and add to the current event before starting a new event.
382                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 we have any leftover entries from the bucket after splitting it up, put them back in the current bucket
395            // and add that bucket to the current client payload.
396            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 we have buckets for this client payload (whether we split or not), add them back to the current client payload.
403        if !current_client_stats_buckets.is_empty() {
404            current_client_payloads.push(client_payload.with_stats(current_client_stats_buckets));
405        }
406    }
407
408    // Stick the remaining entries into a final `TraceStats` event, if any.
409    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                            // We've reached the end of our input stream, so mark ourselves for a final flush and reset the
464                            // interval so it ticks immediately on the next loop iteration.
465                            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
479/// Returns the current time as nanoseconds since Unix epoch.
480fn now_nanos() -> u64 {
481    SystemTime::now()
482        .duration_since(UNIX_EPOCH)
483        .unwrap_or_default()
484        .as_nanos() as u64
485}
486
487/// Extracts process tags from trace, checking both span and trace attributes.
488fn 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    /// Helper to align timestamp to bucket boundary
525    fn align_ts(ts: u64, bsize: u64) -> u64 {
526        ts - ts % bsize
527    }
528
529    /// Creates a test span with the given parameters.
530    #[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    /// Creates a simple measured span for basic tests
553    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    /// Creates a top-level span (`parent_id` = 0, has `_top_level` metric)
560    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        // Flush and verify we got stats
599        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        // Create a span with 0.5 sample rate (weight = 2.0)
617        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        // The hits should be weighted (approximately 2 due to 0.5 sample rate)
641        let bucket = &stats[0].stats()[0];
642        let grouped = &bucket.stats()[0];
643        // With stochastic rounding, hits could be 1 or 2, but with weight 2.0 it should round to 2
644        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        // Add a span
655        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        // ts=0 so that flush always considers buckets not old enough
672        let ts: u64 = 0;
673
674        // Without force flush, should skip the bucket
675        let stats = concentrator.flush(ts, false);
676        assert!(stats.is_empty(), "Non-force flush should return empty");
677
678        // With force flush, should flush buckets regardless of age
679        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        // Create a partial span (has _dd.partial_version metric)
692        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        // Partial spans should be ignored
713        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        // Set oldestTs to allow old buckets
723        let oldest_ts = aligned_now - 2 * BUCKET_DURATION_NS;
724        let mut concentrator = SpanConcentrator::new(true, true, &[], oldest_ts);
725
726        // Build spans spread over time windows
727        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        // Flush all and collect totals
749        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        // Root span (parent_id = 0, top_level)
781        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        // Non-root but top level span (has _top_level but parent_id != 0)
797        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        // Client span (non-root, non-top level, but has span.kind = client)
813        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        // Count grouped stats - should be split by IsTraceRoot
835        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        // We expect 3 grouped stats:
853        // 1. Root span (is_trace_root = true)
854        // 2. Non-root top-level span (is_trace_root = false)
855        // 3. Client span (is_trace_root = false, span.kind = client)
856        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        // Test with compute_stats_by_span_kind DISABLED
866        {
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            // Client span with span.kind=client but no _top_level or _dd.measured
884            // Should NOT produce stats when compute_stats_by_span_kind is disabled
885            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            // When disabled, only top_level span gets stats (client span has no top_level/measured)
907            assert_eq!(count, 1, "Expected 1 stat when span kind check disabled");
908        }
909
910        // Test with compute_stats_by_span_kind ENABLED
911        {
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            // Client span with span.kind=client
929            // SHOULD produce stats when compute_stats_by_span_kind is enabled
930            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            // When enabled, both spans get stats
952            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        // Test without peer tags aggregation enabled
961        {
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            // Without peer tags aggregation, peer_tags should be empty
994            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        // Test with peer tags aggregation enabled
1007        {
1008            // Note: BASE_PEER_TAGS already includes db.instance and db.system
1009            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            // With peer tags aggregation, client span should have peer_tags
1041            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                            // Check that peer tags contain db.instance and db.system
1048                            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        // Test "cold" scenario - all spans in the past should end up in current bucket
1075        {
1076            // Start concentrator at current time (cold start)
1077            let mut concentrator = SpanConcentrator::new(true, true, &[], now);
1078
1079            // Build spans spread over many time windows (all in the past)
1080            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            // Flush multiple times without force
1102            let mut flush_time = now;
1103            let buffer_len = 2; // DEFAULT_BUFFER_LEN
1104
1105            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            // After buffer_len flushes, should get aggregated stats
1112            let stats = concentrator.flush(flush_time, false);
1113            assert!(!stats.is_empty(), "Should flush after buffer fills");
1114
1115            // All spans should be aggregated into one bucket (oldest bucket aggregates old data)
1116            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        // Valid span kinds (case insensitive)
1141        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        // Uppercase
1147        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        // Mixed case
1153        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        // Invalid span kinds
1159        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        // Test with no process tags
1168        {
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        // Test with process tags in first span meta
1176        {
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        // Test with empty process tags
1189        {
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        // Test with empty trace
1205        {
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        // Empty string should return 0
1217        assert_eq!(process_tags_hash(""), 0);
1218
1219        // Same tags should produce same hash
1220        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        // Different tags should produce different hash
1225        let hash3 = process_tags_hash("a:1,b:2");
1226        assert_ne!(hash1, hash3);
1227    }
1228
1229    // Helper to create a ClientGroupedStats for testing
1230    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    // Helper to create a ClientStatsBucket with N stats
1237    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    // Helper to create a ClientStatsPayload with specified buckets
1245    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    // Helper to count total ClientGroupedStats across all payloads in a TraceStats
1253    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        // 50 stats with max 100 - no split needed
1271        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        // Exactly 100 stats with max 100 - no split needed
1283        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        // 250 stats in one bucket with max 100 - should split into 3 TraceStats
1295        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        // Total should be preserved
1306        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        // Two payloads with 60 stats each, max 100 - should combine into 2 TraceStats
1313        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        // First TraceStats: 60 from payload1 + 40 from payload2 = 100
1320        assert_eq!(count_grouped_stats(&result[0]), 100);
1321        // Second TraceStats: remaining 20 from payload2
1322        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        // One payload with two buckets (70 + 80 = 150 stats), max 100
1328        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        // All split payloads should have the same metadata
1354        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        // Create a bucket with specific start/duration/time_shift
1372        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        // All split buckets should have the same start/duration/time_shift
1381        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        // 10,000 stats with max 4000 - should produce 3 TraceStats
1406        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    // Property test strategies
1418
1419    /// Strategy to generate arbitrary ClientGroupedStats.
1420    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    /// Strategy to generate a bucket with 0..=max_stats_per_bucket grouped stats.
1429    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    /// Strategy to generate a payload with 1..=max_buckets buckets.
1435    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    /// Strategy to generate test inputs for the split function.
1441    ///
1442    /// Parameters:
1443    /// - `num_payloads`: 1..=10
1444    /// - `num_buckets_per_payload`: 1..=5
1445    /// - `num_stats_per_bucket`: 0..=500
1446    /// - `max_entries_per_event`: 1..=1000 (never zero)
1447    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    // Mirrors Go agent behavior: for OTLP traces the version on the root span's meta["version"]
1455    // (set with span-over-resource precedence by otel_span_to_dd_span) should win over the
1456    // resource-level version carried in payload.app_version.
1457    // See: pkg/trace/version/version.go:GetAppVersionFromTrace and pkg/trace/api/otlp.go.
1458    #[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        // Root span carries a span-level version attribute (set by otel_span_to_dd_span with
1471        // span-over-resource precedence).
1472        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        // Simulate OTLP resource extraction: payload.app_version comes from resource attrs only.
1481        trace.payload.app_version = MetaString::from("resource-v1");
1482
1483        let key = transform.build_payload_key(&trace, "");
1484
1485        // The Go agent resolves version from root_span.meta["version"] first, which carries
1486        // span-over-resource precedence. The span attribute ("span-v2") must win.
1487        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            // Property 1: No TraceStats should exceed max_entries_per_event
1503            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            // Property 2: Total stats should be preserved
1514            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}