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}