saluki_components/sources/dogstatsd/
origin.rs

1use std::sync::Arc;
2
3use saluki_core::data_model::{
4    origin::{OriginTagCardinality, OriginTagsResolver, RawOrigin},
5    tags::SharedTagSet,
6};
7use saluki_env::{
8    workload::{origin::ResolvedOrigin, EntityId},
9    WorkloadProvider,
10};
11use saluki_io::deser::codec::dogstatsd::{EventPacket, MetricPacket, ServiceCheckPacket};
12use tracing::trace;
13
14use super::{replay::CapturedTaggerHandle, tags::WellKnownTags};
15
16const REPLAY_PROCESS_ID_MARKER: u32 = 1u32 << 31;
17
18#[derive(Clone, Debug, Eq, PartialEq)]
19pub(super) enum ProcessOrigin {
20    Pinned(Option<EntityId>),
21    Unpinned(u32),
22    Replay(u32),
23}
24
25impl ProcessOrigin {
26    pub(super) fn container_entity_id(&self) -> Option<&EntityId> {
27        match self {
28            Self::Pinned(entity_id) => entity_id.as_ref(),
29            Self::Unpinned(_) | Self::Replay(_) => None,
30        }
31    }
32}
33
34pub(super) fn mark_replay_process_id(process_id: u32) -> u32 {
35    process_id | REPLAY_PROCESS_ID_MARKER
36}
37
38fn captured_process_id_from_replay(process_id: u32) -> Option<u32> {
39    if process_id & REPLAY_PROCESS_ID_MARKER != 0 {
40        Some(process_id & !REPLAY_PROCESS_ID_MARKER)
41    } else {
42        None
43    }
44}
45
46/// Origin enrichment configuration.
47///
48/// Origin enrichment controls the when and how of enriching metrics ingested via DogStatsD based on various sources of
49/// "origin" information, such as specific metric tags or UDS socket credentials. Enrichment involves adding additional
50/// metric tags that describe the origin of the metric, such as the Kubernetes pod or container.
51#[derive(Clone)]
52#[cfg_attr(test, derive(Debug, PartialEq))]
53pub struct OriginEnrichmentConfiguration {
54    /// Whether or not to enable origin detection.
55    ///
56    /// If disabled, no origin tags will be added to events even if the origin information is detected.
57    pub enabled: bool,
58
59    /// Whether or not a client-provided entity ID should take precedence over automatically detected origin metadata.
60    ///
61    /// When a client-provided entity ID is specified, and an origin process ID has automatically been detected, setting
62    /// this to `true` will cause the origin process ID to be ignored. This only applies when
63    /// `origin_detection_unified` is `false`.
64    pub entity_id_precedence: bool,
65
66    /// The cardinality of tags to enrich metrics with, when the payload doesn't specify one itself.
67    pub tag_cardinality: OriginTagCardinality,
68
69    /// Whether or not to use the unified origin detection behavior.
70    ///
71    /// When set to `true`, all detected entity IDs -- UDS Origin Detection, `dd.internal.entity_id`, container ID from
72    /// DogStatsD payload -- will be used for querying tags to enrich with. When set to `false`, the original precedence
73    /// behavior will be used, which enriches with the entity ID detected via Origin Detection first, and then
74    /// potentially again with either the client-provided entity ID (`dd.internal.entity_id`) or the container ID from
75    /// the DogStatsD payload, with the client-provided entity ID taking precedence. An entity ID detected via Origin
76    /// Detection is only used when no client-provided entity ID was present, or when `entity_id_precedence` is
77    /// `false`.
78    pub origin_detection_unified: bool,
79
80    /// Whether or not to opt out of origin detection for DogStatsD metrics.
81    ///
82    /// When set to `true`, and the metric explicitly denotes a cardinality of `"none"`, origin enrichment will be
83    /// skipped. This is only applicable to DogStatsD metrics when unified origin detection behavior isn't enabled.
84    pub origin_detection_optout: bool,
85
86    /// Whether or not to parse client-provided origin fields from DogStatsD payloads.
87    ///
88    /// When enabled, the `c:` (Local Data), `e:` (External Data), and `card:` (Cardinality) protocol fields are
89    /// parsed and used for origin enrichment.
90    pub origin_detection_client: bool,
91}
92
93#[cfg(test)]
94impl OriginEnrichmentConfiguration {
95    /// Creates a fixture configuration with origin detection disabled, for tests that exercise enrichment behavior
96    /// rather than configuration.
97    pub(super) fn for_test() -> Self {
98        Self {
99            enabled: false,
100            entity_id_precedence: false,
101            tag_cardinality: OriginTagCardinality::Low,
102            origin_detection_unified: false,
103            origin_detection_optout: true,
104            origin_detection_client: false,
105        }
106    }
107}
108
109#[derive(Clone)]
110pub(super) struct DogStatsDOriginTagResolver {
111    config: OriginEnrichmentConfiguration,
112    workload_provider: Arc<dyn WorkloadProvider + Send + Sync>,
113    captured_tagger: CapturedTaggerHandle,
114}
115
116impl DogStatsDOriginTagResolver {
117    pub fn new(
118        config: OriginEnrichmentConfiguration, workload_provider: Arc<dyn WorkloadProvider + Send + Sync>,
119        captured_tagger: CapturedTaggerHandle,
120    ) -> Self {
121        Self {
122            config,
123            workload_provider,
124            captured_tagger,
125        }
126    }
127
128    fn collect_origin_tags_with_process_entity(
129        &self, origin: &ResolvedOrigin, process_entity_id: Option<&EntityId>,
130    ) -> SharedTagSet {
131        let mut collected_tags = SharedTagSet::default();
132
133        if !self.config.enabled {
134            return collected_tags;
135        }
136
137        // Examine the various possible entity ID values, and based on their state, use one or more of them to grab any
138        // enriched tags attached to the entities. We evalulate a number of possible entity IDs:
139        //
140        // - sender entity (resolved from UDS socket credentials when the packet is received)
141        // - Local Data-based container ID (extracted from special "container ID" extension in DogStatsD protocol; also known as Local Data)
142        // - Local Data-based pod UID (extracted from `dd.internal.entity_id` tag)
143        // - External Data-based container ID (derived from pod UID and container name in External Data)
144        // - External Data-based pod UID (raw pod UID from External Data)
145        let maybe_process_entity_id = process_entity_id;
146        let maybe_local_container_id = origin.local_data();
147        let maybe_local_pod_uid = origin.pod_uid();
148        let maybe_external_container_id = origin.resolved_external_data().map(|red| red.container_entity_id());
149        let maybe_external_pod_uid = origin.resolved_external_data().map(|red| red.pod_entity_id());
150
151        let tag_cardinality = origin.cardinality().unwrap_or(self.config.tag_cardinality);
152
153        if !self.config.origin_detection_unified {
154            if self.config.origin_detection_optout && tag_cardinality == OriginTagCardinality::None {
155                trace!("Skipping origin enrichment for DogStatsD metric with cardinality 'none'.");
156                return collected_tags;
157            }
158
159            // If we discovered an entity ID from socket credentials, and no client-provided entity ID was provided (or
160            // it was, but entity ID precedence is disabled), then try to get tags for the detected entity ID.
161            if let Some(entity_id) = maybe_process_entity_id {
162                if maybe_local_pod_uid.is_none() || !self.config.entity_id_precedence {
163                    if let Some(tags) = self.workload_provider.get_tags_for_entity(entity_id, tag_cardinality) {
164                        collected_tags.extend_from_shared(&tags);
165                    } else {
166                        trace!(
167                            ?entity_id,
168                            cardinality = tag_cardinality.as_str(),
169                            "No tags found for entity."
170                        );
171                    }
172                }
173            }
174
175            // If we have a client-provided entity ID, try to get tags for the entity based on those. A client-provided
176            // entity ID takes precedence over the container ID.
177            let maybe_entity_id = maybe_local_pod_uid.or(maybe_local_container_id);
178            if let Some(entity_id) = maybe_entity_id {
179                if let Some(tags) = self.workload_provider.get_tags_for_entity(entity_id, tag_cardinality) {
180                    collected_tags.extend_from_shared(&tags);
181                } else {
182                    trace!(
183                        ?entity_id,
184                        cardinality = tag_cardinality.as_str(),
185                        "No tags found for entity."
186                    );
187                }
188            }
189        } else {
190            if tag_cardinality == OriginTagCardinality::None {
191                trace!("Skipping origin enrichment for metric with cardinality 'none'.");
192                return collected_tags;
193            }
194
195            // Evaluate all available entity IDs in order of priority: Local Data-based container ID, sender entity,
196            // External Data-based container ID, Local Data-based pod UID, and External Data-based pod UID.
197            //
198            // As soon as the first set of tags for an entity ID is found, we skip the remaining entity IDs.
199            let maybe_entity_ids = &[
200                maybe_local_container_id,
201                maybe_process_entity_id,
202                maybe_external_container_id,
203                maybe_local_pod_uid,
204                maybe_external_pod_uid,
205            ];
206            for entity_id in maybe_entity_ids.iter().flatten() {
207                if let Some(tags) = self.workload_provider.get_tags_for_entity(entity_id, tag_cardinality) {
208                    if !tags.is_empty() {
209                        collected_tags.extend_from_shared(&tags);
210                        break;
211                    }
212                } else {
213                    trace!(
214                        ?entity_id,
215                        cardinality = tag_cardinality.as_str(),
216                        "No tags found for entity."
217                    );
218                }
219            }
220        }
221
222        collected_tags
223    }
224
225    fn collect_origin_tags(&self, origin: &ResolvedOrigin) -> SharedTagSet {
226        self.collect_origin_tags_with_process_entity(origin, origin.process_id())
227    }
228
229    pub(super) fn resolve_origin_tags_with_process_origin(
230        &self, mut origin: RawOrigin<'_>, process_origin: Option<&ProcessOrigin>,
231    ) -> SharedTagSet {
232        match process_origin {
233            Some(ProcessOrigin::Replay(process_id)) => {
234                origin.set_process_id(mark_replay_process_id(*process_id));
235                self.resolve_origin_tags(origin)
236            }
237            Some(ProcessOrigin::Unpinned(process_id)) => {
238                origin.set_process_id(*process_id);
239                self.resolve_origin_tags(origin)
240            }
241            Some(ProcessOrigin::Pinned(process_entity_id)) => {
242                let resolved_origin = self
243                    .workload_provider
244                    .get_resolved_origin(origin.clone())
245                    .unwrap_or_else(|| ResolvedOrigin::from_parts(origin.cardinality(), None, None, None, None));
246                self.collect_origin_tags_with_process_entity(&resolved_origin, process_entity_id.as_ref())
247            }
248            None => self.resolve_origin_tags(origin),
249        }
250    }
251}
252
253impl OriginTagsResolver for DogStatsDOriginTagResolver {
254    fn resolve_origin_tags(&self, origin: RawOrigin<'_>) -> SharedTagSet {
255        // Replay traffic is tagged by setting the marker bit on the origin process ID. It bypasses the live enrichment
256        // pipeline entirely: the captured `TaggerState` already contains fully-resolved tags per entity, so the
257        // resolver's entity-ID-walking logic has nothing to add.
258        if let Some(captured_process_id) = origin.process_id().and_then(captured_process_id_from_replay) {
259            if let Some(store) = self.captured_tagger.current() {
260                let cardinality = origin.cardinality().unwrap_or(self.config.tag_cardinality);
261                return store
262                    .lookup(captured_process_id as i32, cardinality)
263                    .unwrap_or_default();
264            }
265            trace!(?origin, "Replay-flagged origin but no captured tagger available.");
266            return SharedTagSet::default();
267        }
268
269        match self.workload_provider.get_resolved_origin(origin.clone()) {
270            Some(resolved_origin) => self.collect_origin_tags(&resolved_origin),
271            None => {
272                trace!(?origin, "No resolved origin found for origin.");
273                SharedTagSet::default()
274            }
275        }
276    }
277}
278
279/// Builds an `RawOrigin` object from the given metric packet.
280pub fn origin_from_metric_packet<'packet, 'tags>(
281    packet: &'tags MetricPacket<'packet>, well_known_tags: &'tags WellKnownTags<'tags>,
282) -> RawOrigin<'tags>
283where
284    'packet: 'tags,
285{
286    let cardinality = packet.cardinality.or(well_known_tags.cardinality);
287
288    let mut origin = RawOrigin::default();
289    origin.set_pod_uid(well_known_tags.pod_uid);
290    origin.set_local_data(packet.local_data);
291    origin.set_external_data(packet.external_data);
292    origin.set_cardinality(cardinality);
293    origin
294}
295
296/// Builds an `RawOrigin` object from the given event packet.
297pub fn origin_from_event_packet<'packet, 'tags>(
298    packet: &'tags EventPacket<'packet>, well_known_tags: &'tags WellKnownTags<'tags>,
299) -> RawOrigin<'tags>
300where
301    'packet: 'tags,
302{
303    let cardinality = packet.cardinality.or(well_known_tags.cardinality);
304
305    let mut origin = RawOrigin::default();
306    origin.set_pod_uid(well_known_tags.pod_uid);
307    origin.set_local_data(packet.local_data);
308    origin.set_external_data(packet.external_data);
309    origin.set_cardinality(cardinality);
310    origin
311}
312
313/// Builds an `RawOrigin` object from the given service check packet.
314pub fn origin_from_service_check_packet<'packet, 'tags>(
315    packet: &'tags ServiceCheckPacket<'packet>, well_known_tags: &'tags WellKnownTags<'tags>,
316) -> RawOrigin<'tags>
317where
318    'packet: 'tags,
319{
320    let cardinality = packet.cardinality.or(well_known_tags.cardinality);
321
322    let mut origin = RawOrigin::default();
323    origin.set_pod_uid(well_known_tags.pod_uid);
324    origin.set_local_data(packet.local_data);
325    origin.set_external_data(packet.external_data);
326    origin.set_cardinality(cardinality);
327    origin
328}
329
330#[cfg(test)]
331mod tests {
332    use std::collections::HashMap;
333
334    use saluki_core::data_model::{
335        event::{metric::MetricValues, service_check::CheckStatus},
336        tags::{RawTags, TagSet},
337    };
338    use saluki_env::workload::{origin::ResolvedExternalData, providers::TestWorkloadProvider, EntityId};
339    use stringtheory::MetaString;
340
341    use super::*;
342
343    static EID_PID: EntityId = EntityId::ContainerPid(12345);
344    static EID_LOCAL_CID: EntityId = EntityId::Container(MetaString::from_static("local-cid"));
345    static EID_EXTERNAL_CID_VALID: EntityId = EntityId::Container(MetaString::from_static("external-cid"));
346    static EID_EXTERNAL_CID_INVALID: EntityId = EntityId::Container(MetaString::from_static("invalid-external-cid"));
347    static EID_LOCAL_POD: EntityId = EntityId::PodUid(MetaString::from_static("local-pod-uid"));
348    static EID_EXTERNAL_POD: EntityId = EntityId::PodUid(MetaString::from_static("external-pod-uid"));
349
350    fn single_tag(tag: &str) -> SharedTagSet {
351        let mut tag_set = TagSet::default();
352        tag_set.insert_tag(tag);
353        tag_set.into_shared()
354    }
355
356    fn tags_for_entity(entity_id: &EntityId) -> SharedTagSet {
357        if entity_id == &EID_PID {
358            single_tag("tag_source:pid")
359        } else if entity_id == &EID_LOCAL_CID {
360            single_tag("tag_source:local-cid")
361        } else if entity_id == &EID_EXTERNAL_CID_VALID {
362            single_tag("tag_source:external-cid")
363        } else if entity_id == &EID_LOCAL_POD {
364            single_tag("tag_source:local-pod")
365        } else if entity_id == &EID_EXTERNAL_POD {
366            single_tag("tag_source:external-pod")
367        } else {
368            SharedTagSet::default()
369        }
370    }
371
372    fn origin(
373        maybe_process_id: Option<&EntityId>, maybe_local_container_id: Option<&EntityId>,
374        maybe_local_pod_uid: Option<&EntityId>, maybe_external_data: Option<&ResolvedExternalData>,
375    ) -> ResolvedOrigin {
376        ResolvedOrigin::from_parts(
377            None,
378            maybe_process_id.cloned(),
379            maybe_local_container_id.cloned(),
380            maybe_local_pod_uid.cloned(),
381            maybe_external_data.cloned(),
382        )
383    }
384
385    fn build_tags_resolver_with_default_tags(config: OriginEnrichmentConfiguration) -> DogStatsDOriginTagResolver {
386        let mut workload_provider = TestWorkloadProvider::new();
387        workload_provider.add_entity_shared_tags(EID_PID.clone(), tags_for_entity(&EID_PID));
388        workload_provider.add_entity_shared_tags(EID_LOCAL_CID.clone(), tags_for_entity(&EID_LOCAL_CID));
389        workload_provider
390            .add_entity_shared_tags(EID_EXTERNAL_CID_VALID.clone(), tags_for_entity(&EID_EXTERNAL_CID_VALID));
391        workload_provider.add_entity_shared_tags(EID_LOCAL_POD.clone(), tags_for_entity(&EID_LOCAL_POD));
392        workload_provider.add_entity_shared_tags(EID_EXTERNAL_POD.clone(), tags_for_entity(&EID_EXTERNAL_POD));
393
394        let erased_workload_provider = Arc::new(workload_provider);
395
396        DogStatsDOriginTagResolver::new(config, erased_workload_provider, CapturedTaggerHandle::new())
397    }
398
399    #[test]
400    fn metric_cardinality_precedence() {
401        // Tests that the cardinality specified in a metric packet (`|card:high`, etc) takes precedence over the cardinality
402        // specified via the deprecated `dd.internal.card` tag.
403        let raw_tags_input = "dd.internal.card:high";
404        let raw_tags = RawTags::new(raw_tags_input, usize::MAX, usize::MAX);
405
406        let well_known_tags = WellKnownTags::from_raw_tags(&raw_tags);
407        assert_eq!(well_known_tags.cardinality, Some(OriginTagCardinality::High));
408
409        let packet_with_card = MetricPacket {
410            metric_name: "test_metric",
411            tags: raw_tags.clone(),
412            values: MetricValues::counter(1.0),
413            num_points: 1,
414            timestamp: None,
415            local_data: None,
416            external_data: None,
417            cardinality: Some(OriginTagCardinality::Low),
418            unit: None,
419        };
420
421        let packet_without_card = MetricPacket {
422            metric_name: "test_metric",
423            tags: raw_tags.clone(),
424            values: MetricValues::counter(1.0),
425            num_points: 1,
426            timestamp: None,
427            local_data: None,
428            external_data: None,
429            cardinality: None,
430            unit: None,
431        };
432
433        let with_card_origin = origin_from_metric_packet(&packet_with_card, &well_known_tags);
434        assert_ne!(packet_with_card.cardinality, well_known_tags.cardinality);
435        assert_eq!(with_card_origin.cardinality(), packet_with_card.cardinality);
436
437        let without_card_origin = origin_from_metric_packet(&packet_without_card, &well_known_tags);
438        assert_ne!(packet_without_card.cardinality, well_known_tags.cardinality);
439        assert_eq!(without_card_origin.cardinality(), well_known_tags.cardinality);
440    }
441
442    #[test]
443    fn event_cardinality_precedence() {
444        // Tests that the cardinality specified in an event packet (`|card:high`, etc) takes precedence over the cardinality
445        // specified via the deprecated `dd.internal.card` tag.
446        let raw_tags_input = "dd.internal.card:low";
447        let raw_tags = RawTags::new(raw_tags_input, usize::MAX, usize::MAX);
448
449        let well_known_tags = WellKnownTags::from_raw_tags(&raw_tags);
450        assert_eq!(well_known_tags.cardinality, Some(OriginTagCardinality::Low));
451
452        let packet_with_card = EventPacket {
453            title: MetaString::empty(),
454            text: MetaString::empty(),
455            timestamp: None,
456            hostname: None,
457            aggregation_key: None,
458            priority: None,
459            alert_type: None,
460            source_type_name: None,
461            tags: raw_tags.clone(),
462            local_data: None,
463            external_data: None,
464            cardinality: Some(OriginTagCardinality::Orchestrator),
465        };
466
467        let packet_without_card = EventPacket {
468            title: MetaString::empty(),
469            text: MetaString::empty(),
470            timestamp: None,
471            hostname: None,
472            aggregation_key: None,
473            priority: None,
474            alert_type: None,
475            source_type_name: None,
476            tags: raw_tags.clone(),
477            local_data: None,
478            external_data: None,
479            cardinality: None,
480        };
481
482        let with_card_origin = origin_from_event_packet(&packet_with_card, &well_known_tags);
483        assert_ne!(packet_with_card.cardinality, well_known_tags.cardinality);
484        assert_eq!(with_card_origin.cardinality(), packet_with_card.cardinality);
485
486        let without_card_origin = origin_from_event_packet(&packet_without_card, &well_known_tags);
487        assert_ne!(packet_without_card.cardinality, well_known_tags.cardinality);
488        assert_eq!(without_card_origin.cardinality(), well_known_tags.cardinality);
489    }
490
491    #[test]
492    fn service_check_cardinality_precedence() {
493        // Tests that the cardinality specified in an event packet (`|card:high`, etc) takes precedence over the cardinality
494        // specified via the deprecated `dd.internal.card` tag.
495        let raw_tags_input = "dd.internal.card:orchestrator";
496        let raw_tags = RawTags::new(raw_tags_input, usize::MAX, usize::MAX);
497
498        let well_known_tags = WellKnownTags::from_raw_tags(&raw_tags);
499        assert_eq!(well_known_tags.cardinality, Some(OriginTagCardinality::Orchestrator));
500
501        let packet_with_card = ServiceCheckPacket {
502            name: MetaString::empty(),
503            status: CheckStatus::Ok,
504            timestamp: None,
505            hostname: None,
506            message: None,
507            tags: raw_tags.clone(),
508            local_data: None,
509            external_data: None,
510            cardinality: Some(OriginTagCardinality::Low),
511        };
512
513        let packet_without_card = ServiceCheckPacket {
514            name: MetaString::empty(),
515            status: CheckStatus::Ok,
516            timestamp: None,
517            hostname: None,
518            message: None,
519            tags: raw_tags.clone(),
520            local_data: None,
521            external_data: None,
522            cardinality: None,
523        };
524
525        let with_card_origin = origin_from_service_check_packet(&packet_with_card, &well_known_tags);
526        assert_ne!(packet_with_card.cardinality, well_known_tags.cardinality);
527        assert_eq!(with_card_origin.cardinality(), packet_with_card.cardinality);
528
529        let without_card_origin = origin_from_service_check_packet(&packet_without_card, &well_known_tags);
530        assert_ne!(packet_without_card.cardinality, well_known_tags.cardinality);
531        assert_eq!(without_card_origin.cardinality(), well_known_tags.cardinality);
532    }
533
534    #[test]
535    fn origin_detection_legacy_precedence() {
536        let mut pid_plus_local_pod_tags = tags_for_entity(&EID_PID);
537        pid_plus_local_pod_tags.extend_from_shared(&tags_for_entity(&EID_LOCAL_POD));
538
539        let mut pid_plus_local_cid_tags = tags_for_entity(&EID_PID);
540        pid_plus_local_cid_tags.extend_from_shared(&tags_for_entity(&EID_LOCAL_CID));
541
542        let cases = [
543            // We only have the process ID, so entity ID precedence should be irrelevant.
544            (
545                false,
546                origin(Some(&EID_PID), None, None, None),
547                tags_for_entity(&EID_PID),
548            ),
549            (
550                true,
551                origin(Some(&EID_PID), None, None, None),
552                tags_for_entity(&EID_PID),
553            ),
554            // We have both the process ID and local pod UID, but entity ID precedence is disabled, so we get the
555            // process ID and local pod UID tags.
556            (
557                false,
558                origin(Some(&EID_PID), None, Some(&EID_LOCAL_POD), None),
559                pid_plus_local_pod_tags.clone(),
560            ),
561            // We have both the process ID and local pod UID, but entity ID precedence is enabled, so we should only get
562            // the local pod UID tags.
563            (
564                true,
565                origin(Some(&EID_PID), None, Some(&EID_LOCAL_POD), None),
566                tags_for_entity(&EID_LOCAL_POD),
567            ),
568            // We have the process ID, local container ID, and local pod UID, but entity ID precedence is disabled, so
569            // we should get the process ID and local pod UID tags.
570            (
571                false,
572                origin(Some(&EID_PID), Some(&EID_LOCAL_CID), Some(&EID_LOCAL_POD), None),
573                pid_plus_local_pod_tags,
574            ),
575            // We have the process ID, local container ID, and local pod UID, but entity ID precedence is enabled, so we
576            // should only get the local pod UID tags.
577            (
578                true,
579                origin(Some(&EID_PID), Some(&EID_LOCAL_CID), Some(&EID_LOCAL_POD), None),
580                tags_for_entity(&EID_LOCAL_POD),
581            ),
582            // We only have the process ID and local container ID, so entity ID precedence should be irrelevant, and so
583            // we should get the process ID and local container ID tags.
584            (
585                false,
586                origin(Some(&EID_PID), Some(&EID_LOCAL_CID), None, None),
587                pid_plus_local_cid_tags.clone(),
588            ),
589            (
590                true,
591                origin(Some(&EID_PID), Some(&EID_LOCAL_CID), None, None),
592                pid_plus_local_cid_tags,
593            ),
594        ];
595
596        for (entity_id_precedence, resolved_origin, expected_tags) in cases {
597            let tag_resolver_config = OriginEnrichmentConfiguration {
598                enabled: true,
599                entity_id_precedence,
600                tag_cardinality: OriginTagCardinality::High,
601                origin_detection_unified: false,
602                origin_detection_optout: false,
603                origin_detection_client: false,
604            };
605
606            let origin_tags_resolver = build_tags_resolver_with_default_tags(tag_resolver_config);
607
608            let actual_tags = origin_tags_resolver.collect_origin_tags(&resolved_origin);
609            assert_eq!(
610                actual_tags, expected_tags,
611                "failed to resolve the expected tags for origin {:?}",
612                resolved_origin
613            );
614        }
615    }
616
617    #[test]
618    fn origin_detection_unified_precedence() {
619        // We craft a "valid" and "invalid" variant for External Data, where the invalid one has a container ID with no tags
620        // assigned to it, which lets us exercise the tags resolver logic for when we have a pod UID through External Data,
621        // but not a container ID (or a container ID with no tags attached).
622        //
623        // We have to do it this way because we pass both External Data-based entity IDs through `ResolvedExternalData` when
624        // creating `ResolvedOrigin`, so we can't pass them separately.
625        let ext_data_valid = ResolvedExternalData::new(EID_EXTERNAL_POD.clone(), EID_EXTERNAL_CID_VALID.clone());
626        let ext_data_invalid = ResolvedExternalData::new(EID_EXTERNAL_POD.clone(), EID_EXTERNAL_CID_INVALID.clone());
627
628        let tag_resolver_config = OriginEnrichmentConfiguration {
629            enabled: true,
630            entity_id_precedence: false,
631            tag_cardinality: OriginTagCardinality::High,
632            origin_detection_unified: true,
633            origin_detection_optout: false,
634            origin_detection_client: false,
635        };
636
637        let origin_tags_resolver = build_tags_resolver_with_default_tags(tag_resolver_config);
638
639        // We craft our test cases to ensure that we always take the tags of the highest precedence entity ID available,
640        // and don't take any other tags.
641        let cases = [
642            // Cases where we're only setting a single entity ID. This is the happy path.
643            (origin(Some(&EID_PID), None, None, None), tags_for_entity(&EID_PID)),
644            (
645                origin(None, Some(&EID_LOCAL_CID), None, None),
646                tags_for_entity(&EID_LOCAL_CID),
647            ),
648            (
649                origin(None, None, Some(&EID_LOCAL_POD), None),
650                tags_for_entity(&EID_LOCAL_POD),
651            ),
652            (
653                origin(None, None, None, Some(&ext_data_valid)),
654                tags_for_entity(&EID_EXTERNAL_CID_VALID),
655            ),
656            (
657                origin(None, None, None, Some(&ext_data_invalid)),
658                tags_for_entity(&EID_EXTERNAL_POD),
659            ),
660            // Cases where we have multiple entity IDs to choose from. We work our way backwards here.
661            (
662                origin(
663                    Some(&EID_PID),
664                    Some(&EID_LOCAL_CID),
665                    Some(&EID_LOCAL_POD),
666                    Some(&ext_data_valid),
667                ),
668                tags_for_entity(&EID_LOCAL_CID),
669            ),
670            (
671                origin(Some(&EID_PID), None, Some(&EID_LOCAL_POD), Some(&ext_data_valid)),
672                tags_for_entity(&EID_PID),
673            ),
674            (
675                origin(None, None, Some(&EID_LOCAL_POD), Some(&ext_data_valid)),
676                tags_for_entity(&EID_EXTERNAL_CID_VALID),
677            ),
678            (
679                origin(None, None, Some(&EID_LOCAL_POD), Some(&ext_data_invalid)),
680                tags_for_entity(&EID_LOCAL_POD),
681            ),
682            (
683                origin(None, None, None, Some(&ext_data_invalid)),
684                tags_for_entity(&EID_EXTERNAL_POD),
685            ),
686        ];
687
688        for (resolved_origin, expected_tags) in cases {
689            let actual_tags = origin_tags_resolver.collect_origin_tags(&resolved_origin);
690            assert_eq!(
691                actual_tags, expected_tags,
692                "failed to resolve the expected tags for origin {:?}",
693                resolved_origin
694            );
695        }
696    }
697
698    #[test]
699    fn origin_detection_disabled() {
700        // When origin detection is disabled, no tags should be resolved even if we do have mapped tags for the given
701        // resolved origin.
702        let tag_resolver_config = OriginEnrichmentConfiguration::for_test();
703        assert!(!tag_resolver_config.enabled);
704
705        let origin_tags_resolver = build_tags_resolver_with_default_tags(tag_resolver_config);
706
707        let resolved_origin = origin(Some(&EID_PID), None, None, None);
708        let actual_tags = origin_tags_resolver.collect_origin_tags(&resolved_origin);
709        assert!(actual_tags.is_empty());
710    }
711
712    #[test]
713    fn resolve_origin_tags_dispatches_to_captured_store_when_replay_flag_set() {
714        use datadog_protos::agent::{Entity as ProtoEntity, TaggerState};
715
716        // Build a captured store with a single entity keyed by PID 7777.
717        let mut entities = HashMap::new();
718        entities.insert(
719            "container_id://captured-container".to_string(),
720            ProtoEntity {
721                low_cardinality_tags: vec!["env:captured".into(), "service:replayed".into()],
722                ..Default::default()
723            },
724        );
725        let mut pid_map = HashMap::new();
726        pid_map.insert(7777, "container_id://captured-container".to_string());
727        let state = TaggerState {
728            state: entities,
729            pid_map,
730            duration: 0,
731        };
732        let captured_tagger = CapturedTaggerHandle::new();
733        captured_tagger.set_current(Some(super::super::replay::CapturedTaggerStore::from_tagger_state(
734            state,
735        )));
736
737        // Build a resolver whose live workload provider is empty; if the replay path is wired wrong, we'd get
738        // no tags. If it's wired right, we get the captured tags.
739        let config = OriginEnrichmentConfiguration {
740            enabled: true,
741            tag_cardinality: OriginTagCardinality::Low,
742            ..OriginEnrichmentConfiguration::for_test()
743        };
744        let live = Arc::new(TestWorkloadProvider::new());
745        let resolver = DogStatsDOriginTagResolver::new(config, live, captured_tagger);
746
747        let mut origin = RawOrigin::default();
748        origin.set_process_id(mark_replay_process_id(7777));
749        origin.set_cardinality(OriginTagCardinality::Low);
750
751        let tags = resolver.resolve_origin_tags(origin);
752        let tag_strs: Vec<String> = tags.into_iter().map(|t| t.as_str().to_string()).collect();
753        assert!(tag_strs.contains(&"env:captured".to_string()));
754        assert!(tag_strs.contains(&"service:replayed".to_string()));
755    }
756
757    #[test]
758    fn resolve_origin_tags_returns_empty_when_replay_flag_set_but_no_captured_store() {
759        // If the replay flag is set but no captured store is current, the resolver returns an empty tag set
760        // (no fallback to the live tagger). This guards against accidentally serving live tags for replay packets.
761        let config = OriginEnrichmentConfiguration {
762            enabled: true,
763            tag_cardinality: OriginTagCardinality::Low,
764            ..OriginEnrichmentConfiguration::for_test()
765        };
766        let live = Arc::new(TestWorkloadProvider::new());
767        let resolver = DogStatsDOriginTagResolver::new(config, live, CapturedTaggerHandle::new());
768
769        let mut origin = RawOrigin::default();
770        origin.set_process_id(mark_replay_process_id(7777));
771
772        let tags = resolver.resolve_origin_tags(origin);
773        assert!(
774            tags.is_empty(),
775            "replay path with no captured store must return empty tags"
776        );
777    }
778
779    #[test]
780    fn resolve_origin_tags_live_path_resolves_tags_via_workload_provider() {
781        // Exercises the real, non-replay production entry point: `resolve_origin_tags` asks the workload provider to
782        // resolve the raw origin, then walks the resolved entity IDs to collect tags. Every other test either drives
783        // the internal `collect_origin_tags` directly or takes the replay bypass, so this is the only coverage of the
784        // live `get_resolved_origin` -> `collect_origin_tags` path through the trait method.
785        let config = OriginEnrichmentConfiguration {
786            enabled: true,
787            tag_cardinality: OriginTagCardinality::High,
788            ..OriginEnrichmentConfiguration::for_test()
789        };
790
791        let mut workload_provider = TestWorkloadProvider::new();
792        workload_provider.add_entity_shared_tags(EntityId::ContainerPid(4242), single_tag("tag_source:live-pid"));
793        let resolver =
794            DogStatsDOriginTagResolver::new(config, Arc::new(workload_provider), CapturedTaggerHandle::new());
795
796        // A plain process ID (no replay marker bit) resolves to `EntityId::ContainerPid(4242)` via the live path.
797        let mut origin = RawOrigin::default();
798        origin.set_process_id(4242);
799
800        let tags = resolver.resolve_origin_tags(origin);
801        assert_eq!(tags, single_tag("tag_source:live-pid"));
802    }
803
804    #[test]
805    fn resolve_origin_tags_uses_pinned_process_entity() {
806        let pinned_entity = EntityId::Container(MetaString::from_static("original-container"));
807        let config = OriginEnrichmentConfiguration {
808            enabled: true,
809            tag_cardinality: OriginTagCardinality::High,
810            ..OriginEnrichmentConfiguration::for_test()
811        };
812        let workload_provider =
813            TestWorkloadProvider::with_entity(pinned_entity.clone(), &["tag_source:pinned-container"]);
814        let resolver =
815            DogStatsDOriginTagResolver::new(config, Arc::new(workload_provider), CapturedTaggerHandle::new());
816        let process_origin = ProcessOrigin::Pinned(Some(pinned_entity));
817
818        let tags = resolver.resolve_origin_tags_with_process_origin(RawOrigin::default(), Some(&process_origin));
819
820        assert_eq!(tags, single_tag("tag_source:pinned-container"));
821    }
822
823    #[test]
824    fn resolve_origin_tags_live_path_returns_empty_when_origin_unresolved() {
825        // When the workload provider can't resolve the origin (here, an empty origin resolves to `None`), the live
826        // path returns an empty tag set rather than falling through to anything else.
827        let config = OriginEnrichmentConfiguration {
828            enabled: true,
829            tag_cardinality: OriginTagCardinality::High,
830            ..OriginEnrichmentConfiguration::for_test()
831        };
832        let resolver = DogStatsDOriginTagResolver::new(
833            config,
834            Arc::new(TestWorkloadProvider::new()),
835            CapturedTaggerHandle::new(),
836        );
837
838        let tags = resolver.resolve_origin_tags(RawOrigin::default());
839        assert!(tags.is_empty(), "unresolved origin must produce no tags");
840    }
841}