1use std::{num::NonZeroUsize, sync::Arc, time::Duration};
2
3use saluki_common::{collections::PrehashedHashSet, hash::NoopU64BuildHasher};
4use saluki_error::{generic_error, GenericError};
5use saluki_metrics::{static_metrics, Counter, Gauge};
6use stringtheory::{
7 interning::{GenericMapInterner, Interner as _},
8 CheapMetaString, MetaString,
9};
10use tokio::time::sleep;
11use tracing::debug;
12
13use super::{
14 hash::{hash_context_with_host_and_seen, ContextKey, TagSetKey},
15 Context, ContextInner,
16};
17use crate::{
18 cache::{weight::ItemCountWeighter, Cache, CacheBuilder},
19 data_model::{
20 origin::{OriginTagsResolver, RawOrigin},
21 tags::{SharedTagSet, TagSet},
22 },
23};
24
25const DEFAULT_CONTEXT_RESOLVER_CACHED_CONTEXTS_LIMIT: NonZeroUsize = NonZeroUsize::new(500_000).unwrap();
27
28const DEFAULT_CONTEXT_RESOLVER_INTERNER_CAPACITY_BYTES: NonZeroUsize = NonZeroUsize::new(2 * 1024 * 1024).unwrap();
30
31const SEEN_HASHSET_INITIAL_CAPACITY: usize = 128;
32
33type ContextCache = Cache<ContextKey, Context, ItemCountWeighter, NoopU64BuildHasher>;
34type TagSetCache = Cache<TagSetKey, SharedTagSet, ItemCountWeighter, NoopU64BuildHasher>;
35
36#[static_metrics(prefix = context_resolver, labels(resolver_id))]
37#[derive(Clone)]
38struct Telemetry {
39 interner_capacity_bytes: Gauge,
40 interner_len_bytes: Gauge,
41 interner_entries: Gauge,
42 #[metric(level = debug)]
43 intern_fallback_total: Counter,
44 #[metric(level = debug)]
45 resolved_existing_context_total: Counter,
46 #[metric(level = debug)]
47 resolved_new_context_total: Counter,
48 active_contexts: Gauge,
49 #[metric(level = debug)]
50 resolved_existing_tagset_total: Counter,
51 #[metric(level = debug)]
52 resolved_new_tagset_total: Counter,
53}
54
55pub struct ContextResolverBuilder {
61 name: String,
62 caching_enabled: bool,
63 cached_contexts_limit: Option<NonZeroUsize>,
64 idle_context_expiration: Option<Duration>,
65 interner_capacity_bytes: Option<NonZeroUsize>,
66 allow_heap_allocations: Option<bool>,
67 tags_resolver: Option<TagsResolver>,
68 interner: Option<GenericMapInterner>,
69 origin_tags_resolver: Option<Arc<dyn OriginTagsResolver>>,
70 telemetry_enabled: bool,
71}
72
73impl ContextResolverBuilder {
74 pub fn from_name<S: Into<String>>(name: S) -> Result<Self, GenericError> {
84 let name = name.into();
85 if name.is_empty() {
86 return Err(generic_error!("resolver name must not be empty"));
87 }
88
89 Ok(Self {
90 name,
91 caching_enabled: true,
92 cached_contexts_limit: None,
93 idle_context_expiration: None,
94 interner_capacity_bytes: None,
95 allow_heap_allocations: None,
96 tags_resolver: None,
97 interner: None,
98 origin_tags_resolver: None,
99 telemetry_enabled: true,
100 })
101 }
102
103 pub fn without_caching(mut self) -> Self {
118 self.caching_enabled = false;
119 self.idle_context_expiration = None;
120 self
121 }
122
123 pub fn with_cached_contexts_limit(mut self, limit: usize) -> Self {
138 match NonZeroUsize::new(limit) {
139 Some(limit) => {
140 self.cached_contexts_limit = Some(limit);
141 self
142 }
143 None => self.without_caching(),
144 }
145 }
146
147 pub fn with_idle_context_expiration(mut self, time_to_idle: Duration) -> Self {
156 self.idle_context_expiration = Some(time_to_idle);
157 self
158 }
159
160 pub fn with_interner_capacity_bytes(mut self, capacity: NonZeroUsize) -> Self {
176 self.interner_capacity_bytes = Some(capacity);
177 self
178 }
179
180 pub fn with_heap_allocations(mut self, allow: bool) -> Self {
189 self.allow_heap_allocations = Some(allow);
190 self
191 }
192
193 pub fn with_tags_resolver(mut self, resolver: Option<TagsResolver>) -> Self {
197 self.tags_resolver = resolver;
198 self
199 }
200
201 pub fn without_telemetry(mut self) -> Self {
210 self.telemetry_enabled = false;
211 self
212 }
213
214 pub fn with_interner(mut self, interner: GenericMapInterner) -> Self {
218 self.interner = Some(interner);
219 self
220 }
221
222 pub fn for_tests() -> Self {
235 ContextResolverBuilder::from_name("noop")
236 .expect("resolver name not empty")
237 .with_cached_contexts_limit(usize::MAX)
238 .with_interner_capacity_bytes(NonZeroUsize::new(1).expect("not zero"))
239 .with_heap_allocations(true)
240 .with_tags_resolver(Some(TagsResolverBuilder::for_tests().build()))
241 .without_telemetry()
242 }
243
244 pub fn build(self) -> ContextResolver {
246 let interner_capacity_bytes = self
247 .interner_capacity_bytes
248 .unwrap_or(DEFAULT_CONTEXT_RESOLVER_INTERNER_CAPACITY_BYTES);
249
250 let interner = match self.interner {
251 Some(interner) => interner,
252 None => GenericMapInterner::new(interner_capacity_bytes),
253 };
254
255 let cached_context_limit = self
256 .cached_contexts_limit
257 .unwrap_or(DEFAULT_CONTEXT_RESOLVER_CACHED_CONTEXTS_LIMIT);
258
259 let allow_heap_allocations = self.allow_heap_allocations.unwrap_or(true);
260
261 let telemetry = Telemetry::new(&self.name);
262 telemetry
263 .interner_capacity_bytes()
264 .set(interner.capacity_bytes() as f64);
265
266 let context_cache = CacheBuilder::from_identifier(format!("{}/contexts", self.name))
268 .expect("cache identifier cannot possibly be empty")
269 .with_capacity(cached_context_limit)
270 .with_time_to_idle(self.idle_context_expiration)
271 .with_hasher::<NoopU64BuildHasher>()
272 .with_telemetry(self.telemetry_enabled)
273 .build();
274
275 let tags_resolver = match self.tags_resolver {
277 Some(tags_resolver) => tags_resolver,
278 None => TagsResolverBuilder::new(format!("{}/tags", self.name), interner.clone())
279 .expect("tags resolver name not empty")
280 .with_cached_tagsets_limit(cached_context_limit.get())
281 .with_idle_tagsets_expiration(self.idle_context_expiration.unwrap_or_default())
282 .with_heap_allocations(allow_heap_allocations)
283 .with_origin_tags_resolver(self.origin_tags_resolver.clone())
284 .build(),
285 };
286
287 if self.telemetry_enabled {
288 tokio::spawn(drive_telemetry(interner.clone(), telemetry.clone()));
289 }
290
291 ContextResolver {
292 telemetry,
293 interner,
294 caching_enabled: self.caching_enabled,
295 context_cache,
296 hash_seen_buffer: PrehashedHashSet::with_capacity_and_hasher(
297 SEEN_HASHSET_INITIAL_CAPACITY,
298 NoopU64BuildHasher,
299 ),
300 allow_heap_allocations,
301 tags_resolver,
302 }
303 }
304}
305
306pub struct ContextResolver {
330 telemetry: Telemetry,
331 interner: GenericMapInterner,
332 caching_enabled: bool,
333 context_cache: ContextCache,
334 hash_seen_buffer: PrehashedHashSet<u64>,
335 allow_heap_allocations: bool,
336 tags_resolver: TagsResolver,
337}
338
339impl ContextResolver {
340 fn intern<S>(&self, s: S) -> Option<MetaString>
341 where
342 S: AsRef<str> + CheapMetaString,
343 {
344 s.try_cheap_clone()
347 .or_else(|| self.interner.try_intern(s.as_ref()).map(MetaString::from))
348 .or_else(|| {
349 self.allow_heap_allocations.then(|| {
350 saluki_antithesis::sometimes!(
354 true,
355 "context string interner spilled to the heap (unbounded under default config)"
356 );
357 self.telemetry.intern_fallback_total().increment(1);
358 MetaString::from(s.as_ref())
359 })
360 })
361 }
362
363 fn create_context_key_with_host<N, H, I, I2, T, T2>(
364 &mut self, name: N, host: Option<&H>, tags: I, origin_tags: I2,
365 ) -> (ContextKey, TagSetKey)
366 where
367 N: AsRef<str>,
368 H: AsRef<str>,
369 I: IntoIterator<Item = T>,
370 T: AsRef<str>,
371 I2: IntoIterator<Item = T2>,
372 T2: AsRef<str>,
373 {
374 hash_context_with_host_and_seen(
375 name.as_ref(),
376 host.map(AsRef::as_ref),
377 tags,
378 origin_tags,
379 &mut self.hash_seen_buffer,
380 )
381 }
382
383 fn create_context<N, H>(
384 &self, key: ContextKey, name: N, host: Option<H>, context_tags: SharedTagSet, origin_tags: SharedTagSet,
385 ) -> Option<Context>
386 where
387 N: AsRef<str> + CheapMetaString,
388 H: AsRef<str> + CheapMetaString,
389 {
390 let context_name = self.intern(name)?;
392 let context_host = match host {
393 Some(host) => Some(self.intern(host)?),
394 None => None,
395 };
396
397 self.telemetry.resolved_new_context_total().increment(1);
398 self.telemetry.active_contexts().increment(1);
399
400 Some(Context::from_inner(ContextInner::from_parts(
401 key,
402 context_name,
403 context_host,
404 context_tags.into(),
405 origin_tags.into(),
406 self.telemetry.active_contexts().clone(),
407 )))
408 }
409
410 pub fn resolve<N, I, T>(&mut self, name: N, tags: I, maybe_origin: Option<RawOrigin<'_>>) -> Option<Context>
419 where
420 N: AsRef<str> + CheapMetaString,
421 I: IntoIterator<Item = T> + Clone,
422 T: AsRef<str> + CheapMetaString,
423 {
424 let origin_tags = self.tags_resolver.resolve_origin_tags(maybe_origin);
426
427 self.resolve_inner(name, tags, origin_tags)
428 }
429
430 pub fn resolve_with_origin_tags<N, I, T>(
449 &mut self, name: N, tags: I, origin_tags: impl Into<SharedTagSet>,
450 ) -> Option<Context>
451 where
452 N: AsRef<str> + CheapMetaString,
453 I: IntoIterator<Item = T> + Clone,
454 T: AsRef<str> + CheapMetaString,
455 {
456 self.resolve_inner(name, tags, origin_tags.into())
457 }
458
459 pub fn resolve_with_host_and_origin_tags<N, H, I, T>(
461 &mut self, name: N, host: H, tags: I, origin_tags: impl Into<SharedTagSet>,
462 ) -> Option<Context>
463 where
464 N: AsRef<str> + CheapMetaString,
465 H: AsRef<str> + CheapMetaString,
466 I: IntoIterator<Item = T> + Clone,
467 T: AsRef<str> + CheapMetaString,
468 {
469 self.resolve_inner_with_host(name, Some(host), tags, origin_tags.into())
470 }
471
472 pub fn resolve_with_optional_host_and_origin_tags<N, H, I, T>(
474 &mut self, name: N, host: Option<H>, tags: I, origin_tags: impl Into<SharedTagSet>,
475 ) -> Option<Context>
476 where
477 N: AsRef<str> + CheapMetaString,
478 H: AsRef<str> + CheapMetaString,
479 I: IntoIterator<Item = T> + Clone,
480 T: AsRef<str> + CheapMetaString,
481 {
482 self.resolve_inner_with_host(name, host, tags, origin_tags.into())
483 }
484
485 pub fn resolve_with_host<N, H, I, T>(
489 &mut self, name: N, host: H, tags: I, maybe_origin: Option<RawOrigin<'_>>,
490 ) -> Option<Context>
491 where
492 N: AsRef<str> + CheapMetaString,
493 H: AsRef<str> + CheapMetaString,
494 I: IntoIterator<Item = T> + Clone,
495 T: AsRef<str> + CheapMetaString,
496 {
497 self.resolve_with_optional_host(name, Some(host), tags, maybe_origin)
498 }
499
500 pub fn resolve_with_optional_host<N, H, I, T>(
502 &mut self, name: N, host: Option<H>, tags: I, maybe_origin: Option<RawOrigin<'_>>,
503 ) -> Option<Context>
504 where
505 N: AsRef<str> + CheapMetaString,
506 H: AsRef<str> + CheapMetaString,
507 I: IntoIterator<Item = T> + Clone,
508 T: AsRef<str> + CheapMetaString,
509 {
510 let origin_tags = self.tags_resolver.resolve_origin_tags(maybe_origin);
511
512 self.resolve_inner_with_host(name, host, tags, origin_tags)
513 }
514
515 fn resolve_inner<N, I, T>(&mut self, name: N, tags: I, origin_tags: SharedTagSet) -> Option<Context>
516 where
517 N: AsRef<str> + CheapMetaString,
518 I: IntoIterator<Item = T> + Clone,
519 T: AsRef<str> + CheapMetaString,
520 {
521 self.resolve_inner_with_host(name, None::<&str>, tags, origin_tags)
522 }
523
524 fn resolve_inner_with_host<N, H, I, T>(
525 &mut self, name: N, host: Option<H>, tags: I, origin_tags: SharedTagSet,
526 ) -> Option<Context>
527 where
528 N: AsRef<str> + CheapMetaString,
529 H: AsRef<str> + CheapMetaString,
530 I: IntoIterator<Item = T> + Clone,
531 T: AsRef<str> + CheapMetaString,
532 {
533 let (context_key, tagset_key) =
534 self.create_context_key_with_host(&name, host.as_ref(), tags.clone(), &origin_tags);
535
536 if !self.caching_enabled {
538 let tag_set = self.tags_resolver.create_tag_set(tags).unwrap_or_default();
539
540 let context = self.create_context(context_key, name, host, tag_set, origin_tags)?;
541
542 debug!(?context_key, ?context, "Resolved new non-cached context.");
543 return Some(context);
544 }
545
546 match self.context_cache.get(&context_key) {
547 Some(context) => {
548 self.telemetry.resolved_existing_context_total().increment(1);
549 Some(context)
550 }
551 None => {
552 let tag_set = match self.tags_resolver.get_tag_set(tagset_key) {
554 Some(tag_set) => {
555 self.telemetry.resolved_existing_tagset_total().increment(1);
556 tag_set
557 }
558 None => {
559 let tag_set = self.tags_resolver.create_tag_set(tags.clone()).unwrap_or_default();
561
562 self.tags_resolver.insert_tag_set(tagset_key, tag_set.clone());
563
564 tag_set
565 }
566 };
567
568 let context = self.create_context(context_key, name, host, tag_set, origin_tags)?;
569 self.context_cache.insert(context_key, context.clone());
570
571 debug!(?context_key, ?context, "Resolved new context.");
572 Some(context)
573 }
574 }
575 }
576}
577
578impl Clone for ContextResolver {
579 fn clone(&self) -> Self {
580 Self {
581 telemetry: self.telemetry.clone(),
582 interner: self.interner.clone(),
583 caching_enabled: self.caching_enabled,
584 context_cache: self.context_cache.clone(),
585 hash_seen_buffer: PrehashedHashSet::with_capacity_and_hasher(
586 SEEN_HASHSET_INITIAL_CAPACITY,
587 NoopU64BuildHasher,
588 ),
589 allow_heap_allocations: self.allow_heap_allocations,
590 tags_resolver: self.tags_resolver.clone(),
591 }
592 }
593}
594
595async fn drive_telemetry(interner: GenericMapInterner, telemetry: Telemetry) {
596 loop {
597 sleep(Duration::from_secs(1)).await;
598
599 telemetry.interner_entries().set(interner.len() as f64);
600 telemetry
601 .interner_capacity_bytes()
602 .set(interner.capacity_bytes() as f64);
603 telemetry.interner_len_bytes().set(interner.len_bytes() as f64);
604 }
605}
606
607pub struct TagsResolverBuilder {
609 name: String,
610 caching_enabled: bool,
611 cached_tagset_limit: Option<NonZeroUsize>,
612 idle_tagset_expiration: Option<Duration>,
613 allow_heap_allocations: Option<bool>,
614 origin_tags_resolver: Option<Arc<dyn OriginTagsResolver>>,
615 telemetry_enabled: bool,
616 interner: GenericMapInterner,
617}
618
619impl TagsResolverBuilder {
620 pub fn new<S: Into<String>>(name: S, interner: GenericMapInterner) -> Result<Self, GenericError> {
622 let name = name.into();
623 if name.is_empty() {
624 return Err(generic_error!("resolver name must not be empty"));
625 }
626
627 Ok(Self {
628 name,
629 caching_enabled: true,
630 cached_tagset_limit: None,
631 idle_tagset_expiration: None,
632 allow_heap_allocations: None,
633 origin_tags_resolver: None,
634 telemetry_enabled: true,
635 interner,
636 })
637 }
638
639 pub fn with_interner(mut self, interner: GenericMapInterner) -> Self {
645 self.interner = interner;
646 self
647 }
648
649 pub fn without_caching(mut self) -> Self {
664 self.caching_enabled = false;
665 self.idle_tagset_expiration = None;
666 self
667 }
668
669 pub fn with_cached_tagsets_limit(mut self, limit: usize) -> Self {
684 match NonZeroUsize::new(limit) {
685 Some(limit) => {
686 self.cached_tagset_limit = Some(limit);
687 self
688 }
689 None => self.without_caching(),
690 }
691 }
692
693 pub fn with_idle_tagsets_expiration(mut self, time_to_idle: Duration) -> Self {
702 self.idle_tagset_expiration = Some(time_to_idle);
703 self
704 }
705
706 pub fn with_heap_allocations(mut self, allow: bool) -> Self {
715 self.allow_heap_allocations = Some(allow);
716 self
717 }
718
719 pub fn with_origin_tags_resolver(mut self, resolver: Option<Arc<dyn OriginTagsResolver>>) -> Self {
732 self.origin_tags_resolver = resolver;
733 self
734 }
735
736 pub fn without_telemetry(mut self) -> Self {
745 self.telemetry_enabled = false;
746 self
747 }
748
749 pub fn build(self) -> TagsResolver {
751 let cached_tagsets_limit = self
752 .cached_tagset_limit
753 .unwrap_or(DEFAULT_CONTEXT_RESOLVER_CACHED_CONTEXTS_LIMIT);
754
755 let allow_heap_allocations = self.allow_heap_allocations.unwrap_or(true);
756
757 let telemetry = Telemetry::new(self.name.clone());
758 telemetry
759 .interner_capacity_bytes()
760 .set(self.interner.capacity_bytes() as f64);
761
762 let tagset_cache = CacheBuilder::from_identifier(format!("{}/tagsets", self.name))
763 .expect("cache identifier cannot possibly be empty")
764 .with_capacity(cached_tagsets_limit)
765 .with_time_to_idle(self.idle_tagset_expiration)
766 .with_hasher::<NoopU64BuildHasher>()
767 .with_telemetry(self.telemetry_enabled)
768 .build();
769
770 TagsResolver {
771 telemetry,
772 interner: self.interner,
773 caching_enabled: self.caching_enabled,
774 tagset_cache,
775 origin_tags_resolver: self.origin_tags_resolver,
776 allow_heap_allocations,
777 }
778 }
779
780 pub fn for_tests() -> Self {
793 TagsResolverBuilder::new("noop", GenericMapInterner::new(NonZeroUsize::new(1).expect("not zero")))
794 .expect("resolver name not empty")
795 .with_cached_tagsets_limit(usize::MAX)
796 .with_heap_allocations(true)
797 .without_telemetry()
798 }
799}
800
801pub struct TagsResolver {
803 telemetry: Telemetry,
804 interner: GenericMapInterner,
805 caching_enabled: bool,
806 tagset_cache: TagSetCache,
807 origin_tags_resolver: Option<Arc<dyn OriginTagsResolver>>,
808 allow_heap_allocations: bool,
809}
810
811impl TagsResolver {
812 fn intern<S>(&self, s: S) -> Option<MetaString>
813 where
814 S: AsRef<str> + CheapMetaString,
815 {
816 s.try_cheap_clone()
819 .or_else(|| self.interner.try_intern(s.as_ref()).map(MetaString::from))
820 .or_else(|| {
821 self.allow_heap_allocations.then(|| {
822 saluki_antithesis::sometimes!(
826 true,
827 "tag string interner spilled to the heap (unbounded under default config)"
828 );
829 self.telemetry.intern_fallback_total().increment(1);
830 MetaString::from(s.as_ref())
831 })
832 })
833 }
834
835 pub fn create_tag_set<I, T>(&mut self, tags: I) -> Option<SharedTagSet>
842 where
843 I: IntoIterator<Item = T>,
844 T: AsRef<str> + CheapMetaString,
845 {
846 let mut tag_set = TagSet::default();
847 for tag in tags {
848 let tag = self.intern(tag)?;
849 tag_set.insert_tag(tag);
850 }
851
852 self.telemetry.resolved_new_tagset_total().increment(1);
853
854 Some(tag_set.into_shared())
855 }
856
857 pub fn resolve_origin_tags(&self, maybe_origin: Option<RawOrigin<'_>>) -> SharedTagSet {
861 self.origin_tags_resolver
862 .as_ref()
863 .and_then(|resolver| maybe_origin.map(|origin| resolver.resolve_origin_tags(origin)))
864 .unwrap_or_default()
865 }
866
867 fn get_tag_set(&self, key: TagSetKey) -> Option<SharedTagSet> {
868 self.tagset_cache.get(&key)
869 }
870
871 fn insert_tag_set(&self, key: TagSetKey, tag_set: SharedTagSet) {
872 self.tagset_cache.insert(key, tag_set);
873 }
874}
875
876impl Clone for TagsResolver {
877 fn clone(&self) -> Self {
878 Self {
879 telemetry: self.telemetry.clone(),
880 interner: self.interner.clone(),
881 caching_enabled: self.caching_enabled,
882 tagset_cache: self.tagset_cache.clone(),
883 origin_tags_resolver: self.origin_tags_resolver.clone(),
884 allow_heap_allocations: self.allow_heap_allocations,
885 }
886 }
887}
888
889#[cfg(test)]
890mod tests {
891 use metrics::{SharedString, Unit};
892 use metrics_util::{
893 debugging::{DebugValue, DebuggingRecorder},
894 CompositeKey,
895 };
896 use saluki_common::hash::hash_single_fast;
897
898 use super::*;
899 use crate::data_model::tags::Tag;
900
901 fn get_gauge_value(metrics: &[(CompositeKey, Option<Unit>, Option<SharedString>, DebugValue)], key: &str) -> f64 {
902 metrics
903 .iter()
904 .find(|(k, _, _, _)| k.key().name() == key)
905 .map(|(_, _, _, value)| match value {
906 DebugValue::Gauge(value) => value.into_inner(),
907 other => panic!("expected a gauge, got: {:?}", other),
908 })
909 .unwrap_or_else(|| panic!("no metric found with key: {}", key))
910 }
911
912 struct DummyOriginTagsResolver;
913
914 impl OriginTagsResolver for DummyOriginTagsResolver {
915 fn resolve_origin_tags(&self, origin: RawOrigin<'_>) -> SharedTagSet {
916 let origin_key = hash_single_fast(origin);
917
918 let mut tags = TagSet::default();
919 tags.insert_tag(format!("origin_key:{}", origin_key));
920 tags.into_shared()
921 }
922 }
923
924 #[test]
925 fn basic() {
926 let mut resolver = ContextResolverBuilder::for_tests().build();
927
928 let name = "metric_name";
930 let tags1: [&str; 0] = [];
931 let tags2 = ["tag1"];
932
933 assert_ne!(&tags1[..], &tags2[..]);
934
935 let context1 = resolver
936 .resolve(name, &tags1[..], None)
937 .expect("should not fail to resolve");
938 let context2 = resolver
939 .resolve(name, &tags2[..], None)
940 .expect("should not fail to resolve");
941
942 assert_ne!(context1, context2);
945 assert!(!context1.ptr_eq(&context2));
946
947 let context1_redo = resolver
949 .resolve(name, &tags1[..], None)
950 .expect("should not fail to resolve");
951 let context2_redo = resolver
952 .resolve(name, &tags2[..], None)
953 .expect("should not fail to resolve");
954
955 assert_ne!(context1_redo, context2_redo);
956 assert_eq!(context1, context1_redo);
957 assert_eq!(context2, context2_redo);
958 assert!(context1.ptr_eq(&context1_redo));
959 assert!(context2.ptr_eq(&context2_redo));
960 }
961
962 #[test]
963 fn tag_order() {
964 let mut resolver = ContextResolverBuilder::for_tests().build();
965
966 let name = "metric_name";
968 let tags1 = ["tag1", "tag2"];
969 let tags2 = ["tag2", "tag1"];
970
971 assert_ne!(&tags1[..], &tags2[..]);
972
973 let context1 = resolver
974 .resolve(name, &tags1[..], None)
975 .expect("should not fail to resolve");
976 let context2 = resolver
977 .resolve(name, &tags2[..], None)
978 .expect("should not fail to resolve");
979
980 assert_eq!(context1, context2);
983 assert!(context1.ptr_eq(&context2));
984 }
985
986 #[test]
987 fn host_affects_identity_but_not_visible_tags() {
988 let mut resolver = ContextResolverBuilder::for_tests().build();
989
990 let context1 = resolver
991 .resolve_with_host("metric_name", "host-a", &[] as &[&str], None)
992 .expect("should not fail to resolve");
993 let context2 = resolver
994 .resolve_with_host("metric_name", "host-b", &[] as &[&str], None)
995 .expect("should not fail to resolve");
996 let context1_redo = resolver
997 .resolve_with_host("metric_name", "host-a", &[] as &[&str], None)
998 .expect("should not fail to resolve");
999
1000 assert_ne!(context1, context2);
1001 assert_eq!(context1, context1_redo);
1002 assert!(context1.ptr_eq(&context1_redo));
1003 assert_eq!(context1.host(), Some("host-a"));
1004 assert_eq!(context2.host(), Some("host-b"));
1005 assert!(context1.tags().is_empty());
1006 assert!(context2.tags().is_empty());
1007
1008 let mut uncached_resolver = ContextResolverBuilder::for_tests().without_caching().build();
1009 let uncached1 = uncached_resolver
1010 .resolve_with_host("metric_name", "host-a", &[] as &[&str], None)
1011 .expect("should not fail to resolve");
1012 let uncached2 = uncached_resolver
1013 .resolve_with_host("metric_name", "host-b", &[] as &[&str], None)
1014 .expect("should not fail to resolve");
1015
1016 assert_ne!(uncached1, uncached2);
1017 assert!(!uncached1.ptr_eq(&uncached2));
1018 }
1019
1020 #[test]
1021 fn host_survives_rewrites() {
1022 let mut resolver = ContextResolverBuilder::for_tests().build();
1023
1024 let context1 = resolver
1025 .resolve_with_host("metric_name", "host-a", &["env:prod"][..], None)
1026 .expect("should not fail to resolve");
1027 let context2 = resolver
1028 .resolve_with_host("metric_name", "host-b", &["env:prod"][..], None)
1029 .expect("should not fail to resolve");
1030
1031 assert_ne!(context1, context2);
1032 let service_tag_set = TagSet::from_iter([Tag::from("service:api")]);
1033 assert_ne!(context1.with_name("renamed"), context2.with_name("renamed"));
1034 assert_ne!(
1035 context1.with_tags(service_tag_set.clone()),
1036 context2.with_tags(service_tag_set)
1037 );
1038
1039 let mut context1_filtered = context1.clone();
1040 let mut context2_filtered = context2.clone();
1041 let mut state1 = crate::data_model::event::metric::context::TagSetMutViewState::new();
1042 let mut state2 = crate::data_model::event::metric::context::TagSetMutViewState::new();
1043 {
1044 let mut view = context1_filtered.tags_mut_view(&mut state1);
1045 view.retain_tags(|_| false);
1046 view.finish();
1047 }
1048 {
1049 let mut view = context2_filtered.tags_mut_view(&mut state2);
1050 view.retain_tags(|_| false);
1051 view.finish();
1052 }
1053
1054 assert_ne!(context1_filtered, context2_filtered);
1055 assert!(context1_filtered.tags().is_empty());
1056 assert!(context2_filtered.tags().is_empty());
1057 assert_eq!(context1_filtered.host(), Some("host-a"));
1058 assert_eq!(context2_filtered.host(), Some("host-b"));
1059 }
1060
1061 #[test]
1062 fn active_contexts() {
1063 let recorder = DebuggingRecorder::new();
1064 let snapshotter = recorder.snapshotter();
1065
1066 let context = metrics::with_local_recorder(&recorder, || {
1068 let mut resolver = ContextResolverBuilder::for_tests().build();
1069 resolver
1070 .resolve("name", &["tag"][..], None)
1071 .expect("should not fail to resolve")
1072 });
1073
1074 let metrics_before = snapshotter.snapshot().into_vec();
1076 let active_contexts = get_gauge_value(&metrics_before, Telemetry::active_contexts_name());
1077 assert_eq!(active_contexts, 1.0);
1078
1079 drop(context);
1081 let metrics_after = snapshotter.snapshot().into_vec();
1082 let active_contexts = get_gauge_value(&metrics_after, Telemetry::active_contexts_name());
1083 assert_eq!(active_contexts, -1.0);
1084 }
1085
1086 #[test]
1087 fn duplicate_tags() {
1088 let mut resolver = ContextResolverBuilder::for_tests().build();
1089
1090 let name = "metric_name";
1092 let tags1 = ["tag1"];
1093 let tags1_duplicated = ["tag1", "tag1"];
1094 let tags2 = ["tag2"];
1095 let tags2_duplicated = ["tag2", "tag2"];
1096
1097 let context1 = resolver
1098 .resolve(name, &tags1[..], None)
1099 .expect("should not fail to resolve");
1100 let context1_duplicated = resolver
1101 .resolve(name, &tags1_duplicated[..], None)
1102 .expect("should not fail to resolve");
1103 let context2 = resolver
1104 .resolve(name, &tags2[..], None)
1105 .expect("should not fail to resolve");
1106 let context2_duplicated = resolver
1107 .resolve(name, &tags2_duplicated[..], None)
1108 .expect("should not fail to resolve");
1109
1110 assert_eq!(context1, context1_duplicated);
1112 assert_eq!(context2, context2_duplicated);
1113
1114 assert_ne!(context1, context2);
1123 assert_ne!(context1_duplicated, context2_duplicated);
1124 assert_ne!(context1, context2_duplicated);
1125 assert_ne!(context2, context1_duplicated);
1126 }
1127
1128 #[test]
1129 fn differing_origins_with_without_resolver() {
1130 let mut resolver = ContextResolverBuilder::for_tests().build();
1133
1134 let name = "metric_name";
1135 let tags = ["tag1"];
1136 let mut origin1 = RawOrigin::default();
1137 origin1.set_local_data("container1");
1138 let mut origin2 = RawOrigin::default();
1139 origin2.set_local_data("container2");
1140
1141 let context1 = resolver
1142 .resolve(name, &tags[..], Some(origin1.clone()))
1143 .expect("should not fail to resolve");
1144 let context2 = resolver
1145 .resolve(name, &tags[..], Some(origin2.clone()))
1146 .expect("should not fail to resolve");
1147
1148 assert_eq!(context1, context2);
1149
1150 let tags_resolver = TagsResolverBuilder::for_tests()
1151 .with_origin_tags_resolver(Some(Arc::new(DummyOriginTagsResolver)))
1152 .build();
1153 let mut resolver = ContextResolverBuilder::for_tests()
1157 .with_tags_resolver(Some(tags_resolver))
1158 .build();
1159
1160 let context1 = resolver
1161 .resolve(name, &tags[..], Some(origin1))
1162 .expect("should not fail to resolve");
1163 let context2 = resolver
1164 .resolve(name, &tags[..], Some(origin2))
1165 .expect("should not fail to resolve");
1166
1167 assert_ne!(context1, context2);
1168 }
1169
1170 #[test]
1171 fn caching_disabled() {
1172 let tags_resolver = TagsResolverBuilder::for_tests()
1173 .with_origin_tags_resolver(Some(Arc::new(DummyOriginTagsResolver)))
1174 .build();
1175 let mut resolver = ContextResolverBuilder::for_tests()
1176 .without_caching()
1177 .with_tags_resolver(Some(tags_resolver))
1178 .build();
1179
1180 let name = "metric_name";
1181 let tags = ["tag1"];
1182 let mut origin1 = RawOrigin::default();
1183 origin1.set_local_data("container1");
1184
1185 let context1 = resolver
1187 .resolve(name, &tags[..], Some(origin1.clone()))
1188 .expect("should not fail to resolve");
1189 assert_eq!(resolver.context_cache.len(), 0);
1190
1191 let context2 = resolver
1193 .resolve(name, &tags[..], Some(origin1))
1194 .expect("should not fail to resolve");
1195 assert_eq!(resolver.context_cache.len(), 0);
1196
1197 assert_eq!(context1, context2);
1200 assert!(!context1.ptr_eq(&context2));
1201 }
1202
1203 #[test]
1204 fn cheaply_cloneable_name_and_tags() {
1205 const BIG_TAG_ONE: &str = "long-tag-that-cannot-be-inlined-just-to-be-doubly-sure-on-top-of-being-static";
1206 const BIG_TAG_TWO: &str = "another-long-boye-that-we-are-also-sure-wont-be-inlined-and-we-stand-on-that";
1207
1208 let mut resolver = ContextResolverBuilder::for_tests()
1210 .with_interner_capacity_bytes(NonZeroUsize::new(1024).expect("not zero"))
1211 .build();
1212
1213 let name = MetaString::from_static("long-metric-name-that-shouldnt-be-inlined-and-should-end-up-interned");
1215 let tags = [
1216 MetaString::from_static(BIG_TAG_ONE),
1217 MetaString::from_static(BIG_TAG_TWO),
1218 ];
1219 assert!(tags[0].is_cheaply_cloneable());
1220 assert!(tags[1].is_cheaply_cloneable());
1221
1222 assert_eq!(resolver.interner.len(), 0);
1225 assert_eq!(resolver.interner.len_bytes(), 0);
1226
1227 let context = resolver
1228 .resolve(&name, &tags[..], None)
1229 .expect("should not fail to resolve");
1230 assert_eq!(resolver.interner.len(), 0);
1231 assert_eq!(resolver.interner.len_bytes(), 0);
1232
1233 assert_eq!(context.name(), &name);
1235
1236 let context_tags = context.tags();
1237 assert_eq!(context_tags.len(), 2);
1238 assert!(context_tags.has_tag(&tags[0]));
1239 assert!(context_tags.has_tag(&tags[1]));
1240 }
1241}