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