1use std::sync::atomic::{AtomicU64, Ordering};
4
5use crate::{
6 CapacityDimension, DatabaseDisposition, Health, LifecycleState, OutboundResult, Stage,
7 StageOutcome,
8};
9
10pub const LATENCY_BUCKET_UPPER_MS: [u64; 11] =
11 [1, 5, 10, 25, 50, 100, 250, 500, 1_000, 5_000, u64::MAX];
12
13pub(crate) struct Metrics {
14 requests: [AtomicU64; 4],
15 stage_latency: [[AtomicU64; 11]; 7],
16 capacity_accepted: [AtomicU64; 5],
17 capacity_rejected: [AtomicU64; 5],
18 capacity_used: [AtomicU64; 5],
19 database: [AtomicU64; 3],
20 outbound: [AtomicU64; 4],
21 logger_dropped: AtomicU64,
22 logger_output_failed: AtomicU64,
23 lifecycle: AtomicU64,
24 health: AtomicU64,
25 ready: AtomicU64,
26}
27
28impl Default for Metrics {
29 fn default() -> Self {
30 Self {
31 requests: std::array::from_fn(|_| AtomicU64::new(0)),
32 stage_latency: std::array::from_fn(|_| std::array::from_fn(|_| AtomicU64::new(0))),
33 capacity_accepted: std::array::from_fn(|_| AtomicU64::new(0)),
34 capacity_rejected: std::array::from_fn(|_| AtomicU64::new(0)),
35 capacity_used: std::array::from_fn(|_| AtomicU64::new(0)),
36 database: std::array::from_fn(|_| AtomicU64::new(0)),
37 outbound: std::array::from_fn(|_| AtomicU64::new(0)),
38 logger_dropped: AtomicU64::new(0),
39 logger_output_failed: AtomicU64::new(0),
40 lifecycle: AtomicU64::new(LifecycleState::Starting as u64),
41 health: AtomicU64::new(Health::Healthy as u64),
42 ready: AtomicU64::new(0),
43 }
44 }
45}
46
47impl Metrics {
48 pub(crate) fn stage_finished(&self, stage: Stage, outcome: StageOutcome, elapsed_ms: u64) {
49 let bucket = LATENCY_BUCKET_UPPER_MS
50 .iter()
51 .position(|upper| elapsed_ms <= *upper)
52 .unwrap_or(LATENCY_BUCKET_UPPER_MS.len() - 1);
53 increment(&self.stage_latency[stage as usize][bucket], 1);
54 if stage == Stage::Response {
55 increment(&self.requests[outcome as usize], 1);
56 }
57 }
58
59 pub(crate) fn capacity(&self, dimension: Option<CapacityDimension>, used: u64, rejected: bool) {
60 let index = dimension.map_or(4, |value| value as usize);
61 let counters = if rejected {
62 &self.capacity_rejected
63 } else {
64 &self.capacity_accepted
65 };
66 increment(&counters[index], 1);
67 self.capacity_used[index].store(used, Ordering::Relaxed);
68 }
69
70 pub(crate) fn database(&self, disposition: DatabaseDisposition) {
71 increment(&self.database[disposition as usize], 1);
72 }
73
74 pub(crate) fn outbound(&self, result: OutboundResult) {
75 increment(&self.outbound[result as usize], 1);
76 }
77
78 pub(crate) fn logger_output(&self, output_failed: bool) {
79 if output_failed {
80 increment(&self.logger_output_failed, 1);
81 }
82 }
83
84 pub(crate) fn logger_dropped(&self, count: u64) {
85 increment(&self.logger_dropped, count);
86 }
87
88 pub(crate) fn lifecycle(&self, state: LifecycleState, health: Health) {
89 self.lifecycle.store(state as u64, Ordering::Relaxed);
90 self.health.store(health as u64, Ordering::Relaxed);
91 self.ready.store(
92 u64::from(state == LifecycleState::Running && health == Health::Healthy),
93 Ordering::Release,
94 );
95 }
96
97 pub(crate) fn snapshot(&self) -> MetricsSnapshot {
98 MetricsSnapshot {
99 requests: load_array(&self.requests),
100 stage_latency: std::array::from_fn(|stage| load_array(&self.stage_latency[stage])),
101 capacity_accepted: load_array(&self.capacity_accepted),
102 capacity_rejected: load_array(&self.capacity_rejected),
103 capacity_used: load_array(&self.capacity_used),
104 database: load_array(&self.database),
105 outbound: load_array(&self.outbound),
106 logger_dropped: self.logger_dropped.load(Ordering::Relaxed),
107 logger_output_failed: self.logger_output_failed.load(Ordering::Relaxed),
108 lifecycle: self.lifecycle.load(Ordering::Relaxed),
109 health: self.health.load(Ordering::Relaxed),
110 ready: self.ready.load(Ordering::Acquire) != 0,
111 }
112 }
113}
114
115fn load_array<const N: usize>(values: &[AtomicU64; N]) -> [u64; N] {
116 std::array::from_fn(|index| values[index].load(Ordering::Relaxed))
117}
118
119fn increment(value: &AtomicU64, count: u64) {
120 let _ = value.fetch_update(Ordering::Relaxed, Ordering::Relaxed, |current| {
121 Some(current.saturating_add(count))
122 });
123}
124
125#[derive(Clone, Copy, Debug, Eq, PartialEq)]
127pub struct MetricsSnapshot {
128 requests: [u64; 4],
129 stage_latency: [[u64; 11]; 7],
130 capacity_accepted: [u64; 5],
131 capacity_rejected: [u64; 5],
132 capacity_used: [u64; 5],
133 database: [u64; 3],
134 outbound: [u64; 4],
135 logger_dropped: u64,
136 logger_output_failed: u64,
137 lifecycle: u64,
138 health: u64,
139 ready: bool,
140}
141
142impl MetricsSnapshot {
143 pub fn requests(&self, outcome: StageOutcome) -> u64 {
144 self.requests[outcome as usize]
145 }
146 pub fn stage_latency(&self, stage: Stage) -> &[u64; 11] {
147 &self.stage_latency[stage as usize]
148 }
149 pub fn capacity_accepted(&self, dimension: Option<CapacityDimension>) -> u64 {
150 self.capacity_accepted[dimension.map_or(4, |value| value as usize)]
151 }
152 pub fn capacity_rejected(&self, dimension: Option<CapacityDimension>) -> u64 {
153 self.capacity_rejected[dimension.map_or(4, |value| value as usize)]
154 }
155 pub fn capacity_used(&self, dimension: Option<CapacityDimension>) -> u64 {
156 self.capacity_used[dimension.map_or(4, |value| value as usize)]
157 }
158 pub fn database(&self, disposition: DatabaseDisposition) -> u64 {
159 self.database[disposition as usize]
160 }
161 pub fn outbound(&self, result: OutboundResult) -> u64 {
162 self.outbound[result as usize]
163 }
164 pub const fn logger_dropped(&self) -> u64 {
165 self.logger_dropped
166 }
167 pub const fn logger_output_failed(&self) -> u64 {
168 self.logger_output_failed
169 }
170 pub fn lifecycle(&self) -> LifecycleState {
171 decode_lifecycle(self.lifecycle)
172 }
173 pub fn health(&self) -> Health {
174 decode_health(self.health)
175 }
176 pub const fn ready(&self) -> bool {
177 self.ready
178 }
179}
180
181fn decode_lifecycle(value: u64) -> LifecycleState {
182 match value {
183 1 => LifecycleState::Running,
184 2 => LifecycleState::Draining,
185 3 => LifecycleState::Stopped,
186 _ => LifecycleState::Starting,
187 }
188}
189
190fn decode_health(value: u64) -> Health {
191 match value {
192 1 => Health::Degraded,
193 2 => Health::Failed,
194 _ => Health::Healthy,
195 }
196}
197
198#[cfg(test)]
199mod tests {
200 use std::{mem::size_of, sync::Arc, thread};
201
202 use super::*;
203
204 #[test]
205 fn concurrent_updates_are_non_blocking_and_exact() {
206 let metrics = Arc::new(Metrics::default());
207 let mut workers = Vec::new();
208 for _ in 0..8 {
209 let metrics = metrics.clone();
210 workers.push(thread::spawn(move || {
211 for _ in 0..10_000 {
212 metrics.stage_finished(Stage::Response, StageOutcome::Success, 7);
213 }
214 }));
215 }
216 for worker in workers {
217 worker.join().unwrap();
218 }
219 let snapshot = metrics.snapshot();
220 assert_eq!(snapshot.requests(StageOutcome::Success), 80_000);
221 assert_eq!(snapshot.stage_latency(Stage::Response)[2], 80_000);
222 }
223
224 #[test]
225 fn storage_and_snapshot_have_compile_time_fixed_size() {
226 assert!(size_of::<Metrics>() <= 1_024);
227 assert!(size_of::<MetricsSnapshot>() <= 1_024);
228 assert_eq!(LATENCY_BUCKET_UPPER_MS.last(), Some(&u64::MAX));
229 }
230}