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