saluki_components/transforms/trace_sampler/
mod.rs

1//! Trace sampling transform.
2//!
3//! This transform implements agent-side head sampling for traces, supporting:
4//! - Probabilistic sampling based on trace ID
5//! - User-set priority preservation
6//! - Error-based sampling as a safety net
7//! - OTLP trace ingestion with proper sampling decision handling
8//!
9//! TODO:
10//!
11//! - add trace metrics: datadog-agent/pkg/trace/sampler/metrics.go
12//! - adding missing samplers (priority, nopriority)
13//! - add error tracking standalone mode
14
15use std::sync::LazyLock;
16
17use agent_data_plane_config::domains;
18use async_trait::async_trait;
19use saluki_core::accounting::{MemoryBounds, MemoryBoundsBuilder};
20use saluki_core::{
21    components::{transforms::*, BuildContext},
22    data_model::event::{
23        trace::{AttributeValue, Span, Trace},
24        Event, EventType,
25    },
26    topology::OutputDefinition,
27};
28use saluki_error::GenericError;
29use stringtheory::MetaString;
30use tokio::select;
31use tracing::debug;
32
33mod catalog;
34mod core_sampler;
35mod errors;
36mod priority_sampler;
37mod probabilistic;
38mod rare_sampler;
39mod score_sampler;
40mod signature;
41mod telemetry;
42
43use self::probabilistic::PROB_RATE_KEY;
44use self::telemetry::{DecisionWindow, WINDOW};
45use crate::common::datadog::{
46    compute_top_level, get_root_span_index, get_trace_env, sample_by_rate, DECISION_MAKER_MANUAL,
47    DECISION_MAKER_PROBABILISTIC, OTEL_TRACE_ID_META_KEY, SAMPLING_PRIORITY_METRIC_KEY, TAG_DECISION_MAKER, TAG_ORIGIN,
48};
49
50// Sampling priority constants (matching datadog-agent)
51const PRIORITY_AUTO_DROP: i32 = 0;
52const PRIORITY_AUTO_KEEP: i32 = 1;
53const PRIORITY_USER_KEEP: i32 = 2;
54
55// Single Span Sampling and Analytics Events keys
56const KEY_SPAN_SAMPLING_MECHANISM: &str = "_dd.span_sampling.mechanism";
57const KEY_ANALYZED_SPANS: &str = "_dd.analyzed";
58
59fn normalize_sampling_rate(rate: f64) -> f64 {
60    if rate <= 0.0 || rate >= 1.0 {
61        1.0
62    } else {
63        rate
64    }
65}
66
67/// Configuration for the trace sampler transform.
68#[derive(Debug)]
69pub struct TraceSamplerConfiguration {
70    probabilistic_sampler_enabled: bool,
71    probabilistic_hash_seed: u32,
72    probabilistic_full_trace_id: bool,
73    sampling_percentage: f64,
74    error_sampling_enabled: bool,
75    error_tracking_standalone: bool,
76    errors_per_second: f64,
77    target_traces_per_second: f64,
78    extra_sample_rate: f64,
79    max_catalog_entries: usize,
80    default_env: MetaString,
81    rare_sampler_enabled: bool,
82    rare_sampler_tps: f64,
83    rare_sampler_cooldown_secs: f64,
84    rare_sampler_cardinality: usize,
85    otlp_sampling_rate: f64,
86    compute_top_level_by_span_kind: bool,
87}
88
89impl TraceSamplerConfiguration {
90    /// Creates a new `TraceSamplerConfiguration` from the resolved traces domain.
91    ///
92    /// The OTLP trace settings live in their own domain, so they arrive as a separate slice rather
93    /// than through the traces domain.
94    pub fn from_configuration(traces: &domains::traces::Domain, otlp_traces: &domains::otlp::Traces) -> Self {
95        let otlp_sampling_rate = normalize_sampling_rate(otlp_traces.probabilistic_sampler_sampling_percentage / 100.0);
96        Self {
97            probabilistic_sampler_enabled: traces.probabilistic_sampler.enabled,
98            probabilistic_hash_seed: traces.probabilistic_sampler.hash_seed,
99            probabilistic_full_trace_id: traces
100                .features
101                .contains(&domains::traces::ApmFeature::ProbabilisticSamplerFullTraceId),
102            sampling_percentage: traces.probabilistic_sampler.sampling_percentage,
103            error_sampling_enabled: traces.error_sampling_enabled,
104            error_tracking_standalone: traces.error_tracking_standalone_enabled,
105            errors_per_second: traces.errors_per_second,
106            target_traces_per_second: traces.target_traces_per_second,
107            extra_sample_rate: traces.extra_sample_rate,
108            max_catalog_entries: traces.max_catalog_entries,
109            default_env: MetaString::from(traces.default_env.clone()),
110            rare_sampler_enabled: traces.enable_rare_sampler,
111            rare_sampler_tps: traces.rare_sampler.tps,
112            rare_sampler_cooldown_secs: traces.rare_sampler.cooldown,
113            rare_sampler_cardinality: traces.rare_sampler.cardinality,
114            otlp_sampling_rate,
115            compute_top_level_by_span_kind: otlp_traces.enable_compute_top_level_by_span_kind,
116        }
117    }
118}
119
120#[async_trait]
121impl TransformBuilder for TraceSamplerConfiguration {
122    fn input_event_type(&self) -> EventType {
123        EventType::Trace
124    }
125
126    fn outputs(&self) -> &[OutputDefinition<EventType>] {
127        static OUTPUTS: LazyLock<Vec<OutputDefinition<EventType>>> = LazyLock::new(|| {
128            vec![
129                OutputDefinition::default_output(EventType::Trace),
130                OutputDefinition::named_output("metrics", EventType::Metric),
131            ]
132        });
133        &OUTPUTS
134    }
135
136    async fn build(&self, _context: BuildContext) -> Result<Box<dyn Transform + Send>, GenericError> {
137        // TODO: Need to support remote configuration changing these at runtime
138        // See https://github.com/DataDog/saluki/issues/1326
139        let telemetry = DecisionWindow::new();
140
141        let sampler = TraceSampler {
142            sampling_rate: self.sampling_percentage / 100.0,
143            error_sampling_enabled: self.error_sampling_enabled,
144            error_tracking_standalone: self.error_tracking_standalone,
145            probabilistic_sampler_enabled: self.probabilistic_sampler_enabled,
146            probabilistic: probabilistic::ProbabilisticSampler::new(
147                self.probabilistic_hash_seed,
148                self.probabilistic_full_trace_id,
149            ),
150            otlp_sampling_rate: self.otlp_sampling_rate,
151            error_sampler: errors::ErrorsSampler::new(self.errors_per_second, self.extra_sample_rate),
152            priority_sampler: priority_sampler::PrioritySampler::new(
153                self.default_env.clone(),
154                self.extra_sample_rate,
155                self.target_traces_per_second,
156                self.max_catalog_entries,
157            ),
158            no_priority_sampler: score_sampler::NoPrioritySampler::new(
159                self.target_traces_per_second,
160                self.extra_sample_rate,
161            ),
162            rare_sampler: rare_sampler::RareSampler::new(
163                self.rare_sampler_enabled,
164                self.rare_sampler_tps,
165                std::time::Duration::from_secs_f64(self.rare_sampler_cooldown_secs),
166                self.rare_sampler_cardinality,
167                telemetry.counters().clone(),
168            ),
169            telemetry,
170            compute_top_level_by_span_kind: self.compute_top_level_by_span_kind,
171        };
172
173        Ok(Box::new(sampler))
174    }
175}
176
177impl MemoryBounds for TraceSamplerConfiguration {
178    fn specify_bounds(&self, builder: &mut MemoryBoundsBuilder) {
179        builder.minimum().with_single_value::<TraceSampler>("component struct");
180
181        // The catalog holds a hash map slot plus one LRU slab entry per tracked signature. A
182        // configured 0 selects the default cap, matching the catalog's own resolution.
183        let catalog_capacity = if self.max_catalog_entries == 0 {
184            catalog::MAX_CATALOG_ENTRIES
185        } else {
186            self.max_catalog_entries
187        };
188        builder
189            .minimum()
190            .with_map::<signature::ServiceSignature, u32>("priority sampler catalog map", catalog_capacity)
191            .with_map::<signature::ServiceSignature, signature::Signature>(
192                "priority sampler catalog entries",
193                catalog_capacity,
194            );
195
196        // Decision keys span the same (service, env) universe as the catalog, across a handful
197        // of samplers. The window resets every ten seconds, so steady-state growth follows active
198        // combinations; within a window, growth is capped at the telemetry's decision-key limit,
199        // with excess combinations rolled into one bucket per sampler.
200        builder
201            .minimum()
202            .with_map::<telemetry::DecisionKey, telemetry::DecisionCounts>("sampler decision window", catalog_capacity);
203    }
204}
205
206pub struct TraceSampler {
207    sampling_rate: f64,
208    error_tracking_standalone: bool,
209    error_sampling_enabled: bool,
210    probabilistic_sampler_enabled: bool,
211    probabilistic: probabilistic::ProbabilisticSampler,
212    otlp_sampling_rate: f64,
213    compute_top_level_by_span_kind: bool,
214    error_sampler: errors::ErrorsSampler,
215    priority_sampler: priority_sampler::PrioritySampler,
216    no_priority_sampler: score_sampler::NoPrioritySampler,
217    rare_sampler: rare_sampler::RareSampler,
218    telemetry: DecisionWindow,
219}
220
221/// The outcome of evaluating a trace against the configured samplers.
222struct SamplerOutcome {
223    /// Whether the trace should be kept.
224    keep: bool,
225    /// The final sampling priority.
226    priority: i32,
227    /// The decision maker tag value for the deciding sampler, if any.
228    decision_maker: &'static str,
229    /// The index of the root span the samplers evaluated against, when one was found.
230    root_span_idx: Option<usize>,
231    /// The sampler that decided the trace.
232    sampler: telemetry::SamplerName,
233}
234
235impl SamplerOutcome {
236    const fn new(
237        keep: bool, priority: i32, decision_maker: &'static str, root_span_idx: Option<usize>,
238        sampler: telemetry::SamplerName,
239    ) -> Self {
240        Self {
241            keep,
242            priority,
243            decision_maker,
244            root_span_idx,
245            sampler,
246        }
247    }
248}
249
250impl TraceSampler {
251    /// Find the root span index of a trace.
252    fn get_root_span_index(&self, trace: &Trace) -> Option<usize> {
253        // The shared root finder, so every consumer anchors metadata to identical root spans.
254        get_root_span_index(trace.spans())
255    }
256
257    /// Check for user-set sampling priority in trace
258    fn get_user_priority(&self, trace: &Trace, root_span_idx: usize) -> Option<i32> {
259        // First check trace-level sampling priority (last-seen priority from OTLP ingest)
260        if let Some(priority) = trace.priority {
261            return Some(priority);
262        }
263
264        if trace.spans().is_empty() {
265            return None;
266        }
267
268        // Fall back to checking spans (for compatibility with non-OTLP traces)
269        // Prefer the root span (common case), but fall back to scanning all spans to be robust to ordering.
270        if let Some(root) = trace.spans().get(root_span_idx) {
271            if let Some(p) = root
272                .attributes
273                .get(SAMPLING_PRIORITY_METRIC_KEY)
274                .and_then(AttributeValue::as_num)
275            {
276                return Some(p as i32);
277            }
278        }
279        let spans = trace.spans();
280        spans.iter().find_map(|span| {
281            span.attributes
282                .get(SAMPLING_PRIORITY_METRIC_KEY)
283                .and_then(AttributeValue::as_num)
284                .map(|p| p as i32)
285        })
286    }
287
288    /// Returns `true` if the given trace ID should be probabilistically sampled.
289    fn sample_probabilistic(&self, trace_id_high: u64, trace_id_low: u64) -> bool {
290        self.probabilistic
291            .sample(trace_id_high, trace_id_low, self.sampling_rate)
292    }
293
294    fn is_otlp_trace(&self, trace: &Trace, root_span_idx: usize) -> bool {
295        trace
296            .spans()
297            .get(root_span_idx)
298            .map(|span| {
299                span.attributes
300                    .contains_key(&MetaString::from_static(OTEL_TRACE_ID_META_KEY))
301            })
302            .unwrap_or(false)
303    }
304
305    /// Returns `true` if the trace contains a span with an error.
306    fn trace_contains_error(&self, trace: &Trace, consider_exception_span_events: bool) -> bool {
307        trace.spans().iter().any(|span| {
308            span.error() != 0 || (consider_exception_span_events && self.span_contains_exception_span_event(span))
309        })
310    }
311
312    /// Returns `true` if the span has exception span events.
313    ///
314    /// This checks for the `_dd.span_events.has_exception` meta field set to `"true"`.
315    fn span_contains_exception_span_event(&self, span: &Span) -> bool {
316        if let Some(has_exception) = span
317            .attributes
318            .get("_dd.span_events.has_exception")
319            .and_then(AttributeValue::as_string)
320        {
321            return has_exception == "true";
322        }
323        false
324    }
325
326    /// Computes the OTLP pre-sampling priority and decision maker for a trace, mirroring
327    /// `OTLPReceiver.createChunks` in DDA which runs before `runSamplersV1`.
328    ///
329    /// Returns `Some((priority, dm))` for OTLP traces when the probabilistic sampler is disabled,
330    /// or `None` if pre-sampling doesn't apply.
331    ///
332    /// See: https://github.com/DataDog/datadog-agent/blob/be33ac1490c4a34602cbc65a211406b73ad6d00b/pkg/trace/api/otlp.go#L561-L585
333    fn otlp_pre_sample(&mut self, trace: &mut Trace, root_span_idx: usize) -> Option<(i32, &'static str)> {
334        if self.probabilistic_sampler_enabled || !self.is_otlp_trace(trace, root_span_idx) {
335            return None;
336        }
337        let (priority, dm) = if let Some(user_priority) = self.get_user_priority(trace, root_span_idx) {
338            (user_priority, DECISION_MAKER_MANUAL)
339        } else {
340            let root_trace_id = trace.trace_id_low;
341            if sample_by_rate(root_trace_id, self.otlp_sampling_rate) {
342                (PRIORITY_AUTO_KEEP, DECISION_MAKER_PROBABILISTIC)
343            } else {
344                (PRIORITY_AUTO_DROP, DECISION_MAKER_PROBABILISTIC)
345            }
346        };
347        if priority == PRIORITY_AUTO_KEEP {
348            if let Some(root_span) = trace.spans_mut().get_mut(root_span_idx) {
349                root_span.attributes.remove(PROB_RATE_KEY);
350            }
351        }
352        Some((priority, dm))
353    }
354
355    /// Apply analyzed span sampling to the trace.
356    ///
357    /// Returns `true` if the trace was modified.
358    fn analyzed_span_sampling(&self, trace: &mut Trace) -> bool {
359        let retained = trace.retain_spans(|_, span| span.attributes.contains_key(KEY_ANALYZED_SPANS));
360        if retained > 0 {
361            trace.dropped_trace = false;
362            trace.priority = Some(PRIORITY_USER_KEEP);
363            trace.otlp_sampling_rate = Some(self.sampling_rate);
364            true
365        } else {
366            false
367        }
368    }
369
370    /// Returns `true` if the given trace has any analyzed spans.
371    fn has_analyzed_spans(&self, trace: &Trace) -> bool {
372        trace
373            .spans()
374            .iter()
375            .any(|span| span.attributes.contains_key(KEY_ANALYZED_SPANS))
376    }
377
378    /// Apply Single Span Sampling to the trace
379    /// Returns true if the trace was modified
380    fn single_span_sampling(&self, trace: &mut Trace) -> bool {
381        let retained = trace.retain_spans(|_, span| span.attributes.contains_key(KEY_SPAN_SAMPLING_MECHANISM));
382        if retained > 0 {
383            trace.dropped_trace = false;
384            trace.priority = Some(PRIORITY_USER_KEEP);
385            trace.otlp_sampling_rate = Some(self.sampling_rate);
386            true
387        } else {
388            false
389        }
390    }
391
392    /// Evaluates the given trace against all configured samplers.
393    fn run_samplers(&mut self, trace: &mut Trace) -> SamplerOutcome {
394        let outcome = self.run_samplers_inner(trace);
395
396        // Unclaimed decisions carry no sampler identity, so recording them would inflate a
397        // series that exists only to expose attribution gaps; empty and rootless traces are the
398        // only unclaimed cases.
399        if outcome.sampler != telemetry::SamplerName::Unknown {
400            let (service, env) = self.decision_dimensions(trace, outcome.root_span_idx);
401            self.telemetry
402                .record_decision(outcome.keep, outcome.sampler, outcome.priority, service, env);
403        }
404
405        outcome
406    }
407
408    /// Reads the service and env dimensions a decision is recorded under.
409    fn decision_dimensions(&self, trace: &Trace, root_span_idx: Option<usize>) -> (MetaString, MetaString) {
410        match root_span_idx {
411            Some(root_span_idx) => {
412                let service = trace
413                    .spans()
414                    .get(root_span_idx)
415                    .map(|span| MetaString::from(span.service()))
416                    .unwrap_or_default();
417                let env = get_trace_env(trace, root_span_idx).cloned().unwrap_or_default();
418                (service, env)
419            }
420            None => (MetaString::empty(), MetaString::empty()),
421        }
422    }
423
424    fn run_samplers_inner(&mut self, trace: &mut Trace) -> SamplerOutcome {
425        // The name starts unclaimed; whichever sampler decides the trace claims it, so every
426        // decision is recorded under exactly one sampler.
427        let mut sampler_name = telemetry::SamplerName::Unknown;
428
429        // logic taken from: https://github.com/DataDog/datadog-agent/blob/main/pkg/trace/agent/agent.go#L1066
430        // Empty trace check
431        if trace.spans().is_empty() {
432            return SamplerOutcome::new(false, PRIORITY_AUTO_DROP, "", None, sampler_name);
433        }
434
435        let Some(root_span_idx) = self.get_root_span_index(trace) else {
436            return SamplerOutcome::new(false, PRIORITY_AUTO_DROP, "", None, sampler_name);
437        };
438
439        // ETS: only sample traces containing errors (including exception span events); skip all other samplers.
440        // logic taken from: https://github.com/DataDog/datadog-agent/blob/be33ac1490c4a34602cbc65a211406b73ad6d00/pkg/trace/agent/agent.go#L1068
441        if self.error_tracking_standalone {
442            sampler_name = telemetry::SamplerName::Error;
443            let otlp_pre_sample = self.otlp_pre_sample(trace, root_span_idx);
444            if self.trace_contains_error(trace, true) {
445                let now = std::time::SystemTime::now();
446                let keep = self.error_sampler.sample_error(now, trace, root_span_idx);
447                let default_priority = if keep { PRIORITY_AUTO_KEEP } else { PRIORITY_AUTO_DROP };
448                let (priority, dm) = otlp_pre_sample.unwrap_or((default_priority, ""));
449                return SamplerOutcome::new(keep, priority, dm, Some(root_span_idx), sampler_name);
450            }
451            let (pre_priority, pre_dm) = otlp_pre_sample.unwrap_or((PRIORITY_AUTO_DROP, ""));
452            return SamplerOutcome::new(false, pre_priority, pre_dm, Some(root_span_idx), sampler_name);
453        }
454
455        // Run the rare sampler early, before all other samplers. This mirrors the Go agent behavior
456        // where the rare sampler runs first to catch traces that would otherwise be dropped entirely.
457        // logic taken from: https://github.com/DataDog/datadog-agent/blob/main/pkg/trace/agent/agent.go#L1078
458        let rare = self.rare_sampler.sample(trace, root_span_idx);
459
460        // Modern path: ProbabilisticSamplerEnabled = true
461        if self.probabilistic_sampler_enabled {
462            sampler_name = telemetry::SamplerName::Probabilistic;
463            let mut prob_keep = false;
464            let mut decision_maker = "";
465
466            if rare {
467                // Rare sampler wins over probabilistic sampling.
468                sampler_name = telemetry::SamplerName::Rare;
469                prob_keep = true;
470            } else {
471                if self.sample_probabilistic(trace.trace_id_high, trace.trace_id_low) {
472                    decision_maker = DECISION_MAKER_PROBABILISTIC;
473                    prob_keep = true;
474
475                    if let Some(root_span) = trace.spans_mut().get_mut(root_span_idx) {
476                        root_span.attributes.insert(
477                            MetaString::from_static(PROB_RATE_KEY),
478                            AttributeValue::Float(self.sampling_rate),
479                        );
480                    }
481                } else if self.error_sampling_enabled && self.trace_contains_error(trace, false) {
482                    sampler_name = telemetry::SamplerName::Error;
483                    let now = std::time::SystemTime::now();
484                    prob_keep = self.error_sampler.sample_error(now, trace, root_span_idx);
485                }
486            }
487
488            let priority = if prob_keep {
489                PRIORITY_AUTO_KEEP
490            } else {
491                PRIORITY_AUTO_DROP
492            };
493
494            return SamplerOutcome::new(prob_keep, priority, decision_maker, Some(root_span_idx), sampler_name);
495        }
496
497        // Read once here, where the samplers below consume it; every path above returns without
498        // needing it.
499        let now = std::time::SystemTime::now();
500        let user_priority = self.get_user_priority(trace, root_span_idx);
501        if let Some(priority) = user_priority {
502            sampler_name = telemetry::SamplerName::Priority;
503            if Self::trace_chunk_contains_probability_sampling(trace) {
504                // Traces already sampled probabilistically upstream count under the probabilistic
505                // sampler; the priority sampler still decides them.
506                sampler_name = telemetry::SamplerName::Probabilistic;
507            }
508            if priority < PRIORITY_AUTO_DROP {
509                // Manual drop: short-circuit and skip other samplers.
510                return SamplerOutcome::new(false, priority, "", Some(root_span_idx), sampler_name);
511            }
512
513            if rare {
514                sampler_name = telemetry::SamplerName::Rare;
515                return SamplerOutcome::new(true, priority, "", Some(root_span_idx), sampler_name);
516            }
517
518            if self.priority_sampler.sample(now, trace, root_span_idx, priority, 0.0) {
519                return SamplerOutcome::new(true, priority, "", Some(root_span_idx), sampler_name);
520            }
521        } else if self.is_otlp_trace(trace, root_span_idx) {
522            // Rare check mirrors agent behavior: https://github.com/DataDog/datadog-agent/blob/main/pkg/trace/agent/agent.go#L1129-L1140
523            if rare {
524                sampler_name = telemetry::SamplerName::Rare;
525                return SamplerOutcome::new(true, PRIORITY_AUTO_KEEP, "", Some(root_span_idx), sampler_name);
526            }
527
528            // The OTLP rate decision is probabilistic: its keeps and its drops both record
529            // under the probabilistic sampler, never the no-priority bucket.
530            sampler_name = telemetry::SamplerName::Probabilistic;
531            // some sampling happens upstream in the otlp receiver in the agent: https://github.com/DataDog/datadog-agent/blob/main/pkg/trace/api/otlp.go#L572
532            let root_trace_id = trace.trace_id_low;
533            if sample_by_rate(root_trace_id, self.otlp_sampling_rate) {
534                if let Some(root_span) = trace.spans_mut().get_mut(root_span_idx) {
535                    root_span.attributes.remove(PROB_RATE_KEY);
536                }
537                return SamplerOutcome::new(
538                    true,
539                    PRIORITY_AUTO_KEEP,
540                    DECISION_MAKER_PROBABILISTIC,
541                    Some(root_span_idx),
542                    sampler_name,
543                );
544            }
545        } else {
546            sampler_name = telemetry::SamplerName::NoPriority;
547            if rare {
548                sampler_name = telemetry::SamplerName::Rare;
549                return SamplerOutcome::new(true, PRIORITY_AUTO_KEEP, "", Some(root_span_idx), sampler_name);
550            }
551            if self.no_priority_sampler.sample(now, trace, root_span_idx) {
552                return SamplerOutcome::new(true, PRIORITY_AUTO_KEEP, "", Some(root_span_idx), sampler_name);
553            }
554        }
555
556        if self.error_sampling_enabled && self.trace_contains_error(trace, false) {
557            sampler_name = telemetry::SamplerName::Error;
558            let keep = self.error_sampler.sample_error(now, trace, root_span_idx);
559            if keep {
560                return SamplerOutcome::new(true, PRIORITY_AUTO_KEEP, "", Some(root_span_idx), sampler_name);
561            }
562        }
563
564        // Default: drop the trace
565        SamplerOutcome::new(false, PRIORITY_AUTO_DROP, "", Some(root_span_idx), sampler_name)
566    }
567
568    /// Returns whether the trace was already sampled probabilistically upstream.
569    ///
570    /// Tracer head sampling and receiver-side sampling both stamp the decision maker, so the
571    /// first span carrying it identifies the trace.
572    fn trace_chunk_contains_probability_sampling(trace: &Trace) -> bool {
573        trace
574            .spans()
575            .iter()
576            .find_map(|span| {
577                span.attributes
578                    .get(TAG_DECISION_MAKER)
579                    .and_then(AttributeValue::as_string)
580            })
581            .is_some_and(|dm| dm == DECISION_MAKER_PROBABILISTIC)
582    }
583
584    /// Fills in trace-level metadata from the spans.
585    ///
586    /// Origin comes from the root span's `_dd.origin` tag, and the decision maker is the first
587    /// span in chunk order carrying `_dd.p.dm` - both only when not already set. This runs here,
588    /// downstream of span filtering, so the metadata anchors to the final span population rather
589    /// than the one that arrived at decode time.
590    fn backfill_trace_metadata(&self, trace: &mut Trace, root_span_idx: usize) {
591        if trace.origin.is_empty() {
592            if let Some(origin) = trace.spans()[root_span_idx]
593                .attributes
594                .get(TAG_ORIGIN)
595                .and_then(AttributeValue::as_string)
596            {
597                trace.origin = origin.clone();
598            }
599        }
600
601        if trace.decision_maker.is_none() {
602            let first_decision_maker = trace.spans().iter().find_map(|span| {
603                span.attributes
604                    .get(TAG_DECISION_MAKER)
605                    .and_then(AttributeValue::as_string)
606            });
607            if let Some(dm) = first_decision_maker {
608                trace.decision_maker = Some(dm.clone());
609            }
610        }
611
612        if !self.compute_top_level_by_span_kind {
613            compute_top_level(trace.spans_mut());
614        }
615    }
616
617    /// Apply sampling metadata to the trace in-place.
618    ///
619    /// The `root_span_id` parameter identifies which span should receive the sampling metadata.
620    /// This avoids recalculating the root span since it was already found in `run_samplers`.
621    fn apply_sampling_metadata(
622        &self, trace: &mut Trace, keep: bool, priority: i32, decision_maker: &'static str, root_span_idx: usize,
623    ) {
624        let is_otlp = self.is_otlp_trace(trace, root_span_idx);
625        // Chunk-level decision maker: the sampler's value when it decided, otherwise the first
626        // carrier promoted during backfill (falling back to the root's own tag).
627        let promoted_decision_maker = if decision_maker.is_empty() {
628            trace.decision_maker.clone()
629        } else {
630            None
631        };
632        let root_span_value = match trace.spans_mut().get_mut(root_span_idx) {
633            Some(span) => span,
634            None => return,
635        };
636
637        let decision_maker_meta = if decision_maker.is_empty() {
638            promoted_decision_maker.or_else(|| {
639                root_span_value
640                    .attributes
641                    .get(TAG_DECISION_MAKER)
642                    .and_then(AttributeValue::as_string)
643                    .cloned()
644            })
645        } else {
646            Some(MetaString::from_static(decision_maker))
647        };
648
649        // The span-level decision maker is stamped only when a sampler made the call: promotion
650        // targets trace chunk metadata, and each span's own tag is left in place because
651        // downstream systems rely on the span keeping its original value.
652        //
653        // With the APM-level probabilistic sampler on OTLP traces, `_dd.p.dm` is written to
654        // trace chunk tags only (not span meta); on the legacy sampling path it is written to
655        // both. The span meta write is skipped only when both conditions hold; the DM value
656        // still flows through trace fields to the encoder.
657        if priority > 0 && !decision_maker.is_empty() && !(is_otlp && self.probabilistic_sampler_enabled) {
658            root_span_value.attributes.insert(
659                MetaString::from_static(TAG_DECISION_MAKER),
660                AttributeValue::String(MetaString::from_static(decision_maker)),
661            );
662        }
663
664        // Now set sampling metadata directly on the trace.
665        trace.dropped_trace = !keep;
666        trace.priority = Some(priority);
667        trace.decision_maker = if priority > 0 { decision_maker_meta } else { None };
668        trace.otlp_sampling_rate = Some(if is_otlp {
669            self.otlp_sampling_rate
670        } else {
671            self.sampling_rate
672        });
673    }
674
675    fn process_trace(&mut self, trace: &mut Trace) -> bool {
676        // keep is a boolean that indicates if the trace should be kept or dropped
677        // priority is the sampling priority
678        // decision_maker is the tag that indicates the decision maker (probabilistic, error, etc.)
679        // root_span_idx is the index of the root span of the trace
680        let SamplerOutcome {
681            keep,
682            priority,
683            decision_maker,
684            root_span_idx,
685            ..
686        } = self.run_samplers(trace);
687
688        // Backfill trace-derived metadata before any forwarding path (keep, ETS, or
689        // single-span/analytics sampling), so every forwarded trace carries metadata derived
690        // from its final span population.
691        if let Some(root_idx) = root_span_idx {
692            self.backfill_trace_metadata(trace, root_idx);
693        }
694
695        // Apply sampling metadata and forward if kept, or if ETS (dropped non-error traces are
696        // forwarded with DroppedTrace=true, suppressing SSS/analytics).
697        if keep || self.error_tracking_standalone {
698            if let Some(root_idx) = root_span_idx {
699                self.apply_sampling_metadata(trace, keep, priority, decision_maker, root_idx);
700            }
701            return true;
702        }
703
704        // logic taken from here: https://github.com/DataDog/datadog-agent/blob/main/pkg/trace/agent/agent.go#L980-L990
705        // try single span sampling (keeps spans marked for sampling when trace would be dropped)
706        let modified = self.single_span_sampling(trace);
707        if !modified {
708            // Fall back to analytics events if no SSS spans
709            if self.analyzed_span_sampling(trace) {
710                return true;
711            }
712        } else if self.has_analyzed_spans(trace) {
713            // Warn about both SSS and analytics events
714            debug!(
715                "Detected both analytics events AND single span sampling in the same trace. Single span sampling wins because App Analytics is deprecated."
716            );
717            return true;
718        }
719
720        // If we modified the trace with SSS, send it
721        if modified {
722            return true;
723        }
724
725        // Neither SSS nor analytics events found, drop the trace
726        debug!("Dropping trace with priority {}", priority);
727        false
728    }
729}
730
731#[async_trait]
732impl Transform for TraceSampler {
733    async fn run(mut self: Box<Self>, mut context: TransformContext) -> Result<(), GenericError> {
734        let mut health = context.take_health_handle();
735        health.mark_ready();
736        debug!("Trace sampler transform started.");
737
738        let mut window = tokio::time::interval_at(tokio::time::Instant::now() + WINDOW, WINDOW);
739
740        loop {
741            select! {
742                _ = health.live() => continue,
743                _ = window.tick() => {
744                    // Report the sampling telemetry through the dedicated metrics output, so the
745                    // metrics flow into the metrics pipeline and on to the backend. The signature
746                    // counts are read here, at report time, so a quiet window cannot republish
747                    // stale gauges.
748                    let events = self.telemetry.take_window_events(
749                        self.priority_sampler.tracked_signature_count(),
750                        self.no_priority_sampler.tracked_signature_count(),
751                        self.error_sampler.tracked_signature_count(),
752                    );
753                    context.dispatcher().buffered_named("metrics")?.send_all(events).await?;
754                }
755                maybe_events = context.events().next() => match maybe_events {
756                    Some(mut events) => {
757                        events.remove_if(|event| match event {
758                            Event::Trace(trace) => !self.process_trace(trace),
759                            _ => false,
760                        });
761
762                        context.dispatcher().buffered()?.send_all(events).await?;
763                    }
764                    None => {
765                        // The input stream has ended, so report the final window before stopping.
766                        let events = self.telemetry.take_window_events(
767                            self.priority_sampler.tracked_signature_count(),
768                            self.no_priority_sampler.tracked_signature_count(),
769                            self.error_sampler.tracked_signature_count(),
770                        );
771                        context.dispatcher().buffered_named("metrics")?.send_all(events).await?;
772                        break;
773                    }
774                }
775            }
776        }
777
778        debug!("Trace sampler transform stopped.");
779        Ok(())
780    }
781}
782
783#[cfg(test)]
784mod tests {
785    use std::collections::HashMap;
786
787    use saluki_core::data_model::event::trace::{AttributeValue, Span as DdSpan, Trace};
788    const PRIORITY_USER_DROP: i32 = -1;
789
790    use super::*;
791    fn create_test_sampler() -> TraceSampler {
792        TraceSampler {
793            sampling_rate: 1.0,
794            error_sampling_enabled: true,
795            error_tracking_standalone: false,
796            probabilistic_sampler_enabled: true,
797            probabilistic: probabilistic::ProbabilisticSampler::new(0, false),
798            otlp_sampling_rate: 1.0,
799            error_sampler: errors::ErrorsSampler::new(10.0, 1.0),
800            priority_sampler: priority_sampler::PrioritySampler::new(MetaString::from("agent-env"), 1.0, 10.0, 5000),
801            no_priority_sampler: score_sampler::NoPrioritySampler::new(10.0, 1.0),
802            rare_sampler: rare_sampler::RareSampler::new(
803                false,
804                5.0,
805                std::time::Duration::from_secs(300),
806                200,
807                telemetry::SamplerCounters::new(),
808            ),
809            telemetry: telemetry::DecisionWindow::new(),
810            compute_top_level_by_span_kind: false,
811        }
812    }
813
814    fn create_test_span(span_id: u64, error: i32) -> DdSpan {
815        DdSpan::new(
816            MetaString::from("test-service"),
817            MetaString::from("test-operation"),
818            MetaString::from("test-resource"),
819            MetaString::from("test-type"),
820            span_id,
821            0,    // parent_id
822            0,    // start
823            1000, // duration
824            error,
825        )
826    }
827
828    fn create_test_span_with_metrics(span_id: u64, metrics: HashMap<String, f64>) -> DdSpan {
829        let attrs: saluki_common::collections::FastHashMap<MetaString, AttributeValue> = metrics
830            .into_iter()
831            .map(|(k, v)| (MetaString::from(k), AttributeValue::Float(v)))
832            .collect();
833        create_test_span(span_id, 0).with_attributes(attrs)
834    }
835
836    #[allow(dead_code)]
837    fn create_test_span_with_meta(span_id: u64, meta: HashMap<String, String>) -> DdSpan {
838        let attrs: saluki_common::collections::FastHashMap<MetaString, AttributeValue> = meta
839            .into_iter()
840            .map(|(k, v)| (MetaString::from(k), AttributeValue::String(MetaString::from(v))))
841            .collect();
842        create_test_span(span_id, 0).with_attributes(attrs)
843    }
844
845    fn create_test_trace(spans: Vec<DdSpan>) -> Trace {
846        Trace::new(spans)
847    }
848
849    fn create_span_with_parent_and_service(
850        span_id: u64, parent_id: u64, service: &str, meta: &[(&str, &str)],
851    ) -> DdSpan {
852        let mut attrs = saluki_common::collections::FastHashMap::default();
853        for (k, v) in meta {
854            attrs.insert(MetaString::from(*k), AttributeValue::String(MetaString::from(*v)));
855        }
856        DdSpan::new(
857            MetaString::from(service),
858            MetaString::from("operation"),
859            MetaString::from("resource"),
860            MetaString::from("type"),
861            span_id,
862            parent_id,
863            0,
864            1000,
865            0,
866        )
867        .with_attributes(attrs)
868    }
869
870    #[test]
871    fn apply_sampling_metadata_preserves_backfilled_decision_maker() {
872        let sampler = create_test_sampler();
873        let root = create_test_span(1, 0);
874        let mut trace = create_test_trace(vec![root]);
875        trace.decision_maker = Some(MetaString::from("-8"));
876
877        sampler.apply_sampling_metadata(&mut trace, true, PRIORITY_AUTO_KEEP, "", 0);
878
879        assert_eq!(
880            trace.decision_maker.as_deref(),
881            Some("-8"),
882            "the chunk-level decision maker keeps the promoted value"
883        );
884        assert!(
885            !trace.spans()[0].attributes.contains_key(TAG_DECISION_MAKER),
886            "a promoted decision maker is never stamped onto the root span; span-level _dd.p.dm is written only when a sampler decides"
887        );
888    }
889
890    #[test]
891    fn apply_sampling_metadata_keeps_root_decision_maker_without_sampler_decision() {
892        let sampler = create_test_sampler();
893        let mut attrs = saluki_common::collections::FastHashMap::default();
894        attrs.insert(
895            MetaString::from(TAG_DECISION_MAKER),
896            AttributeValue::String(MetaString::from("-9")),
897        );
898        let root = create_test_span(1, 0).with_attributes(attrs);
899        let mut trace = create_test_trace(vec![root]);
900        // Promoted from an earlier child span during backfill.
901        trace.decision_maker = Some(MetaString::from("-8"));
902
903        sampler.apply_sampling_metadata(&mut trace, true, PRIORITY_AUTO_KEEP, "", 0);
904
905        assert_eq!(
906            trace.decision_maker.as_deref(),
907            Some("-8"),
908            "the chunk-level decision maker keeps the promoted value"
909        );
910        assert_eq!(
911            trace.spans()[0]
912                .attributes
913                .get(TAG_DECISION_MAKER)
914                .and_then(AttributeValue::as_string)
915                .map(|s| s.as_ref()),
916            Some("-9"),
917            "the root span's own decision maker is never overwritten by a promoted value"
918        );
919    }
920
921    #[test]
922    fn backfill_promotes_origin_from_root_span() {
923        let sampler = create_test_sampler();
924        // The root is reported last here.
925        let child = create_span_with_parent_and_service(2, 1, "svc", &[]);
926        let root = create_span_with_parent_and_service(1, 0, "svc", &[("_dd.origin", "lambda")]);
927        let mut trace = create_test_trace(vec![child, root]);
928
929        let root_idx = sampler.get_root_span_index(&trace).unwrap();
930        sampler.backfill_trace_metadata(&mut trace, root_idx);
931
932        assert_eq!(trace.origin.as_ref(), "lambda");
933    }
934
935    #[test]
936    fn backfill_ignores_origin_from_non_root_span() {
937        let sampler = create_test_sampler();
938        let root = create_span_with_parent_and_service(1, 0, "svc", &[]);
939        let child = create_span_with_parent_and_service(2, 1, "svc", &[("_dd.origin", "rum")]);
940        let mut trace = create_test_trace(vec![root, child]);
941
942        let root_idx = sampler.get_root_span_index(&trace).unwrap();
943        sampler.backfill_trace_metadata(&mut trace, root_idx);
944
945        assert!(trace.origin.as_ref().is_empty());
946    }
947
948    #[test]
949    fn backfill_promotes_first_decision_maker() {
950        let sampler = create_test_sampler();
951        let root = create_span_with_parent_and_service(1, 0, "svc", &[]);
952        let early = create_span_with_parent_and_service(2, 1, "svc", &[("_dd.p.dm", "-8")]);
953        let late = create_span_with_parent_and_service(3, 1, "svc", &[("_dd.p.dm", "-9")]);
954        let mut trace = create_test_trace(vec![root, early, late]);
955
956        let root_idx = sampler.get_root_span_index(&trace).unwrap();
957        sampler.backfill_trace_metadata(&mut trace, root_idx);
958
959        assert_eq!(
960            trace.decision_maker.as_deref(),
961            Some("-8"),
962            "first span carrying the tag wins - not a root-based rule"
963        );
964    }
965
966    #[test]
967    fn backfill_marks_top_level_fallback_when_span_kind_computation_disabled() {
968        let sampler = create_test_sampler();
969        // Roots, orphans, and service boundaries get marked; interior same-service spans don't.
970        let root = create_span_with_parent_and_service(1, 0, "svc-a", &[]);
971        let child = create_span_with_parent_and_service(2, 1, "svc-a", &[]);
972        let orphan = create_span_with_parent_and_service(3, 0xEE, "svc-a", &[]);
973        let other_service = create_span_with_parent_and_service(5, 1, "svc-b", &[]);
974        let local_entry = create_span_with_parent_and_service(6, 5, "svc-a", &[]);
975        let mut trace = create_test_trace(vec![root, child, orphan, other_service, local_entry]);
976
977        let root_idx = sampler.get_root_span_index(&trace).unwrap();
978        sampler.backfill_trace_metadata(&mut trace, root_idx);
979
980        let spans = trace.spans();
981        let top_level = |sid: u64| {
982            spans
983                .iter()
984                .find(|s| s.span_id() == sid)
985                .unwrap()
986                .attributes
987                .get(crate::common::datadog::TOP_LEVEL_KEY)
988                .and_then(AttributeValue::as_num)
989        };
990        assert_eq!(top_level(1), Some(1.0), "root spans are marked");
991        assert_eq!(
992            top_level(3),
993            Some(1.0),
994            "orphans whose parent is missing from the chunk are marked"
995        );
996        assert_eq!(top_level(5), Some(1.0), "spans entering a different service are marked");
997        assert_eq!(
998            top_level(6),
999            Some(1.0),
1000            "spans crossing back into the resource service are marked (local root)"
1001        );
1002        assert_eq!(
1003            top_level(2),
1004            None,
1005            "same-service children with their parent present are not marked"
1006        );
1007    }
1008
1009    #[test]
1010    fn backfill_skips_top_level_when_span_kind_computation_enabled() {
1011        let sampler = TraceSampler {
1012            compute_top_level_by_span_kind: true,
1013            ..create_test_sampler()
1014        };
1015        let root = create_span_with_parent_and_service(1, 0, "svc", &[]);
1016        let child = create_span_with_parent_and_service(2, 1, "svc", &[]);
1017        let orphan = create_span_with_parent_and_service(3, 0xEE, "svc", &[]);
1018        let mut trace = create_test_trace(vec![root, child, orphan]);
1019
1020        let root_idx = sampler.get_root_span_index(&trace).unwrap();
1021        sampler.backfill_trace_metadata(&mut trace, root_idx);
1022
1023        for span in trace.spans() {
1024            assert!(
1025                !span.attributes.contains_key(crate::common::datadog::TOP_LEVEL_KEY),
1026                "no span should carry a top-level mark with the fallback disabled"
1027            );
1028        }
1029    }
1030
1031    #[test]
1032    fn process_trace_backfills_metadata_for_forwarded_traces() {
1033        let mut sampler = create_test_sampler();
1034        sampler.sampling_rate = 1.0; // keeps the trace, so it is forwarded
1035        sampler.probabilistic_sampler_enabled = true;
1036
1037        let child = create_span_with_parent_and_service(2, 1, "svc", &[]);
1038        let root = create_span_with_parent_and_service(1, 0, "svc", &[("_dd.origin", "lambda")]);
1039        let mut trace = create_test_trace(vec![child, root]);
1040
1041        let forwarded = sampler.process_trace(&mut trace);
1042        assert!(forwarded);
1043        assert_eq!(
1044            trace.origin.as_ref(),
1045            "lambda",
1046            "origin is backfilled for forwarded traces"
1047        );
1048        assert!(
1049            trace
1050                .spans()
1051                .iter()
1052                .any(|s| s.attributes.contains_key(crate::common::datadog::TOP_LEVEL_KEY)),
1053            "top-level marks are set for forwarded traces"
1054        );
1055    }
1056
1057    #[test]
1058    fn apply_sampling_metadata_sampler_decision_overrides_backfill() {
1059        let sampler = create_test_sampler();
1060        let root = create_test_span(1, 0);
1061        let mut trace = create_test_trace(vec![root]);
1062        trace.decision_maker = Some(MetaString::from("-8"));
1063
1064        sampler.apply_sampling_metadata(&mut trace, true, PRIORITY_AUTO_KEEP, DECISION_MAKER_PROBABILISTIC, 0);
1065
1066        assert_eq!(trace.decision_maker.as_deref(), Some(DECISION_MAKER_PROBABILISTIC));
1067        assert_eq!(
1068            trace.spans()[0]
1069                .attributes
1070                .get(TAG_DECISION_MAKER)
1071                .and_then(AttributeValue::as_string)
1072                .map(|s| s.as_ref()),
1073            Some(DECISION_MAKER_PROBABILISTIC)
1074        );
1075    }
1076
1077    #[test]
1078    fn user_priority_detection() {
1079        let sampler = create_test_sampler();
1080
1081        // Test trace with user-set priority = 2 (UserKeep)
1082        let mut metrics = HashMap::new();
1083        metrics.insert(SAMPLING_PRIORITY_METRIC_KEY.to_string(), 2.0);
1084        let span = create_test_span_with_metrics(1, metrics);
1085        let trace = create_test_trace(vec![span]);
1086        let root_idx = sampler.get_root_span_index(&trace).unwrap();
1087
1088        assert_eq!(sampler.get_user_priority(&trace, root_idx), Some(2));
1089
1090        // Test trace with user-set priority = -1 (UserDrop)
1091        let mut metrics = HashMap::new();
1092        metrics.insert(SAMPLING_PRIORITY_METRIC_KEY.to_string(), -1.0);
1093        let span = create_test_span_with_metrics(1, metrics);
1094        let trace = create_test_trace(vec![span]);
1095        let root_idx = sampler.get_root_span_index(&trace).unwrap();
1096
1097        assert_eq!(sampler.get_user_priority(&trace, root_idx), Some(-1));
1098
1099        // Test trace without user priority
1100        let span = create_test_span(1, 0);
1101        let trace = create_test_trace(vec![span]);
1102        let root_idx = sampler.get_root_span_index(&trace).unwrap();
1103
1104        assert_eq!(sampler.get_user_priority(&trace, root_idx), None);
1105    }
1106
1107    #[test]
1108    fn trace_level_priority_takes_precedence() {
1109        let sampler = create_test_sampler();
1110
1111        // Test trace-level priority overrides span priorities (last-seen priority)
1112        // Create spans with different priorities - root has 0, later span has 2
1113        let mut metrics_root = HashMap::new();
1114        metrics_root.insert(SAMPLING_PRIORITY_METRIC_KEY.to_string(), 0.0);
1115        let root_span = create_test_span_with_metrics(1, metrics_root);
1116
1117        let mut metrics_later = HashMap::new();
1118        metrics_later.insert(SAMPLING_PRIORITY_METRIC_KEY.to_string(), 1.0);
1119        let later_span = create_test_span_with_metrics(2, metrics_later).with_parent_id(1);
1120
1121        let mut trace = create_test_trace(vec![root_span, later_span]);
1122        let root_idx = sampler.get_root_span_index(&trace).unwrap();
1123
1124        // Without trace-level priority, should get priority from root (0)
1125        assert_eq!(sampler.get_user_priority(&trace, root_idx), Some(0));
1126
1127        // Now set trace-level priority to 2 (simulating last-seen priority from OTLP translator)
1128        trace.priority = Some(2);
1129
1130        // Trace-level priority should take precedence
1131        assert_eq!(sampler.get_user_priority(&trace, root_idx), Some(2));
1132
1133        // Test that trace-level priority is used even when no span has priority
1134        let span_no_priority = create_test_span(3, 0);
1135        let mut trace_only_trace_level = create_test_trace(vec![span_no_priority]);
1136        trace_only_trace_level.priority = Some(1);
1137        let root_idx = sampler.get_root_span_index(&trace_only_trace_level).unwrap();
1138
1139        assert_eq!(sampler.get_user_priority(&trace_only_trace_level, root_idx), Some(1));
1140    }
1141
1142    #[test]
1143    fn manual_keep_with_trace_level_priority() {
1144        let mut sampler = create_test_sampler();
1145        sampler.probabilistic_sampler_enabled = false; // Use legacy path that checks user priority
1146
1147        // Test that manual keep (priority = 2) works via trace-level priority
1148        let span = create_test_span(1, 0);
1149        let mut trace = create_test_trace(vec![span]);
1150        trace.priority = Some(PRIORITY_USER_KEEP);
1151
1152        let SamplerOutcome {
1153            keep,
1154            priority,
1155            decision_maker,
1156            ..
1157        } = sampler.run_samplers(&mut trace);
1158        assert!(keep);
1159        assert_eq!(priority, PRIORITY_USER_KEEP);
1160        assert_eq!(decision_maker, "");
1161
1162        // Test manual drop (priority = -1) via trace-level priority
1163        let span = create_test_span(1, 0);
1164        let mut trace = create_test_trace(vec![span]);
1165        trace.priority = Some(PRIORITY_USER_DROP);
1166
1167        let SamplerOutcome { keep, priority, .. } = sampler.run_samplers(&mut trace);
1168        assert!(!keep); // Should not keep when user drops
1169        assert_eq!(priority, PRIORITY_USER_DROP);
1170
1171        // Test that priority = 1 (auto keep) via trace-level is also respected
1172        let span = create_test_span(1, 0);
1173        let mut trace = create_test_trace(vec![span]);
1174        trace.priority = Some(PRIORITY_AUTO_KEEP);
1175
1176        let SamplerOutcome {
1177            keep,
1178            priority,
1179            decision_maker,
1180            ..
1181        } = sampler.run_samplers(&mut trace);
1182        assert!(keep);
1183        assert_eq!(priority, PRIORITY_AUTO_KEEP);
1184        assert_eq!(decision_maker, "");
1185    }
1186
1187    #[test]
1188    fn probabilistic_sampling_known_decisions() {
1189        // The bucketed probabilistic sampler is fully deterministic: it hashes the trace ID into one of 0x4000
1190        // buckets and keeps the trace when `bucket < (rate * 0x4000)`. These cases pin the exact keep/drop decision
1191        // for known trace IDs at known rates, so a regression in the hash, the bucket mask, or the comparison is
1192        // caught (a determinism-only check would not catch any of those).
1193        //
1194        // Expected values were computed directly from the FNV-1a bucket math in `ProbabilisticSampler::sample`.
1195        // For reference, the trace IDs below hash to these buckets (out of 0x4000 = 16384):
1196        // 0x1234567890ABCDEF -> 1764, 0x0 -> 9301, u64::MAX -> 12365.
1197        struct Case {
1198            trace_id: u64,
1199            rate: f64,
1200            expected_keep: bool,
1201        }
1202
1203        let cases = [
1204            // rate 1.0 keeps every trace (the maximum bucket, 16383, is always below 16384).
1205            Case {
1206                trace_id: 0x1234567890ABCDEF,
1207                rate: 1.0,
1208                expected_keep: true,
1209            },
1210            // rate 0.0 drops every trace (no bucket is below 0).
1211            Case {
1212                trace_id: 0x1234567890ABCDEF,
1213                rate: 0.0,
1214                expected_keep: false,
1215            },
1216            // Same trace ID (bucket 1764) flips from drop to keep as the rate crosses its bucket ratio (~0.108).
1217            Case {
1218                trace_id: 0x1234567890ABCDEF,
1219                rate: 0.10,
1220                expected_keep: false,
1221            },
1222            Case {
1223                trace_id: 0x1234567890ABCDEF,
1224                rate: 0.20,
1225                expected_keep: true,
1226            },
1227            // Trace ID 0 (bucket 9301) straddles rate 0.5 (scaled bucket 8192) vs 0.6 (scaled bucket 9830).
1228            Case {
1229                trace_id: 0,
1230                rate: 0.50,
1231                expected_keep: false,
1232            },
1233            Case {
1234                trace_id: 0,
1235                rate: 0.60,
1236                expected_keep: true,
1237            },
1238            // u64::MAX (bucket 12365) straddles rate 0.5 vs 0.8.
1239            Case {
1240                trace_id: u64::MAX,
1241                rate: 0.50,
1242                expected_keep: false,
1243            },
1244            Case {
1245                trace_id: u64::MAX,
1246                rate: 0.80,
1247                expected_keep: true,
1248            },
1249        ];
1250
1251        for case in cases {
1252            let mut sampler = create_test_sampler();
1253            sampler.sampling_rate = case.rate;
1254            assert_eq!(
1255                sampler.sample_probabilistic(0, case.trace_id),
1256                case.expected_keep,
1257                "trace_id={:#018x} rate={}",
1258                case.trace_id,
1259                case.rate
1260            );
1261        }
1262    }
1263
1264    #[test]
1265    fn probabilistic_sampling_is_deterministic() {
1266        // Determinism is a documented property of `ProbabilisticSampler::sample` (same trace ID + rate always yields
1267        // the same decision). This is intentionally a determinism-only check; correctness is covered by
1268        // `test_probabilistic_sampling_known_decisions`.
1269        let sampler = create_test_sampler();
1270        let trace_id = 0x1234567890ABCDEF_u64;
1271        assert_eq!(
1272            sampler.sample_probabilistic(0, trace_id),
1273            sampler.sample_probabilistic(0, trace_id)
1274        );
1275    }
1276
1277    #[test]
1278    fn probabilistic_sampling_hash_seed_changes_decisions() {
1279        // A seed re-buckets every trace ID: 0x1234567890ABCDEF moves from bucket 1764 with seed 0 to
1280        // bucket 7338 with seed 22, so rates that straddle only one of the buckets flip decisions.
1281        let mut seeded = create_test_sampler();
1282        seeded.probabilistic = probabilistic::ProbabilisticSampler::new(22, false);
1283
1284        // Bucket 7338 (7338/16384 = 0.4479): dropped below 0.45, kept above it.
1285        seeded.sampling_rate = 0.40;
1286        assert!(!seeded.sample_probabilistic(0, 0x1234567890ABCDEF));
1287        seeded.sampling_rate = 0.50;
1288        assert!(seeded.sample_probabilistic(0, 0x1234567890ABCDEF));
1289
1290        // At rate 0.20 the zero-seed sampler keeps this trace (bucket 1764, pinned in
1291        // `probabilistic_sampling_known_decisions`) while the seeded sampler drops it.
1292        seeded.sampling_rate = 0.20;
1293        assert!(!seeded.sample_probabilistic(0, 0x1234567890ABCDEF));
1294        let mut unseeded = create_test_sampler();
1295        unseeded.sampling_rate = 0.20;
1296        assert!(unseeded.sample_probabilistic(0, 0x1234567890ABCDEF));
1297    }
1298
1299    #[test]
1300    fn probabilistic_sampling_full_trace_id_mode_hashes_both_halves() {
1301        let high = 0xAABBCCDD00112233_u64;
1302        let low = 0x1234567890ABCDEF_u64;
1303
1304        // Legacy mode buckets on the low half alone, so the high half cannot move the decision.
1305        let mut legacy = create_test_sampler();
1306        legacy.probabilistic = probabilistic::ProbabilisticSampler::new(22, false);
1307        legacy.sampling_rate = 0.45;
1308        assert_eq!(
1309            legacy.sample_probabilistic(0, low),
1310            legacy.sample_probabilistic(high, low)
1311        );
1312
1313        // Full mode hashes the 16-byte big-endian ID, high half first: (high, low) lands in bucket
1314        // 454 (kept at 0.05, dropped at 0.02), moving the decision relative to legacy mode.
1315        let mut full = create_test_sampler();
1316        full.probabilistic = probabilistic::ProbabilisticSampler::new(22, true);
1317        full.sampling_rate = 0.02;
1318        assert!(!full.sample_probabilistic(high, low));
1319        full.sampling_rate = 0.05;
1320        assert!(full.sample_probabilistic(high, low));
1321
1322        // A zero high half is still a different input than legacy mode: the halves swap byte
1323        // ranges, bucket 14666 versus 7338 (kept at 0.90, dropped at 0.89).
1324        full.sampling_rate = 0.89;
1325        assert!(!full.sample_probabilistic(0, low));
1326        full.sampling_rate = 0.90;
1327        assert!(full.sample_probabilistic(0, low));
1328    }
1329
1330    #[test]
1331    fn from_configuration_reads_the_full_trace_id_feature() {
1332        let traces = domains::traces::Domain {
1333            features: vec![domains::traces::ApmFeature::ProbabilisticSamplerFullTraceId],
1334            ..Default::default()
1335        };
1336        let config = TraceSamplerConfiguration::from_configuration(&traces, &domains::otlp::Traces::default());
1337        assert!(config.probabilistic_full_trace_id);
1338
1339        // Unrecognized flags are carried, not rejected, and enable nothing.
1340        let traces = domains::traces::Domain {
1341            features: vec![domains::traces::ApmFeature::Other("table_names".to_owned())],
1342            ..Default::default()
1343        };
1344        let config = TraceSamplerConfiguration::from_configuration(&traces, &domains::otlp::Traces::default());
1345        assert!(!config.probabilistic_full_trace_id);
1346    }
1347
1348    #[test]
1349    fn every_trace_is_counted_under_exactly_one_sampler() {
1350        // With the probabilistic sampler enabled at rate 1.0, every trace is decided (and kept) by
1351        // it: the sum of seen across all decision keys must equal the number of traces, with no
1352        // double counting across samplers.
1353        let mut sampler = create_test_sampler();
1354
1355        for i in 0..10 {
1356            let span = create_test_span(i, 0);
1357            let mut trace = create_test_trace(vec![span]);
1358            sampler.run_samplers(&mut trace);
1359        }
1360
1361        let decisions = sampler.telemetry.snapshot_decisions();
1362        let total_seen: u64 = decisions.values().map(|counts| counts.seen).sum();
1363        let total_kept: u64 = decisions.values().map(|counts| counts.kept).sum();
1364        assert_eq!(total_seen, 10);
1365        assert_eq!(total_kept, 10);
1366        assert!(decisions
1367            .keys()
1368            .all(|key| key.sampler == telemetry::SamplerName::Probabilistic));
1369        assert!(decisions.keys().all(|key| key.service == "test-service"));
1370    }
1371
1372    #[test]
1373    fn priority_decisions_are_counted_under_the_priority_sampler_with_priority_tag() {
1374        // A trace with a user-set sampling priority is decided by the priority sampler, and its
1375        // decision key carries that priority; no other sampler claims it.
1376        let mut sampler = create_test_sampler();
1377        sampler.probabilistic_sampler_enabled = false;
1378
1379        let mut metrics = HashMap::new();
1380        metrics.insert(SAMPLING_PRIORITY_METRIC_KEY.to_string(), 2.0);
1381        let span = create_test_span_with_metrics(1, metrics);
1382        let mut trace = create_test_trace(vec![span]);
1383        sampler.run_samplers(&mut trace);
1384
1385        let decisions = sampler.telemetry.snapshot_decisions();
1386        assert_eq!(decisions.len(), 1);
1387        let (key, counts) = decisions.iter().next().unwrap();
1388        assert_eq!(key.sampler, telemetry::SamplerName::Priority);
1389        assert_eq!(key.priority, Some(telemetry::Priority::ManualKeep));
1390        assert_eq!(counts.seen, 1);
1391        assert_eq!(counts.kept, 1);
1392    }
1393
1394    #[test]
1395    fn upstream_probability_sampling_reattributes_to_the_probabilistic_sampler() {
1396        // A trace already sampled probabilistically upstream carries the probabilistic decision
1397        // maker; its decision counts under the probabilistic sampler, without the priority
1398        // sampler's sampling_priority or target_env tags.
1399        let mut sampler = create_test_sampler();
1400        sampler.probabilistic_sampler_enabled = false;
1401
1402        let mut metrics = HashMap::new();
1403        metrics.insert(SAMPLING_PRIORITY_METRIC_KEY.to_string(), 2.0);
1404        let mut span = create_test_span_with_metrics(1, metrics);
1405        span.attributes.insert(
1406            MetaString::from_static(TAG_DECISION_MAKER),
1407            AttributeValue::String(MetaString::from_static(DECISION_MAKER_PROBABILISTIC)),
1408        );
1409        let mut trace = create_test_trace(vec![span]);
1410        sampler.run_samplers(&mut trace);
1411
1412        let decisions = sampler.telemetry.snapshot_decisions();
1413        assert_eq!(decisions.len(), 1);
1414        let (key, _) = decisions.iter().next().unwrap();
1415        assert_eq!(key.sampler, telemetry::SamplerName::Probabilistic);
1416        assert_eq!(key.priority, None);
1417    }
1418
1419    #[test]
1420    fn empty_traces_record_no_decision() {
1421        // Unclaimed decisions carry no sampler identity, so they never inflate the unknown series;
1422        // a real attribution gap remains the only thing that could produce one.
1423        let mut sampler = create_test_sampler();
1424        sampler.probabilistic_sampler_enabled = false;
1425
1426        let mut trace = create_test_trace(vec![]);
1427        sampler.run_samplers(&mut trace);
1428        assert!(sampler.telemetry.snapshot_decisions().is_empty());
1429    }
1430
1431    #[test]
1432    fn error_detection() {
1433        let sampler = create_test_sampler();
1434
1435        // Test trace with error field set
1436        let span_with_error = create_test_span(1, 1);
1437        let trace = create_test_trace(vec![span_with_error]);
1438        assert!(sampler.trace_contains_error(&trace, false));
1439
1440        // Test trace without error
1441        let span_without_error = create_test_span(1, 0);
1442        let trace = create_test_trace(vec![span_without_error]);
1443        assert!(!sampler.trace_contains_error(&trace, false));
1444    }
1445
1446    #[test]
1447    fn sampling_priority_order() {
1448        // Test modern path: error sampler overrides probabilistic drop
1449        let mut sampler = create_test_sampler();
1450        sampler.sampling_rate = 0.5; // 50% sampling rate
1451        sampler.probabilistic_sampler_enabled = true;
1452
1453        // Create trace with error that would be dropped by probabilistic
1454        // Using a trace ID that we know will be dropped at 50% rate
1455        let span_with_error = create_test_span(1, 1);
1456        let mut trace = create_test_trace(vec![span_with_error]);
1457        trace.trace_id_low = u64::MAX - 1;
1458
1459        let SamplerOutcome {
1460            keep,
1461            priority,
1462            decision_maker,
1463            ..
1464        } = sampler.run_samplers(&mut trace);
1465        assert!(keep);
1466        assert_eq!(priority, PRIORITY_AUTO_KEEP);
1467        assert_eq!(decision_maker, ""); // Error sampler doesn't set decision_maker
1468
1469        // Test legacy path: user priority is respected
1470        let mut sampler = create_test_sampler();
1471        sampler.probabilistic_sampler_enabled = false; // Use legacy path
1472
1473        let mut metrics = HashMap::new();
1474        metrics.insert(SAMPLING_PRIORITY_METRIC_KEY.to_string(), 2.0);
1475        let span = create_test_span_with_metrics(1, metrics);
1476        let mut trace = create_test_trace(vec![span]);
1477
1478        let SamplerOutcome {
1479            keep,
1480            priority,
1481            decision_maker,
1482            ..
1483        } = sampler.run_samplers(&mut trace);
1484        assert!(keep);
1485        assert_eq!(priority, 2); // UserKeep
1486        assert_eq!(decision_maker, "");
1487    }
1488
1489    #[test]
1490    fn empty_trace_handling() {
1491        let mut sampler = create_test_sampler();
1492        let mut trace = create_test_trace(vec![]);
1493
1494        let SamplerOutcome { keep, priority, .. } = sampler.run_samplers(&mut trace);
1495        assert!(!keep);
1496        assert_eq!(priority, PRIORITY_AUTO_DROP);
1497    }
1498
1499    #[test]
1500    fn root_span_detection() {
1501        let sampler = create_test_sampler();
1502
1503        // Test 1: Root span with parent_id = 0 (common case)
1504        let root_span = DdSpan::new(
1505            MetaString::from("service"),
1506            MetaString::from("operation"),
1507            MetaString::from("resource"),
1508            MetaString::from("type"),
1509            1,
1510            0, // parent_id = 0 indicates root
1511            0,
1512            1000,
1513            0,
1514        );
1515        let child_span = DdSpan::new(
1516            MetaString::from("service"),
1517            MetaString::from("child_op"),
1518            MetaString::from("resource"),
1519            MetaString::from("type"),
1520            2,
1521            1, // parent_id = 1 (points to root)
1522            100,
1523            500,
1524            0,
1525        );
1526        // Put root span second to test that we find it even when not first
1527        let trace = create_test_trace(vec![child_span.clone(), root_span.clone()]);
1528        let root_idx = sampler.get_root_span_index(&trace).unwrap();
1529        assert_eq!(trace.spans()[root_idx].span_id(), 1);
1530
1531        // Test 2: Orphaned span (parent not in trace)
1532        let orphan_span = DdSpan::new(
1533            MetaString::from("service"),
1534            MetaString::from("orphan"),
1535            MetaString::from("resource"),
1536            MetaString::from("type"),
1537            3,
1538            999, // parent_id = 999 (doesn't exist in trace)
1539            200,
1540            300,
1541            0,
1542        );
1543        let trace = create_test_trace(vec![orphan_span]);
1544        let root_idx = sampler.get_root_span_index(&trace).unwrap();
1545        assert_eq!(trace.spans()[root_idx].span_id(), 3);
1546
1547        // Test 3: Multiple root candidates: should return the last one found (index 1)
1548        let span1 = create_test_span(1, 0);
1549        let span2 = create_test_span(2, 0);
1550        let trace = create_test_trace(vec![span1, span2]);
1551        // Both have parent_id = 0, should return the last one found (span_id = 2)
1552        let root_idx = sampler.get_root_span_index(&trace).unwrap();
1553        assert_eq!(trace.spans()[root_idx].span_id(), 2);
1554    }
1555
1556    #[test]
1557    fn single_span_sampling() {
1558        let mut sampler = create_test_sampler();
1559
1560        // Test 1: Trace with SSS tags should be kept even when probabilistic would drop it
1561        sampler.sampling_rate = 0.0; // 0% sampling rate - should drop everything
1562        sampler.probabilistic_sampler_enabled = true;
1563
1564        // Create span with SSS metric
1565        let mut attrs_map = saluki_common::collections::FastHashMap::default();
1566        attrs_map.insert(
1567            MetaString::from(KEY_SPAN_SAMPLING_MECHANISM),
1568            AttributeValue::Float(8.0),
1569        );
1570        let sss_span = create_test_span(1, 0).with_attributes(attrs_map.clone());
1571
1572        // Create regular span without SSS
1573        let regular_span = create_test_span(2, 0);
1574
1575        let mut trace = create_test_trace(vec![sss_span.clone(), regular_span]);
1576
1577        // Apply SSS
1578        let modified = sampler.single_span_sampling(&mut trace);
1579        assert!(modified);
1580        assert_eq!(trace.spans().len(), 1); // Only SSS span kept
1581        assert_eq!(trace.spans()[0].span_id(), 1); // It's the SSS span
1582
1583        // Check that trace has been marked as kept with high priority
1584        assert_eq!(trace.priority, Some(PRIORITY_USER_KEEP));
1585
1586        // Test 2: Trace without SSS tags should not be modified
1587        let trace_without_sss = create_test_trace(vec![create_test_span(3, 0)]);
1588        let mut trace_copy = trace_without_sss.clone();
1589        let modified = sampler.single_span_sampling(&mut trace_copy);
1590        assert!(!modified);
1591        assert_eq!(trace_copy.spans().len(), trace_without_sss.spans().len());
1592    }
1593
1594    #[test]
1595    fn analytics_events() {
1596        let sampler = create_test_sampler();
1597
1598        // Test 1: Trace with analyzed spans
1599        let mut attrs_map = saluki_common::collections::FastHashMap::default();
1600        attrs_map.insert(MetaString::from(KEY_ANALYZED_SPANS), AttributeValue::Float(1.0));
1601        let analyzed_span = create_test_span(1, 0).with_attributes(attrs_map.clone());
1602        let regular_span = create_test_span(2, 0);
1603
1604        let mut trace = create_test_trace(vec![analyzed_span.clone(), regular_span]);
1605
1606        let analyzed_span_ids: Vec<u64> = trace
1607            .spans()
1608            .iter()
1609            .filter(|span| span.attributes.contains_key(KEY_ANALYZED_SPANS))
1610            .map(|span| span.span_id())
1611            .collect();
1612        assert_eq!(analyzed_span_ids, vec![1]);
1613
1614        assert!(sampler.has_analyzed_spans(&trace));
1615        let modified = sampler.analyzed_span_sampling(&mut trace);
1616        assert!(modified);
1617        assert_eq!(trace.spans().len(), 1);
1618        assert_eq!(trace.spans()[0].span_id(), 1);
1619        assert_eq!(trace.priority, Some(PRIORITY_USER_KEEP));
1620
1621        // Test 2: Trace without analyzed spans
1622        let trace_no_analytics = create_test_trace(vec![create_test_span(3, 0)]);
1623        let mut trace_no_analytics_copy = trace_no_analytics.clone();
1624        let analyzed_span_ids: Vec<u64> = trace_no_analytics
1625            .spans()
1626            .iter()
1627            .filter(|span| span.attributes.contains_key(KEY_ANALYZED_SPANS))
1628            .map(|span| span.span_id())
1629            .collect();
1630        assert!(analyzed_span_ids.is_empty());
1631        assert!(!sampler.has_analyzed_spans(&trace_no_analytics));
1632        let modified = sampler.analyzed_span_sampling(&mut trace_no_analytics_copy);
1633        assert!(!modified);
1634        assert_eq!(trace_no_analytics_copy.spans().len(), trace_no_analytics.spans().len());
1635    }
1636
1637    #[test]
1638    fn probabilistic_sampling_with_prob_rate_key() {
1639        let mut sampler = create_test_sampler();
1640        sampler.sampling_rate = 0.75; // 75% sampling rate
1641        sampler.probabilistic_sampler_enabled = true;
1642
1643        // Use a trace ID that we know will be sampled
1644        let trace_id = 12345_u64;
1645        let root_span = DdSpan::new(
1646            MetaString::from("service"),
1647            MetaString::from("operation"),
1648            MetaString::from("resource"),
1649            MetaString::from("type"),
1650            1,
1651            0, // parent_id = 0 indicates root
1652            0,
1653            1000,
1654            0,
1655        );
1656        let mut trace = create_test_trace(vec![root_span]);
1657        trace.trace_id_low = trace_id;
1658
1659        let SamplerOutcome {
1660            keep,
1661            priority,
1662            decision_maker,
1663            root_span_idx,
1664            ..
1665        } = sampler.run_samplers(&mut trace);
1666
1667        if keep && decision_maker == DECISION_MAKER_PROBABILISTIC {
1668            // If sampled probabilistically, check that probRateKey was already added
1669            assert_eq!(priority, PRIORITY_AUTO_KEEP);
1670            assert_eq!(decision_maker, DECISION_MAKER_PROBABILISTIC); // probabilistic sampling marker
1671
1672            // Check that the root span already has the probRateKey (it should have been added in run_samplers)
1673            let root_idx = root_span_idx.unwrap_or(0);
1674            let root_span = &trace.spans()[root_idx];
1675            assert!(root_span.attributes.contains_key(PROB_RATE_KEY));
1676            assert_eq!(
1677                root_span
1678                    .attributes
1679                    .get(PROB_RATE_KEY)
1680                    .and_then(AttributeValue::as_float),
1681                Some(0.75)
1682            );
1683
1684            // Test that apply_sampling_metadata still works correctly for other metadata
1685            let mut trace_with_metadata = trace.clone();
1686            sampler.apply_sampling_metadata(&mut trace_with_metadata, keep, priority, decision_maker, root_idx);
1687
1688            // Check that decision maker tag was added
1689            let modified_root = &trace_with_metadata.spans()[root_idx];
1690            assert!(modified_root.attributes.contains_key(TAG_DECISION_MAKER));
1691            assert_eq!(
1692                modified_root
1693                    .attributes
1694                    .get(TAG_DECISION_MAKER)
1695                    .and_then(AttributeValue::as_string),
1696                Some(&MetaString::from(DECISION_MAKER_PROBABILISTIC))
1697            );
1698        }
1699    }
1700
1701    // ── Rare-sampler interaction tests ──────────────────────────────────────────
1702    // Adapted from datadog-agent/pkg/trace/agent/agent_test.go TestSampling cases:
1703    // "rare-sampler-catch-unsampled", "rare-sampler-catch-sampled",
1704    // "rare-sampler-disabled", and related probabilistic path interactions.
1705
1706    /// Create a top-level span eligible for rare sampling.
1707    ///
1708    /// The rare sampler only considers spans that have `_top_level=1` or `_dd.measured=1`.
1709    /// This helper sets `_top_level=1` so that the rare sampler can consider the span.
1710    fn create_top_level_span(span_id: u64) -> DdSpan {
1711        let mut attrs = saluki_common::collections::FastHashMap::default();
1712        attrs.insert(MetaString::from("_top_level"), AttributeValue::Float(1.0));
1713        create_test_span(span_id, 0).with_attributes(attrs)
1714    }
1715
1716    /// Create a `TraceSampler` with the rare sampler enabled and a very high TPS limit so it
1717    /// freely samples first occurrences, plus a long TTL so second occurrences stay within TTL.
1718    fn create_sampler_with_rare_enabled() -> TraceSampler {
1719        TraceSampler {
1720            rare_sampler: rare_sampler::RareSampler::new(
1721                true,
1722                1000.0,
1723                std::time::Duration::from_secs(300),
1724                200,
1725                telemetry::SamplerCounters::new(),
1726            ),
1727            ..create_test_sampler()
1728        }
1729    }
1730
1731    /// Adapted from Go "rare-sampler-catch-unsampled":
1732    ///
1733    /// Rare is enabled + probabilistic would drop → rare catches it (first occurrence).
1734    #[test]
1735    fn rare_sampler_catches_unsampled_trace() {
1736        let mut sampler = create_sampler_with_rare_enabled();
1737        sampler.sampling_rate = 0.0; // probabilistic drops everything
1738        sampler.probabilistic_sampler_enabled = true;
1739
1740        let span = create_top_level_span(1);
1741        let mut trace = create_test_trace(vec![span]);
1742
1743        let SamplerOutcome {
1744            keep,
1745            priority,
1746            decision_maker,
1747            ..
1748        } = sampler.run_samplers(&mut trace);
1749        assert!(keep, "rare sampler should catch first occurrence");
1750        assert_eq!(priority, PRIORITY_AUTO_KEEP);
1751        assert_eq!(decision_maker, "", "rare sampler does not set _dd.p.dm");
1752    }
1753
1754    /// Adapted from Go "rare-sampler-catch-sampled" (first trace):
1755    ///
1756    /// Rare is enabled, first occurrence—trace is kept and `_dd.rare` is set on the span.
1757    #[test]
1758    fn rare_sampler_sets_rare_metric_on_first_occurrence() {
1759        let mut sampler = create_sampler_with_rare_enabled();
1760        sampler.sampling_rate = 0.0;
1761        sampler.probabilistic_sampler_enabled = true;
1762
1763        let span = create_top_level_span(1);
1764        let mut trace = create_test_trace(vec![span]);
1765
1766        let SamplerOutcome {
1767            keep,
1768            root_span_idx: root_idx,
1769            ..
1770        } = sampler.run_samplers(&mut trace);
1771        assert!(keep);
1772        let root = &trace.spans()[root_idx.unwrap()];
1773        assert_eq!(
1774            root.attributes
1775                .get(rare_sampler::RARE_KEY)
1776                .and_then(AttributeValue::as_float),
1777            Some(1.0),
1778            "_dd.rare should be 1 on first occurrence"
1779        );
1780    }
1781
1782    /// Adapted from Go "rare-sampler-catch-sampled" (second trace same signature):
1783    ///
1784    /// Within the TTL, the same signature is no longer "rare" and rare doesn't re-sample it.
1785    /// With probabilistic at 0%, the trace should be dropped.
1786    #[test]
1787    fn rare_sampler_does_not_resample_within_ttl() {
1788        let mut sampler = create_sampler_with_rare_enabled();
1789        sampler.sampling_rate = 0.0;
1790        sampler.probabilistic_sampler_enabled = true;
1791
1792        // First trace: rare catches it.
1793        let span1 = create_top_level_span(1);
1794        let mut trace1 = create_test_trace(vec![span1]);
1795        let SamplerOutcome { keep: keep1, .. } = sampler.run_samplers(&mut trace1);
1796        assert!(keep1, "first occurrence should be kept by rare sampler");
1797
1798        // Second trace: same signature (same service/operation/resource on the top-level span),
1799        // still within TTL → rare won't catch it; probabilistic at 0% drops it.
1800        let span2 = create_top_level_span(2);
1801        let mut trace2 = create_test_trace(vec![span2]);
1802        let SamplerOutcome {
1803            keep: keep2,
1804            priority: priority2,
1805            ..
1806        } = sampler.run_samplers(&mut trace2);
1807        assert!(!keep2, "second occurrence within TTL should be dropped");
1808        assert_eq!(priority2, PRIORITY_AUTO_DROP);
1809    }
1810
1811    /// Adapted from Go "rare-sampler-disabled":
1812    ///
1813    /// Rare is disabled + probabilistic at 0% → trace is dropped.
1814    #[test]
1815    fn rare_sampler_disabled_does_not_catch_unsampled() {
1816        let mut sampler = create_test_sampler(); // rare disabled by default
1817        sampler.sampling_rate = 0.0;
1818        sampler.probabilistic_sampler_enabled = true;
1819
1820        let span = create_top_level_span(1);
1821        let mut trace = create_test_trace(vec![span]);
1822
1823        let SamplerOutcome { keep, priority, .. } = sampler.run_samplers(&mut trace);
1824        assert!(!keep, "rare disabled should not catch the trace");
1825        assert_eq!(priority, PRIORITY_AUTO_DROP);
1826    }
1827
1828    /// Rare + non-probabilistic path (priority path): rare catches `PriorityAutoDrop` on first
1829    /// occurrence, preserving the tracer-set priority rather than upgrading to AutoKeep.
1830    #[test]
1831    fn rare_sampler_catches_priority_auto_drop_in_legacy_path() {
1832        let mut sampler = create_sampler_with_rare_enabled();
1833        sampler.probabilistic_sampler_enabled = false;
1834
1835        let mut attrs = saluki_common::collections::FastHashMap::default();
1836        attrs.insert(MetaString::from("_top_level"), AttributeValue::Float(1.0));
1837        attrs.insert(
1838            MetaString::from(SAMPLING_PRIORITY_METRIC_KEY),
1839            AttributeValue::Float(PRIORITY_AUTO_DROP as f64),
1840        );
1841        let span = create_test_span(1, 0).with_attributes(attrs);
1842        let mut trace = create_test_trace(vec![span]);
1843
1844        let SamplerOutcome {
1845            keep,
1846            priority,
1847            decision_maker,
1848            ..
1849        } = sampler.run_samplers(&mut trace);
1850        assert!(keep, "rare sampler should catch PriorityAutoDrop on first occurrence");
1851        assert_eq!(priority, PRIORITY_AUTO_DROP, "tracer-set priority should be preserved");
1852        assert_eq!(decision_maker, "");
1853    }
1854
1855    /// Rare + non-probabilistic path (priority path): UserKeep priority is preserved, not
1856    /// downgraded to AutoKeep. Mirrors Go agent behavior at agent.go#L1129-1131.
1857    #[test]
1858    fn rare_sampler_preserves_user_keep_priority_in_legacy_path() {
1859        let mut sampler = create_sampler_with_rare_enabled();
1860        sampler.probabilistic_sampler_enabled = false;
1861
1862        let mut attrs = saluki_common::collections::FastHashMap::default();
1863        attrs.insert(MetaString::from("_top_level"), AttributeValue::Float(1.0));
1864        attrs.insert(
1865            MetaString::from(SAMPLING_PRIORITY_METRIC_KEY),
1866            AttributeValue::Float(2.0),
1867        ); // UserKeep
1868        let span = create_test_span(1, 0).with_attributes(attrs);
1869        let mut trace = create_test_trace(vec![span]);
1870
1871        let SamplerOutcome { keep, priority, .. } = sampler.run_samplers(&mut trace);
1872        assert!(keep);
1873        assert_eq!(priority, 2, "UserKeep priority must not be downgraded to AutoKeep");
1874    }
1875
1876    /// Probabilistic path with 100% rate and rare disabled: keep with `_dd.p.dm = "-9"`.
1877    #[test]
1878    fn probabilistic_100_percent_keeps_trace_with_decision_maker() {
1879        let mut sampler = create_test_sampler(); // rare disabled
1880        sampler.sampling_rate = 1.0;
1881        sampler.probabilistic_sampler_enabled = true;
1882
1883        let span = create_top_level_span(1);
1884        let mut trace = create_test_trace(vec![span]);
1885
1886        let SamplerOutcome {
1887            keep,
1888            priority,
1889            decision_maker,
1890            ..
1891        } = sampler.run_samplers(&mut trace);
1892        assert!(keep);
1893        assert_eq!(priority, PRIORITY_AUTO_KEEP);
1894        assert_eq!(decision_maker, DECISION_MAKER_PROBABILISTIC);
1895    }
1896
1897    /// Probabilistic path with 0% rate and rare disabled: drop.
1898    #[test]
1899    fn probabilistic_0_percent_drops_trace() {
1900        let mut sampler = create_test_sampler(); // rare disabled
1901        sampler.sampling_rate = 0.0;
1902        sampler.probabilistic_sampler_enabled = true;
1903        sampler.error_sampling_enabled = false;
1904
1905        let span = create_top_level_span(1);
1906        let mut trace = create_test_trace(vec![span]);
1907
1908        let SamplerOutcome { keep, priority, .. } = sampler.run_samplers(&mut trace);
1909        assert!(!keep);
1910        assert_eq!(priority, PRIORITY_AUTO_DROP);
1911    }
1912
1913    /// Rare sampler should catch OTLP traces without a sampling priority on their first occurrence,
1914    /// matching the Go agent behavior: https://github.com/DataDog/datadog-agent/blob/main/pkg/trace/agent/agent.go#L1129-L1140
1915    #[test]
1916    fn otlp_rate_decisions_record_under_the_probabilistic_sampler() {
1917        // Keeps and drops from the OTLP rate decision both belong to the probabilistic sampler;
1918        // the drop must not fall through to the no-priority bucket.
1919        let mut sampler = create_test_sampler();
1920        sampler.probabilistic_sampler_enabled = false;
1921        sampler.error_sampling_enabled = false;
1922        sampler.otlp_sampling_rate = 0.0;
1923
1924        let mut span = create_top_level_span(1);
1925        span.attributes.insert(
1926            MetaString::from_static(OTEL_TRACE_ID_META_KEY),
1927            AttributeValue::String(MetaString::from("00000000000000000000000000000001")),
1928        );
1929        let mut trace = create_test_trace(vec![span]);
1930
1931        let SamplerOutcome { keep, priority, .. } = sampler.run_samplers(&mut trace);
1932        assert!(!keep);
1933        assert_eq!(priority, PRIORITY_AUTO_DROP);
1934
1935        let decisions = sampler.telemetry.snapshot_decisions();
1936        let (key, counts) = decisions.iter().next().unwrap();
1937        assert_eq!(key.sampler, telemetry::SamplerName::Probabilistic);
1938        assert_eq!(counts.seen, 1);
1939        assert_eq!(counts.kept, 0);
1940    }
1941
1942    #[test]
1943    fn rare_sampler_catches_otlp_no_priority_trace() {
1944        let mut sampler = create_sampler_with_rare_enabled();
1945        sampler.probabilistic_sampler_enabled = false;
1946        sampler.error_sampling_enabled = false;
1947        sampler.otlp_sampling_rate = 0.0;
1948
1949        let mut span = create_top_level_span(1);
1950        span.attributes.insert(
1951            MetaString::from_static(OTEL_TRACE_ID_META_KEY),
1952            AttributeValue::String(MetaString::from("00000000000000000000000000000001")),
1953        );
1954        let mut trace = create_test_trace(vec![span]);
1955
1956        let SamplerOutcome {
1957            keep,
1958            priority,
1959            decision_maker,
1960            root_span_idx: root_idx,
1961            ..
1962        } = sampler.run_samplers(&mut trace);
1963        assert!(
1964            keep,
1965            "rare sampler should keep OTLP trace with no priority on first occurrence"
1966        );
1967        assert_eq!(priority, PRIORITY_AUTO_KEEP);
1968        assert_eq!(decision_maker, "");
1969        assert_eq!(
1970            trace.spans()[root_idx.unwrap()]
1971                .attributes
1972                .get(rare_sampler::RARE_KEY)
1973                .and_then(AttributeValue::as_float),
1974            Some(1.0),
1975            "_dd.rare should be set to 1 on first occurrence"
1976        );
1977    }
1978
1979    /// Adapted from Go "probabilistic-rare-100":
1980    ///
1981    /// Rare fires before probabilistic is consulted, so even at 100% sampling rate the decision
1982    /// maker tag isn't set—the trace is attributed to rare, not probabilistic.
1983    #[test]
1984    fn rare_wins_over_probabilistic_no_decision_maker_tag() {
1985        let mut sampler = create_sampler_with_rare_enabled();
1986        sampler.sampling_rate = 1.0;
1987        sampler.probabilistic_sampler_enabled = true;
1988
1989        let span = create_top_level_span(1);
1990        let mut trace = create_test_trace(vec![span]);
1991
1992        let SamplerOutcome {
1993            keep,
1994            priority,
1995            decision_maker,
1996            ..
1997        } = sampler.run_samplers(&mut trace);
1998        assert!(keep);
1999        assert_eq!(priority, PRIORITY_AUTO_KEEP);
2000        assert_eq!(decision_maker, "", "rare takes precedence—_dd.p.dm must not be set");
2001    }
2002
2003    /// Adapted from Go "error-sampled-prio-unsampled":
2004    ///
2005    /// Rare fires before the error sampler is reached. A no-priority error trace on its first
2006    /// occurrence is kept by rare, not by the error sampler.
2007    #[test]
2008    fn rare_catches_error_trace_before_error_sampler() {
2009        let mut sampler = create_sampler_with_rare_enabled();
2010        sampler.probabilistic_sampler_enabled = false;
2011        sampler.error_sampling_enabled = true;
2012
2013        let span = create_top_level_span(1);
2014        let error_span = create_test_span(2, 1); // error=1
2015        let mut trace = create_test_trace(vec![span, error_span]);
2016
2017        let SamplerOutcome {
2018            keep,
2019            priority,
2020            decision_maker,
2021            ..
2022        } = sampler.run_samplers(&mut trace);
2023        assert!(keep, "rare should catch the trace before the error sampler");
2024        assert_eq!(priority, PRIORITY_AUTO_KEEP);
2025        assert_eq!(decision_maker, "");
2026    }
2027
2028    /// Adapted from Go manual-drop short-circuit behavior:
2029    ///
2030    /// UserDrop (-1) priority is checked before rare runs in the priority path. A UserDrop trace
2031    /// must be dropped even when rare is enabled and would otherwise match.
2032    #[test]
2033    fn manual_drop_short_circuits_before_rare() {
2034        let mut sampler = create_sampler_with_rare_enabled();
2035        sampler.probabilistic_sampler_enabled = false;
2036
2037        let mut attrs = saluki_common::collections::FastHashMap::default();
2038        attrs.insert(MetaString::from("_top_level"), AttributeValue::Float(1.0));
2039        attrs.insert(
2040            MetaString::from(SAMPLING_PRIORITY_METRIC_KEY),
2041            AttributeValue::Float(-1.0),
2042        ); // UserDrop
2043        let span = create_test_span(1, 0).with_attributes(attrs);
2044        let mut trace = create_test_trace(vec![span]);
2045
2046        let SamplerOutcome { keep, priority, .. } = sampler.run_samplers(&mut trace);
2047        assert!(!keep, "UserDrop must be dropped even when rare would match");
2048        assert_eq!(priority, -1);
2049    }
2050
2051    // ── Error Tracking Standalone tests ─────────────────────────────────────────
2052    // Adapted from datadog-agent/pkg/trace/agent/agent.go runSamplers ETS block.
2053
2054    fn create_sampler_with_ets() -> TraceSampler {
2055        TraceSampler {
2056            error_tracking_standalone: true,
2057            ..create_test_sampler()
2058        }
2059    }
2060
2061    /// ETS enabled + trace with error → kept by error sampler.
2062    #[test]
2063    fn ets_keeps_trace_with_error() {
2064        let mut sampler = create_sampler_with_ets();
2065
2066        let span = create_test_span(1, 1); // error=1
2067        let mut trace = create_test_trace(vec![span]);
2068
2069        let SamplerOutcome {
2070            keep,
2071            priority,
2072            decision_maker,
2073            ..
2074        } = sampler.run_samplers(&mut trace);
2075        assert!(keep, "ETS should keep traces with errors");
2076        assert_eq!(priority, PRIORITY_AUTO_KEEP);
2077        assert_eq!(decision_maker, "", "ETS does not set a decision maker");
2078    }
2079
2080    /// ETS enabled + trace without error → dropped; rare/probabilistic/priority not consulted.
2081    #[test]
2082    fn ets_drops_trace_without_error() {
2083        let mut sampler = create_sampler_with_ets();
2084
2085        let span = create_test_span(1, 0); // error=0
2086        let mut trace = create_test_trace(vec![span]);
2087
2088        let SamplerOutcome { keep, priority, .. } = sampler.run_samplers(&mut trace);
2089        assert!(!keep, "ETS should drop traces without errors");
2090        assert_eq!(priority, PRIORITY_AUTO_DROP);
2091    }
2092
2093    /// ETS enabled + non-error trace → forwarded with DroppedTrace=true; SSS/analytics suppressed.
2094    #[test]
2095    fn ets_forwards_dropped_trace_with_dropped_flag() {
2096        let mut sampler = create_sampler_with_ets();
2097
2098        // Span with SSS metric — would trigger single span sampling in non-ETS mode.
2099        let mut attrs = saluki_common::collections::FastHashMap::default();
2100        attrs.insert(
2101            MetaString::from(KEY_SPAN_SAMPLING_MECHANISM),
2102            AttributeValue::Float(8.0),
2103        );
2104        let span = create_test_span(1, 0).with_attributes(attrs);
2105        let mut trace = create_test_trace(vec![span]);
2106
2107        let forwarded = sampler.process_trace(&mut trace);
2108        assert!(forwarded, "ETS should forward non-error traces to intake");
2109        assert!(trace.dropped_trace, "non-error ETS trace should have DroppedTrace=true");
2110    }
2111
2112    /// ETS enabled + trace with exception span event → kept (exception events count as errors in ETS).
2113    #[test]
2114    fn ets_keeps_trace_with_exception_span_event() {
2115        let mut sampler = create_sampler_with_ets();
2116
2117        // Span with error=0 but exception span event metadata.
2118        let mut attrs = saluki_common::collections::FastHashMap::default();
2119        attrs.insert(
2120            MetaString::from("_dd.span_events.has_exception"),
2121            AttributeValue::String(MetaString::from("true")),
2122        );
2123        let span = create_test_span(1, 0).with_attributes(attrs);
2124        let mut trace = create_test_trace(vec![span]);
2125
2126        let SamplerOutcome { keep, .. } = sampler.run_samplers(&mut trace);
2127        assert!(keep, "ETS should treat exception span events as errors");
2128    }
2129
2130    /// ETS disabled → normal sampling path (probabilistic) is used.
2131    #[test]
2132    fn ets_disabled_uses_normal_sampling() {
2133        let mut sampler = create_test_sampler(); // ETS disabled
2134        sampler.sampling_rate = 1.0;
2135        sampler.probabilistic_sampler_enabled = true;
2136
2137        let span = create_test_span(1, 0); // no error
2138        let mut trace = create_test_trace(vec![span]);
2139
2140        let SamplerOutcome {
2141            keep, decision_maker, ..
2142        } = sampler.run_samplers(&mut trace);
2143        assert!(keep, "normal probabilistic sampling should keep the trace");
2144        assert_eq!(decision_maker, DECISION_MAKER_PROBABILISTIC);
2145    }
2146
2147    // ── ETS + OTLP pre-sampling tests ────────────────────────────────────────────
2148    // These mirror DDA's OTLPReceiver.createChunks behavior which pre-assigns
2149    // priority/dm before runSamplersV1, so ETS sees those values even when it
2150    // short-circuits. See: pkg/trace/api/otlp.go#L561-L585.
2151
2152    fn create_otlp_test_span(span_id: u64, error: i32) -> DdSpan {
2153        let mut attrs = saluki_common::collections::FastHashMap::default();
2154        attrs.insert(
2155            MetaString::from_static(OTEL_TRACE_ID_META_KEY),
2156            AttributeValue::String(MetaString::from("0000000000000000deadbeefcafebabe")),
2157        );
2158        create_test_span(span_id, error).with_attributes(attrs)
2159    }
2160
2161    fn create_sampler_with_ets_legacy() -> TraceSampler {
2162        TraceSampler {
2163            error_tracking_standalone: true,
2164            probabilistic_sampler_enabled: false,
2165            otlp_sampling_rate: 1.0,
2166            ..create_test_sampler()
2167        }
2168    }
2169
2170    /// ETS + OTLP non-error trace (legacy sampler path): pre-sampling sets priority=AutoKeep and dm=-9.
2171    /// Mirrors DDA `OTLPReceiver` assigning priority=1 + dm=-9 before ETS returns early.
2172    #[test]
2173    fn ets_otlp_non_error_gets_presample_priority_and_dm() {
2174        let mut sampler = create_sampler_with_ets_legacy();
2175
2176        let span = create_otlp_test_span(1, 0); // no error
2177        let mut trace = create_test_trace(vec![span]);
2178
2179        let SamplerOutcome {
2180            keep,
2181            priority,
2182            decision_maker: dm,
2183            ..
2184        } = sampler.run_samplers(&mut trace);
2185        assert!(!keep, "ETS should drop non-error OTLP traces");
2186        assert_eq!(
2187            priority, PRIORITY_AUTO_KEEP,
2188            "OTLP pre-sampling sets priority=AutoKeep even for ETS-dropped traces"
2189        );
2190        assert_eq!(dm, DECISION_MAKER_PROBABILISTIC, "OTLP pre-sampling sets dm=-9");
2191    }
2192
2193    /// ETS + OTLP error trace (legacy sampler path): pre-sampling sets priority=AutoKeep and dm=-9.
2194    #[test]
2195    fn ets_otlp_error_gets_presample_priority_and_dm() {
2196        let mut sampler = create_sampler_with_ets_legacy();
2197
2198        let span = create_otlp_test_span(1, 1); // error=1
2199        let mut trace = create_test_trace(vec![span]);
2200
2201        let SamplerOutcome {
2202            keep,
2203            priority,
2204            decision_maker: dm,
2205            ..
2206        } = sampler.run_samplers(&mut trace);
2207        assert!(keep, "ETS should keep error OTLP traces");
2208        assert_eq!(priority, PRIORITY_AUTO_KEEP, "OTLP pre-sampling sets priority=AutoKeep");
2209        assert_eq!(dm, DECISION_MAKER_PROBABILISTIC, "OTLP pre-sampling sets dm=-9");
2210    }
2211
2212    /// ETS + OTLP + probabilistic_sampler_enabled=true: `OTLPReceiver` defers, no pre-sampling.
2213    /// DDA's `OTLPReceiver` sets PriorityNone and skips when ProbabilisticSamplerEnabled.
2214    #[test]
2215    fn ets_otlp_probabilistic_path_skips_presample() {
2216        let mut sampler = create_sampler_with_ets_legacy();
2217        sampler.probabilistic_sampler_enabled = true; // override to prob path
2218
2219        let span = create_otlp_test_span(1, 0); // no error
2220        let mut trace = create_test_trace(vec![span]);
2221
2222        let SamplerOutcome {
2223            keep,
2224            priority,
2225            decision_maker: dm,
2226            ..
2227        } = sampler.run_samplers(&mut trace);
2228        assert!(!keep, "ETS should drop non-error traces");
2229        assert_eq!(
2230            priority, PRIORITY_AUTO_DROP,
2231            "no pre-sampling when probabilistic path active"
2232        );
2233        assert_eq!(dm, "", "no dm when probabilistic path active");
2234    }
2235
2236    /// ETS + non-OTLP trace (legacy sampler path): behavior unchanged—no pre-sampling.
2237    #[test]
2238    fn ets_non_otlp_unaffected_by_presample() {
2239        let mut sampler = create_sampler_with_ets_legacy();
2240
2241        let span = create_test_span(1, 0); // no error, no OTLP meta
2242        let mut trace = create_test_trace(vec![span]);
2243
2244        let SamplerOutcome {
2245            keep,
2246            priority,
2247            decision_maker: dm,
2248            ..
2249        } = sampler.run_samplers(&mut trace);
2250        assert!(!keep, "ETS should drop non-error non-OTLP traces");
2251        assert_eq!(priority, PRIORITY_AUTO_DROP, "non-OTLP traces use default ETS priority");
2252        assert_eq!(dm, "", "non-OTLP traces get no dm");
2253    }
2254
2255    /// ETS + OTLP trace with user-set priority: dm="-4" (manual sampling), matching DDA.
2256    #[test]
2257    fn ets_otlp_user_priority_gets_manual_dm() {
2258        let mut sampler = create_sampler_with_ets_legacy();
2259
2260        let mut span = create_otlp_test_span(1, 0);
2261        span.attributes.insert(
2262            MetaString::from(SAMPLING_PRIORITY_METRIC_KEY),
2263            AttributeValue::Float(2.0),
2264        ); // UserKeep
2265        let mut trace = create_test_trace(vec![span]);
2266
2267        let SamplerOutcome {
2268            keep,
2269            priority,
2270            decision_maker: dm,
2271            ..
2272        } = sampler.run_samplers(&mut trace);
2273        assert!(!keep, "ETS drops non-error traces regardless of user priority");
2274        assert_eq!(priority, PRIORITY_USER_KEEP, "user priority is preserved");
2275        assert_eq!(dm, DECISION_MAKER_MANUAL, "user-set priority gets dm=-4");
2276    }
2277}