saluki_env/workload/
origin.rs1use 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
15const 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#[derive(Clone, Debug, Eq, Hash, PartialEq)]
21pub struct ResolvedExternalData {
22 pod_entity_id: EntityId,
23 container_entity_id: EntityId,
24}
25
26impl ResolvedExternalData {
27 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 pub fn pod_entity_id(&self) -> &EntityId {
37 &self.pod_entity_id
38 }
39
40 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#[derive(Clone, Debug, Hash, Eq, PartialEq)]
62pub struct ResolvedOrigin {
63 inner: Arc<ResolvedOriginInner>,
64}
65
66impl ResolvedOrigin {
67 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 pub fn cardinality(&self) -> Option<OriginTagCardinality> {
85 self.inner.cardinality
86 }
87
88 pub fn process_id(&self) -> Option<&EntityId> {
90 self.inner.process_id.as_ref()
91 }
92
93 pub fn local_data(&self) -> Option<&EntityId> {
95 self.inner.local_data.as_ref()
96 }
97
98 pub fn pod_uid(&self) -> Option<&EntityId> {
100 self.inner.pod_uid.as_ref()
101 }
102
103 pub fn resolved_external_data(&self) -> Option<&ResolvedExternalData> {
105 self.inner.resolved_external_data.as_ref()
106 }
107}
108
109#[derive(Clone)]
111pub struct OriginResolver {
112 ed_resolver: ExternalDataStoreResolver,
113 origin_cache: Cache<u64, ResolvedOrigin>,
114}
115
116impl OriginResolver {
117 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 pub fn get_resolved_origin(&self, origin: RawOrigin<'_>) -> Option<ResolvedOrigin> {
147 if origin.is_empty() {
149 return None;
150 }
151
152 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 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 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 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 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 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}