Skip to main content

pb_mapper_client/
diagnostics.rs

1//! Bounded, credential-free recovery counters exposed separately from lifecycle status.
2
3use crate::endpoint::RelayEndpoint;
4use serde::Serialize;
5use std::sync::{Arc, Mutex};
6use tokio::time::{Duration, Instant};
7
8/// The phase most recently entered by a tunnel worker.
9#[derive(Clone, Copy, Debug, Serialize, PartialEq, Eq)]
10#[serde(rename_all = "snake_case")]
11pub enum RecoveryPhase {
12    Starting,
13    Capacity,
14    Resolving,
15    Connecting,
16    Handshake,
17    Ready,
18    Backoff,
19    Stopped,
20}
21
22impl RecoveryPhase {
23    /// Stable diagnostic label, matching the serialized representation.
24    pub fn as_str(self) -> &'static str {
25        match self {
26            Self::Starting => "starting",
27            Self::Capacity => "capacity",
28            Self::Resolving => "resolving",
29            Self::Connecting => "connecting",
30            Self::Handshake => "handshake",
31            Self::Ready => "ready",
32            Self::Backoff => "backoff",
33            Self::Stopped => "stopped",
34        }
35    }
36}
37
38/// Why a recovery attempt failed, without credentials or application data.
39#[derive(Clone, Copy, Debug, Serialize, PartialEq, Eq)]
40#[serde(rename_all = "snake_case")]
41pub enum RecoveryFailure {
42    Dns,
43    Transport,
44    Timeout,
45    ServiceUnavailable,
46    Rejected,
47    NetworkChanged,
48}
49
50impl RecoveryFailure {
51    /// Stable diagnostic label, matching the serialized representation.
52    pub fn as_str(self) -> &'static str {
53        match self {
54            Self::Dns => "dns",
55            Self::Transport => "transport",
56            Self::Timeout => "timeout",
57            Self::ServiceUnavailable => "service_unavailable",
58            Self::Rejected => "rejected",
59            Self::NetworkChanged => "network_changed",
60        }
61    }
62}
63
64/// A point-in-time view of this tunnel and its process-shared relay budget.
65#[derive(Clone, Debug, Serialize)]
66pub struct TunnelDiagnostics {
67    pub sdk_version: &'static str,
68    /// Most recent worker phase; use tunnel status for aggregate readiness.
69    pub last_attempt_phase: RecoveryPhase,
70    pub attempts: u64,
71    pub consecutive_failures: u64,
72    pub last_failure: Option<RecoveryFailure>,
73    /// Time since the last authenticated relay reply, even a negative reply.
74    pub last_success_age_ms: Option<u64>,
75    /// Last answered setup latency, excluding admission queue time.
76    pub last_setup_latency_ms: Option<u64>,
77    pub next_retry_in_ms: Option<u64>,
78    pub dns_age_ms: Option<u64>,
79    pub network_generation: u64,
80    /// Shared across this process for the same configured relay address.
81    pub active_control_setups: usize,
82    pub active_data_setups: usize,
83}
84
85struct State {
86    phase: RecoveryPhase,
87    attempts: u64,
88    failures: u64,
89    failure: Option<RecoveryFailure>,
90    success: Option<Instant>,
91    latency: Option<Duration>,
92    retry: Option<Instant>,
93}
94
95#[derive(Clone)]
96pub(crate) struct Diagnostics(Arc<Mutex<State>>);
97
98impl Default for Diagnostics {
99    fn default() -> Self {
100        Self(Arc::new(Mutex::new(State {
101            phase: RecoveryPhase::Starting,
102            attempts: 0,
103            failures: 0,
104            failure: None,
105            success: None,
106            latency: None,
107            retry: None,
108        })))
109    }
110}
111
112impl Diagnostics {
113    pub(crate) fn heard(&self) {
114        self.0.lock().unwrap_or_else(|e| e.into_inner()).success = Some(Instant::now());
115    }
116    pub(crate) fn responded(&self, elapsed: Duration) {
117        let mut state = self.0.lock().unwrap_or_else(|e| e.into_inner());
118        state.success = Some(Instant::now());
119        state.latency = Some(elapsed);
120    }
121    pub(crate) fn phase(&self, phase: RecoveryPhase) {
122        self.0.lock().unwrap_or_else(|e| e.into_inner()).phase = phase;
123    }
124    pub(crate) fn attempt(&self) {
125        let mut state = self.0.lock().unwrap_or_else(|e| e.into_inner());
126        state.attempts = state.attempts.saturating_add(1);
127        state.retry = None;
128        state.phase = RecoveryPhase::Capacity;
129    }
130    pub(crate) fn succeeded(&self, elapsed: Duration) {
131        let mut state = self.0.lock().unwrap_or_else(|e| e.into_inner());
132        state.phase = RecoveryPhase::Ready;
133        state.failures = 0;
134        state.success = Some(Instant::now());
135        state.latency = Some(elapsed);
136        state.retry = None;
137    }
138    pub(crate) fn failed(&self, failure: RecoveryFailure, retry: Duration) -> bool {
139        let mut state = self.0.lock().unwrap_or_else(|e| e.into_inner());
140        state.failures = state.failures.saturating_add(1);
141        let emit = state.failure != Some(failure) || state.failures.is_power_of_two();
142        state.phase = RecoveryPhase::Backoff;
143        state.failure = Some(failure);
144        state.retry = Some(Instant::now() + retry);
145        emit
146    }
147    pub(crate) fn snapshot(&self, endpoint: &RelayEndpoint) -> TunnelDiagnostics {
148        let state = self.0.lock().unwrap_or_else(|e| e.into_inner());
149        let ms = |time: Duration| time.as_millis().min(u128::from(u64::MAX)) as u64;
150        let (active_control_setups, active_data_setups) = endpoint.active_setups();
151        TunnelDiagnostics {
152            sdk_version: env!("CARGO_PKG_VERSION"),
153            last_attempt_phase: state.phase,
154            attempts: state.attempts,
155            consecutive_failures: state.failures,
156            last_failure: state.failure,
157            last_success_age_ms: state.success.map(|time| ms(time.elapsed())),
158            last_setup_latency_ms: state.latency.map(ms),
159            next_retry_in_ms: state
160                .retry
161                .map(|time| ms(time.saturating_duration_since(Instant::now()))),
162            dns_age_ms: endpoint.dns_age().map(ms),
163            network_generation: endpoint.network_generation(),
164            active_control_setups,
165            active_data_setups,
166        }
167    }
168}