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#[derive(Clone)]
52#[cfg_attr(test, derive(Debug, PartialEq))]
53pub struct OriginEnrichmentConfiguration {
54 pub enabled: bool,
58
59 pub entity_id_precedence: bool,
65
66 pub tag_cardinality: OriginTagCardinality,
68
69 pub origin_detection_unified: bool,
79
80 pub origin_detection_optout: bool,
85
86 pub origin_detection_client: bool,
91}
92
93#[cfg(test)]
94impl OriginEnrichmentConfiguration {
95 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 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 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 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 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 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
279pub 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
296pub 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
313pub 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 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 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 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 (
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 (
557 false,
558 origin(Some(&EID_PID), None, Some(&EID_LOCAL_POD), None),
559 pid_plus_local_pod_tags.clone(),
560 ),
561 (
564 true,
565 origin(Some(&EID_PID), None, Some(&EID_LOCAL_POD), None),
566 tags_for_entity(&EID_LOCAL_POD),
567 ),
568 (
571 false,
572 origin(Some(&EID_PID), Some(&EID_LOCAL_CID), Some(&EID_LOCAL_POD), None),
573 pid_plus_local_pod_tags,
574 ),
575 (
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 (
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 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 let cases = [
642 (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 (
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 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 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 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 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 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 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 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}