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 subc_protocol::error_codes::TARGET_FLOW_UNSUPPORTED,
28 "module_removed",
29 "module_no_protocol",
30 "unknown_module",
31 "module_reloading",
32 "op_not_allowed",
33 "bad_consumer_identity",
34 "capability_forbidden",
35 "admission_facts_not_permitted",
36 "admission_facts_target_not_allowed",
37 "route_limit",
38 "forwarding_error",
39 "module_timeout",
40 ROUTE_OPEN_REFUSED_BREAKER_OPEN,
41 "module_rejected",
42 subc_protocol::error_codes::SCOPE_EPOCH_REQUIRED,
43 subc_protocol::error_codes::SCOPE_NOT_SYNCED,
44 subc_protocol::error_codes::SCOPE_NOT_LIVE,
45 subc_protocol::error_codes::SCOPE_ENDED,
46 subc_protocol::error_codes::SCOPE_NOT_CARRIER,
47 subc_protocol::error_codes::SCOPE_CHANGED,
48 subc_protocol::error_codes::INVALID_REQUEST,
49];
50
51pub(crate) const ROUTE_OPEN_REFUSED_BREAKER_OPEN: &str = "module_timeout_breaker_open";
63
64pub(crate) const ROUTE_OPEN_REFUSED_DECLARED_NOT_READY: &str = "module_warming_declared_not_ready";
70
71pub(crate) const ROUTE_OPEN_REFUSED_REQUIRED_CAPABILITY_UNPROVIDED: &str =
78 "module_warming_required_capability_unprovided";
79
80#[derive(Debug, Clone, Default)]
82pub struct ConnectedClients {
83 count: Arc<AtomicU64>,
84}
85
86impl ConnectedClients {
87 pub fn new() -> Self {
88 Self::default()
89 }
90
91 pub fn count(&self) -> u64 {
92 self.count.load(Ordering::SeqCst)
93 }
94
95 pub(crate) fn open(&self, connection_id: ConnectionId) -> ConnectedClientGuard {
96 let previous = self.count.fetch_add(1, Ordering::SeqCst);
97 let current = previous + 1;
98 debug!(
106 connection_id = connection_id.get(),
107 connected_clients = current,
108 previous_connected_clients = previous,
109 "authenticated connection count changed"
110 );
111 ConnectedClientGuard {
112 clients: self.clone(),
113 connection_id,
114 }
115 }
116}
117
118pub(crate) struct ConnectedClientGuard {
119 clients: ConnectedClients,
120 connection_id: ConnectionId,
121}
122
123impl Drop for ConnectedClientGuard {
124 fn drop(&mut self) {
125 let previous = self.clients.count.fetch_sub(1, Ordering::SeqCst);
126 let current = previous.saturating_sub(1);
127 debug!(
129 connection_id = self.connection_id.get(),
130 connected_clients = current,
131 previous_connected_clients = previous,
132 "authenticated connection count changed"
133 );
134 }
135}
136
137#[derive(Debug, Clone, Default)]
139pub struct DaemonCounters {
140 module_frames_dropped_no_route: Arc<AtomicU64>,
145 module_frames_dropped_released_route: Arc<AtomicU64>,
151 module_frames_dropped_released_route_by_module: Arc<Mutex<HashMap<String, u64>>>,
152 module_orphan_route_goodbyes_sent: Arc<AtomicU64>,
157 module_frames_dropped_no_route_by_module: Arc<Mutex<HashMap<String, u64>>>,
160 route_open_refused_by_code: Arc<Mutex<HashMap<String, u64>>>,
161 route_open_accepted_by_principal: Arc<Mutex<HashMap<String, u64>>>,
162 module_frames_dropped_no_route_window: Arc<Mutex<DropWindow>>,
163 module_requests_dropped_stale_route: Arc<AtomicU64>,
164 client_frames_dropped_stale_route: Arc<AtomicU64>,
165 client_egress_close_delivery_failed: Arc<AtomicU64>,
166 goodbye_relay_client_failed: Arc<AtomicU64>,
167 goodbye_relay_module_dropped: Arc<AtomicU64>,
168 goodbye_relay_module_dropped_by_module: Arc<Mutex<HashMap<String, u64>>>,
169 route_released_epoch_fenced: Arc<AtomicU64>,
170 route_release_stale_skipped: Arc<AtomicU64>,
171 drains_with_undeclared_gauge: Arc<AtomicU64>,
172}
173
174#[derive(Debug)]
177struct DropWindow {
178 started_at: tokio::time::Instant,
179 buckets: VecDeque<DropBucket>,
180}
181
182#[derive(Debug)]
183struct DropBucket {
184 minute: u64,
185 count: u64,
186}
187
188impl Default for DropWindow {
189 fn default() -> Self {
190 Self {
191 started_at: tokio::time::Instant::now(),
192 buckets: VecDeque::new(),
193 }
194 }
195}
196
197impl DropWindow {
198 const MINUTE: Duration = Duration::from_secs(60);
199 const BUCKETS: u64 = 10;
200
201 fn record(&mut self, now: tokio::time::Instant) {
202 let minute = self.minute_at(now);
203 self.prune_before(minute);
204 match self.buckets.back_mut() {
205 Some(bucket) if bucket.minute == minute => bucket.count += 1,
206 _ => self.buckets.push_back(DropBucket { minute, count: 1 }),
207 }
208 }
209
210 fn count_last_10m(&mut self, now: tokio::time::Instant) -> u64 {
211 let minute = self.minute_at(now);
212 self.prune_before(minute);
213 self.buckets.iter().map(|bucket| bucket.count).sum()
214 }
215
216 fn nonzero_minutes_last_10m(&mut self, now: tokio::time::Instant) -> u64 {
217 let minute = self.minute_at(now);
218 self.prune_before(minute);
219 self.buckets.len() as u64
220 }
221
222 fn minute_at(&self, now: tokio::time::Instant) -> u64 {
223 now.saturating_duration_since(self.started_at).as_secs() / Self::MINUTE.as_secs()
224 }
225
226 fn prune_before(&mut self, current_minute: u64) {
227 while self
228 .buckets
229 .front()
230 .is_some_and(|bucket| current_minute.saturating_sub(bucket.minute) >= Self::BUCKETS)
231 {
232 self.buckets.pop_front();
233 }
234 }
235}
236
237impl DaemonCounters {
238 pub fn new() -> Self {
239 Self::default()
240 }
241
242 pub fn snapshot(&self) -> Value {
245 let mut snapshot = serde_json::Map::new();
246 snapshot.insert(
247 "module_frames_dropped_no_route".into(),
248 self.module_frames_dropped_no_route
249 .load(Ordering::Relaxed)
250 .into(),
251 );
252 let mut drop_window = self
253 .module_frames_dropped_no_route_window
254 .lock()
255 .expect("drop-rate window mutex poisoned");
256 let now = tokio::time::Instant::now();
257 snapshot.insert(
258 "module_frames_dropped_no_route_last_10m".into(),
259 drop_window.count_last_10m(now).into(),
260 );
261 snapshot.insert(
262 "module_frames_dropped_no_route_nonzero_minutes_last_10m".into(),
263 drop_window.nonzero_minutes_last_10m(now).into(),
264 );
265 insert_nonempty_counts(
266 &mut snapshot,
267 "module_frames_dropped_no_route_by_module",
268 &self.module_frames_dropped_no_route_by_module,
269 );
270 snapshot.insert(
271 "module_frames_dropped_released_route".into(),
272 self.module_frames_dropped_released_route
273 .load(Ordering::Relaxed)
274 .into(),
275 );
276 insert_nonempty_counts(
277 &mut snapshot,
278 "module_frames_dropped_released_route_by_module",
279 &self.module_frames_dropped_released_route_by_module,
280 );
281 snapshot.insert(
282 "module_orphan_route_goodbyes_sent".into(),
283 self.module_orphan_route_goodbyes_sent
284 .load(Ordering::Relaxed)
285 .into(),
286 );
287 insert_nonempty_counts(
288 &mut snapshot,
289 "route_open_refused_by_code",
290 &self.route_open_refused_by_code,
291 );
292 insert_nonempty_counts(
293 &mut snapshot,
294 "route_open_accepted_by_principal",
295 &self.route_open_accepted_by_principal,
296 );
297 snapshot.insert(
298 "module_requests_dropped_stale_route".into(),
299 self.module_requests_dropped_stale_route
300 .load(Ordering::Relaxed)
301 .into(),
302 );
303 snapshot.insert(
304 "client_frames_dropped_stale_route".into(),
305 self.client_frames_dropped_stale_route
306 .load(Ordering::Relaxed)
307 .into(),
308 );
309 snapshot.insert(
310 "client_egress_close_delivery_failed".into(),
311 self.client_egress_close_delivery_failed
312 .load(Ordering::Relaxed)
313 .into(),
314 );
315 snapshot.insert(
316 "goodbye_relay_client_failed".into(),
317 self.goodbye_relay_client_failed
318 .load(Ordering::Relaxed)
319 .into(),
320 );
321 snapshot.insert(
322 "goodbye_relay_module_dropped".into(),
323 self.goodbye_relay_module_dropped
324 .load(Ordering::Relaxed)
325 .into(),
326 );
327 insert_nonempty_counts(
328 &mut snapshot,
329 "goodbye_relay_module_dropped_by_module",
330 &self.goodbye_relay_module_dropped_by_module,
331 );
332 snapshot.insert(
333 "route_released_epoch_fenced".into(),
334 self.route_released_epoch_fenced
335 .load(Ordering::Relaxed)
336 .into(),
337 );
338 snapshot.insert(
339 "route_release_stale_skipped".into(),
340 self.route_release_stale_skipped
341 .load(Ordering::Relaxed)
342 .into(),
343 );
344 snapshot.insert(
345 "drains_with_undeclared_gauge".into(),
346 self.drains_with_undeclared_gauge
347 .load(Ordering::Relaxed)
348 .into(),
349 );
350 Value::Object(snapshot)
351 }
352
353 pub(crate) fn increment_module_frames_dropped_no_route(&self, module_id: Option<&str>) {
354 self.module_frames_dropped_no_route
355 .fetch_add(1, Ordering::Relaxed);
356 if let Some(module_id) = module_id {
357 increment_keyed_count(&self.module_frames_dropped_no_route_by_module, module_id);
358 }
359 self.module_frames_dropped_no_route_window
360 .lock()
361 .expect("drop-rate window mutex poisoned")
362 .record(tokio::time::Instant::now());
363 }
364
365 pub(crate) fn increment_module_frames_dropped_released_route(&self, module_id: Option<&str>) {
369 self.module_frames_dropped_released_route
370 .fetch_add(1, Ordering::Relaxed);
371 if let Some(module_id) = module_id {
372 increment_keyed_count(
373 &self.module_frames_dropped_released_route_by_module,
374 module_id,
375 );
376 }
377 }
378
379 pub(crate) fn increment_module_orphan_route_goodbyes_sent(&self) {
380 self.module_orphan_route_goodbyes_sent
381 .fetch_add(1, Ordering::Relaxed);
382 }
383
384 pub(crate) fn increment_route_open_refused(&self, code: &'static str) {
385 debug_assert!(ROUTE_OPEN_REFUSAL_COUNTER_CODES.contains(&code));
386 increment_keyed_count(&self.route_open_refused_by_code, code);
387 }
388
389 pub(crate) fn increment_route_open_accepted(&self, principal: &str) {
399 increment_keyed_count(&self.route_open_accepted_by_principal, principal);
400 }
401
402 pub(crate) fn increment_module_requests_dropped_stale_route(&self) {
403 self.module_requests_dropped_stale_route
404 .fetch_add(1, Ordering::Relaxed);
405 }
406
407 pub(crate) fn increment_client_frames_dropped_stale_route(&self) {
408 self.client_frames_dropped_stale_route
409 .fetch_add(1, Ordering::Relaxed);
410 }
411
412 pub(crate) fn increment_client_egress_close_delivery_failed(&self) {
413 self.client_egress_close_delivery_failed
414 .fetch_add(1, Ordering::Relaxed);
415 }
416
417 pub(crate) fn increment_goodbye_relay_client_failed(&self) {
418 self.goodbye_relay_client_failed
419 .fetch_add(1, Ordering::Relaxed);
420 }
421
422 pub(crate) fn increment_goodbye_relay_module_dropped(&self, module_id: Option<&str>) {
423 self.goodbye_relay_module_dropped
424 .fetch_add(1, Ordering::Relaxed);
425 if let Some(module_id) = module_id {
426 increment_keyed_count(&self.goodbye_relay_module_dropped_by_module, module_id);
427 }
428 }
429
430 pub(crate) fn increment_route_released_epoch_fenced(&self) {
431 self.route_released_epoch_fenced
432 .fetch_add(1, Ordering::Relaxed);
433 }
434
435 pub(crate) fn increment_route_release_stale_skipped(&self) {
436 self.route_release_stale_skipped
437 .fetch_add(1, Ordering::Relaxed);
438 }
439
440 pub(crate) fn increment_drains_with_undeclared_gauge(&self) {
441 self.drains_with_undeclared_gauge
442 .fetch_add(1, Ordering::Relaxed);
443 }
444}
445
446fn increment_keyed_count(counts: &Mutex<HashMap<String, u64>>, key: &str) {
447 *counts
448 .lock()
449 .expect("keyed counter mutex poisoned")
450 .entry(key.to_string())
451 .or_default() += 1;
452}
453
454fn insert_nonempty_counts(
455 snapshot: &mut serde_json::Map<String, Value>,
456 key: &str,
457 counts: &Mutex<HashMap<String, u64>>,
458) {
459 let counts: MutexGuard<'_, HashMap<String, u64>> =
460 counts.lock().expect("keyed counter mutex poisoned");
461 if !counts.is_empty() {
462 snapshot.insert(key.to_string(), json!(&*counts));
463 }
464}
465
466#[cfg(test)]
467mod tests {
468 use super::*;
469
470 #[test]
471 fn counter_snapshot_includes_zero_rate_and_omits_empty_module_maps() {
472 let counters = DaemonCounters::new();
473 let snapshot = counters.snapshot();
474
475 assert_eq!(snapshot["module_frames_dropped_no_route_last_10m"], 0);
476 assert_eq!(
477 snapshot["module_frames_dropped_no_route_nonzero_minutes_last_10m"],
478 0
479 );
480 assert!(snapshot
481 .get("module_frames_dropped_no_route_by_module")
482 .is_none());
483 assert!(snapshot
484 .get("goodbye_relay_module_dropped_by_module")
485 .is_none());
486 }
487
488 #[tokio::test(start_paused = true)]
489 async fn module_frame_drops_are_attributed_to_the_emitting_module() {
490 let counters = DaemonCounters::new();
491 counters.increment_module_frames_dropped_no_route(Some("alpha"));
492 counters.increment_module_frames_dropped_no_route(Some("alpha"));
493
494 let snapshot = counters.snapshot();
495 assert_eq!(snapshot["module_frames_dropped_no_route"], 2);
496 assert_eq!(
497 snapshot["module_frames_dropped_no_route_by_module"],
498 json!({ "alpha": 2 })
499 );
500 assert_eq!(snapshot["module_frames_dropped_no_route_last_10m"], 2);
501 }
502
503 #[tokio::test(start_paused = true)]
504 async fn frame_drop_rate_ages_out_after_ten_minute_buckets() {
505 let counters = DaemonCounters::new();
506 counters.increment_module_frames_dropped_no_route(Some("alpha"));
507
508 tokio::time::advance(Duration::from_secs(9 * 60)).await;
509 assert_eq!(
510 counters.snapshot()["module_frames_dropped_no_route_last_10m"],
511 1
512 );
513
514 tokio::time::advance(Duration::from_secs(60)).await;
515 assert_eq!(
516 counters.snapshot()["module_frames_dropped_no_route_last_10m"],
517 0
518 );
519 }
520
521 #[tokio::test(start_paused = true)]
522 async fn frame_drop_window_counts_only_nonzero_minutes() {
523 let counters = DaemonCounters::new();
524 for minute in 0..10 {
525 if minute > 0 {
526 tokio::time::advance(Duration::from_secs(60)).await;
527 }
528 if minute != 4 {
529 counters.increment_module_frames_dropped_no_route(Some("alpha"));
530 }
531 }
532
533 assert_eq!(
534 counters.snapshot()["module_frames_dropped_no_route_nonzero_minutes_last_10m"],
535 9
536 );
537 }
538
539 #[test]
540 fn goodbye_relay_drops_are_attributed_to_the_target_module() {
541 let counters = DaemonCounters::new();
542 counters.increment_goodbye_relay_module_dropped(Some("alpha"));
543
544 assert_eq!(
545 counters.snapshot()["goodbye_relay_module_dropped_by_module"],
546 json!({ "alpha": 1 })
547 );
548 }
549}