Skip to main content

macp_runtime/
metrics.rs

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
41/// Maximum number of distinct mode names tracked in metrics.
42/// Beyond this limit, metrics are aggregated into an "_overflow" bucket.
43const MAX_MODE_CARDINALITY: usize = 1000;
44const OVERFLOW_MODE: &str = "_overflow";
45
46pub struct RuntimeMetrics {
47    per_mode: RwLock<HashMap<String, Arc<ModeMetrics>>>,
48    /// Replay/snapshot divergences observed during startup recovery (D7).
49    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 at cardinality limit, aggregate into overflow bucket
130        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    /// Render this snapshot as Prometheus text-format lines for one mode.
199    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    /// Assert one Prometheus text-format sample line: either
241    /// `name{labels} value` or `name value`, with a numeric value and a
242    /// well-formed metric name.
243    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        // A metrics registry with no recorded activity snapshots to nothing.
378        let m = RuntimeMetrics::new();
379        assert!(m.snapshot().is_empty());
380
381        // An all-zero snapshot renders every sample with value 0.
382        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        // Beyond the cardinality limit, new modes fold into "_overflow".
414        m.record_message_accepted("mode-beyond-limit-a");
415        m.record_message_accepted("mode-beyond-limit-b");
416        // Existing modes keep their own bucket.
417        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}