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    subc_protocol::error_codes::INVALID_REQUEST,
48];
49
50/// Counter key for a `route.open` refused by the per-module bind-relay breaker
51/// before any relay was attempted.
52///
53/// The frame the caller receives carries `module_timeout`, because both SDKs
54/// already classify that as retryable with capped backoff and inventing a new
55/// wire code would need a change in each of them. The COUNTER is deliberately a
56/// different key: "this module burned the full bind budget" and "this module is
57/// being refused in microseconds because it already did that repeatedly" are
58/// the two states an operator most needs to tell apart, and they are
59/// indistinguishable from the client side, where both look like one retryable
60/// error that the next attempt may well satisfy.
61pub(crate) const ROUTE_OPEN_REFUSED_BREAKER_OPEN: &str = "module_timeout_breaker_open";
62
63/// Counter key for a registered module that declared itself not ready.
64///
65/// The caller still receives `module_warming`, but operators must be able to
66/// distinguish declared readiness from a supervised process that has not
67/// registered yet.
68pub(crate) const ROUTE_OPEN_REFUSED_DECLARED_NOT_READY: &str = "module_warming_declared_not_ready";
69
70/// Counter key for a registered, declared-ready module held not-ready because
71/// a capability it declares `need: required` has no registered provider.
72///
73/// The caller still receives `module_warming`; a separate key lets an operator
74/// tell "the module says it is warming" from "the module is waiting on a
75/// provider that is not running", which point at different fixes.
76pub(crate) const ROUTE_OPEN_REFUSED_REQUIRED_CAPABILITY_UNPROVIDED: &str =
77    "module_warming_required_capability_unprovided";
78
79/// Shared count of authenticated socket connections accepted by the daemon.
80#[derive(Debug, Clone, Default)]
81pub struct ConnectedClients {
82    count: Arc<AtomicU64>,
83}
84
85impl ConnectedClients {
86    pub fn new() -> Self {
87        Self::default()
88    }
89
90    pub fn count(&self) -> u64 {
91        self.count.load(Ordering::SeqCst)
92    }
93
94    pub(crate) fn open(&self, connection_id: ConnectionId) -> ConnectedClientGuard {
95        let previous = self.count.fetch_add(1, Ordering::SeqCst);
96        let current = previous + 1;
97        // DEBUG, not INFO: this pair fired on every connect and disconnect and was
98        // measured at 43-74% of the daemon's own log (issue #114), burying the
99        // lines an operator opens the file for. The count itself is not lost: it
100        // is served live as `connected_clients` on server.describe (`ck daemon`),
101        // and a connection that opens a route still names its connection_id on the
102        // `route.open accepted` line, so per-connection forensics survive where
103        // they matter. Raise the filter (`CK_LOG=subc=debug`) to see every churn.
104        debug!(
105            connection_id = connection_id.get(),
106            connected_clients = current,
107            previous_connected_clients = previous,
108            "authenticated connection count changed"
109        );
110        ConnectedClientGuard {
111            clients: self.clone(),
112            connection_id,
113        }
114    }
115}
116
117pub(crate) struct ConnectedClientGuard {
118    clients: ConnectedClients,
119    connection_id: ConnectionId,
120}
121
122impl Drop for ConnectedClientGuard {
123    fn drop(&mut self) {
124        let previous = self.clients.count.fetch_sub(1, Ordering::SeqCst);
125        let current = previous.saturating_sub(1);
126        // DEBUG for the same reason as the open side above.
127        debug!(
128            connection_id = self.connection_id.get(),
129            connected_clients = current,
130            previous_connected_clients = previous,
131            "authenticated connection count changed"
132        );
133    }
134}
135
136/// Lock-free counters for route lifecycle drops and delivery failures.
137#[derive(Debug, Clone, Default)]
138pub struct DaemonCounters {
139    /// Every non-request frame a module sent on a (channel, epoch) the daemon
140    /// holds no bound route for, dropped. Counts both kinds below: orphan
141    /// traffic on a route the daemon released, and frames on a (channel,
142    /// epoch) that was never allocated on that module connection.
143    module_frames_dropped_no_route: Arc<AtomicU64>,
144    /// The orphan-traffic part of `module_frames_dropped_no_route`: the daemon
145    /// had allocated that (channel, epoch) on the sending module connection and
146    /// has since released it, so the module is still holding a route the
147    /// daemon closed. `module_frames_dropped_no_route` minus this is the count
148    /// of frames on a (channel, epoch) that never existed on the connection.
149    module_frames_dropped_released_route: Arc<AtomicU64>,
150    module_frames_dropped_released_route_by_module: Arc<Mutex<HashMap<String, u64>>>,
151    /// Route GOODBYEs the daemon enqueued to a module in answer to its frames
152    /// on a (channel, epoch) the daemon holds no route for, telling the module
153    /// to drop the route. Rate-limited per module connection and channel, so
154    /// this counts answers, not orphan frames.
155    module_orphan_route_goodbyes_sent: Arc<AtomicU64>,
156    // Per-module maps and the rate window are daemon-lifetime diagnostics only:
157    // they deliberately reset on restart instead of becoming durable daemon state.
158    module_frames_dropped_no_route_by_module: Arc<Mutex<HashMap<String, u64>>>,
159    route_open_refused_by_code: Arc<Mutex<HashMap<String, u64>>>,
160    route_open_accepted_by_principal: Arc<Mutex<HashMap<String, u64>>>,
161    module_frames_dropped_no_route_window: Arc<Mutex<DropWindow>>,
162    module_requests_dropped_stale_route: Arc<AtomicU64>,
163    client_frames_dropped_stale_route: Arc<AtomicU64>,
164    client_egress_close_delivery_failed: Arc<AtomicU64>,
165    goodbye_relay_client_failed: Arc<AtomicU64>,
166    goodbye_relay_module_dropped: Arc<AtomicU64>,
167    goodbye_relay_module_dropped_by_module: Arc<Mutex<HashMap<String, u64>>>,
168    route_released_epoch_fenced: Arc<AtomicU64>,
169    route_release_stale_skipped: Arc<AtomicU64>,
170    drains_with_undeclared_gauge: Arc<AtomicU64>,
171}
172
173/// Ten one-minute buckets make sustained module-to-client route drops visible
174/// without retaining one record for every dropped frame.
175#[derive(Debug)]
176struct DropWindow {
177    started_at: tokio::time::Instant,
178    buckets: VecDeque<DropBucket>,
179}
180
181#[derive(Debug)]
182struct DropBucket {
183    minute: u64,
184    count: u64,
185}
186
187impl Default for DropWindow {
188    fn default() -> Self {
189        Self {
190            started_at: tokio::time::Instant::now(),
191            buckets: VecDeque::new(),
192        }
193    }
194}
195
196impl DropWindow {
197    const MINUTE: Duration = Duration::from_secs(60);
198    const BUCKETS: u64 = 10;
199
200    fn record(&mut self, now: tokio::time::Instant) {
201        let minute = self.minute_at(now);
202        self.prune_before(minute);
203        match self.buckets.back_mut() {
204            Some(bucket) if bucket.minute == minute => bucket.count += 1,
205            _ => self.buckets.push_back(DropBucket { minute, count: 1 }),
206        }
207    }
208
209    fn count_last_10m(&mut self, now: tokio::time::Instant) -> u64 {
210        let minute = self.minute_at(now);
211        self.prune_before(minute);
212        self.buckets.iter().map(|bucket| bucket.count).sum()
213    }
214
215    fn nonzero_minutes_last_10m(&mut self, now: tokio::time::Instant) -> u64 {
216        let minute = self.minute_at(now);
217        self.prune_before(minute);
218        self.buckets.len() as u64
219    }
220
221    fn minute_at(&self, now: tokio::time::Instant) -> u64 {
222        now.saturating_duration_since(self.started_at).as_secs() / Self::MINUTE.as_secs()
223    }
224
225    fn prune_before(&mut self, current_minute: u64) {
226        while self
227            .buckets
228            .front()
229            .is_some_and(|bucket| current_minute.saturating_sub(bucket.minute) >= Self::BUCKETS)
230        {
231            self.buckets.pop_front();
232        }
233    }
234}
235
236impl DaemonCounters {
237    pub fn new() -> Self {
238        Self::default()
239    }
240
241    /// Returns a JSON snapshot whose stable, additive schema keeps the
242    /// `server.describe` diagnostic endpoint backward-compatible.
243    pub fn snapshot(&self) -> Value {
244        let mut snapshot = serde_json::Map::new();
245        snapshot.insert(
246            "module_frames_dropped_no_route".into(),
247            self.module_frames_dropped_no_route
248                .load(Ordering::Relaxed)
249                .into(),
250        );
251        let mut drop_window = self
252            .module_frames_dropped_no_route_window
253            .lock()
254            .expect("drop-rate window mutex poisoned");
255        let now = tokio::time::Instant::now();
256        snapshot.insert(
257            "module_frames_dropped_no_route_last_10m".into(),
258            drop_window.count_last_10m(now).into(),
259        );
260        snapshot.insert(
261            "module_frames_dropped_no_route_nonzero_minutes_last_10m".into(),
262            drop_window.nonzero_minutes_last_10m(now).into(),
263        );
264        insert_nonempty_counts(
265            &mut snapshot,
266            "module_frames_dropped_no_route_by_module",
267            &self.module_frames_dropped_no_route_by_module,
268        );
269        snapshot.insert(
270            "module_frames_dropped_released_route".into(),
271            self.module_frames_dropped_released_route
272                .load(Ordering::Relaxed)
273                .into(),
274        );
275        insert_nonempty_counts(
276            &mut snapshot,
277            "module_frames_dropped_released_route_by_module",
278            &self.module_frames_dropped_released_route_by_module,
279        );
280        snapshot.insert(
281            "module_orphan_route_goodbyes_sent".into(),
282            self.module_orphan_route_goodbyes_sent
283                .load(Ordering::Relaxed)
284                .into(),
285        );
286        insert_nonempty_counts(
287            &mut snapshot,
288            "route_open_refused_by_code",
289            &self.route_open_refused_by_code,
290        );
291        insert_nonempty_counts(
292            &mut snapshot,
293            "route_open_accepted_by_principal",
294            &self.route_open_accepted_by_principal,
295        );
296        snapshot.insert(
297            "module_requests_dropped_stale_route".into(),
298            self.module_requests_dropped_stale_route
299                .load(Ordering::Relaxed)
300                .into(),
301        );
302        snapshot.insert(
303            "client_frames_dropped_stale_route".into(),
304            self.client_frames_dropped_stale_route
305                .load(Ordering::Relaxed)
306                .into(),
307        );
308        snapshot.insert(
309            "client_egress_close_delivery_failed".into(),
310            self.client_egress_close_delivery_failed
311                .load(Ordering::Relaxed)
312                .into(),
313        );
314        snapshot.insert(
315            "goodbye_relay_client_failed".into(),
316            self.goodbye_relay_client_failed
317                .load(Ordering::Relaxed)
318                .into(),
319        );
320        snapshot.insert(
321            "goodbye_relay_module_dropped".into(),
322            self.goodbye_relay_module_dropped
323                .load(Ordering::Relaxed)
324                .into(),
325        );
326        insert_nonempty_counts(
327            &mut snapshot,
328            "goodbye_relay_module_dropped_by_module",
329            &self.goodbye_relay_module_dropped_by_module,
330        );
331        snapshot.insert(
332            "route_released_epoch_fenced".into(),
333            self.route_released_epoch_fenced
334                .load(Ordering::Relaxed)
335                .into(),
336        );
337        snapshot.insert(
338            "route_release_stale_skipped".into(),
339            self.route_release_stale_skipped
340                .load(Ordering::Relaxed)
341                .into(),
342        );
343        snapshot.insert(
344            "drains_with_undeclared_gauge".into(),
345            self.drains_with_undeclared_gauge
346                .load(Ordering::Relaxed)
347                .into(),
348        );
349        Value::Object(snapshot)
350    }
351
352    pub(crate) fn increment_module_frames_dropped_no_route(&self, module_id: Option<&str>) {
353        self.module_frames_dropped_no_route
354            .fetch_add(1, Ordering::Relaxed);
355        if let Some(module_id) = module_id {
356            increment_keyed_count(&self.module_frames_dropped_no_route_by_module, module_id);
357        }
358        self.module_frames_dropped_no_route_window
359            .lock()
360            .expect("drop-rate window mutex poisoned")
361            .record(tokio::time::Instant::now());
362    }
363
364    /// Count a dropped module frame whose (channel, epoch) the daemon had
365    /// allocated on that connection and since released. Called in addition to
366    /// [`Self::increment_module_frames_dropped_no_route`], never instead of it.
367    pub(crate) fn increment_module_frames_dropped_released_route(&self, module_id: Option<&str>) {
368        self.module_frames_dropped_released_route
369            .fetch_add(1, Ordering::Relaxed);
370        if let Some(module_id) = module_id {
371            increment_keyed_count(
372                &self.module_frames_dropped_released_route_by_module,
373                module_id,
374            );
375        }
376    }
377
378    pub(crate) fn increment_module_orphan_route_goodbyes_sent(&self) {
379        self.module_orphan_route_goodbyes_sent
380            .fetch_add(1, Ordering::Relaxed);
381    }
382
383    pub(crate) fn increment_route_open_refused(&self, code: &'static str) {
384        debug_assert!(ROUTE_OPEN_REFUSAL_COUNTER_CODES.contains(&code));
385        increment_keyed_count(&self.route_open_refused_by_code, code);
386    }
387
388    /// Count an accepted route.open by the principal the daemon stamped.
389    ///
390    /// THE KEY SPACE IS CLOSED BY CONSTRUCTION, unlike the refusal counter which
391    /// needs an explicit allowlist: a principal is `direct` or
392    /// `reserved:<module_id>`, and a module id was already refused at HELLO
393    /// unless it is a single path component free of control characters. So an
394    /// untrusted string cannot expand this map without first passing module-id
395    /// validation, and the bound is the number of modules rather than the number
396    /// of distinct strings a caller can invent.
397    pub(crate) fn increment_route_open_accepted(&self, principal: &str) {
398        increment_keyed_count(&self.route_open_accepted_by_principal, principal);
399    }
400
401    pub(crate) fn increment_module_requests_dropped_stale_route(&self) {
402        self.module_requests_dropped_stale_route
403            .fetch_add(1, Ordering::Relaxed);
404    }
405
406    pub(crate) fn increment_client_frames_dropped_stale_route(&self) {
407        self.client_frames_dropped_stale_route
408            .fetch_add(1, Ordering::Relaxed);
409    }
410
411    pub(crate) fn increment_client_egress_close_delivery_failed(&self) {
412        self.client_egress_close_delivery_failed
413            .fetch_add(1, Ordering::Relaxed);
414    }
415
416    pub(crate) fn increment_goodbye_relay_client_failed(&self) {
417        self.goodbye_relay_client_failed
418            .fetch_add(1, Ordering::Relaxed);
419    }
420
421    pub(crate) fn increment_goodbye_relay_module_dropped(&self, module_id: Option<&str>) {
422        self.goodbye_relay_module_dropped
423            .fetch_add(1, Ordering::Relaxed);
424        if let Some(module_id) = module_id {
425            increment_keyed_count(&self.goodbye_relay_module_dropped_by_module, module_id);
426        }
427    }
428
429    pub(crate) fn increment_route_released_epoch_fenced(&self) {
430        self.route_released_epoch_fenced
431            .fetch_add(1, Ordering::Relaxed);
432    }
433
434    pub(crate) fn increment_route_release_stale_skipped(&self) {
435        self.route_release_stale_skipped
436            .fetch_add(1, Ordering::Relaxed);
437    }
438
439    pub(crate) fn increment_drains_with_undeclared_gauge(&self) {
440        self.drains_with_undeclared_gauge
441            .fetch_add(1, Ordering::Relaxed);
442    }
443}
444
445fn increment_keyed_count(counts: &Mutex<HashMap<String, u64>>, key: &str) {
446    *counts
447        .lock()
448        .expect("keyed counter mutex poisoned")
449        .entry(key.to_string())
450        .or_default() += 1;
451}
452
453fn insert_nonempty_counts(
454    snapshot: &mut serde_json::Map<String, Value>,
455    key: &str,
456    counts: &Mutex<HashMap<String, u64>>,
457) {
458    let counts: MutexGuard<'_, HashMap<String, u64>> =
459        counts.lock().expect("keyed counter mutex poisoned");
460    if !counts.is_empty() {
461        snapshot.insert(key.to_string(), json!(&*counts));
462    }
463}
464
465#[cfg(test)]
466mod tests {
467    use super::*;
468
469    #[test]
470    fn counter_snapshot_includes_zero_rate_and_omits_empty_module_maps() {
471        let counters = DaemonCounters::new();
472        let snapshot = counters.snapshot();
473
474        assert_eq!(snapshot["module_frames_dropped_no_route_last_10m"], 0);
475        assert_eq!(
476            snapshot["module_frames_dropped_no_route_nonzero_minutes_last_10m"],
477            0
478        );
479        assert!(snapshot
480            .get("module_frames_dropped_no_route_by_module")
481            .is_none());
482        assert!(snapshot
483            .get("goodbye_relay_module_dropped_by_module")
484            .is_none());
485    }
486
487    #[tokio::test(start_paused = true)]
488    async fn module_frame_drops_are_attributed_to_the_emitting_module() {
489        let counters = DaemonCounters::new();
490        counters.increment_module_frames_dropped_no_route(Some("alpha"));
491        counters.increment_module_frames_dropped_no_route(Some("alpha"));
492
493        let snapshot = counters.snapshot();
494        assert_eq!(snapshot["module_frames_dropped_no_route"], 2);
495        assert_eq!(
496            snapshot["module_frames_dropped_no_route_by_module"],
497            json!({ "alpha": 2 })
498        );
499        assert_eq!(snapshot["module_frames_dropped_no_route_last_10m"], 2);
500    }
501
502    #[tokio::test(start_paused = true)]
503    async fn frame_drop_rate_ages_out_after_ten_minute_buckets() {
504        let counters = DaemonCounters::new();
505        counters.increment_module_frames_dropped_no_route(Some("alpha"));
506
507        tokio::time::advance(Duration::from_secs(9 * 60)).await;
508        assert_eq!(
509            counters.snapshot()["module_frames_dropped_no_route_last_10m"],
510            1
511        );
512
513        tokio::time::advance(Duration::from_secs(60)).await;
514        assert_eq!(
515            counters.snapshot()["module_frames_dropped_no_route_last_10m"],
516            0
517        );
518    }
519
520    #[tokio::test(start_paused = true)]
521    async fn frame_drop_window_counts_only_nonzero_minutes() {
522        let counters = DaemonCounters::new();
523        for minute in 0..10 {
524            if minute > 0 {
525                tokio::time::advance(Duration::from_secs(60)).await;
526            }
527            if minute != 4 {
528                counters.increment_module_frames_dropped_no_route(Some("alpha"));
529            }
530        }
531
532        assert_eq!(
533            counters.snapshot()["module_frames_dropped_no_route_nonzero_minutes_last_10m"],
534            9
535        );
536    }
537
538    #[test]
539    fn goodbye_relay_drops_are_attributed_to_the_target_module() {
540        let counters = DaemonCounters::new();
541        counters.increment_goodbye_relay_module_dropped(Some("alpha"));
542
543        assert_eq!(
544            counters.snapshot()["goodbye_relay_module_dropped_by_module"],
545            json!({ "alpha": 1 })
546        );
547    }
548}