1use std::{
2 sync::atomic::{AtomicU8, AtomicU64, Ordering},
3 time::Duration,
4};
5
6#[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 pub const fn samples(self) -> u64 {
18 self.samples
19 }
20
21 pub const fn failures(self) -> u64 {
23 self.failures
24 }
25
26 pub const fn total(self) -> Duration {
28 self.total
29 }
30
31 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}