Skip to main content

platform/monitoring/
metrics_store.rs

1//! Three-tier in-memory ring buffer for per-service metrics.
2//!
3//! Each service gets three retention tiers:
4//! - Tier 0: 1-min granularity, 15 samples (15 min window)
5//! - Tier 1: 15-min granularity, 16 samples (4 h window)
6//! - Tier 2: 4-hour granularity, 18 samples (72 h window)
7//!
8//! A background sampler task reads atomic counters once per minute and
9//! pushes samples into tier 0.  Rollups cascade automatically.
10
11use std::collections::{HashMap, VecDeque};
12use std::sync::Arc;
13use std::sync::atomic::{AtomicU32, AtomicU64, Ordering};
14use std::time::Duration;
15
16use tokio::sync::{Mutex, RwLock};
17
18// ---------------------------------------------------------------------------
19// MetricSample
20// ---------------------------------------------------------------------------
21
22/// One sample point for a single service.
23#[derive(Clone, Debug, serde::Serialize)]
24pub struct MetricSample {
25    /// Unix epoch seconds.
26    pub ts: i64,
27    /// Current open connections (gauge).
28    pub active_conns: u32,
29    /// Request count within the interval (delta).
30    pub requests: u64,
31    /// Failed request count within the interval (delta).
32    pub failed_requests: u64,
33    /// 95th-percentile latency in milliseconds.
34    pub latency_p95_ms: f64,
35}
36
37// ---------------------------------------------------------------------------
38// ServiceCounters
39// ---------------------------------------------------------------------------
40
41/// Atomic counters that services increment in their hot paths.
42///
43/// The 1-minute sampler calls [`snapshot_and_reset`] to drain deltas and
44/// produce a [`MetricSample`].
45pub struct ServiceCounters {
46    /// Current open connections (gauge — not reset on snapshot).
47    pub active_conns: AtomicU32,
48    /// Monotonic total requests since last snapshot.
49    pub total_requests: AtomicU64,
50    /// Monotonic failed requests since last snapshot.
51    pub failed_requests: AtomicU64,
52    /// Latencies collected during the current interval, drained on snapshot.
53    latencies: Mutex<Vec<f64>>,
54}
55
56impl ServiceCounters {
57    pub fn new() -> Self {
58        Self {
59            active_conns: AtomicU32::new(0),
60            total_requests: AtomicU64::new(0),
61            failed_requests: AtomicU64::new(0),
62            latencies: Mutex::new(Vec::new()),
63        }
64    }
65
66    /// Increment the active-connections gauge.
67    pub fn inc_conns(&self) {
68        self.active_conns.fetch_add(1, Ordering::Relaxed);
69    }
70
71    /// Decrement the active-connections gauge.
72    pub fn dec_conns(&self) {
73        self.active_conns.fetch_sub(1, Ordering::Relaxed);
74    }
75
76    /// Record one completed request.
77    ///
78    /// Increments `total_requests` unconditionally, and `failed_requests` when
79    /// `success` is false.  The latency value is pushed into a buffer that
80    /// will be drained when the next sample is taken.
81    pub async fn record_request(&self, success: bool, latency_ms: f64) {
82        self.total_requests.fetch_add(1, Ordering::Relaxed);
83        if !success {
84            self.failed_requests.fetch_add(1, Ordering::Relaxed);
85        }
86        self.latencies.lock().await.push(latency_ms);
87    }
88
89    /// Take a point-in-time snapshot and reset the delta counters.
90    ///
91    /// - `active_conns` is read but **not** reset (it is a gauge).
92    /// - `total_requests` and `failed_requests` are swapped to zero.
93    /// - Latencies are drained and the p95 value is computed.
94    pub async fn snapshot_and_reset(&self) -> MetricSample {
95        let active_conns = self.active_conns.load(Ordering::Relaxed);
96        let requests = self.total_requests.swap(0, Ordering::Relaxed);
97        let failed_requests = self.failed_requests.swap(0, Ordering::Relaxed);
98
99        let mut lats = {
100            let mut guard = self.latencies.lock().await;
101            std::mem::take(&mut *guard)
102        };
103
104        let latency_p95_ms = compute_p95(&mut lats);
105
106        MetricSample {
107            ts: chrono::Utc::now().timestamp(),
108            active_conns,
109            requests,
110            failed_requests,
111            latency_p95_ms,
112        }
113    }
114}
115
116impl Default for ServiceCounters {
117    fn default() -> Self {
118        Self::new()
119    }
120}
121
122impl std::fmt::Debug for ServiceCounters {
123    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
124        f.debug_struct("ServiceCounters")
125            .field("active_conns", &self.active_conns.load(Ordering::Relaxed))
126            .field(
127                "total_requests",
128                &self.total_requests.load(Ordering::Relaxed),
129            )
130            .field(
131                "failed_requests",
132                &self.failed_requests.load(Ordering::Relaxed),
133            )
134            .finish()
135    }
136}
137
138/// Compute the 95th-percentile value from an unsorted slice.
139/// Returns 0.0 for an empty slice.
140fn compute_p95(values: &mut [f64]) -> f64 {
141    if values.is_empty() {
142        return 0.0;
143    }
144    values.sort_by(|a, b| a.partial_cmp(b).unwrap_or(std::cmp::Ordering::Equal));
145    let idx = ((values.len() as f64) * 0.95).ceil() as usize;
146    let idx = idx.min(values.len()) - 1;
147    values[idx]
148}
149
150// ---------------------------------------------------------------------------
151// Ring
152// ---------------------------------------------------------------------------
153
154/// Fixed-capacity ring buffer backed by a `VecDeque`.
155struct Ring {
156    buf: VecDeque<MetricSample>,
157    cap: usize,
158}
159
160impl Ring {
161    fn new(cap: usize) -> Self {
162        Self {
163            buf: VecDeque::with_capacity(cap),
164            cap,
165        }
166    }
167
168    /// Push a sample, evicting the oldest entry when at capacity.
169    fn push(&mut self, sample: MetricSample) {
170        if self.buf.len() == self.cap {
171            self.buf.pop_front();
172        }
173        self.buf.push_back(sample);
174    }
175
176    fn to_vec(&self) -> Vec<MetricSample> {
177        self.buf.iter().cloned().collect()
178    }
179}
180
181// ---------------------------------------------------------------------------
182// ServiceRings
183// ---------------------------------------------------------------------------
184
185/// Three retention tiers for one service, plus staging accumulators for
186/// tier roll-ups.
187struct ServiceRings {
188    tier0: Ring, // cap=15, 1-min granularity
189    tier1: Ring, // cap=16, 15-min granularity
190    tier2: Ring, // cap=18, 4-hour granularity
191    /// Accumulator for tier0 -> tier1 rollup (collects 15 samples).
192    tier1_acc: Vec<MetricSample>,
193    /// Accumulator for tier1 -> tier2 rollup (collects 16 samples).
194    tier2_acc: Vec<MetricSample>,
195}
196
197impl ServiceRings {
198    fn new() -> Self {
199        Self {
200            tier0: Ring::new(15),
201            tier1: Ring::new(16),
202            tier2: Ring::new(18),
203            tier1_acc: Vec::with_capacity(15),
204            tier2_acc: Vec::with_capacity(16),
205        }
206    }
207}
208
209// ---------------------------------------------------------------------------
210// Rollup helpers
211// ---------------------------------------------------------------------------
212
213/// Aggregate a batch of samples into a single rolled-up sample.
214///
215/// Rules:
216/// - `active_conns` -> mean
217/// - `requests` -> sum
218/// - `failed_requests` -> sum
219/// - `latency_p95_ms` -> max
220fn aggregate(samples: &[MetricSample]) -> MetricSample {
221    debug_assert!(!samples.is_empty());
222
223    let ts = samples.last().map(|s| s.ts).unwrap_or(0);
224
225    let conns_sum: u64 = samples.iter().map(|s| s.active_conns as u64).sum();
226    let active_conns = (conns_sum as f64 / samples.len() as f64).round() as u32;
227
228    let requests: u64 = samples.iter().map(|s| s.requests).sum();
229    let failed_requests: u64 = samples.iter().map(|s| s.failed_requests).sum();
230
231    let latency_p95_ms = samples
232        .iter()
233        .map(|s| s.latency_p95_ms)
234        .fold(0.0_f64, f64::max);
235
236    MetricSample {
237        ts,
238        active_conns,
239        requests,
240        failed_requests,
241        latency_p95_ms,
242    }
243}
244
245// ---------------------------------------------------------------------------
246// MetricsStore
247// ---------------------------------------------------------------------------
248
249/// Top-level per-service metrics store, keyed by service type (`ResourceType` as `i32`).
250#[derive(Clone)]
251pub struct MetricsStore {
252    inner: Arc<RwLock<HashMap<i32, ServiceRings>>>,
253}
254
255impl MetricsStore {
256    pub fn new() -> Self {
257        Self {
258            inner: Arc::new(RwLock::new(HashMap::new())),
259        }
260    }
261
262    /// Push a 1-minute sample for `service_type` into tier 0 and cascade
263    /// rollups into tier 1 / tier 2 as needed.
264    pub async fn push_sample(&self, service_type: i32, sample: MetricSample) {
265        let mut map = self.inner.write().await;
266        let rings = map.entry(service_type).or_insert_with(ServiceRings::new);
267
268        // Push into tier 0.
269        rings.tier0.push(sample.clone());
270
271        // Accumulate for tier 1 rollup.
272        rings.tier1_acc.push(sample);
273
274        // Every 15 tier-0 samples -> roll up into one tier-1 sample.
275        if rings.tier1_acc.len() == 15 {
276            let rolled = aggregate(&rings.tier1_acc);
277            rings.tier1_acc.clear();
278
279            rings.tier1.push(rolled.clone());
280
281            // Accumulate for tier 2 rollup.
282            rings.tier2_acc.push(rolled);
283
284            // Every 16 tier-1 samples -> roll up into one tier-2 sample.
285            if rings.tier2_acc.len() == 16 {
286                let rolled2 = aggregate(&rings.tier2_acc);
287                rings.tier2_acc.clear();
288                rings.tier2.push(rolled2);
289            }
290        }
291    }
292
293    /// Query samples at the given tier for a service.
294    ///
295    /// Returns an empty vec if no data exists for `service_type` or `tier`
296    /// is out of range.
297    pub async fn query(&self, service_type: i32, tier: u8) -> Vec<MetricSample> {
298        let map = self.inner.read().await;
299        let Some(rings) = map.get(&service_type) else {
300            return Vec::new();
301        };
302        match tier {
303            0 => rings.tier0.to_vec(),
304            1 => rings.tier1.to_vec(),
305            2 => rings.tier2.to_vec(),
306            _ => Vec::new(),
307        }
308    }
309
310    /// Spawn a background sampler task that periodically snapshots each
311    /// service's counters and pushes the samples into the store.
312    ///
313    /// The task runs every `interval` (typically 60 s) until the runtime
314    /// shuts down.
315    pub fn start_sampler(
316        self,
317        counters: Arc<HashMap<i32, Arc<ServiceCounters>>>,
318        interval: Duration,
319    ) {
320        tokio::spawn(async move {
321            let mut tick = tokio::time::interval(interval);
322            loop {
323                tick.tick().await;
324                for (&svc_type, ctr) in counters.iter() {
325                    let sample = ctr.snapshot_and_reset().await;
326                    self.push_sample(svc_type, sample).await;
327                }
328            }
329        });
330    }
331}
332
333impl std::fmt::Debug for MetricsStore {
334    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
335        f.debug_struct("MetricsStore").finish()
336    }
337}
338
339impl Default for MetricsStore {
340    fn default() -> Self {
341        Self::new()
342    }
343}
344
345#[cfg(test)]
346mod tests {
347    use super::*;
348
349    #[test]
350    fn compute_p95_basic() {
351        // 20 values: 1..=20 -> p95 index = ceil(20*0.95)-1 = 19-1 = 18 -> value 19
352        let mut vals: Vec<f64> = (1..=20).map(|v| v as f64).collect();
353        assert!((compute_p95(&mut vals) - 19.0).abs() < f64::EPSILON);
354    }
355
356    #[test]
357    fn compute_p95_empty() {
358        assert!((compute_p95(&mut []) - 0.0).abs() < f64::EPSILON);
359    }
360
361    #[test]
362    fn compute_p95_single() {
363        assert!((compute_p95(&mut [42.0]) - 42.0).abs() < f64::EPSILON);
364    }
365
366    #[test]
367    fn ring_evicts_oldest() {
368        let mut ring = Ring::new(3);
369        for i in 0..5 {
370            ring.push(MetricSample {
371                ts: i,
372                active_conns: 0,
373                requests: 0,
374                failed_requests: 0,
375                latency_p95_ms: 0.0,
376            });
377        }
378        let v = ring.to_vec();
379        assert_eq!(v.len(), 3);
380        assert_eq!(v[0].ts, 2);
381        assert_eq!(v[2].ts, 4);
382    }
383
384    #[test]
385    fn aggregate_applies_rules() {
386        let samples = vec![
387            MetricSample {
388                ts: 100,
389                active_conns: 10,
390                requests: 50,
391                failed_requests: 2,
392                latency_p95_ms: 3.5,
393            },
394            MetricSample {
395                ts: 200,
396                active_conns: 20,
397                requests: 60,
398                failed_requests: 3,
399                latency_p95_ms: 7.1,
400            },
401        ];
402        let agg = aggregate(&samples);
403        assert_eq!(agg.ts, 200); // latest
404        assert_eq!(agg.active_conns, 15); // mean(10,20) = 15
405        assert_eq!(agg.requests, 110); // sum
406        assert_eq!(agg.failed_requests, 5); // sum
407        assert!((agg.latency_p95_ms - 7.1).abs() < f64::EPSILON); // max
408    }
409
410    #[tokio::test]
411    async fn push_and_query() {
412        let store = MetricsStore::new();
413        let sample = MetricSample {
414            ts: 1000,
415            active_conns: 5,
416            requests: 100,
417            failed_requests: 1,
418            latency_p95_ms: 2.0,
419        };
420        store.push_sample(1, sample).await;
421
422        let tier0 = store.query(1, 0).await;
423        assert_eq!(tier0.len(), 1);
424        assert_eq!(tier0[0].ts, 1000);
425
426        // tier 1 should be empty (need 15 samples to trigger rollup)
427        assert!(store.query(1, 1).await.is_empty());
428    }
429
430    #[tokio::test]
431    async fn tier0_to_tier1_rollup() {
432        let store = MetricsStore::new();
433        for i in 0..15 {
434            store
435                .push_sample(
436                    1,
437                    MetricSample {
438                        ts: i * 60,
439                        active_conns: 10,
440                        requests: 100,
441                        failed_requests: 1,
442                        latency_p95_ms: 5.0,
443                    },
444                )
445                .await;
446        }
447        let tier1 = store.query(1, 1).await;
448        assert_eq!(tier1.len(), 1);
449        assert_eq!(tier1[0].active_conns, 10); // mean of uniform = same
450        assert_eq!(tier1[0].requests, 1500); // 15 * 100
451        assert_eq!(tier1[0].failed_requests, 15); // 15 * 1
452        assert!((tier1[0].latency_p95_ms - 5.0).abs() < f64::EPSILON);
453    }
454
455    #[tokio::test]
456    async fn service_counters_snapshot() {
457        let ctr = ServiceCounters::new();
458        ctr.inc_conns();
459        ctr.inc_conns();
460        ctr.record_request(true, 1.0).await;
461        ctr.record_request(false, 10.0).await;
462        ctr.record_request(true, 5.0).await;
463
464        let snap = ctr.snapshot_and_reset().await;
465        assert_eq!(snap.active_conns, 2);
466        assert_eq!(snap.requests, 3);
467        assert_eq!(snap.failed_requests, 1);
468        // p95 of [1.0, 5.0, 10.0] -> ceil(3*0.95)=3, idx=2 -> 10.0
469        assert!((snap.latency_p95_ms - 10.0).abs() < f64::EPSILON);
470
471        // After snapshot, deltas should be reset
472        let snap2 = ctr.snapshot_and_reset().await;
473        assert_eq!(snap2.requests, 0);
474        assert_eq!(snap2.failed_requests, 0);
475        assert!((snap2.latency_p95_ms - 0.0).abs() < f64::EPSILON);
476        // Gauge stays
477        assert_eq!(snap2.active_conns, 2);
478    }
479}