saluki_core/observability/metrics/
mod.rs

1//! Internal metrics support.
2//!
3//! Includes the [`MetricsStream`] broadcast of internally emitted metrics events and the
4//! [`Reflector`]-based [`AggregatedMetricsState`] view that downstream callers query for
5//! Prometheus exposition.
6
7use 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
70/// Number of consecutive flush intervals that a metric must go without a live handle and without being re-registered
71/// before it is evicted from the registry.
72///
73/// This grace period prevents metrics that are emitted through `metrics` macros without caching the returned handle --
74/// and that therefore re-register on every call -- from being churned in and out of the registry while they are still
75/// being emitted. Naturally, this only applies to metrics which re-register within `IDLE_EVICT_INTERVALS *
76/// FLUSH_INTERVAL`, which is 3 seconds by default... but that's a tradeoff we're making here to avoid unbounded growth.
77const IDLE_EVICT_INTERVALS: u32 = 3;
78
79static RECEIVER_STATE: OnceLock<Arc<State>> = OnceLock::new();
80
81/// A batch of metric updates produced by a single flush.
82#[derive(Clone, Default)]
83pub struct MetricsSnapshot {
84    /// Updates to active metrics.
85    pub upserts: Vec<Event>,
86
87    /// Metrics which are no longer active and should be removed.
88    pub evictions: Vec<Context>,
89}
90
91/// A [`MetricsSnapshot`] that can be cheaply cloned and shared between multiple consumers.
92pub 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    /// Consumes the accumulated delta since the last flush, resetting the counter to zero.
137    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    // Tracks whether or not the gauge has been written since the last flush, since the gauge _could_
155    // be modified in a way where the value seen by two consecutive flushes is the same, even though
156    // it _was_ modified between the two and should be considered non-idle.
157    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    /// Reads the current gauge value.
172    fn load(&self) -> f64 {
173        f64::from_bits(self.inner.value.load(Ordering::Relaxed))
174    }
175
176    /// Returns whether the gauge was written since the last flush, clearing the flag.
177    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
199/// Storage for a single histogram, shared between the registry and every caller-held [`Histogram`].
200struct 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    /// Drains all recorded samples since the last flush into `out`, clearing the histogram.
215    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/// A registry for all internal metrics.
227///
228/// Optimized for simplicity and the ability to efficiently track metric state and evict idle metrics. Not optimized for
229/// high concurrency with regards to registration, so metric handles should always be held for as long as possible, when
230/// possible.
231#[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    /// Returns a handle to the counter for `key`, creating it at `level` if it doesn't yet exist.
245    ///
246    /// The level is fixed when the metric is first created; re-registration only refreshes the idle
247    /// counter (the activity signal that keeps actively emitted metrics from being evicted).
248    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    /// Returns a handle to the gauge for `key`, creating it at `level` if it doesn't yet exist.
268    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    /// Returns a handle to the histogram for `key`, creating it at `level` if it doesn't yet exist.
283    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
298/// Handle to the metrics filter.
299///
300/// Allows for overriding the current metrics filter level, which influences which metrics are emitted to downstream
301/// receivers.
302pub struct FilterHandle {
303    state: Arc<State>,
304}
305
306impl FilterHandle {
307    /// Overrides the current metrics filter level.
308    pub fn override_filter(&self, level: Level) {
309        *self.state.current_level.lock().unwrap() = level;
310    }
311
312    /// Resets the metrics filter level to the default that was configured when the metrics subsystem was initialized.
313    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
397/// Deletes an internal counter from the metrics registry.
398///
399/// If the counter exists and metrics snapshots are being consumed, this sends any pending counter
400/// delta before the eviction update so downstream consumers do not lose the final increment.
401pub 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
446/// Internal metrics stream
447///
448/// Used to receive periodic snapshots of the internal metrics registry, which contains all metrics that are currently
449/// active within the process.
450pub struct MetricsStream {
451    inner: ReusableBoxFuture<
452        'static,
453        (
454            Result<SharedMetricsSnapshot, RecvError>,
455            Receiver<SharedMetricsSnapshot>,
456        ),
457    >,
458}
459
460impl MetricsStream {
461    /// Creates a new `MetricsStream` that receives updates from the internal metrics registry.
462    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            // Poll the receiver, and rearm the future once we actually resolve it.
476            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
537/// Performs a single flush pass over the registry: consumes each metric's value, evicts metrics that
538/// are no longer referenced and have gone idle, and broadcasts the resulting metric values and
539/// evictions to downstream consumers.
540fn 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        // Counters: a counter is never idle so long as it has a non-zero delta during a flush _or_ has
558        // an active reference (strong count > 1).
559        //
560        // The strong-count check precedes `consume()`, with an Acquire fence on the orphaned path:
561        // observing `strong_count == 1` synchronizes-with the dropping thread's release, which
562        // guarantees the following `consume()` observes that thread's final increment.
563        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        // Gauges: a gauge is never idle so long as it was touched ("dirty") prior to a flush _or_ has
583        // an active reference (strong count > 1).
584        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        // Histograms: a histogram is never idle so long as it has recorded samples prior to a flush _or_
605        // has an active reference (strong count > 1).
606        //
607        // The strong-count check precedes the drain, with an Acquire fence on the orphaned path, so the
608        // drain observes a dropping thread's final `record()`.
609        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 every evicted key we collected, take its cached context (if it exists) and add it to the
632    // snapshot's evictions list.
633    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            // Set up our context resolver without caching, since we will be caching the contexts ourselves.
674            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        // Check the cache first.
691        if let Some(context) = self.key_context_cache.get(&key) {
692            return context.clone();
693        }
694
695        // We don't have the context cached, so we need to resolve it.
696        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    /// Removes a key's cached context, returning it if present.
711    ///
712    /// Returns `Some` only if the key had been resolved before (that is, the metric was emitted at
713    /// least once), which is what tells the flush loop whether a downstream eviction needs to be sent.
714    fn take(&mut self, key: &Key) -> Option<Context> {
715        self.key_context_cache.remove(key)
716    }
717}
718
719/// Initializes the metrics subsystem with the given metrics prefix and default filter level.
720///
721/// `default_level` sets the initial filter level for emitted metrics, and is also what the filter is restored to when
722/// [`FilterHandle::reset_filter`] is invoked. Metrics whose level is more verbose than this default are filtered out
723/// until a runtime override is applied via [`FilterHandle::override_filter`].
724///
725/// Returns a [`FilterHandle`] for adjusting the runtime metrics filter, plus a [`MetricsFlusherWorker`]
726/// that must be added to a [`Supervisor`][crate::runtime::Supervisor] in order to drive the periodic
727/// flush loop. Internal metrics aren't propagated to subscribers until the worker is running.
728///
729/// # Errors
730///
731/// If a global recorder was already installed, an error will be returned.
732pub 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
742/// A worker that periodically flushes the internal metrics registry to broadcast subscribers.
743///
744/// Wraps the internal flush loop and runs it under a [`Supervisor`][crate::runtime::Supervisor]. Must
745/// only be added to a supervisor after [`initialize_metrics`] has been called -- the flush loop
746/// reads from the global recorder state set by that call.
747pub 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/// Feeds `upserts` through a fresh [`AggregatedMetricsProcessor`] and returns the resulting
768/// [`AggregatedMetricsState`].
769#[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    // We keep a live receiver in each test so `flush_once` observes a listener and emits updates.
794    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    // Mimics `MetricsRecorder::register_counter` against a raw key (no prefixing needed in tests).
813    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        // While the handle is held, the metric is retained and its delta is emitted.
979        flush_once(&state, &mut resolver, &mut flush_state);
980        assert!(counter_present(&state, "dropped_counter"));
981
982        // Drop the last handle; the metric must be reclaimed within `idle_evict_intervals` flushes.
983        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        // The eviction is propagated downstream so the aggregated view drops it too.
991        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        // Idle well past the eviction threshold without dropping the handle.
1004        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        // Models the uncached-macro pattern: a fresh handle is registered (and dropped) every interval.
1018        // Re-registration resets the idle counter, so the metric is never churned out while in use.
1019        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        // Once it stops being re-registered, it is reclaimed within the idle threshold.
1032        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        // A gauge re-set to the same value every interval (e.g. a backoff gauge that sits at 0.0) must
1041        // stay alive. Activity is driven by re-registration, not value comparison.
1042        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        // Arm eviction by letting the idle counter reach the threshold while the handle is held.
1064        flush_once(&state, &mut resolver, &mut flush_state);
1065
1066        // Increment and drop the handle in the same inter-flush window. The next flush sees a non-zero
1067        // delta, so it must emit it and defer eviction by one interval (so the delta isn't lost).
1068        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        // The emitted delta must reach the downstream aggregated state.
1077        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        // The following interval (no delta) finally evicts it.
1085        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        // Arm eviction by letting the idle counter reach the threshold while the handle is held.
1099        flush_once(&state, &mut resolver, &mut flush_state);
1100
1101        // Update to a NEW value and drop the handle in the same inter-flush window. Even though the
1102        // gauge is now orphaned and idle, the write must not be lost: the next flush sees the dirty flag,
1103        // emits the value, and defers eviction by one interval.
1104        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        // The final value must reach the downstream aggregated state.
1113        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        // The following interval (no write) finally evicts it.
1121        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        // Regression test for the histogram eviction race fixed in commit 3752c5ad1b
1128        // ("rework internal metrics registry to support evicting expired/idle metrics", #1947).
1129        //
1130        // The histogram branch of `flush_once` must read the Arc strong count first (Acquire-fencing
1131        // on the orphaned path) and only THEN drain the sample bucket, exactly like the counter and
1132        // gauge paths. If a final `record()` and the handle drop both land in the inter-flush window
1133        // where eviction is armed, the flush that observes the orphaned handle must still see and emit
1134        // that sample -- and defer eviction by one interval -- rather than dropping the samples along
1135        // with the evicted registry entry. That commit shipped the counter regression coverage
1136        // (`final_counter_delta_reaches_downstream_before_eviction` plus the loom model) but never the
1137        // symmetric histogram coverage; this is it.
1138        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        // Arm eviction by letting the idle counter reach the threshold while the handle is held.
1145        flush_once(&state, &mut resolver, &mut flush_state);
1146
1147        // Record a sample and drop the handle in the same inter-flush window. The next flush observes
1148        // the recorded sample, so it must emit it and defer eviction by one interval.
1149        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        // The recorded sample must reach the downstream aggregated state.
1158        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        // The following interval (no samples) finally evicts it.
1179        flush_once(&state, &mut resolver, &mut flush_state);
1180        assert!(!histogram_present(&state, "final_histogram"));
1181    }
1182}
1183
1184// Loom model of the counter eviction path in `flush_once`.
1185//
1186// In `flush_once`, the registry `maps` lock is held during the drain, but the increment/drop path
1187// does not take that lock: `Counter::increment` writes straight through the `Arc<Handle>`, and
1188// dropping a caller's `Counter` only decrements the Arc strong count. A component holding a cached
1189// handle can therefore land a final increment and drop the handle while a flush is mid-pass over it.
1190//
1191// The model transcribes the counter branch with loom primitives instead of exercising the real types,
1192// which loom cannot instrument here:
1193//
1194// - The strong count is an explicit `AtomicUsize`. loom's `Arc::strong_count` does not reflect a
1195//   concurrent decrement from another thread (it models Arc's drop synchronization, not the observable
1196//   count value), so a flush reading it never sees the orphaned (`== 1`) state mid-race. The explicit
1197//   atomic matches `std`'s `Arc`: clone increments (Relaxed), drop decrements with `Release`,
1198//   observation is a `Relaxed` load -- the same shape as how `flush_once` reads `Arc::strong_count`.
1199// - The real `Handle` is wrapped by the `metrics` crate's `Counter`/`CounterFn`, which construct
1200//   `std::sync::Arc` internally. loom only instruments atomics and `Arc`s swapped to `loom::sync::*`
1201//   behind the `loom` cfg, not those inside a third-party crate.
1202//
1203// `flush_decision_consume_first` and `flush_decision_strong_count_first` mirror the two possible
1204// orderings of `flush_once`'s counter branch and must stay in lockstep with it.
1205#[cfg(all(test, feature = "loom"))]
1206mod loom_tests {
1207    use loom::sync::atomic::{fence, AtomicU64, AtomicUsize, Ordering};
1208    use loom::sync::Arc;
1209
1210    /// The shared state behind a counter handle: the accumulator the flush loop drains, and the Arc
1211    /// strong count it consults to decide whether the handle is orphaned.
1212    struct Shared {
1213        value: AtomicU64,
1214        strong: AtomicUsize,
1215    }
1216
1217    /// A component holding a cached handle lands one final increment, then drops the handle.
1218    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); // Arc clone drop
1221    }
1222
1223    /// Consumes the value, then checks the strong count (value read before the liveness check).
1224    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    /// Checks the strong count first; if orphaned, Acquire-fences to synchronize with the dropping
1231    /// thread's `Release`, then consumes the value.
1232    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            // strong count starts at 2: one ref in the registry map, one held by the caller.
1246            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            // The flush task evaluates this counter while the component races.
1257            let (delta, evicted) = decision(&shared);
1258            caller.join().unwrap();
1259
1260            // Invariant: a counter may only be evicted once everything it accumulated has been
1261            // emitted. If the flush evicts it, the delta it emitted downstream must already include
1262            // the final increment -- otherwise those counts are silently discarded with the handle.
1263            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    // Consuming the value before checking the strong count is lossy: loom finds an interleaving where
1275    // the handle is evicted while a final increment goes unemitted. `#[should_panic]` asserts loom
1276    // reaches that interleaving.
1277    #[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    // Checking the strong count first (Acquire-fencing when orphaned) before consuming preserves every
1284    // increment across all interleavings loom explores.
1285    #[test]
1286    fn strong_count_first_ordering_preserves_a_final_increment() {
1287        model(flush_decision_strong_count_first);
1288    }
1289
1290    // The histogram branch of `flush_once` has the same race shape as the counter branch, and got the
1291    // same fix in commit 3752c5ad1b (#1947): `Histogram::record` pushes straight through the
1292    // `Arc<Handle>`'s sample bucket, and dropping a caller's `Histogram` only decrements the Arc strong
1293    // count, so a component can land a final sample and drop its handle while a flush is mid-pass over
1294    // that histogram. Reading the strong count first (Acquire-fencing when orphaned) before draining
1295    // guarantees the drain observes the dropping thread's final `record()`.
1296    //
1297    // As with the counter model, loom cannot instrument the real storage: histogram samples live in
1298    // `metrics_util`'s `AtomicBucket`, which loom does not rewrite behind the `loom` cfg. We transcribe
1299    // "number of un-drained samples" into the same loom `AtomicU64` accumulator (record => fetch_add(1),
1300    // drain => swap(0)); that is precisely the ordering that decides whether the drain sees the final
1301    // sample, so the two decision functions above -- `flush_decision_consume_first` (drain-first) and
1302    // `flush_decision_strong_count_first` -- model the histogram drain unchanged, just with a
1303    // single-sample final write.
1304    fn histogram_model(decision: fn(&Arc<Shared>) -> (u64, bool)) {
1305        loom::model(move || {
1306            const FINAL_SAMPLES: u64 = 1;
1307
1308            // strong count starts at 2: one ref in the registry map, one held by the caller.
1309            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            // A component holding a cached histogram handle records one final sample, then drops it.
1316            let caller = loom::thread::spawn(move || {
1317                caller_shared.value.fetch_add(FINAL_SAMPLES, Ordering::Relaxed); // Histogram::record
1318                caller_shared.strong.fetch_sub(1, Ordering::Release); // Arc clone drop
1319            });
1320
1321            // The flush task evaluates this histogram while the component races.
1322            let (drained, evicted) = decision(&shared);
1323            caller.join().unwrap();
1324
1325            // Invariant: a histogram may only be evicted once every recorded sample has been drained
1326            // (and thus emitted). If the flush evicts it, the drain it performed must have observed the
1327            // final sample -- otherwise that sample is discarded with the evicted registry entry.
1328            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    // Draining the sample bucket before checking the strong count is lossy: loom finds an interleaving
1340    // where the histogram is evicted while a final recorded sample goes undrained. `#[should_panic]`
1341    // asserts loom reaches that interleaving.
1342    #[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    // Checking the strong count first (Acquire-fencing when orphaned) before draining preserves the
1349    // final recorded sample across all interleavings loom explores.
1350    #[test]
1351    fn histogram_strong_count_first_ordering_preserves_a_final_sample() {
1352        histogram_model(flush_decision_strong_count_first);
1353    }
1354}