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