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}