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}
49
50impl RuntimeMetrics {
51 pub fn new() -> Self {
52 Self {
53 per_mode: RwLock::new(HashMap::new()),
54 }
55 }
56
57 pub fn record_session_start(&self, mode: &str) {
58 self.get_or_create(mode)
59 .sessions_started
60 .fetch_add(1, Ordering::Relaxed);
61 }
62
63 pub fn record_message_accepted(&self, mode: &str) {
64 self.get_or_create(mode)
65 .messages_accepted
66 .fetch_add(1, Ordering::Relaxed);
67 }
68
69 pub fn record_message_rejected(&self, mode: &str) {
70 self.get_or_create(mode)
71 .messages_rejected
72 .fetch_add(1, Ordering::Relaxed);
73 }
74
75 pub fn record_session_resolved(&self, mode: &str) {
76 self.get_or_create(mode)
77 .sessions_resolved
78 .fetch_add(1, Ordering::Relaxed);
79 }
80
81 pub fn record_session_expired(&self, mode: &str) {
82 self.get_or_create(mode)
83 .sessions_expired
84 .fetch_add(1, Ordering::Relaxed);
85 }
86
87 pub fn record_session_cancelled(&self, mode: &str) {
88 self.get_or_create(mode)
89 .sessions_cancelled
90 .fetch_add(1, Ordering::Relaxed);
91 }
92
93 pub fn record_session_suspended(&self, mode: &str) {
94 self.get_or_create(mode)
95 .sessions_suspended
96 .fetch_add(1, Ordering::Relaxed);
97 }
98
99 pub fn record_session_resumed(&self, mode: &str) {
100 self.get_or_create(mode)
101 .sessions_resumed
102 .fetch_add(1, Ordering::Relaxed);
103 }
104
105 pub fn record_commitment_accepted(&self, mode: &str) {
106 self.get_or_create(mode)
107 .commitments_accepted
108 .fetch_add(1, Ordering::Relaxed);
109 }
110
111 pub fn record_commitment_rejected(&self, mode: &str) {
112 self.get_or_create(mode)
113 .commitments_rejected
114 .fetch_add(1, Ordering::Relaxed);
115 }
116
117 fn get_or_create(&self, mode: &str) -> Arc<ModeMetrics> {
118 {
119 let guard = self.per_mode.read().unwrap_or_else(|e| e.into_inner());
120 if let Some(metrics) = guard.get(mode) {
121 return Arc::clone(metrics);
122 }
123 }
124
125 let mut guard = self.per_mode.write().unwrap_or_else(|e| e.into_inner());
126 if guard.len() >= MAX_MODE_CARDINALITY && !guard.contains_key(mode) {
128 return Arc::clone(
129 guard
130 .entry(OVERFLOW_MODE.to_string())
131 .or_insert_with(|| Arc::new(ModeMetrics::default())),
132 );
133 }
134 Arc::clone(
135 guard
136 .entry(mode.to_string())
137 .or_insert_with(|| Arc::new(ModeMetrics::default())),
138 )
139 }
140
141 pub fn snapshot(&self) -> Vec<(String, MetricsSnapshot)> {
142 let guard = self.per_mode.read().unwrap_or_else(|e| e.into_inner());
143 guard
144 .iter()
145 .map(|(mode, m)| {
146 (
147 mode.clone(),
148 MetricsSnapshot {
149 messages_accepted: m.messages_accepted.load(Ordering::Relaxed),
150 messages_rejected: m.messages_rejected.load(Ordering::Relaxed),
151 sessions_started: m.sessions_started.load(Ordering::Relaxed),
152 sessions_resolved: m.sessions_resolved.load(Ordering::Relaxed),
153 sessions_expired: m.sessions_expired.load(Ordering::Relaxed),
154 sessions_cancelled: m.sessions_cancelled.load(Ordering::Relaxed),
155 commitments_accepted: m.commitments_accepted.load(Ordering::Relaxed),
156 commitments_rejected: m.commitments_rejected.load(Ordering::Relaxed),
157 },
158 )
159 })
160 .collect()
161 }
162}
163
164impl Default for RuntimeMetrics {
165 fn default() -> Self {
166 Self::new()
167 }
168}
169
170#[derive(Debug)]
171pub struct MetricsSnapshot {
172 pub messages_accepted: u64,
173 pub messages_rejected: u64,
174 pub sessions_started: u64,
175 pub sessions_resolved: u64,
176 pub sessions_expired: u64,
177 pub sessions_cancelled: u64,
178 pub commitments_accepted: u64,
179 pub commitments_rejected: u64,
180}