iroh_db/
telemetry.rs

1use std::{
2    sync::atomic::{AtomicU8, AtomicU64, Ordering},
3    time::Duration,
4};
5
6/// Aggregate latency evidence without high-cardinality labels.
7#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
8pub struct LatencyMetrics {
9    samples: u64,
10    failures: u64,
11    total: Duration,
12    maximum: Duration,
13}
14
15impl LatencyMetrics {
16    /// Number of completed operations observed.
17    pub const fn samples(self) -> u64 {
18        self.samples
19    }
20
21    /// Number of observed operations that returned an error.
22    pub const fn failures(self) -> u64 {
23        self.failures
24    }
25
26    /// Sum of all observed operation durations.
27    pub const fn total(self) -> Duration {
28        self.total
29    }
30
31    /// Largest observed operation duration.
32    pub const fn maximum(self) -> Duration {
33        self.maximum
34    }
35}
36
37#[derive(Debug, Default)]
38struct LatencyRecorder {
39    samples: AtomicU64,
40    failures: AtomicU64,
41    total_nanos: AtomicU64,
42    maximum_nanos: AtomicU64,
43}
44
45impl LatencyRecorder {
46    fn observe(&self, elapsed: Duration, failed: bool) {
47        let nanos = u64::try_from(elapsed.as_nanos()).unwrap_or(u64::MAX);
48        saturating_add(&self.samples, 1);
49        if failed {
50            saturating_add(&self.failures, 1);
51        }
52        saturating_add(&self.total_nanos, nanos);
53        self.maximum_nanos.fetch_max(nanos, Ordering::Relaxed);
54    }
55
56    fn snapshot(&self) -> LatencyMetrics {
57        LatencyMetrics {
58            samples: self.samples.load(Ordering::Relaxed),
59            failures: self.failures.load(Ordering::Relaxed),
60            total: Duration::from_nanos(self.total_nanos.load(Ordering::Relaxed)),
61            maximum: Duration::from_nanos(self.maximum_nanos.load(Ordering::Relaxed)),
62        }
63    }
64}
65
66#[derive(Debug, Default)]
67pub(crate) struct Telemetry {
68    local_apply: LatencyRecorder,
69    query: LatencyRecorder,
70    index_records_scanned: AtomicU64,
71    reconciliation_commits_scanned: AtomicU64,
72    equivocation_quarantines: AtomicU64,
73    schema_rebases: AtomicU64,
74    sync_rounds: AtomicU64,
75    sync_failures: AtomicU64,
76    sync_received: AtomicU64,
77    sync_sent: AtomicU64,
78    snapshots_created: AtomicU64,
79    snapshots_installed: AtomicU64,
80    repairs_succeeded: AtomicU64,
81    repairs_failed: AtomicU64,
82    last_repair_consistency: AtomicU8,
83}
84
85impl Telemetry {
86    pub(crate) fn observe_local_apply(&self, elapsed: Duration, failed: bool) {
87        self.local_apply.observe(elapsed, failed);
88    }
89
90    pub(crate) fn observe_query(&self, elapsed: Duration, failed: bool) {
91        self.query.observe(elapsed, failed);
92    }
93
94    pub(crate) fn add_index_records_scanned(&self, records: usize) {
95        saturating_add(
96            &self.index_records_scanned,
97            u64::try_from(records).unwrap_or(u64::MAX),
98        );
99    }
100
101    pub(crate) fn add_reconciliation_commits_scanned(&self, commits: usize) {
102        saturating_add(
103            &self.reconciliation_commits_scanned,
104            u64::try_from(commits).unwrap_or(u64::MAX),
105        );
106    }
107
108    pub(crate) fn record_equivocation_quarantine(&self) {
109        saturating_add(&self.equivocation_quarantines, 1);
110    }
111
112    pub(crate) fn record_schema_rebase(&self) {
113        saturating_add(&self.schema_rebases, 1);
114    }
115
116    pub(crate) fn record_sync(&self, received: usize, sent: usize, failed: bool) {
117        saturating_add(&self.sync_rounds, 1);
118        if failed {
119            saturating_add(&self.sync_failures, 1);
120        }
121        saturating_add(
122            &self.sync_received,
123            u64::try_from(received).unwrap_or(u64::MAX),
124        );
125        saturating_add(&self.sync_sent, u64::try_from(sent).unwrap_or(u64::MAX));
126    }
127
128    pub(crate) fn record_snapshot_created(&self) {
129        saturating_add(&self.snapshots_created, 1);
130    }
131
132    pub(crate) fn record_snapshot_installed(&self) {
133        saturating_add(&self.snapshots_installed, 1);
134    }
135
136    pub(crate) fn record_repair(&self, consistent: Option<bool>) {
137        if let Some(consistent) = consistent {
138            saturating_add(&self.repairs_succeeded, 1);
139            self.last_repair_consistency
140                .store(if consistent { 2 } else { 1 }, Ordering::Relaxed);
141        } else {
142            saturating_add(&self.repairs_failed, 1);
143            self.last_repair_consistency.store(0, Ordering::Relaxed);
144        }
145    }
146
147    pub(crate) fn snapshot(&self) -> RuntimeMetrics {
148        RuntimeMetrics {
149            local_apply: self.local_apply.snapshot(),
150            query: self.query.snapshot(),
151            index_records_scanned: self.index_records_scanned.load(Ordering::Relaxed),
152            reconciliation_commits_scanned: self
153                .reconciliation_commits_scanned
154                .load(Ordering::Relaxed),
155            equivocation_quarantines: self.equivocation_quarantines.load(Ordering::Relaxed),
156            schema_rebases: self.schema_rebases.load(Ordering::Relaxed),
157            sync_rounds: self.sync_rounds.load(Ordering::Relaxed),
158            sync_failures: self.sync_failures.load(Ordering::Relaxed),
159            sync_received: self.sync_received.load(Ordering::Relaxed),
160            sync_sent: self.sync_sent.load(Ordering::Relaxed),
161            snapshots_created: self.snapshots_created.load(Ordering::Relaxed),
162            snapshots_installed: self.snapshots_installed.load(Ordering::Relaxed),
163            repairs_succeeded: self.repairs_succeeded.load(Ordering::Relaxed),
164            repairs_failed: self.repairs_failed.load(Ordering::Relaxed),
165            last_repair_derived_consistent: match self
166                .last_repair_consistency
167                .load(Ordering::Relaxed)
168            {
169                1 => Some(false),
170                2 => Some(true),
171                _ => None,
172            },
173        }
174    }
175}
176
177fn saturating_add(counter: &AtomicU64, value: u64) {
178    let _previous = counter.fetch_update(Ordering::Relaxed, Ordering::Relaxed, |current| {
179        Some(current.saturating_add(value))
180    });
181}
182
183#[derive(Debug, Clone, Copy)]
184pub(crate) struct RuntimeMetrics {
185    pub(crate) local_apply: LatencyMetrics,
186    pub(crate) query: LatencyMetrics,
187    pub(crate) index_records_scanned: u64,
188    pub(crate) reconciliation_commits_scanned: u64,
189    pub(crate) equivocation_quarantines: u64,
190    pub(crate) schema_rebases: u64,
191    pub(crate) sync_rounds: u64,
192    pub(crate) sync_failures: u64,
193    pub(crate) sync_received: u64,
194    pub(crate) sync_sent: u64,
195    pub(crate) snapshots_created: u64,
196    pub(crate) snapshots_installed: u64,
197    pub(crate) repairs_succeeded: u64,
198    pub(crate) repairs_failed: u64,
199    pub(crate) last_repair_derived_consistent: Option<bool>,
200}
201
202#[cfg(test)]
203mod tests {
204    use super::*;
205
206    #[test]
207    fn counters_saturate_instead_of_wrapping() {
208        let counter = AtomicU64::new(u64::MAX);
209        saturating_add(&counter, 1);
210        assert_eq!(counter.load(Ordering::Relaxed), u64::MAX);
211    }
212}