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