saluki_env/workload/
origin.rs

1//! Origin detection and resolution.
2
3use std::{num::NonZeroUsize, sync::Arc, time::Duration};
4
5use saluki_common::{
6    cache::{Cache, CacheBuilder},
7    hash::hash_single_fast,
8};
9use saluki_context::origin::{OriginTagCardinality, RawOrigin};
10use tracing::trace;
11
12use super::stores::ExternalDataStoreResolver;
13use crate::workload::EntityId;
14
15// SAFETY: This number is obviously non-zero.
16const DEFAULT_ORIGIN_CACHE_ITEM_LIMIT: NonZeroUsize = NonZeroUsize::new(500_000).unwrap();
17const DEFAULT_ORIGIN_CACHE_ITEM_TIME_TO_IDLE: Duration = Duration::from_secs(30);
18
19/// A resolved External Data entry.
20#[derive(Clone, Debug, Eq, Hash, PartialEq)]
21pub struct ResolvedExternalData {
22    pod_entity_id: EntityId,
23    container_entity_id: EntityId,
24}
25
26impl ResolvedExternalData {
27    /// Creates a new `ResolvedExternalData` from the given pod and container entity IDs.
28    pub fn new(pod_entity_id: EntityId, container_entity_id: EntityId) -> Self {
29        Self {
30            pod_entity_id,
31            container_entity_id,
32        }
33    }
34
35    /// Returns a reference to the pod entity ID.
36    pub fn pod_entity_id(&self) -> &EntityId {
37        &self.pod_entity_id
38    }
39
40    /// Returns a reference to the container entity ID.
41    pub fn container_entity_id(&self) -> &EntityId {
42        &self.container_entity_id
43    }
44}
45
46#[derive(Debug, Hash, Eq, PartialEq)]
47struct ResolvedOriginInner {
48    cardinality: Option<OriginTagCardinality>,
49    process_id: Option<EntityId>,
50    local_data: Option<EntityId>,
51    pod_uid: Option<EntityId>,
52    resolved_external_data: Option<ResolvedExternalData>,
53}
54
55/// An resolved representation of `RawOrigin<'a>`
56///
57/// This representation is used to store the pre-calculated entity IDs derived from a borrowed `RawOrigin<'a>` in order
58/// to speed the lookup of origin tags attached to each individual entity ID that comprises an origin.
59///
60/// This type can be cheaply cloned and shared across threads.
61#[derive(Clone, Debug, Hash, Eq, PartialEq)]
62pub struct ResolvedOrigin {
63    inner: Arc<ResolvedOriginInner>,
64}
65
66impl ResolvedOrigin {
67    /// Creates a new `ResolvedOrigin` from the given parts.
68    pub fn from_parts(
69        cardinality: Option<OriginTagCardinality>, process_id: Option<EntityId>, local_data: Option<EntityId>,
70        pod_uid: Option<EntityId>, resolved_external_data: Option<ResolvedExternalData>,
71    ) -> Self {
72        Self {
73            inner: Arc::new(ResolvedOriginInner {
74                cardinality,
75                process_id,
76                local_data,
77                pod_uid,
78                resolved_external_data,
79            }),
80        }
81    }
82
83    /// Returns the cardinality of the origin.
84    pub fn cardinality(&self) -> Option<OriginTagCardinality> {
85        self.inner.cardinality
86    }
87
88    /// Returns the process ID of the origin.
89    pub fn process_id(&self) -> Option<&EntityId> {
90        self.inner.process_id.as_ref()
91    }
92
93    /// Returns the Local Data-based entity ID of the origin.
94    pub fn local_data(&self) -> Option<&EntityId> {
95        self.inner.local_data.as_ref()
96    }
97
98    /// Returns the pod UID of the origin.
99    pub fn pod_uid(&self) -> Option<&EntityId> {
100        self.inner.pod_uid.as_ref()
101    }
102
103    /// Returns the resolved External Data of the origin.
104    pub fn resolved_external_data(&self) -> Option<&ResolvedExternalData> {
105        self.inner.resolved_external_data.as_ref()
106    }
107}
108
109/// Resolves and tracks origins.
110#[derive(Clone)]
111pub struct OriginResolver {
112    ed_resolver: ExternalDataStoreResolver,
113    origin_cache: Cache<u64, ResolvedOrigin>,
114}
115
116impl OriginResolver {
117    /// Creates a new `OriginResolver`.
118    pub fn new(ed_resolver: ExternalDataStoreResolver) -> Self {
119        Self {
120            ed_resolver,
121            origin_cache: CacheBuilder::from_identifier("origin_cache")
122                .expect("identifier cannot be invalid")
123                .with_capacity(DEFAULT_ORIGIN_CACHE_ITEM_LIMIT)
124                .with_time_to_idle(Some(DEFAULT_ORIGIN_CACHE_ITEM_TIME_TO_IDLE))
125                .build(),
126        }
127    }
128
129    fn build_resolved_origin(&self, origin: RawOrigin<'_>) -> ResolvedOrigin {
130        ResolvedOrigin::from_parts(
131            origin.cardinality(),
132            origin.process_id().map(EntityId::ContainerPid),
133            origin.local_data().and_then(EntityId::from_local_data),
134            origin.pod_uid().and_then(EntityId::from_pod_uid),
135            origin
136                .external_data()
137                .and_then(|raw_ed| self.ed_resolver.resolve(raw_ed)),
138        )
139    }
140
141    /// Returns the resolved origin for the given raw origin.
142    ///
143    /// If the raw origin is "empty" -- no origin information is available -- then `None` is returned.
144    ///
145    /// The resolved origin may be cached for speeding up future lookups.
146    pub fn get_resolved_origin(&self, origin: RawOrigin<'_>) -> Option<ResolvedOrigin> {
147        // If there's no origin information at all, then there's nothing to key off of.
148        if origin.is_empty() {
149            return None;
150        }
151
152        // Create the origin key, and populate our cache with the resolved origin if we don't already have it.
153        let origin_key = hash_single_fast(&origin);
154        match self.origin_cache.get(&origin_key) {
155            Some(resolved_origin) => {
156                trace!(?origin_key, "Found origin in cache.");
157                Some(resolved_origin)
158            }
159            None => {
160                trace!(?origin_key, "Origin not found in cache. Resolving.");
161                let resolved_origin = self.build_resolved_origin(origin);
162                self.origin_cache.insert(origin_key, resolved_origin.clone());
163
164                Some(resolved_origin)
165            }
166        }
167    }
168}
169
170#[cfg(test)]
171mod tests {
172    // These tests were restored from a module that was commented out in mid-2025 (commit 6635adbf9f) and never
173    // brought back. In the interim, `OriginResolver` was reworked: it no longer performs process-ID -> container-ID
174    // alias redirection (that behavior moved to `TagStoreQuerier::get_entity_tags`, whose alias handling is now
175    // covered in `stores/tag.rs`), and the old `resolve_origin`/`get_resolved_origin_by_key`/`ResolvedOrigin::container_id`
176    // API was replaced by `get_resolved_origin` returning a cached `ResolvedOrigin`. These tests therefore cover the
177    // current documented contract of `OriginResolver`: field-to-entity resolution and its caching behavior.
178    use std::num::NonZeroUsize;
179
180    use super::*;
181    use crate::workload::stores::ExternalDataStore;
182
183    fn origin_resolver() -> OriginResolver {
184        let external_data_store = ExternalDataStore::with_entity_limit(NonZeroUsize::new(usize::MAX).unwrap());
185        OriginResolver::new(external_data_store.resolver())
186    }
187
188    fn raw_origin(
189        process_id: Option<u32>, local_data: Option<&'static str>, pod_uid: Option<&'static str>,
190    ) -> RawOrigin<'static> {
191        let mut origin = RawOrigin::default();
192        if let Some(process_id) = process_id {
193            origin.set_process_id(process_id);
194        }
195        origin.set_local_data(local_data);
196        origin.set_pod_uid(pod_uid);
197        origin
198    }
199
200    #[tokio::test]
201    async fn get_resolved_origin_returns_none_for_empty_origin() {
202        // An origin with no information at all can't be keyed off of, so resolution yields `None`.
203        let resolver = origin_resolver();
204
205        assert!(resolver.get_resolved_origin(RawOrigin::default()).is_none());
206    }
207
208    #[tokio::test]
209    async fn get_resolved_origin_maps_process_id_and_local_data_to_entities() {
210        // Resolution derives entity IDs directly from the raw origin's fields: the process ID becomes a `ContainerPid`
211        // entity, and the Local Data (here, a bare container ID) becomes a `Container` entity.
212        let resolver = origin_resolver();
213
214        let resolved = resolver
215            .get_resolved_origin(raw_origin(Some(1234), Some("container-a"), None))
216            .expect("non-empty origin should resolve");
217
218        assert_eq!(resolved.process_id(), Some(&EntityId::ContainerPid(1234)));
219        assert_eq!(resolved.local_data(), Some(&EntityId::Container("container-a".into())));
220        assert_eq!(resolved.pod_uid(), None);
221    }
222
223    #[tokio::test]
224    async fn get_resolved_origin_caches_identical_raw_origins() {
225        // Two lookups for equal raw origins must return the very same cached `ResolvedOrigin` instance rather than a
226        // freshly rebuilt one -- the documented caching contract.
227        let resolver = origin_resolver();
228
229        let origin = raw_origin(Some(1234), Some("container-a"), None);
230        let first = resolver.get_resolved_origin(origin.clone()).expect("should resolve");
231        let second = resolver.get_resolved_origin(origin).expect("should resolve");
232
233        assert!(
234            Arc::ptr_eq(&first.inner, &second.inner),
235            "identical raw origins should resolve to the same cached instance"
236        );
237    }
238
239    #[tokio::test]
240    async fn get_resolved_origin_distinguishes_different_local_data() {
241        // Raw origins that differ only in their Local Data resolve to distinct, non-equal resolved origins.
242        let resolver = origin_resolver();
243
244        let resolved_a = resolver
245            .get_resolved_origin(raw_origin(Some(1234), Some("container-a"), None))
246            .expect("should resolve");
247        let resolved_b = resolver
248            .get_resolved_origin(raw_origin(Some(1234), Some("container-b"), None))
249            .expect("should resolve");
250
251        assert_ne!(resolved_a, resolved_b);
252        assert_eq!(
253            resolved_a.local_data(),
254            Some(&EntityId::Container("container-a".into()))
255        );
256        assert_eq!(
257            resolved_b.local_data(),
258            Some(&EntityId::Container("container-b".into()))
259        );
260    }
261}