1use 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#[derive(Clone, Debug, serde::Serialize)]
24pub struct MetricSample {
25 pub ts: i64,
27 pub active_conns: u32,
29 pub requests: u64,
31 pub failed_requests: u64,
33 pub latency_p95_ms: f64,
35}
36
37pub struct ServiceCounters {
46 pub active_conns: AtomicU32,
48 pub total_requests: AtomicU64,
50 pub failed_requests: AtomicU64,
52 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 pub fn inc_conns(&self) {
68 self.active_conns.fetch_add(1, Ordering::Relaxed);
69 }
70
71 pub fn dec_conns(&self) {
73 self.active_conns.fetch_sub(1, Ordering::Relaxed);
74 }
75
76 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 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
138fn 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
150struct 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 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
181struct ServiceRings {
188 tier0: Ring, tier1: Ring, tier2: Ring, tier1_acc: Vec<MetricSample>,
193 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
209fn 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#[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 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 rings.tier0.push(sample.clone());
270
271 rings.tier1_acc.push(sample);
273
274 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 rings.tier2_acc.push(rolled);
283
284 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 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 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 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); assert_eq!(agg.active_conns, 15); assert_eq!(agg.requests, 110); assert_eq!(agg.failed_requests, 5); assert!((agg.latency_p95_ms - 7.1).abs() < f64::EPSILON); }
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 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); assert_eq!(tier1[0].requests, 1500); assert_eq!(tier1[0].failed_requests, 15); 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 assert!((snap.latency_p95_ms - 10.0).abs() < f64::EPSILON);
470
471 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 assert_eq!(snap2.active_conns, 2);
478 }
479}