1use std::{
8 num::NonZeroUsize,
9 pin::Pin,
10 sync::{
11 atomic::{AtomicBool, AtomicU32, Ordering},
12 Arc, LazyLock, Mutex, OnceLock,
13 },
14 task::{self, ready, Poll},
15 time::Duration,
16};
17
18use async_trait::async_trait;
19use futures::Stream;
20use metrics::{
21 atomics::AtomicU64, Counter, CounterFn, Gauge, GaugeFn, Histogram, HistogramFn, Key, KeyName, Level, Metadata,
22 Recorder, SetRecorderError, SharedString, Unit,
23};
24use metrics_util::storage::AtomicBucket;
25use saluki_common::{collections::FastHashMap, sync::shutdown::ShutdownHandle};
26use saluki_error::GenericError;
27use tokio::{
28 select,
29 sync::broadcast::{self, error::RecvError, Receiver},
30};
31use tokio_util::sync::ReusableBoxFuture;
32use tracing::debug;
33
34use crate::{
35 data_model::{
36 event::{
37 metric::{
38 context::{Context, ContextResolver, ContextResolverBuilder},
39 *,
40 },
41 Event,
42 },
43 origin::RawOrigin,
44 tags::{Tag, TagSet},
45 },
46 runtime::{InitializationError, Supervisable, SupervisorFuture},
47};
48
49mod aggregated;
50pub use self::aggregated::{
51 get_shared_metrics_state, initialize_shared_metrics_state, AggregatedMetricValue, AggregatedMetricsProcessor,
52 AggregatedMetricsState, SharedMetricsWorker,
53};
54
55mod histogram;
56pub use self::histogram::AggregatedHistogram;
57
58mod processor;
59pub use self::processor::TelemetryProcessor;
60
61mod reflector;
62pub use self::reflector::{Processor, Reflector, ReflectorWorker};
63
64mod remapper;
65pub use self::remapper::{RemappedMetric, RemapperRule};
66
67const FLUSH_INTERVAL: Duration = Duration::from_secs(1);
68const INTERNAL_METRICS_INTERNER_SIZE: NonZeroUsize = NonZeroUsize::new(16_384).unwrap();
69
70const IDLE_EVICT_INTERVALS: u32 = 3;
78
79static RECEIVER_STATE: OnceLock<Arc<State>> = OnceLock::new();
80
81#[derive(Clone, Default)]
83pub struct MetricsSnapshot {
84 pub upserts: Vec<Event>,
86
87 pub evictions: Vec<Context>,
89}
90
91pub type SharedMetricsSnapshot = Arc<MetricsSnapshot>;
93
94struct Handle<T> {
95 inner: T,
96 level: Level,
97 idle: AtomicU32,
98}
99
100impl<T> Handle<T> {
101 const fn new(inner: T, level: Level) -> Self {
102 Self {
103 inner,
104 level,
105 idle: AtomicU32::new(0),
106 }
107 }
108
109 fn level(&self) -> Level {
110 self.level
111 }
112
113 fn reset_idle(&self) {
114 self.idle.store(0, Ordering::Relaxed);
115 }
116
117 fn bump_idle(&self) -> u32 {
118 self.idle.fetch_add(1, Ordering::Relaxed).saturating_add(1)
119 }
120}
121
122struct CounterInner {
123 value: AtomicU64,
124}
125
126impl Handle<CounterInner> {
127 const fn new_counter(level: Level) -> Self {
128 Self::new(
129 CounterInner {
130 value: AtomicU64::new(0),
131 },
132 level,
133 )
134 }
135
136 fn consume(&self) -> u64 {
138 self.inner.value.swap(0, Ordering::Relaxed)
139 }
140}
141
142impl CounterFn for Handle<CounterInner> {
143 fn increment(&self, value: u64) {
144 CounterFn::increment(&self.inner.value, value);
145 }
146
147 fn absolute(&self, value: u64) {
148 CounterFn::absolute(&self.inner.value, value);
149 }
150}
151
152struct GaugeInner {
153 value: AtomicU64,
154 dirty: AtomicBool,
158}
159
160impl Handle<GaugeInner> {
161 const fn new_gauge(level: Level) -> Self {
162 Self::new(
163 GaugeInner {
164 value: AtomicU64::new(0),
165 dirty: AtomicBool::new(false),
166 },
167 level,
168 )
169 }
170
171 fn load(&self) -> f64 {
173 f64::from_bits(self.inner.value.load(Ordering::Relaxed))
174 }
175
176 fn take_dirty(&self) -> bool {
178 self.inner.dirty.swap(false, Ordering::Relaxed)
179 }
180}
181
182impl GaugeFn for Handle<GaugeInner> {
183 fn increment(&self, value: f64) {
184 GaugeFn::increment(&self.inner.value, value);
185 self.inner.dirty.store(true, Ordering::Relaxed);
186 }
187
188 fn decrement(&self, value: f64) {
189 GaugeFn::decrement(&self.inner.value, value);
190 self.inner.dirty.store(true, Ordering::Relaxed);
191 }
192
193 fn set(&self, value: f64) {
194 GaugeFn::set(&self.inner.value, value);
195 self.inner.dirty.store(true, Ordering::Relaxed);
196 }
197}
198
199struct HistogramInner {
201 value: AtomicBucket<f64>,
202}
203
204impl Handle<HistogramInner> {
205 fn new_histogram(level: Level) -> Self {
206 Self::new(
207 HistogramInner {
208 value: AtomicBucket::new(),
209 },
210 level,
211 )
212 }
213
214 fn drain_into(&self, out: &mut Vec<f64>) {
216 self.inner.value.clear_with(|samples| out.extend(samples));
217 }
218}
219
220impl HistogramFn for Handle<HistogramInner> {
221 fn record(&self, value: f64) {
222 self.inner.value.push(value);
223 }
224}
225
226#[derive(Default)]
232struct MetricsRegistry {
233 maps: Mutex<RegistryMaps>,
234}
235
236#[derive(Default)]
237struct RegistryMaps {
238 counters: FastHashMap<Key, Arc<Handle<CounterInner>>>,
239 gauges: FastHashMap<Key, Arc<Handle<GaugeInner>>>,
240 histograms: FastHashMap<Key, Arc<Handle<HistogramInner>>>,
241}
242
243impl MetricsRegistry {
244 fn get_or_create_counter(&self, key: &Key, level: Level) -> Counter {
249 let mut maps = self.maps.lock().unwrap();
250 let handle = if let Some(handle) = maps.counters.get(key) {
251 handle.reset_idle();
252 Arc::clone(handle)
253 } else {
254 let handle = Arc::new(Handle::new_counter(level));
255 maps.counters.insert(key.clone(), Arc::clone(&handle));
256 handle
257 };
258
259 Counter::from_arc(handle)
260 }
261
262 fn remove_counter(&self, key: &Key) -> Option<Arc<Handle<CounterInner>>> {
263 let mut maps = self.maps.lock().unwrap();
264 maps.counters.remove(key)
265 }
266
267 fn get_or_create_gauge(&self, key: &Key, level: Level) -> Gauge {
269 let mut maps = self.maps.lock().unwrap();
270 let handle = if let Some(handle) = maps.gauges.get(key) {
271 handle.reset_idle();
272 Arc::clone(handle)
273 } else {
274 let handle = Arc::new(Handle::new_gauge(level));
275 maps.gauges.insert(key.clone(), Arc::clone(&handle));
276 handle
277 };
278
279 Gauge::from_arc(handle)
280 }
281
282 fn get_or_create_histogram(&self, key: &Key, level: Level) -> Histogram {
284 let mut maps = self.maps.lock().unwrap();
285 let handle = if let Some(handle) = maps.histograms.get(key) {
286 handle.reset_idle();
287 Arc::clone(handle)
288 } else {
289 let handle = Arc::new(Handle::new_histogram(level));
290 maps.histograms.insert(key.clone(), Arc::clone(&handle));
291 handle
292 };
293
294 Histogram::from_arc(handle)
295 }
296}
297
298pub struct FilterHandle {
303 state: Arc<State>,
304}
305
306impl FilterHandle {
307 pub fn override_filter(&self, level: Level) {
309 *self.state.current_level.lock().unwrap() = level;
310 }
311
312 pub fn reset_filter(&self) {
314 *self.state.current_level.lock().unwrap() = self.state.default_level;
315 }
316}
317
318struct State {
319 registry: MetricsRegistry,
320 flush_tx: broadcast::Sender<SharedMetricsSnapshot>,
321 metrics_prefix: String,
322 default_level: Level,
323 current_level: Mutex<Level>,
324 idle_evict_intervals: u32,
325 flush_interval: Duration,
326}
327
328struct MetricsRecorder {
329 state: Arc<State>,
330}
331
332impl MetricsRecorder {
333 fn new(metrics_prefix: String, default_level: Level) -> Self {
334 let (flush_tx, _) = broadcast::channel(2);
335 Self {
336 state: Arc::new(State {
337 registry: MetricsRegistry::default(),
338 flush_tx,
339 metrics_prefix,
340 default_level,
341 current_level: Mutex::new(default_level),
342 idle_evict_intervals: IDLE_EVICT_INTERVALS,
343 flush_interval: FLUSH_INTERVAL,
344 }),
345 }
346 }
347
348 fn filter_handle(&self) -> FilterHandle {
349 FilterHandle {
350 state: Arc::clone(&self.state),
351 }
352 }
353
354 fn install(self) -> Result<(), SetRecorderError<Self>> {
355 let state = Arc::clone(&self.state);
356 metrics::set_global_recorder(self)?;
357
358 if RECEIVER_STATE.set(state).is_err() {
359 panic!("metrics receiver should never be set prior to global recorder being installed");
360 }
361
362 Ok(())
363 }
364
365 fn prefix_key(&self, key: &Key) -> Key {
366 Key::from_parts(format!("{}.{}", self.state.metrics_prefix, key.name()), key.labels())
367 }
368}
369
370impl Recorder for MetricsRecorder {
371 fn describe_counter(&self, _: KeyName, _: Option<Unit>, _: SharedString) {}
372 fn describe_gauge(&self, _: KeyName, _: Option<Unit>, _: SharedString) {}
373 fn describe_histogram(&self, _: KeyName, _: Option<Unit>, _: SharedString) {}
374
375 fn register_counter(&self, key: &Key, metadata: &Metadata<'_>) -> Counter {
376 let prefixed_key = self.prefix_key(key);
377 self.state
378 .registry
379 .get_or_create_counter(&prefixed_key, *metadata.level())
380 }
381
382 fn register_gauge(&self, key: &Key, metadata: &Metadata<'_>) -> Gauge {
383 let prefixed_key = self.prefix_key(key);
384 self.state
385 .registry
386 .get_or_create_gauge(&prefixed_key, *metadata.level())
387 }
388
389 fn register_histogram(&self, key: &Key, metadata: &Metadata<'_>) -> Histogram {
390 let prefixed_key = self.prefix_key(key);
391 self.state
392 .registry
393 .get_or_create_histogram(&prefixed_key, *metadata.level())
394 }
395}
396
397pub fn delete_counter(key: Key) -> bool {
402 let Some(state) = RECEIVER_STATE.get() else {
403 return false;
404 };
405
406 let prefixed_key = Key::from_parts(format!("{}.{}", state.metrics_prefix, key.name()), key.labels());
407 let Some(snapshot) = delete_counter_from_state(state, &prefixed_key) else {
408 return false;
409 };
410
411 if state.flush_tx.receiver_count() > 0 {
412 let _ = state.flush_tx.send(Arc::new(snapshot));
413 }
414
415 true
416}
417
418fn delete_counter_from_state(state: &State, prefixed_key: &Key) -> Option<MetricsSnapshot> {
419 let handle = state.registry.remove_counter(prefixed_key)?;
420 let delta = handle.consume();
421 let context = context_from_key(prefixed_key);
422 let current_level = *state.current_level.lock().unwrap();
423
424 let mut snapshot = MetricsSnapshot {
425 upserts: Vec::new(),
426 evictions: vec![context.clone()],
427 };
428 if delta > 0 && handle.level() >= current_level {
429 snapshot
430 .upserts
431 .push(Event::Metric(Metric::counter(context, delta as f64)));
432 }
433
434 Some(snapshot)
435}
436
437fn context_from_key(key: &Key) -> Context {
438 let tags = key
439 .labels()
440 .map(|l| Tag::from(format!("{}:{}", l.key(), l.value())))
441 .collect::<TagSet>();
442
443 Context::from_parts(key.name().to_string(), tags)
444}
445
446pub struct MetricsStream {
451 inner: ReusableBoxFuture<
452 'static,
453 (
454 Result<SharedMetricsSnapshot, RecvError>,
455 Receiver<SharedMetricsSnapshot>,
456 ),
457 >,
458}
459
460impl MetricsStream {
461 pub fn register() -> Self {
463 let state = RECEIVER_STATE.get().expect("metrics receiver should be set");
464 Self {
465 inner: ReusableBoxFuture::new(make_rx_future(state.flush_tx.subscribe())),
466 }
467 }
468}
469
470impl Stream for MetricsStream {
471 type Item = SharedMetricsSnapshot;
472
473 fn poll_next(mut self: Pin<&mut Self>, cx: &mut task::Context<'_>) -> Poll<Option<Self::Item>> {
474 loop {
475 let (result, rx) = ready!(self.inner.poll(cx));
477 self.inner.set(make_rx_future(rx));
478
479 match result {
480 Ok(item) => return Poll::Ready(Some(item)),
481 Err(RecvError::Closed) => return Poll::Ready(None),
482 Err(RecvError::Lagged(n)) => {
483 debug!(
484 missed_payloads = n,
485 "Stream lagging behind internal metrics producer. Internal metrics may have been lost."
486 );
487 continue;
488 }
489 }
490 }
491 }
492}
493
494async fn make_rx_future(
495 mut rx: Receiver<SharedMetricsSnapshot>,
496) -> (
497 Result<SharedMetricsSnapshot, RecvError>,
498 Receiver<SharedMetricsSnapshot>,
499) {
500 let result = rx.recv().await;
501 (result, rx)
502}
503
504#[derive(Default)]
505struct FlushState {
506 counter_emits: Vec<(Key, f64)>,
507 gauge_emits: Vec<(Key, f64)>,
508 histogram_emits: Vec<(Key, Vec<f64>)>,
509 evicted_keys: Vec<Key>,
510}
511
512impl FlushState {
513 fn clear(&mut self) {
514 self.counter_emits.clear();
515 self.gauge_emits.clear();
516 self.histogram_emits.clear();
517 self.evicted_keys.clear();
518 }
519}
520
521async fn flush_metrics() {
522 let mut context_resolver = MetricsContextResolver::new(INTERNAL_METRICS_INTERNER_SIZE);
523
524 let state = RECEIVER_STATE.get().expect("metrics receiver should be set");
525
526 let mut flush_interval = tokio::time::interval(state.flush_interval);
527 flush_interval.tick().await;
528
529 let mut flush_state = FlushState::default();
530
531 loop {
532 flush_interval.tick().await;
533 flush_once(state, &mut context_resolver, &mut flush_state);
534 }
535}
536
537fn flush_once(state: &State, context_resolver: &mut MetricsContextResolver, flush_state: &mut FlushState) {
541 let has_listeners = state.flush_tx.receiver_count() > 0;
542 let current_level = *state.current_level.lock().unwrap();
543 let idle_threshold = state.idle_evict_intervals;
544 let mut snapshot = MetricsSnapshot::default();
545
546 flush_state.clear();
547 let FlushState {
548 counter_emits,
549 gauge_emits,
550 histogram_emits,
551 evicted_keys,
552 } = flush_state;
553
554 {
555 let mut maps = state.registry.maps.lock().unwrap();
556
557 maps.counters.retain(|key, handle| {
564 let idle = handle.bump_idle();
565 let orphaned = Arc::strong_count(handle) == 1;
566 if orphaned {
567 std::sync::atomic::fence(Ordering::Acquire);
568 }
569 let delta = handle.consume();
570
571 if orphaned && delta == 0 && idle >= idle_threshold {
572 evicted_keys.push(key.clone());
573 return false;
574 }
575
576 if has_listeners && handle.level() >= current_level {
577 counter_emits.push((key.clone(), delta as f64));
578 }
579 true
580 });
581
582 maps.gauges.retain(|key, handle| {
585 let idle = handle.bump_idle();
586 let orphaned = Arc::strong_count(handle) == 1;
587 if orphaned {
588 std::sync::atomic::fence(Ordering::Acquire);
589 }
590 let value = handle.load();
591 let written = handle.take_dirty();
592
593 if orphaned && !written && idle >= idle_threshold {
594 evicted_keys.push(key.clone());
595 return false;
596 }
597
598 if has_listeners && handle.level() >= current_level {
599 gauge_emits.push((key.clone(), value));
600 }
601 true
602 });
603
604 maps.histograms.retain(|key, handle| {
610 let idle = handle.bump_idle();
611 let orphaned = Arc::strong_count(handle) == 1;
612 if orphaned {
613 std::sync::atomic::fence(Ordering::Acquire);
614 }
615 let mut histogram_samples = Vec::new();
616 handle.drain_into(&mut histogram_samples);
617 let had_samples = !histogram_samples.is_empty();
618
619 if orphaned && !had_samples && idle >= idle_threshold {
620 evicted_keys.push(key.clone());
621 return false;
622 }
623
624 if has_listeners && had_samples && handle.level() >= current_level {
625 histogram_emits.push((key.clone(), histogram_samples));
626 }
627 true
628 });
629 }
630
631 for key in evicted_keys {
634 if let Some(context) = context_resolver.take(key) {
635 if has_listeners {
636 snapshot.evictions.push(context);
637 }
638 }
639 }
640
641 if !has_listeners {
642 return;
643 }
644
645 for (key, value) in counter_emits.drain(..) {
646 let context = context_resolver.resolve_from_key(key);
647 snapshot.upserts.push(Event::Metric(Metric::counter(context, value)));
648 }
649 for (key, value) in gauge_emits.drain(..) {
650 let context = context_resolver.resolve_from_key(key);
651 snapshot.upserts.push(Event::Metric(Metric::gauge(context, value)));
652 }
653 for (key, samples) in histogram_emits.drain(..) {
654 let context = context_resolver.resolve_from_key(key);
655 snapshot
656 .upserts
657 .push(Event::Metric(Metric::histogram(context, &samples[..])));
658 }
659
660 if !snapshot.upserts.is_empty() || !snapshot.evictions.is_empty() {
661 let _ = state.flush_tx.send(Arc::new(snapshot));
662 }
663}
664
665struct MetricsContextResolver {
666 context_resolver: ContextResolver,
667 key_context_cache: FastHashMap<Key, Context>,
668}
669
670impl MetricsContextResolver {
671 fn new(resolver_interner_size_bytes: NonZeroUsize) -> Self {
672 Self {
673 context_resolver: ContextResolverBuilder::from_name("core/internal_metrics")
675 .expect("resolver name is not empty")
676 .with_interner_capacity_bytes(resolver_interner_size_bytes)
677 .without_caching()
678 .build(),
679 key_context_cache: FastHashMap::default(),
680 }
681 }
682
683 fn resolve_from_key(&mut self, key: Key) -> Context {
684 static SELF_ORIGIN_INFO: LazyLock<RawOrigin<'static>> = LazyLock::new(|| {
685 let mut origin_info = RawOrigin::default();
686 origin_info.set_process_id(std::process::id());
687 origin_info
688 });
689
690 if let Some(context) = self.key_context_cache.get(&key) {
692 return context.clone();
693 }
694
695 let tags = key
697 .labels()
698 .map(|l| Tag::from(format!("{}:{}", l.key(), l.value())))
699 .collect::<TagSet>();
700
701 let context = self
702 .context_resolver
703 .resolve(key.name(), &tags, Some(SELF_ORIGIN_INFO.clone()))
704 .expect("resolver should always allow falling back");
705
706 self.key_context_cache.insert(key, context.clone());
707 context
708 }
709
710 fn take(&mut self, key: &Key) -> Option<Context> {
715 self.key_context_cache.remove(key)
716 }
717}
718
719pub async fn initialize_metrics(
733 metrics_prefix: String, default_level: Level,
734) -> Result<(FilterHandle, MetricsFlusherWorker), GenericError> {
735 let recorder = MetricsRecorder::new(metrics_prefix, default_level);
736 let filter_handle = recorder.filter_handle();
737 recorder.install()?;
738
739 Ok((filter_handle, MetricsFlusherWorker))
740}
741
742pub struct MetricsFlusherWorker;
748
749#[async_trait]
750impl Supervisable for MetricsFlusherWorker {
751 fn name(&self) -> &str {
752 "flusher"
753 }
754
755 async fn initialize(&self, process_shutdown: ShutdownHandle) -> Result<SupervisorFuture, InitializationError> {
756 Ok(Box::pin(async move {
757 select! {
758 _ = process_shutdown => {},
759 _ = flush_metrics() => {},
760 }
761
762 Ok(())
763 }))
764 }
765}
766
767#[cfg(test)]
770pub(crate) fn aggregate_upserts(upserts: Vec<Event>) -> AggregatedMetricsState {
771 let processor = AggregatedMetricsProcessor;
772 let state = processor.build_initial_state();
773 processor.process(
774 MetricsSnapshot {
775 upserts,
776 evictions: Vec::new(),
777 },
778 &state,
779 );
780 state
781}
782
783#[cfg(test)]
784mod tests {
785 use metrics::Label;
786 use saluki_common::collections::FastHashSet;
787
788 use super::*;
789
790 const HIGH_CARDINALITY_COUNTER_NAME: &str = "high_cardinality_counter";
791 const HIGH_CARDINALITY_COUNTERS: usize = 64;
792
793 fn make_state(idle_evict_intervals: u32) -> (Arc<State>, Receiver<SharedMetricsSnapshot>) {
795 let (flush_tx, rx) = broadcast::channel(64);
796 let state = Arc::new(State {
797 registry: MetricsRegistry::default(),
798 flush_tx,
799 metrics_prefix: "test".to_string(),
800 default_level: Level::TRACE,
801 current_level: Mutex::new(Level::TRACE),
802 idle_evict_intervals,
803 flush_interval: FLUSH_INTERVAL,
804 });
805 (state, rx)
806 }
807
808 fn resolver() -> MetricsContextResolver {
809 MetricsContextResolver::new(INTERNAL_METRICS_INTERNER_SIZE)
810 }
811
812 fn register_counter(state: &State, name: &'static str) -> Counter {
814 state
815 .registry
816 .get_or_create_counter(&Key::from_name(name), Level::TRACE)
817 }
818
819 fn register_counter_with_origin(state: &State, name: &'static str, origin: &str) -> Counter {
820 let key = Key::from_parts(name, vec![Label::new("origin", origin.to_string())]);
821 state.registry.get_or_create_counter(&key, Level::TRACE)
822 }
823
824 fn register_gauge(state: &State, name: &'static str) -> Gauge {
825 state.registry.get_or_create_gauge(&Key::from_name(name), Level::TRACE)
826 }
827
828 fn register_histogram(state: &State, name: &'static str) -> Histogram {
829 state
830 .registry
831 .get_or_create_histogram(&Key::from_name(name), Level::TRACE)
832 }
833
834 fn counter_present(state: &State, name: &'static str) -> bool {
835 state
836 .registry
837 .maps
838 .lock()
839 .unwrap()
840 .counters
841 .contains_key(&Key::from_name(name))
842 }
843
844 fn counter_count(state: &State) -> usize {
845 state.registry.maps.lock().unwrap().counters.len()
846 }
847
848 fn gauge_present(state: &State, name: &'static str) -> bool {
849 state
850 .registry
851 .maps
852 .lock()
853 .unwrap()
854 .gauges
855 .contains_key(&Key::from_name(name))
856 }
857
858 fn histogram_present(state: &State, name: &'static str) -> bool {
859 state
860 .registry
861 .maps
862 .lock()
863 .unwrap()
864 .histograms
865 .contains_key(&Key::from_name(name))
866 }
867
868 fn drain(rx: &mut Receiver<SharedMetricsSnapshot>) -> Vec<MetricsSnapshot> {
869 let mut out = Vec::new();
870 while let Ok(batch) = rx.try_recv() {
871 out.push((*batch).clone());
872 }
873 out
874 }
875
876 fn evicted_names(snapshots: &[MetricsSnapshot]) -> Vec<String> {
877 snapshots
878 .iter()
879 .flat_map(|snapshot| snapshot.evictions.iter().map(|ctx| ctx.name().to_string()))
880 .collect()
881 }
882
883 fn metric_count(state: &AggregatedMetricsState, name: &str) -> usize {
884 let mut count = 0;
885 state.visit_metrics(|context, _| {
886 if context.name() == name {
887 count += 1;
888 }
889 });
890 count
891 }
892
893 #[test]
894 fn delete_counter_from_state_emits_pending_delta_before_eviction() {
895 let (state, _rx) = make_state(IDLE_EVICT_INTERVALS);
896 let key = Key::from_parts(
897 "deleted_counter",
898 vec![Label::new("origin", "container_id://deleted".to_string())],
899 );
900
901 let counter = state.registry.get_or_create_counter(&key, Level::TRACE);
902 counter.increment(7);
903 drop(counter);
904
905 let snapshot = delete_counter_from_state(&state, &key).expect("counter should be deleted");
906
907 assert!(!state.registry.maps.lock().unwrap().counters.contains_key(&key));
908 assert_eq!(snapshot.evictions.len(), 1);
909 assert_eq!(snapshot.evictions[0].name(), "deleted_counter");
910 assert!(snapshot.evictions[0].tags().has_tag("origin:container_id://deleted"));
911
912 assert_eq!(snapshot.upserts.len(), 1);
913 let Event::Metric(metric) = &snapshot.upserts[0] else {
914 panic!("delete should emit the pending counter delta");
915 };
916 assert_eq!(metric.context(), &snapshot.evictions[0]);
917 let MetricValues::Counter(points) = metric.values() else {
918 panic!("delete should emit a counter metric");
919 };
920 assert_eq!(points.into_iter().map(|(_, value)| value).sum::<f64>(), 7.0);
921 }
922
923 #[tokio::test]
924 async fn evicts_high_cardinality_counters_from_registry_and_aggregated_state() {
925 let (state, mut rx) = make_state(IDLE_EVICT_INTERVALS);
926 let mut resolver = resolver();
927 let mut flush_state = FlushState::default();
928 let agg_processor = AggregatedMetricsProcessor;
929 let agg_state = agg_processor.build_initial_state();
930 let mut expected_eviction_tags = FastHashSet::default();
931 let mut observed_eviction_tags = FastHashSet::default();
932
933 for index in 0..HIGH_CARDINALITY_COUNTERS {
934 let origin = format!("entity://source-{index}");
935 expected_eviction_tags.insert(format!("origin:{origin}"));
936
937 let counter = register_counter_with_origin(&state, HIGH_CARDINALITY_COUNTER_NAME, &origin);
938 counter.increment(1);
939 drop(counter);
940 }
941
942 assert_eq!(counter_count(&state), HIGH_CARDINALITY_COUNTERS);
943
944 for _ in 0..IDLE_EVICT_INTERVALS {
945 flush_once(&state, &mut resolver, &mut flush_state);
946
947 for update in drain(&mut rx) {
948 for context in &update.evictions {
949 if context.name() != HIGH_CARDINALITY_COUNTER_NAME {
950 continue;
951 }
952
953 for tag in &expected_eviction_tags {
954 if context.tags().has_tag(tag) {
955 observed_eviction_tags.insert(tag.clone());
956 }
957 }
958 }
959
960 agg_processor.process(update, &agg_state);
961 }
962 }
963
964 assert_eq!(counter_count(&state), 0);
965 assert_eq!(metric_count(&agg_state, HIGH_CARDINALITY_COUNTER_NAME), 0);
966 assert_eq!(observed_eviction_tags, expected_eviction_tags);
967 }
968
969 #[tokio::test]
970 async fn evicts_counter_within_bounded_time_after_handle_dropped() {
971 let (state, mut rx) = make_state(3);
972 let mut resolver = resolver();
973 let mut flush_state = FlushState::default();
974
975 let counter = register_counter(&state, "dropped_counter");
976 counter.increment(5);
977
978 flush_once(&state, &mut resolver, &mut flush_state);
980 assert!(counter_present(&state, "dropped_counter"));
981
982 drop(counter);
984 let _ = drain(&mut rx);
985 for _ in 0..3 {
986 flush_once(&state, &mut resolver, &mut flush_state);
987 }
988 assert!(!counter_present(&state, "dropped_counter"));
989
990 let updates = drain(&mut rx);
992 assert!(evicted_names(&updates).iter().any(|n| n == "dropped_counter"));
993 }
994
995 #[tokio::test]
996 async fn does_not_evict_while_handle_is_held() {
997 let (state, _rx) = make_state(3);
998 let mut resolver = resolver();
999 let mut flush_state = FlushState::default();
1000
1001 let counter = register_counter(&state, "held_counter");
1002
1003 for _ in 0..8 {
1005 flush_once(&state, &mut resolver, &mut flush_state);
1006 }
1007
1008 assert!(
1009 counter_present(&state, "held_counter"),
1010 "a held metric must never be evicted"
1011 );
1012 drop(counter);
1013 }
1014
1015 #[tokio::test]
1016 async fn reregistration_keeps_uncached_metric_alive() {
1017 let (state, _rx) = make_state(3);
1020 let mut resolver = resolver();
1021 let mut flush_state = FlushState::default();
1022
1023 for _ in 0..8 {
1024 let counter = register_counter(&state, "uncached_counter");
1025 counter.increment(1);
1026 drop(counter);
1027 flush_once(&state, &mut resolver, &mut flush_state);
1028 }
1029 assert!(counter_present(&state, "uncached_counter"));
1030
1031 for _ in 0..3 {
1033 flush_once(&state, &mut resolver, &mut flush_state);
1034 }
1035 assert!(!counter_present(&state, "uncached_counter"));
1036 }
1037
1038 #[tokio::test]
1039 async fn gauge_reset_to_same_value_is_not_evicted() {
1040 let (state, _rx) = make_state(3);
1043 let mut resolver = resolver();
1044 let mut flush_state = FlushState::default();
1045
1046 for _ in 0..8 {
1047 let gauge = register_gauge(&state, "steady_gauge");
1048 gauge.set(0.0);
1049 drop(gauge);
1050 flush_once(&state, &mut resolver, &mut flush_state);
1051 }
1052 assert!(gauge_present(&state, "steady_gauge"));
1053 }
1054
1055 #[tokio::test]
1056 async fn final_counter_delta_reaches_downstream_before_eviction() {
1057 let (state, mut rx) = make_state(1);
1058 let mut resolver = resolver();
1059 let mut flush_state = FlushState::default();
1060
1061 let counter = register_counter(&state, "final_counter");
1062
1063 flush_once(&state, &mut resolver, &mut flush_state);
1065
1066 counter.increment(7);
1069 drop(counter);
1070 flush_once(&state, &mut resolver, &mut flush_state);
1071 assert!(
1072 counter_present(&state, "final_counter"),
1073 "must not evict in an interval that produced a delta"
1074 );
1075
1076 let agg_processor = AggregatedMetricsProcessor;
1078 let agg_state = agg_processor.build_initial_state();
1079 for update in drain(&mut rx) {
1080 agg_processor.process(update, &agg_state);
1081 }
1082 assert_eq!(agg_state.get_aggregated_with_tags("final_counter", &[]), 7.0);
1083
1084 flush_once(&state, &mut resolver, &mut flush_state);
1086 assert!(!counter_present(&state, "final_counter"));
1087 }
1088
1089 #[tokio::test]
1090 async fn gauge_final_value_reaches_downstream_before_eviction() {
1091 let (state, mut rx) = make_state(1);
1092 let mut resolver = resolver();
1093 let mut flush_state = FlushState::default();
1094
1095 let gauge = register_gauge(&state, "final_gauge");
1096 gauge.set(1.0);
1097
1098 flush_once(&state, &mut resolver, &mut flush_state);
1100
1101 gauge.set(42.0);
1105 drop(gauge);
1106 flush_once(&state, &mut resolver, &mut flush_state);
1107 assert!(
1108 gauge_present(&state, "final_gauge"),
1109 "must not evict in an interval that wrote a new value"
1110 );
1111
1112 let agg_processor = AggregatedMetricsProcessor;
1114 let agg_state = agg_processor.build_initial_state();
1115 for update in drain(&mut rx) {
1116 agg_processor.process(update, &agg_state);
1117 }
1118 assert_eq!(agg_state.find_single_with_tags("final_gauge", &[]), Some(42.0));
1119
1120 flush_once(&state, &mut resolver, &mut flush_state);
1122 assert!(!gauge_present(&state, "final_gauge"));
1123 }
1124
1125 #[tokio::test]
1126 async fn histogram_final_samples_reach_downstream_before_eviction() {
1127 let (state, mut rx) = make_state(1);
1139 let mut resolver = resolver();
1140 let mut flush_state = FlushState::default();
1141
1142 let histogram = register_histogram(&state, "final_histogram");
1143
1144 flush_once(&state, &mut resolver, &mut flush_state);
1146
1147 histogram.record(1.5);
1150 drop(histogram);
1151 flush_once(&state, &mut resolver, &mut flush_state);
1152 assert!(
1153 histogram_present(&state, "final_histogram"),
1154 "must not evict in an interval that recorded a sample"
1155 );
1156
1157 let agg_processor = AggregatedMetricsProcessor;
1159 let agg_state = agg_processor.build_initial_state();
1160 for update in drain(&mut rx) {
1161 agg_processor.process(update, &agg_state);
1162 }
1163
1164 let mut downstream_sample_count = None;
1165 agg_state.visit_metrics(|context, value| {
1166 if context.name() == "final_histogram" {
1167 if let AggregatedMetricValue::Histogram(histogram) = value {
1168 downstream_sample_count = Some(histogram.count());
1169 }
1170 }
1171 });
1172 assert_eq!(
1173 downstream_sample_count,
1174 Some(1),
1175 "the final recorded sample must reach downstream before eviction"
1176 );
1177
1178 flush_once(&state, &mut resolver, &mut flush_state);
1180 assert!(!histogram_present(&state, "final_histogram"));
1181 }
1182}
1183
1184#[cfg(all(test, feature = "loom"))]
1206mod loom_tests {
1207 use loom::sync::atomic::{fence, AtomicU64, AtomicUsize, Ordering};
1208 use loom::sync::Arc;
1209
1210 struct Shared {
1213 value: AtomicU64,
1214 strong: AtomicUsize,
1215 }
1216
1217 fn caller_increment_then_drop(shared: &Arc<Shared>, amount: u64) {
1219 shared.value.fetch_add(amount, Ordering::Relaxed);
1220 shared.strong.fetch_sub(1, Ordering::Release); }
1222
1223 fn flush_decision_consume_first(shared: &Arc<Shared>) -> (u64, bool) {
1225 let delta = shared.value.swap(0, Ordering::Relaxed);
1226 let orphaned = shared.strong.load(Ordering::Relaxed) == 1;
1227 (delta, orphaned && delta == 0)
1228 }
1229
1230 fn flush_decision_strong_count_first(shared: &Arc<Shared>) -> (u64, bool) {
1233 let orphaned = shared.strong.load(Ordering::Relaxed) == 1;
1234 if orphaned {
1235 fence(Ordering::Acquire);
1236 }
1237 let delta = shared.value.swap(0, Ordering::Relaxed);
1238 (delta, orphaned && delta == 0)
1239 }
1240
1241 fn model(decision: fn(&Arc<Shared>) -> (u64, bool)) {
1242 loom::model(move || {
1243 const FINAL: u64 = 7;
1244
1245 let shared = Arc::new(Shared {
1247 value: AtomicU64::new(0),
1248 strong: AtomicUsize::new(2),
1249 });
1250 let caller_shared = Arc::clone(&shared);
1251
1252 let caller = loom::thread::spawn(move || {
1253 caller_increment_then_drop(&caller_shared, FINAL);
1254 });
1255
1256 let (delta, evicted) = decision(&shared);
1258 caller.join().unwrap();
1259
1260 if evicted {
1264 assert_eq!(
1265 delta,
1266 FINAL,
1267 "counter evicted while {} counts were never emitted -> permanently lost",
1268 FINAL - delta
1269 );
1270 }
1271 });
1272 }
1273
1274 #[test]
1278 #[should_panic(expected = "permanently lost")]
1279 fn consume_first_ordering_loses_a_final_increment() {
1280 model(flush_decision_consume_first);
1281 }
1282
1283 #[test]
1286 fn strong_count_first_ordering_preserves_a_final_increment() {
1287 model(flush_decision_strong_count_first);
1288 }
1289
1290 fn histogram_model(decision: fn(&Arc<Shared>) -> (u64, bool)) {
1305 loom::model(move || {
1306 const FINAL_SAMPLES: u64 = 1;
1307
1308 let shared = Arc::new(Shared {
1310 value: AtomicU64::new(0),
1311 strong: AtomicUsize::new(2),
1312 });
1313 let caller_shared = Arc::clone(&shared);
1314
1315 let caller = loom::thread::spawn(move || {
1317 caller_shared.value.fetch_add(FINAL_SAMPLES, Ordering::Relaxed); caller_shared.strong.fetch_sub(1, Ordering::Release); });
1320
1321 let (drained, evicted) = decision(&shared);
1323 caller.join().unwrap();
1324
1325 if evicted {
1329 assert_eq!(
1330 drained,
1331 FINAL_SAMPLES,
1332 "histogram evicted while {} recorded sample(s) were never drained -> permanently lost",
1333 FINAL_SAMPLES - drained
1334 );
1335 }
1336 });
1337 }
1338
1339 #[test]
1343 #[should_panic(expected = "permanently lost")]
1344 fn histogram_drain_first_ordering_loses_a_final_sample() {
1345 histogram_model(flush_decision_consume_first);
1346 }
1347
1348 #[test]
1351 fn histogram_strong_count_first_ordering_preserves_a_final_sample() {
1352 histogram_model(flush_decision_strong_count_first);
1353 }
1354}