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