Skip to main content

subc_daemon/
observability.rs

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
15/// The only keys the route.open refusal counter may carry. A module's own
16/// error code on a rejected bind is counted under `module_rejected` and rides
17/// the log event as a separate field: if the module's string were the key, a
18/// module could grow this map for the daemon's lifetime and push terminal
19/// control sequences through `ck daemon` to an operator's screen. The
20/// `&'static str` increment signature plus the debug assertion keep the set
21/// closed at the call site, not just here.
22const 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
52/// Counter key for a `route.open` refused by the per-module bind-relay breaker
53/// before any relay was attempted.
54///
55/// The frame the caller receives carries `module_timeout`, because both SDKs
56/// already classify that as retryable with capped backoff and inventing a new
57/// wire code would need a change in each of them. The COUNTER is deliberately a
58/// different key: "this module burned the full bind budget" and "this module is
59/// being refused in microseconds because it already did that repeatedly" are
60/// the two states an operator most needs to tell apart, and they are
61/// indistinguishable from the client side, where both look like one retryable
62/// error that the next attempt may well satisfy.
63pub(crate) const ROUTE_OPEN_REFUSED_BREAKER_OPEN: &str = "module_timeout_breaker_open";
64
65/// Counter key for a registered module that declared itself not ready.
66///
67/// The caller still receives `module_warming`, but operators must be able to
68/// distinguish declared readiness from a supervised process that has not
69/// registered yet.
70pub(crate) const ROUTE_OPEN_REFUSED_DECLARED_NOT_READY: &str = "module_warming_declared_not_ready";
71
72/// Counter key for a registered, declared-ready module held not-ready because
73/// a capability it declares `need: required` has no registered provider.
74///
75/// The caller still receives `module_warming`; a separate key lets an operator
76/// tell "the module says it is warming" from "the module is waiting on a
77/// provider that is not running", which point at different fixes.
78pub(crate) const ROUTE_OPEN_REFUSED_REQUIRED_CAPABILITY_UNPROVIDED: &str =
79    "module_warming_required_capability_unprovided";
80
81/// Shared count of authenticated socket connections accepted by the daemon.
82#[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, not INFO: this pair fired on every connect and disconnect and was
100        // measured at 43-74% of the daemon's own log (issue #114), burying the
101        // lines an operator opens the file for. The count itself is not lost: it
102        // is served live as `connected_clients` on server.describe (`ck daemon`),
103        // and a connection that opens a route still names its connection_id on the
104        // `route.open accepted` line, so per-connection forensics survive where
105        // they matter. Raise the filter (`CK_LOG=subc=debug`) to see every churn.
106        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 for the same reason as the open side above.
129        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/// Lock-free counters for route lifecycle drops and delivery failures.
139#[derive(Debug, Clone, Default)]
140pub struct DaemonCounters {
141    /// Every non-request frame a module sent on a (channel, epoch) the daemon
142    /// holds no bound route for, dropped. Counts both kinds below: orphan
143    /// traffic on a route the daemon released, and frames on a (channel,
144    /// epoch) that was never allocated on that module connection.
145    module_frames_dropped_no_route: Arc<AtomicU64>,
146    /// The orphan-traffic part of `module_frames_dropped_no_route`: the daemon
147    /// had allocated that (channel, epoch) on the sending module connection and
148    /// has since released it, so the module is still holding a route the
149    /// daemon closed. `module_frames_dropped_no_route` minus this is the count
150    /// of frames on a (channel, epoch) that never existed on the connection.
151    module_frames_dropped_released_route: Arc<AtomicU64>,
152    module_frames_dropped_released_route_by_module: Arc<Mutex<HashMap<String, u64>>>,
153    /// Route GOODBYEs the daemon enqueued to a module in answer to its frames
154    /// on a (channel, epoch) the daemon holds no route for, telling the module
155    /// to drop the route. Rate-limited per module connection and channel, so
156    /// this counts answers, not orphan frames.
157    module_orphan_route_goodbyes_sent: Arc<AtomicU64>,
158    // Per-module maps and the rate window are daemon-lifetime diagnostics only:
159    // they deliberately reset on restart instead of becoming durable daemon state.
160    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/// Ten one-minute buckets make sustained module-to-client route drops visible
176/// without retaining one record for every dropped frame.
177#[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    /// Returns a JSON snapshot whose stable, additive schema keeps the
244    /// `server.describe` diagnostic endpoint backward-compatible.
245    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    /// Count a dropped module frame whose (channel, epoch) the daemon had
367    /// allocated on that connection and since released. Called in addition to
368    /// [`Self::increment_module_frames_dropped_no_route`], never instead of it.
369    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    /// Count an accepted route.open by the principal the daemon stamped.
391    ///
392    /// THE KEY SPACE IS CLOSED BY CONSTRUCTION, unlike the refusal counter which
393    /// needs an explicit allowlist: a principal is `direct` or
394    /// `reserved:<module_id>`, and a module id was already refused at HELLO
395    /// unless it is a single path component free of control characters. So an
396    /// untrusted string cannot expand this map without first passing module-id
397    /// validation, and the bound is the number of modules rather than the number
398    /// of distinct strings a caller can invent.
399    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}