1use std::{
2 collections::{HashMap, VecDeque},
3 sync::{
4 atomic::{AtomicU64, Ordering},
5 Arc, Mutex, MutexGuard,
6 },
7 time::Duration,
8};
9
10use serde_json::{json, Value};
11use tracing::debug;
12
13use crate::registry::ConnectionId;
14
15const ROUTE_OPEN_REFUSAL_COUNTER_CODES: &[&str] = &[
23 "module_warming",
24 ROUTE_OPEN_REFUSED_DECLARED_NOT_READY,
25 ROUTE_OPEN_REFUSED_REQUIRED_CAPABILITY_UNPROVIDED,
26 "target_unavailable",
27 "module_removed",
28 "module_no_protocol",
29 "unknown_module",
30 "module_reloading",
31 "op_not_allowed",
32 "bad_consumer_identity",
33 "capability_forbidden",
34 "admission_facts_not_permitted",
35 "admission_facts_target_not_allowed",
36 "route_limit",
37 "forwarding_error",
38 "module_timeout",
39 ROUTE_OPEN_REFUSED_BREAKER_OPEN,
40 "module_rejected",
41];
42
43pub(crate) const ROUTE_OPEN_REFUSED_BREAKER_OPEN: &str = "module_timeout_breaker_open";
55
56pub(crate) const ROUTE_OPEN_REFUSED_DECLARED_NOT_READY: &str = "module_warming_declared_not_ready";
62
63pub(crate) const ROUTE_OPEN_REFUSED_REQUIRED_CAPABILITY_UNPROVIDED: &str =
70 "module_warming_required_capability_unprovided";
71
72#[derive(Debug, Clone, Default)]
74pub struct ConnectedClients {
75 count: Arc<AtomicU64>,
76}
77
78impl ConnectedClients {
79 pub fn new() -> Self {
80 Self::default()
81 }
82
83 pub fn count(&self) -> u64 {
84 self.count.load(Ordering::SeqCst)
85 }
86
87 pub(crate) fn open(&self, connection_id: ConnectionId) -> ConnectedClientGuard {
88 let previous = self.count.fetch_add(1, Ordering::SeqCst);
89 let current = previous + 1;
90 debug!(
98 connection_id = connection_id.get(),
99 connected_clients = current,
100 previous_connected_clients = previous,
101 "authenticated connection count changed"
102 );
103 ConnectedClientGuard {
104 clients: self.clone(),
105 connection_id,
106 }
107 }
108}
109
110pub(crate) struct ConnectedClientGuard {
111 clients: ConnectedClients,
112 connection_id: ConnectionId,
113}
114
115impl Drop for ConnectedClientGuard {
116 fn drop(&mut self) {
117 let previous = self.clients.count.fetch_sub(1, Ordering::SeqCst);
118 let current = previous.saturating_sub(1);
119 debug!(
121 connection_id = self.connection_id.get(),
122 connected_clients = current,
123 previous_connected_clients = previous,
124 "authenticated connection count changed"
125 );
126 }
127}
128
129#[derive(Debug, Clone, Default)]
131pub struct DaemonCounters {
132 module_frames_dropped_no_route: Arc<AtomicU64>,
133 module_frames_dropped_no_route_by_module: Arc<Mutex<HashMap<String, u64>>>,
136 route_open_refused_by_code: Arc<Mutex<HashMap<String, u64>>>,
137 route_open_accepted_by_principal: Arc<Mutex<HashMap<String, u64>>>,
138 module_frames_dropped_no_route_window: Arc<Mutex<DropWindow>>,
139 module_requests_dropped_stale_route: Arc<AtomicU64>,
140 client_frames_dropped_stale_route: Arc<AtomicU64>,
141 client_egress_close_delivery_failed: Arc<AtomicU64>,
142 goodbye_relay_client_failed: Arc<AtomicU64>,
143 goodbye_relay_module_dropped: Arc<AtomicU64>,
144 goodbye_relay_module_dropped_by_module: Arc<Mutex<HashMap<String, u64>>>,
145 route_released_epoch_fenced: Arc<AtomicU64>,
146 route_release_stale_skipped: Arc<AtomicU64>,
147 drains_with_undeclared_gauge: Arc<AtomicU64>,
148}
149
150#[derive(Debug)]
153struct DropWindow {
154 started_at: tokio::time::Instant,
155 buckets: VecDeque<DropBucket>,
156}
157
158#[derive(Debug)]
159struct DropBucket {
160 minute: u64,
161 count: u64,
162}
163
164impl Default for DropWindow {
165 fn default() -> Self {
166 Self {
167 started_at: tokio::time::Instant::now(),
168 buckets: VecDeque::new(),
169 }
170 }
171}
172
173impl DropWindow {
174 const MINUTE: Duration = Duration::from_secs(60);
175 const BUCKETS: u64 = 10;
176
177 fn record(&mut self, now: tokio::time::Instant) {
178 let minute = self.minute_at(now);
179 self.prune_before(minute);
180 match self.buckets.back_mut() {
181 Some(bucket) if bucket.minute == minute => bucket.count += 1,
182 _ => self.buckets.push_back(DropBucket { minute, count: 1 }),
183 }
184 }
185
186 fn count_last_10m(&mut self, now: tokio::time::Instant) -> u64 {
187 let minute = self.minute_at(now);
188 self.prune_before(minute);
189 self.buckets.iter().map(|bucket| bucket.count).sum()
190 }
191
192 fn nonzero_minutes_last_10m(&mut self, now: tokio::time::Instant) -> u64 {
193 let minute = self.minute_at(now);
194 self.prune_before(minute);
195 self.buckets.len() as u64
196 }
197
198 fn minute_at(&self, now: tokio::time::Instant) -> u64 {
199 now.saturating_duration_since(self.started_at).as_secs() / Self::MINUTE.as_secs()
200 }
201
202 fn prune_before(&mut self, current_minute: u64) {
203 while self
204 .buckets
205 .front()
206 .is_some_and(|bucket| current_minute.saturating_sub(bucket.minute) >= Self::BUCKETS)
207 {
208 self.buckets.pop_front();
209 }
210 }
211}
212
213impl DaemonCounters {
214 pub fn new() -> Self {
215 Self::default()
216 }
217
218 pub fn snapshot(&self) -> Value {
221 let mut snapshot = serde_json::Map::new();
222 snapshot.insert(
223 "module_frames_dropped_no_route".into(),
224 self.module_frames_dropped_no_route
225 .load(Ordering::Relaxed)
226 .into(),
227 );
228 let mut drop_window = self
229 .module_frames_dropped_no_route_window
230 .lock()
231 .expect("drop-rate window mutex poisoned");
232 let now = tokio::time::Instant::now();
233 snapshot.insert(
234 "module_frames_dropped_no_route_last_10m".into(),
235 drop_window.count_last_10m(now).into(),
236 );
237 snapshot.insert(
238 "module_frames_dropped_no_route_nonzero_minutes_last_10m".into(),
239 drop_window.nonzero_minutes_last_10m(now).into(),
240 );
241 insert_nonempty_counts(
242 &mut snapshot,
243 "module_frames_dropped_no_route_by_module",
244 &self.module_frames_dropped_no_route_by_module,
245 );
246 insert_nonempty_counts(
247 &mut snapshot,
248 "route_open_refused_by_code",
249 &self.route_open_refused_by_code,
250 );
251 insert_nonempty_counts(
252 &mut snapshot,
253 "route_open_accepted_by_principal",
254 &self.route_open_accepted_by_principal,
255 );
256 snapshot.insert(
257 "module_requests_dropped_stale_route".into(),
258 self.module_requests_dropped_stale_route
259 .load(Ordering::Relaxed)
260 .into(),
261 );
262 snapshot.insert(
263 "client_frames_dropped_stale_route".into(),
264 self.client_frames_dropped_stale_route
265 .load(Ordering::Relaxed)
266 .into(),
267 );
268 snapshot.insert(
269 "client_egress_close_delivery_failed".into(),
270 self.client_egress_close_delivery_failed
271 .load(Ordering::Relaxed)
272 .into(),
273 );
274 snapshot.insert(
275 "goodbye_relay_client_failed".into(),
276 self.goodbye_relay_client_failed
277 .load(Ordering::Relaxed)
278 .into(),
279 );
280 snapshot.insert(
281 "goodbye_relay_module_dropped".into(),
282 self.goodbye_relay_module_dropped
283 .load(Ordering::Relaxed)
284 .into(),
285 );
286 insert_nonempty_counts(
287 &mut snapshot,
288 "goodbye_relay_module_dropped_by_module",
289 &self.goodbye_relay_module_dropped_by_module,
290 );
291 snapshot.insert(
292 "route_released_epoch_fenced".into(),
293 self.route_released_epoch_fenced
294 .load(Ordering::Relaxed)
295 .into(),
296 );
297 snapshot.insert(
298 "route_release_stale_skipped".into(),
299 self.route_release_stale_skipped
300 .load(Ordering::Relaxed)
301 .into(),
302 );
303 snapshot.insert(
304 "drains_with_undeclared_gauge".into(),
305 self.drains_with_undeclared_gauge
306 .load(Ordering::Relaxed)
307 .into(),
308 );
309 Value::Object(snapshot)
310 }
311
312 pub(crate) fn increment_module_frames_dropped_no_route(&self, module_id: Option<&str>) {
313 self.module_frames_dropped_no_route
314 .fetch_add(1, Ordering::Relaxed);
315 if let Some(module_id) = module_id {
316 increment_keyed_count(&self.module_frames_dropped_no_route_by_module, module_id);
317 }
318 self.module_frames_dropped_no_route_window
319 .lock()
320 .expect("drop-rate window mutex poisoned")
321 .record(tokio::time::Instant::now());
322 }
323
324 pub(crate) fn increment_route_open_refused(&self, code: &'static str) {
325 debug_assert!(ROUTE_OPEN_REFUSAL_COUNTER_CODES.contains(&code));
326 increment_keyed_count(&self.route_open_refused_by_code, code);
327 }
328
329 pub(crate) fn increment_route_open_accepted(&self, principal: &str) {
339 increment_keyed_count(&self.route_open_accepted_by_principal, principal);
340 }
341
342 pub(crate) fn increment_module_requests_dropped_stale_route(&self) {
343 self.module_requests_dropped_stale_route
344 .fetch_add(1, Ordering::Relaxed);
345 }
346
347 pub(crate) fn increment_client_frames_dropped_stale_route(&self) {
348 self.client_frames_dropped_stale_route
349 .fetch_add(1, Ordering::Relaxed);
350 }
351
352 pub(crate) fn increment_client_egress_close_delivery_failed(&self) {
353 self.client_egress_close_delivery_failed
354 .fetch_add(1, Ordering::Relaxed);
355 }
356
357 pub(crate) fn increment_goodbye_relay_client_failed(&self) {
358 self.goodbye_relay_client_failed
359 .fetch_add(1, Ordering::Relaxed);
360 }
361
362 pub(crate) fn increment_goodbye_relay_module_dropped(&self, module_id: Option<&str>) {
363 self.goodbye_relay_module_dropped
364 .fetch_add(1, Ordering::Relaxed);
365 if let Some(module_id) = module_id {
366 increment_keyed_count(&self.goodbye_relay_module_dropped_by_module, module_id);
367 }
368 }
369
370 pub(crate) fn increment_route_released_epoch_fenced(&self) {
371 self.route_released_epoch_fenced
372 .fetch_add(1, Ordering::Relaxed);
373 }
374
375 pub(crate) fn increment_route_release_stale_skipped(&self) {
376 self.route_release_stale_skipped
377 .fetch_add(1, Ordering::Relaxed);
378 }
379
380 pub(crate) fn increment_drains_with_undeclared_gauge(&self) {
381 self.drains_with_undeclared_gauge
382 .fetch_add(1, Ordering::Relaxed);
383 }
384}
385
386fn increment_keyed_count(counts: &Mutex<HashMap<String, u64>>, key: &str) {
387 *counts
388 .lock()
389 .expect("keyed counter mutex poisoned")
390 .entry(key.to_string())
391 .or_default() += 1;
392}
393
394fn insert_nonempty_counts(
395 snapshot: &mut serde_json::Map<String, Value>,
396 key: &str,
397 counts: &Mutex<HashMap<String, u64>>,
398) {
399 let counts: MutexGuard<'_, HashMap<String, u64>> =
400 counts.lock().expect("keyed counter mutex poisoned");
401 if !counts.is_empty() {
402 snapshot.insert(key.to_string(), json!(&*counts));
403 }
404}
405
406#[cfg(test)]
407mod tests {
408 use super::*;
409
410 #[test]
411 fn counter_snapshot_includes_zero_rate_and_omits_empty_module_maps() {
412 let counters = DaemonCounters::new();
413 let snapshot = counters.snapshot();
414
415 assert_eq!(snapshot["module_frames_dropped_no_route_last_10m"], 0);
416 assert_eq!(
417 snapshot["module_frames_dropped_no_route_nonzero_minutes_last_10m"],
418 0
419 );
420 assert!(snapshot
421 .get("module_frames_dropped_no_route_by_module")
422 .is_none());
423 assert!(snapshot
424 .get("goodbye_relay_module_dropped_by_module")
425 .is_none());
426 }
427
428 #[tokio::test(start_paused = true)]
429 async fn module_frame_drops_are_attributed_to_the_emitting_module() {
430 let counters = DaemonCounters::new();
431 counters.increment_module_frames_dropped_no_route(Some("alpha"));
432 counters.increment_module_frames_dropped_no_route(Some("alpha"));
433
434 let snapshot = counters.snapshot();
435 assert_eq!(snapshot["module_frames_dropped_no_route"], 2);
436 assert_eq!(
437 snapshot["module_frames_dropped_no_route_by_module"],
438 json!({ "alpha": 2 })
439 );
440 assert_eq!(snapshot["module_frames_dropped_no_route_last_10m"], 2);
441 }
442
443 #[tokio::test(start_paused = true)]
444 async fn frame_drop_rate_ages_out_after_ten_minute_buckets() {
445 let counters = DaemonCounters::new();
446 counters.increment_module_frames_dropped_no_route(Some("alpha"));
447
448 tokio::time::advance(Duration::from_secs(9 * 60)).await;
449 assert_eq!(
450 counters.snapshot()["module_frames_dropped_no_route_last_10m"],
451 1
452 );
453
454 tokio::time::advance(Duration::from_secs(60)).await;
455 assert_eq!(
456 counters.snapshot()["module_frames_dropped_no_route_last_10m"],
457 0
458 );
459 }
460
461 #[tokio::test(start_paused = true)]
462 async fn frame_drop_window_counts_only_nonzero_minutes() {
463 let counters = DaemonCounters::new();
464 for minute in 0..10 {
465 if minute > 0 {
466 tokio::time::advance(Duration::from_secs(60)).await;
467 }
468 if minute != 4 {
469 counters.increment_module_frames_dropped_no_route(Some("alpha"));
470 }
471 }
472
473 assert_eq!(
474 counters.snapshot()["module_frames_dropped_no_route_nonzero_minutes_last_10m"],
475 9
476 );
477 }
478
479 #[test]
480 fn goodbye_relay_drops_are_attributed_to_the_target_module() {
481 let counters = DaemonCounters::new();
482 counters.increment_goodbye_relay_module_dropped(Some("alpha"));
483
484 assert_eq!(
485 counters.snapshot()["goodbye_relay_module_dropped_by_module"],
486 json!({ "alpha": 1 })
487 );
488 }
489}