1use std::collections::HashMap;
2use std::sync::atomic::{AtomicU64, Ordering};
3use std::sync::{Arc, RwLock};
4
5pub struct ModeMetrics {
6 pub messages_accepted: AtomicU64,
7 pub messages_rejected: AtomicU64,
8 pub sessions_started: AtomicU64,
9 pub sessions_resolved: AtomicU64,
10 pub sessions_expired: AtomicU64,
11 pub sessions_cancelled: AtomicU64,
12 pub sessions_suspended: AtomicU64,
13 pub sessions_resumed: AtomicU64,
14 pub commitments_accepted: AtomicU64,
15 pub commitments_rejected: AtomicU64,
16}
17
18impl ModeMetrics {
19 pub fn new() -> Self {
20 Self {
21 messages_accepted: AtomicU64::new(0),
22 messages_rejected: AtomicU64::new(0),
23 sessions_started: AtomicU64::new(0),
24 sessions_resolved: AtomicU64::new(0),
25 sessions_expired: AtomicU64::new(0),
26 sessions_cancelled: AtomicU64::new(0),
27 sessions_suspended: AtomicU64::new(0),
28 sessions_resumed: AtomicU64::new(0),
29 commitments_accepted: AtomicU64::new(0),
30 commitments_rejected: AtomicU64::new(0),
31 }
32 }
33}
34
35impl Default for ModeMetrics {
36 fn default() -> Self {
37 Self::new()
38 }
39}
40
41const MAX_MODE_CARDINALITY: usize = 1000;
44const OVERFLOW_MODE: &str = "_overflow";
45
46pub struct RuntimeMetrics {
47 per_mode: RwLock<HashMap<String, Arc<ModeMetrics>>>,
48 replay_mismatches: AtomicU64,
50}
51
52impl RuntimeMetrics {
53 pub fn new() -> Self {
54 Self {
55 per_mode: RwLock::new(HashMap::new()),
56 replay_mismatches: AtomicU64::new(0),
57 }
58 }
59
60 pub fn record_session_start(&self, mode: &str) {
61 self.get_or_create(mode)
62 .sessions_started
63 .fetch_add(1, Ordering::Relaxed);
64 }
65
66 pub fn record_message_accepted(&self, mode: &str) {
67 self.get_or_create(mode)
68 .messages_accepted
69 .fetch_add(1, Ordering::Relaxed);
70 }
71
72 pub fn record_message_rejected(&self, mode: &str) {
73 self.get_or_create(mode)
74 .messages_rejected
75 .fetch_add(1, Ordering::Relaxed);
76 }
77
78 pub fn record_session_resolved(&self, mode: &str) {
79 self.get_or_create(mode)
80 .sessions_resolved
81 .fetch_add(1, Ordering::Relaxed);
82 }
83
84 pub fn record_session_expired(&self, mode: &str) {
85 self.get_or_create(mode)
86 .sessions_expired
87 .fetch_add(1, Ordering::Relaxed);
88 }
89
90 pub fn record_session_cancelled(&self, mode: &str) {
91 self.get_or_create(mode)
92 .sessions_cancelled
93 .fetch_add(1, Ordering::Relaxed);
94 }
95
96 pub fn record_session_suspended(&self, mode: &str) {
97 self.get_or_create(mode)
98 .sessions_suspended
99 .fetch_add(1, Ordering::Relaxed);
100 }
101
102 pub fn record_session_resumed(&self, mode: &str) {
103 self.get_or_create(mode)
104 .sessions_resumed
105 .fetch_add(1, Ordering::Relaxed);
106 }
107
108 pub fn record_commitment_accepted(&self, mode: &str) {
109 self.get_or_create(mode)
110 .commitments_accepted
111 .fetch_add(1, Ordering::Relaxed);
112 }
113
114 pub fn record_commitment_rejected(&self, mode: &str) {
115 self.get_or_create(mode)
116 .commitments_rejected
117 .fetch_add(1, Ordering::Relaxed);
118 }
119
120 fn get_or_create(&self, mode: &str) -> Arc<ModeMetrics> {
121 {
122 let guard = self.per_mode.read().unwrap_or_else(|e| e.into_inner());
123 if let Some(metrics) = guard.get(mode) {
124 return Arc::clone(metrics);
125 }
126 }
127
128 let mut guard = self.per_mode.write().unwrap_or_else(|e| e.into_inner());
129 if guard.len() >= MAX_MODE_CARDINALITY && !guard.contains_key(mode) {
131 return Arc::clone(
132 guard
133 .entry(OVERFLOW_MODE.to_string())
134 .or_insert_with(|| Arc::new(ModeMetrics::default())),
135 );
136 }
137 Arc::clone(
138 guard
139 .entry(mode.to_string())
140 .or_insert_with(|| Arc::new(ModeMetrics::default())),
141 )
142 }
143
144 pub fn record_replay_mismatch(&self, count: u64) {
145 self.replay_mismatches.fetch_add(count, Ordering::Relaxed);
146 }
147
148 pub fn replay_mismatches(&self) -> u64 {
149 self.replay_mismatches.load(Ordering::Relaxed)
150 }
151
152 pub fn snapshot(&self) -> Vec<(String, MetricsSnapshot)> {
153 let guard = self.per_mode.read().unwrap_or_else(|e| e.into_inner());
154 guard
155 .iter()
156 .map(|(mode, m)| {
157 (
158 mode.clone(),
159 MetricsSnapshot {
160 messages_accepted: m.messages_accepted.load(Ordering::Relaxed),
161 messages_rejected: m.messages_rejected.load(Ordering::Relaxed),
162 sessions_started: m.sessions_started.load(Ordering::Relaxed),
163 sessions_resolved: m.sessions_resolved.load(Ordering::Relaxed),
164 sessions_expired: m.sessions_expired.load(Ordering::Relaxed),
165 sessions_cancelled: m.sessions_cancelled.load(Ordering::Relaxed),
166 commitments_accepted: m.commitments_accepted.load(Ordering::Relaxed),
167 commitments_rejected: m.commitments_rejected.load(Ordering::Relaxed),
168 sessions_suspended: m.sessions_suspended.load(Ordering::Relaxed),
169 sessions_resumed: m.sessions_resumed.load(Ordering::Relaxed),
170 },
171 )
172 })
173 .collect()
174 }
175}
176
177impl Default for RuntimeMetrics {
178 fn default() -> Self {
179 Self::new()
180 }
181}
182
183#[derive(Debug)]
184pub struct MetricsSnapshot {
185 pub messages_accepted: u64,
186 pub messages_rejected: u64,
187 pub sessions_started: u64,
188 pub sessions_resolved: u64,
189 pub sessions_expired: u64,
190 pub sessions_cancelled: u64,
191 pub commitments_accepted: u64,
192 pub commitments_rejected: u64,
193 pub sessions_suspended: u64,
194 pub sessions_resumed: u64,
195}
196
197impl MetricsSnapshot {
198 pub fn prometheus_lines(&self, mode: &str, out: &mut String) {
200 use std::fmt::Write;
201 let pairs: [(&str, u64); 10] = [
202 ("macp_messages_accepted_total", self.messages_accepted),
203 ("macp_messages_rejected_total", self.messages_rejected),
204 ("macp_sessions_started_total", self.sessions_started),
205 ("macp_sessions_resolved_total", self.sessions_resolved),
206 ("macp_sessions_expired_total", self.sessions_expired),
207 ("macp_sessions_cancelled_total", self.sessions_cancelled),
208 ("macp_sessions_suspended_total", self.sessions_suspended),
209 ("macp_sessions_resumed_total", self.sessions_resumed),
210 ("macp_commitments_accepted_total", self.commitments_accepted),
211 ("macp_commitments_rejected_total", self.commitments_rejected),
212 ];
213 for (name, value) in pairs {
214 let _ = writeln!(out, "{name}{{mode=\"{mode}\"}} {value}");
215 }
216 }
217}
218
219#[cfg(test)]
220mod tests {
221 use super::*;
222
223 const EXPECTED_METRIC_NAMES: [&str; 10] = [
224 "macp_messages_accepted_total",
225 "macp_messages_rejected_total",
226 "macp_sessions_started_total",
227 "macp_sessions_resolved_total",
228 "macp_sessions_expired_total",
229 "macp_sessions_cancelled_total",
230 "macp_sessions_suspended_total",
231 "macp_sessions_resumed_total",
232 "macp_commitments_accepted_total",
233 "macp_commitments_rejected_total",
234 ];
235
236 fn snapshot_map(m: &RuntimeMetrics) -> HashMap<String, MetricsSnapshot> {
237 m.snapshot().into_iter().collect()
238 }
239
240 fn assert_valid_prometheus_line(line: &str) {
244 let (series, value) = line
245 .rsplit_once(' ')
246 .unwrap_or_else(|| panic!("sample line must contain a value: {line:?}"));
247 assert!(
248 value.parse::<u64>().is_ok(),
249 "sample value must be numeric: {line:?}"
250 );
251 let name = match series.split_once('{') {
252 Some((name, labels)) => {
253 assert!(labels.ends_with('}'), "label set must be closed: {line:?}");
254 name
255 }
256 None => series,
257 };
258 assert!(!name.is_empty(), "metric name must be non-empty: {line:?}");
259 let mut chars = name.chars();
260 let first = chars.next().unwrap();
261 assert!(
262 first.is_ascii_alphabetic() || first == '_' || first == ':',
263 "metric name must not start with a digit: {line:?}"
264 );
265 assert!(
266 name.chars()
267 .all(|c| c.is_ascii_alphanumeric() || c == '_' || c == ':'),
268 "metric name has invalid characters: {line:?}"
269 );
270 }
271
272 #[test]
273 fn counters_record_per_mode_increments() {
274 let m = RuntimeMetrics::new();
275
276 m.record_session_start("macp.mode.decision.v1");
277 m.record_message_accepted("macp.mode.decision.v1");
278 m.record_message_accepted("macp.mode.decision.v1");
279 m.record_message_rejected("macp.mode.decision.v1");
280 m.record_session_resolved("macp.mode.decision.v1");
281 m.record_commitment_accepted("macp.mode.decision.v1");
282
283 m.record_session_start("macp.mode.quorum.v1");
284 m.record_session_expired("macp.mode.quorum.v1");
285 m.record_session_cancelled("macp.mode.quorum.v1");
286 m.record_session_suspended("macp.mode.quorum.v1");
287 m.record_session_resumed("macp.mode.quorum.v1");
288 m.record_commitment_rejected("macp.mode.quorum.v1");
289
290 let snap = snapshot_map(&m);
291 assert_eq!(snap.len(), 2, "one entry per distinct mode");
292
293 let d = &snap["macp.mode.decision.v1"];
294 assert_eq!(d.sessions_started, 1);
295 assert_eq!(d.messages_accepted, 2);
296 assert_eq!(d.messages_rejected, 1);
297 assert_eq!(d.sessions_resolved, 1);
298 assert_eq!(d.commitments_accepted, 1);
299 assert_eq!(d.sessions_expired, 0);
300 assert_eq!(d.sessions_cancelled, 0);
301 assert_eq!(d.sessions_suspended, 0);
302 assert_eq!(d.sessions_resumed, 0);
303 assert_eq!(d.commitments_rejected, 0);
304
305 let q = &snap["macp.mode.quorum.v1"];
306 assert_eq!(q.sessions_started, 1);
307 assert_eq!(q.sessions_expired, 1);
308 assert_eq!(q.sessions_cancelled, 1);
309 assert_eq!(q.sessions_suspended, 1);
310 assert_eq!(q.sessions_resumed, 1);
311 assert_eq!(q.commitments_rejected, 1);
312 assert_eq!(q.messages_accepted, 0);
313 assert_eq!(q.messages_rejected, 0);
314 assert_eq!(q.sessions_resolved, 0);
315 assert_eq!(q.commitments_accepted, 0);
316 }
317
318 #[test]
319 fn replay_mismatch_counter_accumulates() {
320 let m = RuntimeMetrics::new();
321 assert_eq!(m.replay_mismatches(), 0);
322 m.record_replay_mismatch(2);
323 m.record_replay_mismatch(3);
324 assert_eq!(m.replay_mismatches(), 5);
325 }
326
327 #[test]
328 fn prometheus_rendering_reflects_recorded_counts() {
329 let m = RuntimeMetrics::new();
330 m.record_message_accepted("macp.mode.decision.v1");
331 m.record_message_accepted("macp.mode.decision.v1");
332 m.record_message_rejected("macp.mode.decision.v1");
333 m.record_session_start("macp.mode.decision.v1");
334
335 let mut body = String::new();
336 for (mode, snap) in m.snapshot() {
337 snap.prometheus_lines(&mode, &mut body);
338 }
339
340 for name in EXPECTED_METRIC_NAMES {
341 assert!(
342 body.contains(name),
343 "rendered output must include {name}:\n{body}"
344 );
345 }
346 assert!(body.contains("macp_messages_accepted_total{mode=\"macp.mode.decision.v1\"} 2"));
347 assert!(body.contains("macp_messages_rejected_total{mode=\"macp.mode.decision.v1\"} 1"));
348 assert!(body.contains("macp_sessions_started_total{mode=\"macp.mode.decision.v1\"} 1"));
349 assert!(body.contains("macp_sessions_resolved_total{mode=\"macp.mode.decision.v1\"} 0"));
350 }
351
352 #[test]
353 fn prometheus_lines_are_valid_exposition_format() {
354 let m = RuntimeMetrics::new();
355 m.record_session_start("macp.mode.task.v1");
356 m.record_commitment_accepted("macp.mode.task.v1");
357 m.record_session_start("ext.multi_round.v1");
358
359 let mut body = String::new();
360 for (mode, snap) in m.snapshot() {
361 snap.prometheus_lines(&mode, &mut body);
362 }
363
364 let lines: Vec<&str> = body.lines().collect();
365 assert_eq!(
366 lines.len(),
367 2 * EXPECTED_METRIC_NAMES.len(),
368 "10 samples per mode"
369 );
370 for line in lines {
371 assert_valid_prometheus_line(line);
372 }
373 }
374
375 #[test]
376 fn zero_state_rendering_does_not_panic() {
377 let m = RuntimeMetrics::new();
379 assert!(m.snapshot().is_empty());
380
381 let zero = MetricsSnapshot {
383 messages_accepted: 0,
384 messages_rejected: 0,
385 sessions_started: 0,
386 sessions_resolved: 0,
387 sessions_expired: 0,
388 sessions_cancelled: 0,
389 commitments_accepted: 0,
390 commitments_rejected: 0,
391 sessions_suspended: 0,
392 sessions_resumed: 0,
393 };
394 let mut body = String::new();
395 zero.prometheus_lines("macp.mode.decision.v1", &mut body);
396 let lines: Vec<&str> = body.lines().collect();
397 assert_eq!(lines.len(), EXPECTED_METRIC_NAMES.len());
398 for line in lines {
399 assert!(
400 line.ends_with(" 0"),
401 "zero-state sample must be 0: {line:?}"
402 );
403 assert_valid_prometheus_line(line);
404 }
405 }
406
407 #[test]
408 fn mode_cardinality_overflow_aggregates_into_overflow_bucket() {
409 let m = RuntimeMetrics::new();
410 for i in 0..MAX_MODE_CARDINALITY {
411 m.record_message_accepted(&format!("mode-{i}"));
412 }
413 m.record_message_accepted("mode-beyond-limit-a");
415 m.record_message_accepted("mode-beyond-limit-b");
416 m.record_message_accepted("mode-0");
418
419 let snap = snapshot_map(&m);
420 assert_eq!(snap.len(), MAX_MODE_CARDINALITY + 1);
421 assert!(!snap.contains_key("mode-beyond-limit-a"));
422 assert_eq!(snap[OVERFLOW_MODE].messages_accepted, 2);
423 assert_eq!(snap["mode-0"].messages_accepted, 2);
424 }
425}