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