1use 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#[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 $(mk_opts.buckets = Some($bk_name);)*
101 $(let _: $crate::metrics::MetricVisibility = $visibility;)?
106 $($(let _: $crate::metrics::MetricTag = $tag;)*)?
109 mk_opts
110 }}
111}
112
113#[derive(Debug, Clone)]
115pub struct MakeCollectorOpts {
116 pub opts: PrometheusOpts,
118 pub buckets: Option<Vec<f64>>,
121}
122
123#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, serde::Serialize)]
129#[serde(rename_all = "snake_case")]
130pub enum MetricVisibility {
131 #[default]
133 Internal,
134 Public,
138}
139
140#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize)]
149#[serde(rename_all = "kebab-case")]
150pub enum MetricTag {
151 Environment,
153 Compute,
155 Source,
157 Sink,
159}
160
161#[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#[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 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
216pub type UIntGauge = GenericGauge<AtomicU64>;
219
220pub type CounterVec = DeleteOnDropWrapper<prometheus::CounterVec>;
222pub type Gauge = DeleteOnDropWrapper<prometheus::Gauge>;
224pub type GaugeVec = DeleteOnDropWrapper<prometheus::GaugeVec>;
226pub type HistogramVec = DeleteOnDropWrapper<prometheus::HistogramVec>;
228pub type IntCounterVec = DeleteOnDropWrapper<prometheus::IntCounterVec>;
230pub type IntGaugeVec = DeleteOnDropWrapper<prometheus::IntGaugeVec>;
232pub type UIntGaugeVec = DeleteOnDropWrapper<raw::UIntGaugeVec>;
234
235use crate::assert_none;
236
237pub use prometheus::{Counter, Histogram, IntCounter, IntGauge};
238
239pub mod raw {
241 use prometheus::core::{AtomicU64, GenericGaugeVec};
242
243 pub type UIntGaugeVec = GenericGaugeVec<AtomicU64>;
246
247 pub use prometheus::{CounterVec, Gauge, GaugeVec, HistogramVec, IntCounterVec, IntGaugeVec};
248}
249
250impl MetricsRegistry {
251 pub fn new() -> Self {
253 MetricsRegistry {
254 inner: Registry::new(),
255 postprocessors: Arc::new(Mutex::new(vec![])),
256 }
257 }
258
259 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 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 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 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 Box::new(())
322 })
323 }
324
325 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 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 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 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
376pub trait MakeCollector: Collector + Clone + 'static {
381 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
456pub 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 pub fn get(&self) -> P::T {
508 (self.f)()
509 }
510}
511
512pub type ComputedGauge = ComputedGenericGauge<AtomicF64>;
514
515pub type ComputedIntGauge = ComputedGenericGauge<AtomicI64>;
517
518pub type ComputedUIntGauge = ComputedGenericGauge<AtomicU64>;
520
521pub trait MetricsFutureExt<F> {
523 fn wall_time(self) -> WallTimeFuture<F, UnspecifiedMetric>;
558
559 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#[must_use = "futures do nothing unless you `.await` or poll them"]
620#[pin_project]
621pub struct WallTimeFuture<F, Metric> {
622 #[pin]
624 fut: F,
625 metric: Metric,
627 start: Option<Instant>,
629 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 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 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 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 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#[must_use = "futures do nothing unless you `.await` or poll them"]
736#[pin_project]
737pub struct ExecTimeFuture<F, Metric> {
738 #[pin]
740 fut: F,
741 metric: Metric,
743 running_duration: Duration,
745 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 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 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 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#[derive(Debug)]
846pub struct UnspecifiedMetric(());
847
848trait 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
867impl DurationMetric for &'_ mut f64 {
870 fn record(&mut self, seconds: f64) {
871 **self = seconds;
872 }
873}
874
875#[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#[cfg(feature = "async")]
1021pub fn describe_runtime_metrics() -> Vec<(String, String, Vec<String>, &'static str)> {
1022 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(), ®istry);
1029 registry
1030 .gather()
1031 .into_iter()
1032 .map(|mf| {
1033 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
1047pub fn remove_children_with_label<V: MetricVec_ + Collector>(vec: &V, name: &str, value: &str) {
1053 let descs = vec.desc();
1054 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 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)] 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(®istry);
1137
1138 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 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 assert_eq!(wall_histogram.get_bucket()[12].cumulative_count(), 0);
1179
1180 let registry = MetricsRegistry::new();
1182 let metrics = Metrics::register_into(®istry);
1183
1184 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 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 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 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 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(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 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 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}