mz_storage_operators/
metrics.rs1use std::sync::atomic::{AtomicU64, Ordering};
13
14use mz_ore::metrics::{IntCounter, UIntGauge};
15
16#[derive(Debug)]
24pub struct BackpressureOperatorMetrics {
25 pub emitted_bytes: IntCounter,
27 pub last_backpressured_bytes: GaugeContribution,
30 pub retired_bytes: IntCounter,
32}
33
34impl BackpressureOperatorMetrics {
35 pub fn new(
36 emitted_bytes: IntCounter,
37 last_backpressured_bytes: UIntGauge,
38 retired_bytes: IntCounter,
39 ) -> Self {
40 BackpressureOperatorMetrics {
41 emitted_bytes,
42 last_backpressured_bytes: GaugeContribution::new(last_backpressured_bytes),
43 retired_bytes,
44 }
45 }
46}
47
48#[derive(Debug)]
51pub struct GaugeContribution {
52 gauge: UIntGauge,
53 contributed: AtomicU64,
54}
55
56impl GaugeContribution {
57 pub fn new(gauge: UIntGauge) -> Self {
58 GaugeContribution {
59 gauge,
60 contributed: AtomicU64::new(0),
61 }
62 }
63
64 pub fn set(&self, value: u64) {
66 let previous = self.contributed.swap(value, Ordering::AcqRel);
67 self.gauge.add(value);
70 self.gauge.sub(previous);
71 }
72}
73
74impl Drop for GaugeContribution {
75 fn drop(&mut self) {
76 self.gauge.sub(self.contributed.load(Ordering::Acquire));
77 }
78}
79
80#[cfg(test)]
81mod tests {
82 use mz_ore::metrics::UIntGauge;
83
84 use super::GaugeContribution;
85
86 #[mz_ore::test]
87 fn gauge_contribution_sums_live_shares_and_withdraws_on_drop() {
88 let gauge = UIntGauge::new("gauge", "help").expect("valid metric");
89 let a = GaugeContribution::new(gauge.clone());
90 let b = GaugeContribution::new(gauge.clone());
91
92 a.set(5);
93 b.set(3);
94 assert_eq!(gauge.get(), 8);
95
96 a.set(2);
97 assert_eq!(gauge.get(), 5);
98
99 drop(a);
100 assert_eq!(gauge.get(), 3);
101
102 b.set(0);
103 assert_eq!(gauge.get(), 0);
104 }
105}