Skip to main content

mz_ore/
metrics.rs

1// Copyright Materialize, Inc. and contributors. All rights reserved.
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License in the LICENSE file at the
6// root of this repository, or online at
7//
8//     http://www.apache.org/licenses/LICENSE-2.0
9//
10// Unless required by applicable law or agreed to in writing, software
11// distributed under the License is distributed on an "AS IS" BASIS,
12// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13// See the License for the specific language governing permissions and
14// limitations under the License.
15
16//! Metrics for materialize systems.
17//!
18//! The idea here is that each subsystem keeps its metrics in a scoped-to-it struct, which gets
19//! registered (once) to the server's (or a test's) prometheus registry.
20//!
21//! Instead of using prometheus's (very verbose) metrics definitions, we rely on type inference to
22//! reduce the verbosity a little bit. A typical subsystem will look like the following:
23//!
24//! ```rust
25//! # use mz_ore::metrics::{MetricsRegistry, IntCounter};
26//! # use mz_ore::metric;
27//! #[derive(Debug, Clone)] // Note that prometheus metrics can safely be cloned
28//! struct Metrics {
29//!     pub bytes_sent: IntCounter,
30//! }
31//!
32//! impl Metrics {
33//!     pub fn register_into(registry: &MetricsRegistry) -> Metrics {
34//!         Metrics {
35//!             bytes_sent: registry.register(metric!(
36//!                 name: "mz_pg_sent_bytes",
37//!                 help: "total number of bytes sent here",
38//!             )),
39//!         }
40//!     }
41//! }
42//! ```
43
44use std::any::Any;
45use std::collections::BTreeMap;
46use std::fmt;
47use std::fmt::{Debug, Formatter};
48use std::future::Future;
49use std::pin::Pin;
50use std::sync::{Arc, Mutex};
51use std::task::{Context, Poll};
52use std::time::{Duration, Instant};
53
54use derivative::Derivative;
55use pin_project::pin_project;
56use prometheus::core::{
57    Atomic, AtomicF64, AtomicI64, AtomicU64, Collector, Desc, GenericCounter, GenericCounterVec,
58    GenericGauge, GenericGaugeVec,
59};
60use prometheus::proto::MetricFamily;
61use prometheus::{HistogramOpts, Registry};
62
63mod delete_on_drop;
64
65pub use delete_on_drop::*;
66pub use prometheus::Opts as PrometheusOpts;
67
68/// Define a metric for use in materialize.
69#[macro_export]
70macro_rules! metric {
71    (
72        name: $name:expr,
73        help: $help:expr
74        $(, subsystem: $subsystem_name:expr)?
75        $(, const_labels: { $($cl_key:expr => $cl_value:expr ),* })?
76        $(, var_labels: [ $($vl_name:expr),* ])?
77        $(, buckets: $bk_name:expr)?
78        $(, visibility: $visibility:expr)?
79        $(, tags: [ $($tag:expr),* $(,)? ])?
80        $(,)?
81    ) => {{
82        let const_labels = (&[
83            $($(
84                ($cl_key.to_string(), $cl_value.to_string()),
85            )*)?
86        ]).into_iter().cloned().collect();
87        let var_labels = vec![
88            $(
89                $($vl_name.into(),)*
90            )?];
91        #[allow(unused_mut)]
92        let mut mk_opts = $crate::metrics::MakeCollectorOpts {
93            opts: $crate::metrics::PrometheusOpts::new($name, $help)
94                $(.subsystem( $subsystem_name ))?
95                .const_labels(const_labels)
96                .variable_labels(var_labels),
97            buckets: None,
98        };
99        // Set buckets if passed
100        $(mk_opts.buckets = Some($bk_name);)*
101        // `visibility` is documentation metadata for the metrics catalog
102        // (`bin/gen-metrics-catalog`).
103        // It has no runtime effect; we only type-check it here so a bad value is
104        // a compile error rather than silently ignored.
105        $(let _: $crate::metrics::MetricVisibility = $visibility;)?
106        // `tags` is documentation metadata for the metrics catalog.
107        // It has no runtime effect.
108        $($(let _: $crate::metrics::MetricTag = $tag;)*)?
109        mk_opts
110    }}
111}
112
113/// Options for MakeCollector. This struct should be instantiated using the metric macro.
114#[derive(Debug, Clone)]
115pub struct MakeCollectorOpts {
116    /// Common Prometheus options
117    pub opts: PrometheusOpts,
118    /// Buckets to be used with Histogram and HistogramVec. Must be set to create Histogram types
119    /// and must not be set for other types.
120    pub buckets: Option<Vec<f64>>,
121}
122
123/// This is documentation metadata: it is set via the optional `visibility:`
124/// field of [`metric!`] and consumed by the metrics catalog
125/// (`bin/gen-metrics-catalog`), which reads it from the source tree to produce
126/// the user-facing metrics reference. It has no effect on the metric at
127/// runtime.
128#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, serde::Serialize)]
129#[serde(rename_all = "snake_case")]
130pub enum MetricVisibility {
131    /// A metric we use for internal development.
132    #[default]
133    Internal,
134    /// A metric we want customers to build dashboards and
135    /// alerts on. We do not guarantee stability for this group
136    /// of metrics.
137    Public,
138}
139
140/// A tag categorizing a metric in the user-facing metrics catalog.
141///
142/// It is set via the optional `tags:` field of [`metric!`] and
143/// consumed by the metrics catalog (`bin/gen-metrics-catalog`).
144/// It has no effect on the metric at runtime.
145///
146/// A metric may carry many tags, or none (the default). A tag names the
147/// grouping the metric is presented under in user-facing documentation.
148#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize)]
149#[serde(rename_all = "kebab-case")]
150pub enum MetricTag {
151    /// SQL Control plane metrics (client connections, availability, catalog).
152    Environment,
153    /// Metrics for compute objects (indexes, materialized views).
154    Compute,
155    /// Metrics for sources.
156    Source,
157    /// Metrics for sinks.
158    Sink,
159}
160
161/// The materialize metrics registry.
162#[derive(Clone, Derivative)]
163#[derivative(Debug)]
164pub struct MetricsRegistry {
165    inner: Registry,
166    #[derivative(Debug = "ignore")]
167    postprocessors: Arc<Mutex<Vec<Box<dyn FnMut(&mut Vec<MetricFamily>) + Send + Sync>>>>,
168}
169
170/// A wrapper for metrics to require delete on drop semantics
171///
172/// The wrapper behaves like regular metrics but only provides functions to create delete-on-drop
173/// variants. This way, no metrics of this type can be leaked.
174///
175/// In situations where the delete-on-drop behavior is not desired or in legacy code, use the raw
176/// variants of the metrics, as defined in [self::raw].
177#[derive(Clone)]
178pub struct DeleteOnDropWrapper<M> {
179    inner: M,
180}
181
182impl<M: MakeCollector + Debug> Debug for DeleteOnDropWrapper<M> {
183    fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
184        self.inner.fmt(f)
185    }
186}
187
188impl<M: Collector> Collector for DeleteOnDropWrapper<M> {
189    fn desc(&self) -> Vec<&Desc> {
190        self.inner.desc()
191    }
192
193    fn collect(&self) -> Vec<MetricFamily> {
194        self.inner.collect()
195    }
196}
197
198impl<M: MakeCollector> MakeCollector for DeleteOnDropWrapper<M> {
199    fn make_collector(opts: MakeCollectorOpts) -> Self {
200        DeleteOnDropWrapper {
201            inner: M::make_collector(opts),
202        }
203    }
204}
205
206impl<M: MetricVecExt> DeleteOnDropWrapper<M> {
207    /// Returns a metric that deletes its labels from this metrics vector when dropped.
208    pub fn get_delete_on_drop_metric<L: PromLabelsExt>(
209        &self,
210        labels: L,
211    ) -> DeleteOnDropMetric<M, L> {
212        self.inner.get_delete_on_drop_metric(labels)
213    }
214}
215
216/// The unsigned integer version of [`Gauge`]. Provides better performance if
217/// metric values are all unsigned integers.
218pub type UIntGauge = GenericGauge<AtomicU64>;
219
220/// Delete-on-drop shadow of Prometheus [prometheus::CounterVec].
221pub type CounterVec = DeleteOnDropWrapper<prometheus::CounterVec>;
222/// Delete-on-drop shadow of Prometheus [prometheus::Gauge].
223pub type Gauge = DeleteOnDropWrapper<prometheus::Gauge>;
224/// Delete-on-drop shadow of Prometheus [prometheus::GaugeVec].
225pub type GaugeVec = DeleteOnDropWrapper<prometheus::GaugeVec>;
226/// Delete-on-drop shadow of Prometheus [prometheus::HistogramVec].
227pub type HistogramVec = DeleteOnDropWrapper<prometheus::HistogramVec>;
228/// Delete-on-drop shadow of Prometheus [prometheus::IntCounterVec].
229pub type IntCounterVec = DeleteOnDropWrapper<prometheus::IntCounterVec>;
230/// Delete-on-drop shadow of Prometheus [prometheus::IntGaugeVec].
231pub type IntGaugeVec = DeleteOnDropWrapper<prometheus::IntGaugeVec>;
232/// Delete-on-drop shadow of Prometheus [raw::UIntGaugeVec].
233pub type UIntGaugeVec = DeleteOnDropWrapper<raw::UIntGaugeVec>;
234
235use crate::assert_none;
236
237pub use prometheus::{Counter, Histogram, IntCounter, IntGauge};
238
239/// Access to non-delete-on-drop vector types
240pub mod raw {
241    use prometheus::core::{AtomicU64, GenericGaugeVec};
242
243    /// The unsigned integer version of [`GaugeVec`].
244    /// Provides better performance if metric values are all unsigned integers.
245    pub type UIntGaugeVec = GenericGaugeVec<AtomicU64>;
246
247    pub use prometheus::{CounterVec, Gauge, GaugeVec, HistogramVec, IntCounterVec, IntGaugeVec};
248}
249
250impl MetricsRegistry {
251    /// Creates a new metrics registry.
252    pub fn new() -> Self {
253        MetricsRegistry {
254            inner: Registry::new(),
255            postprocessors: Arc::new(Mutex::new(vec![])),
256        }
257    }
258
259    /// Register a metric defined with the [`metric`] macro.
260    pub fn register<M>(&self, opts: MakeCollectorOpts) -> M
261    where
262        M: MakeCollector,
263    {
264        let collector = M::make_collector(opts);
265        self.inner.register(Box::new(collector.clone())).unwrap();
266        collector
267    }
268
269    /// Registers a gauge whose value is computed when observed.
270    pub fn register_computed_gauge<P>(
271        &self,
272        opts: MakeCollectorOpts,
273        f: impl Fn() -> P::T + Send + Sync + 'static,
274    ) -> ComputedGenericGauge<P>
275    where
276        P: Atomic + 'static,
277    {
278        let gauge = ComputedGenericGauge {
279            gauge: GenericGauge::make_collector(opts),
280            f: Arc::new(f),
281        };
282        self.inner.register(Box::new(gauge.clone())).unwrap();
283        gauge
284    }
285
286    /// Register a pre-defined prometheus collector.
287    pub fn register_collector<C: 'static + prometheus::core::Collector>(&self, collector: C) {
288        self.inner
289            .register(Box::new(collector))
290            .expect("registering pre-defined metrics collector");
291    }
292
293    /// Register a pre-defined collector and return a handle that unregisters it on drop.
294    ///
295    /// Use this for collectors with a bounded lifetime (e.g. a metric sink that is torn down with
296    /// its dataflow), so the collector's series stop being scraped once the owner drops the handle.
297    /// The returned value is opaque: the caller only needs to hold it for as long as the collector
298    /// should stay registered, then drop it.
299    ///
300    /// `prometheus::Registry::unregister` matches a collector by the id of its `Desc`s, not by
301    /// object identity, so the guard keeps a clone of the collector and hands it back to
302    /// `unregister` on drop.
303    ///
304    /// If a collector with the same descriptor id is already registered this soft-panics and
305    /// returns a handle that owns no registration. A duplicate registration is a logic error, but
306    /// panicking here would run on a worker thread and could crash the process, so it degrades to
307    /// missing series rather than taking down the whole scrape. The contract is that a caller
308    /// replacing a collector must drop the old handle before registering the new one: doing so
309    /// unregisters the old descriptor id first, so the re-registration succeeds cleanly.
310    ///
311    /// A caller that expects such a collision should use
312    /// [`MetricsRegistry::try_register_collector_with_dropper`] instead.
313    pub fn register_collector_with_dropper<C>(&self, collector: C) -> Box<dyn Any + Send + Sync>
314    where
315        C: 'static + prometheus::core::Collector + Clone + Send + Sync,
316    {
317        self.try_register_collector_with_dropper(collector)
318            .unwrap_or_else(|e| {
319                crate::soft_panic_or_log!("collector already registered: {e}");
320                // Nothing was registered, so the handle must not unregister anything on drop.
321                Box::new(())
322            })
323    }
324
325    /// Like [`MetricsRegistry::register_collector_with_dropper`], but returns `Err` on a duplicate
326    /// descriptor id instead of soft-panicking.
327    ///
328    /// For callers that expect a transient collision, where a collector is re-registered before
329    /// its predecessor's teardown has dropped the old handle, and want to retry until the id frees
330    /// up. Nothing is registered on `Err`, so the caller may retry with the same collector.
331    pub fn try_register_collector_with_dropper<C>(
332        &self,
333        collector: C,
334    ) -> Result<Box<dyn Any + Send + Sync>, prometheus::Error>
335    where
336        C: 'static + prometheus::core::Collector + Clone + Send + Sync,
337    {
338        self.inner.register(Box::new(collector.clone()))?;
339        // `prometheus::Registry` is `Arc`-backed, so this clone is cheap and shares the same
340        // underlying registry the collector was registered into.
341        let registry = self.inner.clone();
342        Ok(Box::new(scopeguard::guard(collector, move |c| {
343            let _ = registry.unregister(Box::new(c));
344        })))
345    }
346
347    /// Registers a metric postprocessor.
348    ///
349    /// Postprocessors are invoked on every call to [`MetricsRegistry::gather`]
350    /// in the order that they are registered.
351    pub fn register_postprocessor<F>(&self, f: F)
352    where
353        F: FnMut(&mut Vec<MetricFamily>) + Send + Sync + 'static,
354    {
355        let mut postprocessors = self.postprocessors.lock().expect("lock poisoned");
356        postprocessors.push(Box::new(f));
357    }
358
359    /// Gather all the metrics from the metrics registry for reporting.
360    ///
361    /// This function invokes the postprocessors on all gathered metrics (see
362    /// [`MetricsRegistry::register_postprocessor`]) in the order the
363    /// postprocessors were registered.
364    ///
365    /// See also [`prometheus::Registry::gather`].
366    pub fn gather(&self) -> Vec<MetricFamily> {
367        let mut metrics = self.inner.gather();
368        let mut postprocessors = self.postprocessors.lock().expect("lock poisoned");
369        for postprocessor in &mut *postprocessors {
370            postprocessor(&mut metrics);
371        }
372        metrics
373    }
374}
375
376/// A wrapper for creating prometheus metrics more conveniently.
377///
378/// Together with the [`metric`] macro, this trait is mainly used by [`MetricsRegistry`] and should
379/// not normally be used outside the metric registration flow.
380pub trait MakeCollector: Collector + Clone + 'static {
381    /// Creates a new collector.
382    fn make_collector(opts: MakeCollectorOpts) -> Self;
383}
384
385impl<T> MakeCollector for GenericCounter<T>
386where
387    T: Atomic + 'static,
388{
389    fn make_collector(mk_opts: MakeCollectorOpts) -> Self {
390        assert_none!(mk_opts.buckets);
391        Self::with_opts(mk_opts.opts).expect("defining a counter")
392    }
393}
394
395impl<T> MakeCollector for GenericCounterVec<T>
396where
397    T: Atomic + 'static,
398{
399    fn make_collector(mk_opts: MakeCollectorOpts) -> Self {
400        assert_none!(mk_opts.buckets);
401        let labels: Vec<String> = mk_opts.opts.variable_labels.clone();
402        let label_refs: Vec<&str> = labels.iter().map(String::as_str).collect();
403        Self::new(mk_opts.opts, label_refs.as_slice()).expect("defining a counter vec")
404    }
405}
406
407impl<T> MakeCollector for GenericGauge<T>
408where
409    T: Atomic + 'static,
410{
411    fn make_collector(mk_opts: MakeCollectorOpts) -> Self {
412        assert_none!(mk_opts.buckets);
413        Self::with_opts(mk_opts.opts).expect("defining a gauge")
414    }
415}
416
417impl<T> MakeCollector for GenericGaugeVec<T>
418where
419    T: Atomic + 'static,
420{
421    fn make_collector(mk_opts: MakeCollectorOpts) -> Self {
422        assert_none!(mk_opts.buckets);
423        let labels = mk_opts.opts.variable_labels.clone();
424        let labels = &labels.iter().map(|x| x.as_str()).collect::<Vec<_>>();
425        Self::new(mk_opts.opts, labels).expect("defining a gauge vec")
426    }
427}
428
429impl MakeCollector for Histogram {
430    fn make_collector(mk_opts: MakeCollectorOpts) -> Self {
431        assert!(mk_opts.buckets.is_some());
432        Self::with_opts(HistogramOpts {
433            common_opts: mk_opts.opts,
434            buckets: mk_opts.buckets.unwrap(),
435        })
436        .expect("defining a histogram")
437    }
438}
439
440impl MakeCollector for raw::HistogramVec {
441    fn make_collector(mk_opts: MakeCollectorOpts) -> Self {
442        assert!(mk_opts.buckets.is_some());
443        let labels = mk_opts.opts.variable_labels.clone();
444        let labels = &labels.iter().map(|x| x.as_str()).collect::<Vec<_>>();
445        Self::new(
446            HistogramOpts {
447                common_opts: mk_opts.opts,
448                buckets: mk_opts.buckets.unwrap(),
449            },
450            labels,
451        )
452        .expect("defining a histogram vec")
453    }
454}
455
456/// A [`Gauge`] whose value is computed whenever it is observed.
457pub struct ComputedGenericGauge<P>
458where
459    P: Atomic,
460{
461    gauge: GenericGauge<P>,
462    f: Arc<dyn Fn() -> P::T + Send + Sync>,
463}
464
465impl<P> fmt::Debug for ComputedGenericGauge<P>
466where
467    P: Atomic + fmt::Debug,
468{
469    fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
470        f.debug_struct("ComputedGenericGauge")
471            .field("gauge", &self.gauge)
472            .finish_non_exhaustive()
473    }
474}
475
476impl<P> Clone for ComputedGenericGauge<P>
477where
478    P: Atomic,
479{
480    fn clone(&self) -> ComputedGenericGauge<P> {
481        ComputedGenericGauge {
482            gauge: self.gauge.clone(),
483            f: Arc::clone(&self.f),
484        }
485    }
486}
487
488impl<T> Collector for ComputedGenericGauge<T>
489where
490    T: Atomic,
491{
492    fn desc(&self) -> Vec<&prometheus::core::Desc> {
493        self.gauge.desc()
494    }
495
496    fn collect(&self) -> Vec<MetricFamily> {
497        self.gauge.set((self.f)());
498        self.gauge.collect()
499    }
500}
501
502impl<P> ComputedGenericGauge<P>
503where
504    P: Atomic,
505{
506    /// Computes the current value of the gauge.
507    pub fn get(&self) -> P::T {
508        (self.f)()
509    }
510}
511
512/// A [`ComputedGenericGauge`] for 64-bit floating point numbers.
513pub type ComputedGauge = ComputedGenericGauge<AtomicF64>;
514
515/// A [`ComputedGenericGauge`] for 64-bit signed integers.
516pub type ComputedIntGauge = ComputedGenericGauge<AtomicI64>;
517
518/// A [`ComputedGenericGauge`] for 64-bit unsigned integers.
519pub type ComputedUIntGauge = ComputedGenericGauge<AtomicU64>;
520
521/// Exposes combinators that report metrics related to the execution of a [`Future`] to prometheus.
522pub trait MetricsFutureExt<F> {
523    /// Records the number of seconds it takes a [`Future`] to complete according to "the clock on
524    /// the wall".
525    ///
526    /// More specifically, it records the instant at which the `Future` was first polled, and the
527    /// instant at which the `Future` completes. Then reports the duration between those two
528    /// instances to the provided metric.
529    ///
530    /// # Wall Time vs Execution Time
531    ///
532    /// There is also [`MetricsFutureExt::exec_time`], which measures how long a [`Future`] spent
533    /// executing, instead of how long it took to complete. For example, a network request may have
534    /// a wall time of 1 second, meanwhile it's execution time may have only been 50ms. The 950ms
535    /// delta would be how long the [`Future`] waited for a response from the network.
536    ///
537    /// # Uses
538    ///
539    /// Recording the wall time can be useful for monitoring latency, for example the latency of a
540    /// SQL request.
541    ///
542    /// Note: You must call either [`observe`] to record the execution time to a [`Histogram`] or
543    /// [`inc_by`] to record to a [`Counter`]. The following will not compile:
544    ///
545    /// ```compile_fail
546    /// use mz_ore::metrics::MetricsFutureExt;
547    ///
548    /// # let _ = async {
549    /// async { Ok(()) }
550    ///     .wall_time()
551    ///     .await;
552    /// # };
553    /// ```
554    ///
555    /// [`observe`]: WallTimeFuture::observe
556    /// [`inc_by`]: WallTimeFuture::inc_by
557    fn wall_time(self) -> WallTimeFuture<F, UnspecifiedMetric>;
558
559    /// Records the total number of seconds for which a [`Future`] was executing.
560    ///
561    /// More specifically, every time the `Future` is polled it records how long that individual
562    /// call took, and maintains a running sum until the `Future` completes. Then we report that
563    /// duration to the provided metric.
564    ///
565    /// # Wall Time vs Execution Time
566    ///
567    /// There is also [`MetricsFutureExt::wall_time`], which measures how long a [`Future`] took to
568    /// complete, instead of how long it spent executing. For example, a network request may have
569    /// a wall time of 1 second, meanwhile it's execution time may have only been 50ms. The 950ms
570    /// delta would be how long the [`Future`] waited for a response from the network.
571    ///
572    /// # Uses
573    ///
574    /// Recording execution time can be useful if you want to monitor [`Future`]s that could be
575    /// sensitive to CPU usage. For example, if you have a single logical control thread you'll
576    /// want to make sure that thread never spends too long running a single `Future`. Reporting
577    /// the execution time of `Future`s running on this thread can help ensure there is no
578    /// unexpected blocking.
579    ///
580    /// Note: You must call either [`observe`] to record the execution time to a [`Histogram`] or
581    /// [`inc_by`] to record to a [`Counter`]. The following will not compile:
582    ///
583    /// ```compile_fail
584    /// use mz_ore::metrics::MetricsFutureExt;
585    ///
586    /// # let _ = async {
587    /// async { Ok(()) }
588    ///     .exec_time()
589    ///     .await;
590    /// # };
591    /// ```
592    ///
593    /// [`observe`]: ExecTimeFuture::observe
594    /// [`inc_by`]: ExecTimeFuture::inc_by
595    fn exec_time(self) -> ExecTimeFuture<F, UnspecifiedMetric>;
596}
597
598impl<F: Future> MetricsFutureExt<F> for F {
599    fn wall_time(self) -> WallTimeFuture<F, UnspecifiedMetric> {
600        WallTimeFuture {
601            fut: self,
602            metric: UnspecifiedMetric(()),
603            start: None,
604            filter: None,
605        }
606    }
607
608    fn exec_time(self) -> ExecTimeFuture<F, UnspecifiedMetric> {
609        ExecTimeFuture {
610            fut: self,
611            metric: UnspecifiedMetric(()),
612            running_duration: Duration::from_millis(0),
613            filter: None,
614        }
615    }
616}
617
618/// Future returned by [`MetricsFutureExt::wall_time`].
619#[must_use = "futures do nothing unless you `.await` or poll them"]
620#[pin_project]
621pub struct WallTimeFuture<F, Metric> {
622    /// The inner [`Future`] that we're recording the wall time for.
623    #[pin]
624    fut: F,
625    /// Prometheus metric that we'll report to.
626    metric: Metric,
627    /// [`Instant`] at which the [`Future`] was first polled.
628    start: Option<Instant>,
629    /// Optional filter that determines if we observe the wall time of this [`Future`].
630    filter: Option<Box<dyn FnMut(Duration) -> bool + Send + Sync>>,
631}
632
633impl<F: Debug, M: Debug> fmt::Debug for WallTimeFuture<F, M> {
634    fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
635        f.debug_struct("WallTimeFuture")
636            .field("fut", &self.fut)
637            .field("metric", &self.metric)
638            .field("start", &self.start)
639            .field("filter", &self.filter.is_some())
640            .finish()
641    }
642}
643
644impl<F> WallTimeFuture<F, UnspecifiedMetric> {
645    /// Sets the recored metric to be a [`prometheus::Histogram`].
646    ///
647    /// ```text
648    /// my_future
649    ///     .wall_time()
650    ///     .observe(metrics.slow_queries_hist.with_label_values(&["select"]))
651    /// ```
652    pub fn observe(
653        self,
654        histogram: prometheus::Histogram,
655    ) -> WallTimeFuture<F, prometheus::Histogram> {
656        WallTimeFuture {
657            fut: self.fut,
658            metric: histogram,
659            start: self.start,
660            filter: self.filter,
661        }
662    }
663
664    /// Sets the recored metric to be a [`prometheus::Counter`].
665    ///
666    /// ```text
667    /// my_future
668    ///     .wall_time()
669    ///     .inc_by(metrics.slow_queries.with_label_values(&["select"]))
670    /// ```
671    pub fn inc_by(self, counter: prometheus::Counter) -> WallTimeFuture<F, prometheus::Counter> {
672        WallTimeFuture {
673            fut: self.fut,
674            metric: counter,
675            start: self.start,
676            filter: self.filter,
677        }
678    }
679
680    /// Sets the recorded duration in a specific f64.
681    pub fn set_at(self, place: &mut f64) -> WallTimeFuture<F, &mut f64> {
682        WallTimeFuture {
683            fut: self.fut,
684            metric: place,
685            start: self.start,
686            filter: self.filter,
687        }
688    }
689}
690
691impl<F, M> WallTimeFuture<F, M> {
692    /// Specifies a filter which much return `true` for the wall time to be recorded.
693    ///
694    /// This can be particularly useful if you have a high volume `Future` and you only want to
695    /// record ones that take a long time to complete.
696    pub fn with_filter(
697        mut self,
698        filter: impl FnMut(Duration) -> bool + Send + Sync + 'static,
699    ) -> Self {
700        self.filter = Some(Box::new(filter));
701        self
702    }
703}
704
705impl<F: Future, M: DurationMetric> Future for WallTimeFuture<F, M> {
706    type Output = F::Output;
707
708    fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
709        let this = self.project();
710
711        if this.start.is_none() {
712            *this.start = Some(Instant::now());
713        }
714
715        let result = match this.fut.poll(cx) {
716            Poll::Ready(r) => r,
717            Poll::Pending => return Poll::Pending,
718        };
719        let duration = Instant::now().duration_since(this.start.expect("timer to be started"));
720
721        let pass = this
722            .filter
723            .as_mut()
724            .map(|filter| filter(duration))
725            .unwrap_or(true);
726        if pass {
727            this.metric.record(duration.as_secs_f64())
728        }
729
730        Poll::Ready(result)
731    }
732}
733
734/// Future returned by [`MetricsFutureExt::exec_time`].
735#[must_use = "futures do nothing unless you `.await` or poll them"]
736#[pin_project]
737pub struct ExecTimeFuture<F, Metric> {
738    /// The inner [`Future`] that we're recording the wall time for.
739    #[pin]
740    fut: F,
741    /// Prometheus metric that we'll report to.
742    metric: Metric,
743    /// Total [`Duration`] for which this [`Future`] has been executing.
744    running_duration: Duration,
745    /// Optional filter that determines if we observe the execution time of this [`Future`].
746    filter: Option<Box<dyn FnMut(Duration) -> bool + Send + Sync>>,
747}
748
749impl<F: Debug, M: Debug> fmt::Debug for ExecTimeFuture<F, M> {
750    fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
751        f.debug_struct("ExecTimeFuture")
752            .field("fut", &self.fut)
753            .field("metric", &self.metric)
754            .field("running_duration", &self.running_duration)
755            .field("filter", &self.filter.is_some())
756            .finish()
757    }
758}
759
760impl<F> ExecTimeFuture<F, UnspecifiedMetric> {
761    /// Sets the recored metric to be a [`prometheus::Histogram`].
762    ///
763    /// ```text
764    /// my_future
765    ///     .exec_time()
766    ///     .observe(metrics.slow_queries_hist.with_label_values(&["select"]))
767    /// ```
768    pub fn observe(
769        self,
770        histogram: prometheus::Histogram,
771    ) -> ExecTimeFuture<F, prometheus::Histogram> {
772        ExecTimeFuture {
773            fut: self.fut,
774            metric: histogram,
775            running_duration: self.running_duration,
776            filter: self.filter,
777        }
778    }
779
780    /// Sets the recored metric to be a [`prometheus::Counter`].
781    ///
782    /// ```text
783    /// my_future
784    ///     .exec_time()
785    ///     .inc_by(metrics.slow_queries.with_label_values(&["select"]))
786    /// ```
787    pub fn inc_by(self, counter: prometheus::Counter) -> ExecTimeFuture<F, prometheus::Counter> {
788        ExecTimeFuture {
789            fut: self.fut,
790            metric: counter,
791            running_duration: self.running_duration,
792            filter: self.filter,
793        }
794    }
795}
796
797impl<F, M> ExecTimeFuture<F, M> {
798    /// Specifies a filter which much return `true` for the execution time to be recorded.
799    pub fn with_filter(
800        mut self,
801        filter: impl FnMut(Duration) -> bool + Send + Sync + 'static,
802    ) -> Self {
803        self.filter = Some(Box::new(filter));
804        self
805    }
806}
807
808impl<F: Future, M: DurationMetric> Future for ExecTimeFuture<F, M> {
809    type Output = F::Output;
810
811    fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
812        let this = self.project();
813
814        let start = Instant::now();
815        let result = this.fut.poll(cx);
816        let duration = Instant::now().duration_since(start);
817
818        *this.running_duration = this.running_duration.saturating_add(duration);
819
820        let result = match result {
821            Poll::Ready(result) => result,
822            Poll::Pending => return Poll::Pending,
823        };
824
825        let duration = *this.running_duration;
826        let pass = this
827            .filter
828            .as_mut()
829            .map(|filter| filter(duration))
830            .unwrap_or(true);
831        if pass {
832            this.metric.record(duration.as_secs_f64());
833        }
834
835        Poll::Ready(result)
836    }
837}
838
839/// A type level flag used to ensure callers specify the kind of metric to record for
840/// [`MetricsFutureExt`].
841///
842/// For example, `WallTimeFuture<F, M>` only implements [`Future`] for `M` that implements
843/// `DurationMetric` which [`UnspecifiedMetric`] does not. This forces users at build time to
844/// call [`WallTimeFuture::observe`] or [`WallTimeFuture::inc_by`].
845#[derive(Debug)]
846pub struct UnspecifiedMetric(());
847
848/// A trait makes recording a duration generic over different prometheus metrics. This allows us to
849/// de-dupe the implemenation of [`Future`] for our wrapper Futures like [`WallTimeFuture`] and
850/// [`ExecTimeFuture`] over different kinds of prometheus metrics.
851trait DurationMetric {
852    fn record(&mut self, seconds: f64);
853}
854
855impl DurationMetric for prometheus::Histogram {
856    fn record(&mut self, seconds: f64) {
857        self.observe(seconds)
858    }
859}
860
861impl DurationMetric for prometheus::Counter {
862    fn record(&mut self, seconds: f64) {
863        self.inc_by(seconds)
864    }
865}
866
867// An implementation of `DurationMetric` that lets the user take the recorded
868// value and use it elsewhere.
869impl DurationMetric for &'_ mut f64 {
870    fn record(&mut self, seconds: f64) {
871        **self = seconds;
872    }
873}
874
875/// Register the Tokio runtime's metrics in our metrics registry.
876#[cfg(feature = "async")]
877pub fn register_runtime_metrics(
878    name: &'static str,
879    runtime_metrics: tokio::runtime::RuntimeMetrics,
880    registry: &MetricsRegistry,
881) {
882    macro_rules! register {
883        ($method:ident, $doc:literal) => {
884            let metrics = runtime_metrics.clone();
885            registry.register_computed_gauge::<prometheus::core::AtomicU64>(
886                crate::metric!(
887                    name: concat!("mz_tokio_", stringify!($method)),
888                    help: $doc,
889                    const_labels: {"runtime" => name},
890                ),
891                move || <u64 as crate::cast::CastFrom<_>>::cast_from(metrics.$method()),
892            );
893        };
894    }
895
896    macro_rules! register_per_worker {
897        ($method:ident, $doc:literal) => {
898            let metrics = runtime_metrics.clone();
899            registry.register_computed_gauge::<prometheus::core::AtomicU64>(
900                crate::metric!(
901                    name: concat!("mz_tokio_", stringify!($method)),
902                    help: $doc,
903                    const_labels: {"runtime" => name},
904                ),
905                move || {
906                    (0..metrics.num_workers())
907                        .map(|i| <u64 as crate::cast::CastFrom<_>>::cast_from(metrics.$method(i)))
908                        .sum::<u64>()
909                },
910            );
911        };
912    }
913
914    macro_rules! register_per_worker_duration_secs {
915        ($method:ident, $doc:literal) => {
916            let metrics = runtime_metrics.clone();
917            registry.register_computed_gauge::<prometheus::core::AtomicF64>(
918                crate::metric!(
919                    name: concat!("mz_tokio_", stringify!($method)),
920                    help: $doc,
921                    const_labels: {"runtime" => name},
922                ),
923                move || {
924                    (0..metrics.num_workers())
925                        .map(|i| metrics.$method(i).as_secs_f64())
926                        .sum::<f64>()
927                },
928            );
929        };
930    }
931
932    register!(
933        num_workers,
934        "The number of worker threads used by the runtime."
935    );
936    register!(
937        num_alive_tasks,
938        "The current number of alive tasks in the runtime."
939    );
940    register!(
941        global_queue_depth,
942        "The number of tasks currently scheduled in the runtime's global queue."
943    );
944    register_per_worker_duration_secs!(
945        worker_total_busy_duration,
946        "The amount of time the worker threads have been busy, in seconds."
947    );
948    register_per_worker!(
949        worker_park_count,
950        "The total number of times the worker threads have parked."
951    );
952    register_per_worker!(
953        worker_park_unpark_count,
954        "The total number of times the worker threads have parked and unparked."
955    );
956
957    #[cfg(tokio_unstable)]
958    {
959        register!(
960            num_blocking_threads,
961            "The number of additional threads spawned by the runtime."
962        );
963        register!(
964            num_idle_blocking_threads,
965            "The number of idle threads which have spawned by the runtime for spawn_blocking calls."
966        );
967        register_per_worker!(
968            worker_local_queue_depth,
969            "The number of tasks currently scheduled in the workers' local queues."
970        );
971        register!(
972            blocking_queue_depth,
973            "The number of tasks currently scheduled in the blocking thread pool, spawned using spawn_blocking."
974        );
975        register!(
976            spawned_tasks_count,
977            "The number of tasks spawned in this runtime since it was created."
978        );
979        register!(
980            remote_schedule_count,
981            "The number of tasks scheduled from outside of the runtime."
982        );
983        register!(
984            budget_forced_yield_count,
985            "The number of times that tasks have been forced to yield back to the scheduler after exhausting their task budgets."
986        );
987        register_per_worker!(
988            worker_noop_count,
989            "The number of times the given worker thread unparked but performed no work before parking again."
990        );
991        register_per_worker!(
992            worker_steal_count,
993            "The number of tasks the given worker thread stole from another worker thread."
994        );
995        register_per_worker!(
996            worker_steal_operations,
997            "The number of times the given worker thread stole tasks from another worker thread."
998        );
999        register_per_worker!(
1000            worker_poll_count,
1001            "The number of tasks the given worker thread has polled."
1002        );
1003        register_per_worker!(
1004            worker_local_schedule_count,
1005            "The number of tasks scheduled from within the runtime on the given worker's local queue."
1006        );
1007        register_per_worker!(
1008            worker_overflow_count,
1009            "The number of times the given worker thread saturated its local queue."
1010        );
1011        register_per_worker_duration_secs!(
1012            worker_mean_poll_time,
1013            "The mean duration of task polls in seconds."
1014        );
1015    }
1016}
1017
1018/// Returns the `(name, help, labels, source)` of every Tokio runtime metric
1019/// registered by [`register_runtime_metrics`].
1020#[cfg(feature = "async")]
1021pub fn describe_runtime_metrics() -> Vec<(String, String, Vec<String>, &'static str)> {
1022    // A current-thread runtime is enough to enumerate the metrics; we only read
1023    // their names, help text, and labels, never their values.
1024    let runtime = tokio::runtime::Builder::new_current_thread()
1025        .build()
1026        .expect("building a current-thread runtime");
1027    let registry = MetricsRegistry::new();
1028    register_runtime_metrics("describe", runtime.handle().metrics(), &registry);
1029    registry
1030        .gather()
1031        .into_iter()
1032        .map(|mf| {
1033            // Every series in a family shares the same label keys, so the first
1034            // metric's labels are representative.
1035            let mut labels: Vec<String> = mf
1036                .get_metric()
1037                .first()
1038                .map(|m| m.get_label().iter().map(|l| l.name().to_owned()).collect())
1039                .unwrap_or_default();
1040            labels.sort();
1041            labels.dedup();
1042            (mf.name().to_owned(), mf.help().to_owned(), labels, file!())
1043        })
1044        .collect()
1045}
1046
1047/// Removes every child of `vec` whose label `name` has the value `value`.
1048///
1049/// Prometheus removes children only by their full label tuple, so this learns
1050/// the tuples by collecting the vec, which clones every child once. Meant for
1051/// occasional cleanup such as dropping an object's series, not for hot paths.
1052pub fn remove_children_with_label<V: MetricVec_ + Collector>(vec: &V, name: &str, value: &str) {
1053    let descs = vec.desc();
1054    // A metric vec has exactly one `Desc`.
1055    let Some(desc) = descs.first() else {
1056        return;
1057    };
1058    for family in vec.collect() {
1059        for child in family.get_metric() {
1060            let labels: BTreeMap<&str, &str> = child
1061                .get_label()
1062                .iter()
1063                .map(|pair| (pair.name(), pair.value()))
1064                .collect();
1065            if labels.get(name) != Some(&value) {
1066                continue;
1067            }
1068            let values: Vec<&str> = desc
1069                .variable_labels
1070                .iter()
1071                .map(|label| labels.get(label.as_str()).copied().unwrap_or_default())
1072                .collect();
1073            // `remove_label_values` fails when the series is already gone, which happens
1074            // if another caller removed it between the `collect` snapshot above and now.
1075            // That is the state we wanted, so the error is ignored.
1076            let _ = vec.remove_label_values(&values);
1077        }
1078    }
1079}
1080
1081#[cfg(test)]
1082mod tests {
1083    use std::time::Duration;
1084
1085    use prometheus::core::Collector;
1086    use prometheus::{CounterVec, HistogramVec};
1087
1088    use crate::stats::histogram_seconds_buckets;
1089
1090    use super::{MetricsFutureExt, MetricsRegistry};
1091
1092    struct Metrics {
1093        pub wall_time_hist: HistogramVec,
1094        pub wall_time_cnt: CounterVec,
1095        pub exec_time_hist: HistogramVec,
1096        pub exec_time_cnt: CounterVec,
1097    }
1098
1099    impl Metrics {
1100        pub fn register_into(registry: &MetricsRegistry) -> Self {
1101            Self {
1102                wall_time_hist: registry.register(metric!(
1103                    name: "wall_time_hist",
1104                    help: "help",
1105                    var_labels: ["action"],
1106                    buckets: histogram_seconds_buckets(0.000_128, 8.0),
1107                )),
1108                wall_time_cnt: registry.register(metric!(
1109                    name: "wall_time_cnt",
1110                    help: "help",
1111                    var_labels: ["action"],
1112                )),
1113                exec_time_hist: registry.register(metric!(
1114                    name: "exec_time_hist",
1115                    help: "help",
1116                    var_labels: ["action"],
1117                    buckets: histogram_seconds_buckets(0.000_128, 8.0),
1118                )),
1119                exec_time_cnt: registry.register(metric!(
1120                    name: "exec_time_cnt",
1121                    help: "help",
1122                    var_labels: ["action"],
1123                )),
1124            }
1125        }
1126    }
1127
1128    #[crate::test]
1129    #[cfg_attr(miri, ignore)] // unsupported operation: integer-to-pointer casts and `ptr::from_exposed_addr` are not supported with `-Zmiri-strict-provenance`
1130    fn smoke_test_metrics_future_ext() {
1131        let runtime = tokio::runtime::Builder::new_current_thread()
1132            .enable_time()
1133            .build()
1134            .expect("failed to start runtime");
1135        let registry = MetricsRegistry::new();
1136        let metrics = Metrics::register_into(&registry);
1137
1138        // Record the walltime and execution time of an async sleep.
1139        let async_sleep_future = async {
1140            tokio::time::sleep(tokio::time::Duration::from_secs(2)).await;
1141        };
1142        runtime.block_on(
1143            async_sleep_future
1144                .wall_time()
1145                .observe(metrics.wall_time_hist.with_label_values(&["async_sleep_w"]))
1146                .exec_time()
1147                .observe(metrics.exec_time_hist.with_label_values(&["async_sleep_e"])),
1148        );
1149
1150        let reports = registry.gather();
1151
1152        let exec_family = reports
1153            .iter()
1154            .find(|m| m.name() == "exec_time_hist")
1155            .expect("metric not found");
1156        let exec_metric = exec_family.get_metric();
1157        assert_eq!(exec_metric.len(), 1);
1158        assert_eq!(exec_metric[0].get_label()[0].value(), "async_sleep_e");
1159
1160        let exec_histogram = exec_metric[0].get_histogram();
1161        assert_eq!(exec_histogram.get_sample_count(), 1);
1162        // This future will normally complete very quickly, but it's hard to guarantee any
1163        // particular timing in an arbitrary test environment, so we don't assert on it
1164        // here.
1165
1166        let wall_family = reports
1167            .iter()
1168            .find(|m| m.name() == "wall_time_hist")
1169            .expect("metric not found");
1170        let wall_metric = wall_family.get_metric();
1171        assert_eq!(wall_metric.len(), 1);
1172        assert_eq!(wall_metric[0].get_label()[0].value(), "async_sleep_w");
1173
1174        let wall_histogram = wall_metric[0].get_histogram();
1175        assert_eq!(wall_histogram.get_sample_count(), 1);
1176        // The 13th bucket is 512ms, which the wall time should be longer than, but is also much
1177        // faster than the actual execution time of the async sleep.
1178        assert_eq!(wall_histogram.get_bucket()[12].cumulative_count(), 0);
1179
1180        // Reset the registery to make collecting metrics easier.
1181        let registry = MetricsRegistry::new();
1182        let metrics = Metrics::register_into(&registry);
1183
1184        // Record the walltime and execution time of a thread sleep.
1185        let thread_sleep_future = async {
1186            std::thread::sleep(std::time::Duration::from_secs(1));
1187        };
1188        runtime.block_on(
1189            thread_sleep_future
1190                .wall_time()
1191                .with_filter(|duration| duration < Duration::from_millis(10))
1192                .inc_by(metrics.wall_time_cnt.with_label_values(&["thread_sleep_w"]))
1193                .exec_time()
1194                .inc_by(metrics.exec_time_cnt.with_label_values(&["thread_sleep_e"])),
1195        );
1196
1197        let reports = registry.gather();
1198
1199        let exec_family = reports
1200            .iter()
1201            .find(|m| m.name() == "exec_time_cnt")
1202            .expect("metric not found");
1203        let exec_metric = exec_family.get_metric();
1204        assert_eq!(exec_metric.len(), 1);
1205        assert_eq!(exec_metric[0].get_label()[0].value(), "thread_sleep_e");
1206
1207        let exec_counter = exec_metric[0].get_counter();
1208        // Since we're synchronously sleeping the execution time will be long.
1209        assert!(exec_counter.value() >= 1.0);
1210
1211        let wall_family = reports
1212            .iter()
1213            .find(|m| m.name() == "wall_time_cnt")
1214            .expect("metric not found");
1215        let wall_metric = wall_family.get_metric();
1216        assert_eq!(wall_metric.len(), 1);
1217
1218        let wall_counter = wall_metric[0].get_counter();
1219        // We filtered wall time to < 10ms, so our wall time metric should be filtered out.
1220        assert_eq!(wall_counter.value(), 0.0);
1221    }
1222
1223    #[crate::test]
1224    fn collector_drop_handle_unregisters() {
1225        use prometheus::IntGauge;
1226
1227        let registry = MetricsRegistry::new();
1228        let gauge = IntGauge::new("mz_test_guarded", "help").unwrap();
1229        gauge.set(7);
1230        let before = registry.gather().len();
1231
1232        let handle = registry.register_collector_with_dropper(gauge.clone());
1233        assert_eq!(registry.gather().len(), before + 1);
1234
1235        // Dropping the handle unregisters the collector, so its series stops being scraped.
1236        drop(handle);
1237        assert_eq!(registry.gather().len(), before);
1238    }
1239
1240    #[crate::test]
1241    fn register_drop_then_reregister() {
1242        use prometheus::IntGauge;
1243
1244        // Two collectors sharing a descriptor id: re-registering the second while the first is
1245        // still registered would collide. This locks in the contract that dropping the old handle
1246        // first clears the id so the re-registration succeeds cleanly, the ordering a caller
1247        // replacing a collector (e.g. a metric sink re-rendered on reconciliation) must uphold.
1248        let registry = MetricsRegistry::new();
1249        let old = IntGauge::new("mz_test_reregister", "help").unwrap();
1250        let new = IntGauge::new("mz_test_reregister", "help").unwrap();
1251        let before = registry.gather().len();
1252
1253        let handle = registry.register_collector_with_dropper(old);
1254        assert_eq!(registry.gather().len(), before + 1);
1255
1256        // Drop old before registering new, matching the required ordering.
1257        drop(handle);
1258        let handle = registry.register_collector_with_dropper(new);
1259        assert_eq!(registry.gather().len(), before + 1);
1260
1261        drop(handle);
1262        assert_eq!(registry.gather().len(), before);
1263    }
1264
1265    #[crate::test]
1266    fn try_register_errors_on_duplicate_then_succeeds_after_drop() {
1267        use prometheus::IntGauge;
1268
1269        // The fallible variant reports the collision instead of soft-panicking, and the collision
1270        // clears as soon as the incumbent's handle drops. This is the retry contract the metric
1271        // sink operator leans on when a new incarnation races the old one's teardown.
1272        let registry = MetricsRegistry::new();
1273        let old = IntGauge::new("mz_test_try_register", "help").unwrap();
1274        let new = IntGauge::new("mz_test_try_register", "help").unwrap();
1275        let before = registry.gather().len();
1276
1277        let handle = registry
1278            .try_register_collector_with_dropper(old)
1279            .expect("first registration succeeds");
1280        assert_eq!(registry.gather().len(), before + 1);
1281
1282        let err = registry
1283            .try_register_collector_with_dropper(new.clone())
1284            .err()
1285            .expect("duplicate descriptor id is rejected");
1286        assert!(matches!(err, prometheus::Error::AlreadyReg));
1287        // The failed attempt registered nothing, so the incumbent is still the only series.
1288        assert_eq!(registry.gather().len(), before + 1);
1289
1290        drop(handle);
1291        let handle = registry
1292            .try_register_collector_with_dropper(new)
1293            .expect("registration succeeds once the id is free");
1294        assert_eq!(registry.gather().len(), before + 1);
1295
1296        drop(handle);
1297        assert_eq!(registry.gather().len(), before);
1298    }
1299
1300    #[crate::test]
1301    fn remove_children_with_label_removes_only_matching_children() {
1302        let registry = MetricsRegistry::new();
1303        let vec: HistogramVec = registry.register(metric!(
1304            name: "labeled_hist",
1305            help: "help",
1306            var_labels: ["cluster", "kind"],
1307            buckets: histogram_seconds_buckets(0.000_128, 8.0),
1308        ));
1309        vec.with_label_values(&["u1", "a"]).observe(1.0);
1310        vec.with_label_values(&["u1", "b"]).observe(1.0);
1311        vec.with_label_values(&["u2", "a"]).observe(1.0);
1312
1313        super::remove_children_with_label(&vec, "cluster", "u1");
1314
1315        let remaining: Vec<Vec<String>> = vec
1316            .collect()
1317            .into_iter()
1318            .flat_map(|family| family.get_metric().to_vec())
1319            .map(|child| {
1320                child
1321                    .get_label()
1322                    .iter()
1323                    .map(|pair| pair.value().to_string())
1324                    .collect()
1325            })
1326            .collect();
1327        assert_eq!(remaining, vec![vec!["u2".to_string(), "a".to_string()]]);
1328    }
1329}