1use std::collections::{HashMap, VecDeque};
10use std::time::{Duration, Instant};
11
12const LAT_WINDOW: usize = 256;
15
16#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize)]
26pub struct LatencySummary {
27 pub min_us: i64,
28 pub median_us: i64,
29 pub p95_us: i64,
30 pub max_us: i64,
31 pub samples: usize,
33}
34
35#[derive(Debug, Clone)]
37pub struct KeyStats {
38 pub count: u64,
39 pub bytes: u64,
40 pub rate_hz: f64,
42 pub last_seen: Instant,
43 pub sn_gaps: u64,
46 pub unstamped: u64,
50 last_sn: Option<u32>,
51 lat: VecDeque<i64>,
53}
54
55impl KeyStats {
56 pub fn latency(&self) -> Option<LatencySummary> {
59 if self.lat.is_empty() {
60 return None;
61 }
62 let mut sorted: Vec<i64> = self.lat.iter().copied().collect();
63 sorted.sort_unstable();
64 let at = |q: f64| sorted[((sorted.len() - 1) as f64 * q) as usize];
65 Some(LatencySummary {
66 min_us: sorted[0],
67 median_us: at(0.5),
68 p95_us: at(0.95),
69 max_us: *sorted.last().expect("non-empty"),
70 samples: sorted.len(),
71 })
72 }
73}
74
75#[derive(Debug)]
85pub struct StatsTable {
86 keys: HashMap<String, KeyStats>,
87 max_keys: usize,
88 evicted: u64,
89 unwatched: u64,
90}
91
92impl Default for StatsTable {
93 fn default() -> Self {
94 StatsTable::with_capacity(DEFAULT_MAX_KEYS)
95 }
96}
97
98const TAU: Duration = Duration::from_secs(2);
101
102pub const DEFAULT_MAX_KEYS: usize = 50_000;
106
107const EVICT_FRACTION: usize = 16;
113
114impl StatsTable {
115 pub fn new() -> Self {
116 Self::default()
117 }
118
119 pub fn with_capacity(max_keys: usize) -> Self {
121 StatsTable {
122 keys: HashMap::new(),
123 max_keys: max_keys.max(1),
124 evicted: 0,
125 unwatched: 0,
126 }
127 }
128
129 pub fn evicted(&self) -> u64 {
135 self.evicted
136 }
137
138 pub fn max_keys(&self) -> usize {
140 self.max_keys
141 }
142
143 pub fn unwatched(&self) -> u64 {
152 self.unwatched
153 }
154
155 pub fn retire_unwatched(&mut self, gone: &str, kept: &[String]) -> usize {
160 use zenoh::key_expr::keyexpr;
161 let Ok(gone) = keyexpr::new(gone) else {
167 return 0;
168 };
169 let kept: Vec<&keyexpr> = kept
170 .iter()
171 .filter_map(|k| keyexpr::new(k.as_str()).ok())
172 .collect();
173 let doomed: Vec<String> = self
174 .keys
175 .keys()
176 .filter(|key| match keyexpr::new(key.as_str()) {
177 Ok(ke) => gone.intersects(ke) && !kept.iter().any(|k| k.intersects(ke)),
178 Err(_) => false,
179 })
180 .cloned()
181 .collect();
182 for key in &doomed {
183 self.keys.remove(key);
184 }
185 self.unwatched += doomed.len() as u64;
186 doomed.len()
187 }
188
189 fn evict(&mut self) {
191 let target = self.max_keys - (self.max_keys / EVICT_FRACTION).max(1);
192 let mut seen: Vec<(Instant, String)> = self
193 .keys
194 .iter()
195 .map(|(k, s)| (s.last_seen, k.clone()))
196 .collect();
197 seen.sort_unstable_by_key(|(last_seen, _)| *last_seen);
199 for (_, key) in seen.into_iter().take(self.keys.len() - target) {
200 self.keys.remove(&key);
201 self.evicted += 1;
202 }
203 }
204
205 pub fn record(
209 &mut self,
210 key: &str,
211 payload_len: usize,
212 sn: Option<u32>,
213 now: Instant,
214 latency_us: Option<i64>,
215 ) {
216 if let Some(s) = self.keys.get_mut(key) {
217 let dt = now.saturating_duration_since(s.last_seen).as_secs_f64();
218 if dt > 0.0 {
219 let alpha = 1.0 - (-dt / TAU.as_secs_f64()).exp();
220 let instant_rate = 1.0 / dt;
221 s.rate_hz += alpha * (instant_rate - s.rate_hz);
222 }
223 s.count += 1;
224 s.bytes += payload_len as u64;
225 s.last_seen = now;
226 if let (Some(prev), Some(cur)) = (s.last_sn, sn)
227 && cur > prev + 1
228 {
229 s.sn_gaps += u64::from(cur - prev - 1);
230 }
231 s.last_sn = sn;
232 match latency_us {
233 Some(us) => {
234 if s.lat.len() >= LAT_WINDOW {
235 s.lat.pop_front();
236 }
237 s.lat.push_back(us);
238 }
239 None => s.unstamped += 1,
240 }
241 } else {
242 if self.keys.len() >= self.max_keys {
243 self.evict();
244 }
245 self.keys.insert(
246 key.to_string(),
247 KeyStats {
248 count: 1,
249 bytes: payload_len as u64,
250 rate_hz: 0.0,
251 last_seen: now,
252 sn_gaps: 0,
253 unstamped: u64::from(latency_us.is_none()),
254 last_sn: sn,
255 lat: latency_us.into_iter().collect(),
256 },
257 );
258 }
259 }
260
261 pub fn get(&self, key: &str) -> Option<&KeyStats> {
262 self.keys.get(key)
263 }
264
265 pub fn iter(&self) -> impl Iterator<Item = (&str, &KeyStats)> {
266 self.keys.iter().map(|(k, v)| (k.as_str(), v))
267 }
268
269 pub fn len(&self) -> usize {
270 self.keys.len()
271 }
272
273 pub fn is_empty(&self) -> bool {
274 self.keys.is_empty()
275 }
276
277 pub fn totals(&self) -> (u64, u64, f64) {
279 self.keys.values().fold((0, 0, 0.0), |(c, b, r), s| {
280 (c + s.count, b + s.bytes, r + s.rate_hz)
281 })
282 }
283}
284
285#[cfg(test)]
286mod tests {
287 use super::*;
288
289 #[test]
292 fn the_table_is_bounded() {
293 let mut t = StatsTable::with_capacity(100);
294 let now = Instant::now();
295 for i in 0..1000 {
296 t.record(&format!("demo/k{i}"), 4, None, now, None);
297 }
298 assert!(t.len() <= 100, "len {} exceeds the bound", t.len());
299 assert!(t.evicted() > 0);
300 assert_eq!(t.len() as u64 + t.evicted(), 1000);
302 }
303
304 #[test]
307 fn eviction_drops_the_least_recently_seen() {
308 let mut t = StatsTable::with_capacity(10);
309 let t0 = Instant::now();
310
311 for i in 0..10 {
313 t.record(
314 &format!("old/k{i}"),
315 4,
316 None,
317 t0 + Duration::from_millis(i),
318 None,
319 );
320 }
321 let fresh = t0 + Duration::from_secs(60);
323 t.record("old/k0", 4, None, fresh, None);
324
325 for i in 1..=5 {
328 t.record(
329 &format!("new/k{i}"),
330 4,
331 None,
332 fresh + Duration::from_millis(i),
333 None,
334 );
335 }
336
337 assert!(
338 t.get("old/k0").is_some(),
339 "a key that is still publishing must survive"
340 );
341 assert!(
342 t.get("old/k1").is_none(),
343 "a key that went quiet should have been evicted first"
344 );
345 }
346
347 #[test]
351 fn latency_is_summarised_and_unstamped_is_counted_not_defaulted() {
352 let mut t = StatsTable::new();
353 let now = Instant::now();
354 for us in [1000, -200, 5000, 3000] {
355 t.record("k", 4, None, now, Some(us));
356 }
357 t.record("k", 4, None, now, None);
358 let s = t.get("k").unwrap();
359 assert_eq!(s.unstamped, 1);
360 let lat = s.latency().unwrap();
361 assert_eq!(lat.min_us, -200, "negative skew is shown, not clamped");
362 assert_eq!(lat.max_us, 5000);
363 assert_eq!(lat.samples, 4);
364 assert!(lat.median_us >= -200 && lat.median_us <= 5000);
365
366 t.record("quiet", 4, None, now, None);
368 assert!(t.get("quiet").unwrap().latency().is_none());
369 assert_eq!(t.get("quiet").unwrap().unstamped, 1);
370 }
371
372 #[test]
375 fn repeated_keys_never_trigger_eviction() {
376 let mut t = StatsTable::with_capacity(4);
377 let t0 = Instant::now();
378 for i in 0..1000 {
379 t.record("demo/one", 4, None, t0 + Duration::from_millis(i), None);
380 }
381 assert_eq!(t.len(), 1);
382 assert_eq!(t.evicted(), 0);
383 assert_eq!(t.get("demo/one").unwrap().count, 1000);
384 }
385
386 #[test]
388 fn a_capacity_of_one_still_works() {
389 let mut t = StatsTable::with_capacity(1);
390 let now = Instant::now();
391 t.record("a", 1, None, now, None);
392 t.record("b", 1, None, now, None);
393 assert_eq!(t.len(), 1);
394 assert_eq!(t.evicted(), 1);
395 assert_eq!(StatsTable::with_capacity(0).max_keys(), 1);
397 }
398
399 #[test]
400 fn rates_converge_and_gaps_count() {
401 let mut t = StatsTable::new();
402 let t0 = Instant::now();
403 for i in 0..100u32 {
405 t.record(
406 "v1/h-a/telemetry/x/m",
407 8,
408 Some(i),
409 t0 + Duration::from_millis(100 * u64::from(i)),
410 None,
411 );
412 }
413 let s = t.get("v1/h-a/telemetry/x/m").unwrap();
414 assert_eq!(s.count, 100);
415 assert_eq!(s.bytes, 800);
416 assert!((s.rate_hz - 10.0).abs() < 1.0, "rate {}", s.rate_hz);
417 assert_eq!(s.sn_gaps, 0);
418
419 t.record(
421 "v1/h-a/telemetry/x/m",
422 8,
423 Some(105),
424 t0 + Duration::from_millis(10_100),
425 None,
426 );
427 assert_eq!(t.get("v1/h-a/telemetry/x/m").unwrap().sn_gaps, 5);
428 }
429
430 #[test]
431 fn totals_aggregate() {
432 let mut t = StatsTable::new();
433 let now = Instant::now();
434 t.record("a", 10, None, now, None);
435 t.record("b", 20, None, now, None);
436 let (count, bytes, _) = t.totals();
437 assert_eq!((count, bytes), (2, 30));
438 assert_eq!(t.len(), 2);
439 }
440
441 #[test]
444 fn retire_unwatched_respects_remaining_coverage() {
445 let mut t = StatsTable::new();
446 let now = Instant::now();
447 t.record("v1/h-a/telemetry/x/m1", 4, None, now, None);
448 t.record("v1/h-a/state/x/health", 4, None, now, None);
449 t.record("v1/h-b/telemetry/y/m2", 4, None, now, None);
450
451 let retired = t.retire_unwatched("v1/*/telemetry/**", &["v1/h-a/**".to_string()]);
453 assert_eq!(retired, 1, "only h-b's telemetry loses coverage");
454 assert!(
455 t.get("v1/h-a/telemetry/x/m1").is_some(),
456 "still covered by kept"
457 );
458 assert!(t.get("v1/h-b/telemetry/y/m2").is_none());
459 assert_eq!(t.unwatched(), 1);
460
461 let retired = t.retire_unwatched("**", &[]);
463 assert_eq!(retired, 2);
464 assert_eq!(t.len(), 0);
465 assert_eq!(t.unwatched(), 3);
466 }
467
468 #[test]
471 fn retire_unwatched_tolerates_bad_selectors() {
472 let mut t = StatsTable::new();
473 t.record("a/b", 1, None, Instant::now(), None);
474 assert_eq!(t.retire_unwatched("", &[]), 0);
475 assert_eq!(t.len(), 1);
476 }
477}