confium_coordinator/coordinator/
request_log.rs1use chrono::{DateTime, Utc};
4use serde::{Deserialize, Serialize};
5use std::collections::VecDeque;
6use std::sync::Mutex;
7
8#[derive(Debug, Clone, Serialize, Deserialize)]
10pub struct RequestLogEntry {
11 pub timestamp: DateTime<Utc>,
12 pub request_type: String,
13 pub session_id: Option<String>,
14 pub signer_id: Option<String>,
15 pub duration_us: u64,
16 pub success: bool,
17 pub error: Option<String>,
18 pub bytes_in: u64,
19 pub bytes_out: u64,
20}
21
22pub struct RequestLog {
24 entries: Mutex<VecDeque<RequestLogEntry>>,
25 max_entries: usize,
26}
27
28impl RequestLog {
29 pub fn new(max_entries: usize) -> Self {
30 Self {
31 entries: Mutex::new(VecDeque::with_capacity(max_entries)),
32 max_entries,
33 }
34 }
35
36 pub fn log(&self, entry: RequestLogEntry) {
37 let mut entries = self.entries.lock().unwrap();
38 if entries.len() >= self.max_entries {
39 entries.pop_front();
40 }
41 entries.push_back(entry);
42 }
43
44 pub fn entries(&self) -> Vec<RequestLogEntry> {
45 self.entries.lock().unwrap().iter().cloned().collect()
46 }
47
48 pub fn count(&self) -> usize {
49 self.entries.lock().unwrap().len()
50 }
51
52 pub fn error_count(&self) -> usize {
53 self.entries
54 .lock()
55 .unwrap()
56 .iter()
57 .filter(|e| !e.success)
58 .count()
59 }
60
61 pub fn avg_duration_us(&self) -> f64 {
62 let entries = self.entries.lock().unwrap();
63 if entries.is_empty() {
64 return 0.0;
65 }
66 entries.iter().map(|e| e.duration_us as f64).sum::<f64>() / entries.len() as f64
67 }
68
69 pub fn clear(&self) {
70 self.entries.lock().unwrap().clear();
71 }
72}
73
74pub struct RequestTimer {
76 request_type: String,
77 start: std::time::Instant,
78}
79
80impl RequestTimer {
81 pub fn start(request_type: &str) -> Self {
82 Self {
83 request_type: request_type.into(),
84 start: std::time::Instant::now(),
85 }
86 }
87
88 pub fn finish(self, log: &RequestLog, success: bool, error: Option<String>) {
89 log.log(RequestLogEntry {
90 timestamp: Utc::now(),
91 request_type: self.request_type,
92 session_id: None,
93 signer_id: None,
94 duration_us: self.start.elapsed().as_micros() as u64,
95 success,
96 error,
97 bytes_in: 0,
98 bytes_out: 0,
99 });
100 }
101}
102
103#[cfg(test)]
104mod tests {
105 use super::*;
106
107 #[test]
108 fn empty_log() {
109 let log = RequestLog::new(100);
110 assert_eq!(log.count(), 0);
111 assert_eq!(log.avg_duration_us(), 0.0);
112 }
113
114 #[test]
115 fn log_and_retrieve() {
116 let log = RequestLog::new(100);
117 log.log(RequestLogEntry {
118 timestamp: Utc::now(),
119 request_type: "test".into(),
120 session_id: None,
121 signer_id: None,
122 duration_us: 100,
123 success: true,
124 error: None,
125 bytes_in: 10,
126 bytes_out: 20,
127 });
128 assert_eq!(log.count(), 1);
129 }
130
131 #[test]
132 fn bounded_capacity() {
133 let log = RequestLog::new(3);
134 for i in 0..5 {
135 log.log(RequestLogEntry {
136 timestamp: Utc::now(),
137 request_type: format!("req-{i}"),
138 session_id: None,
139 signer_id: None,
140 duration_us: i * 100,
141 success: true,
142 error: None,
143 bytes_in: 0,
144 bytes_out: 0,
145 });
146 }
147 assert_eq!(log.count(), 3);
148 }
149
150 #[test]
151 fn error_count() {
152 let log = RequestLog::new(100);
153 log.log(RequestLogEntry {
154 timestamp: Utc::now(),
155 request_type: "a".into(),
156 session_id: None,
157 signer_id: None,
158 duration_us: 1,
159 success: true,
160 error: None,
161 bytes_in: 0,
162 bytes_out: 0,
163 });
164 log.log(RequestLogEntry {
165 timestamp: Utc::now(),
166 request_type: "b".into(),
167 session_id: None,
168 signer_id: None,
169 duration_us: 1,
170 success: false,
171 error: Some("fail".into()),
172 bytes_in: 0,
173 bytes_out: 0,
174 });
175 assert_eq!(log.error_count(), 1);
176 }
177
178 #[test]
179 fn avg_duration() {
180 let log = RequestLog::new(100);
181 for d in [100u64, 200, 300] {
182 log.log(RequestLogEntry {
183 timestamp: Utc::now(),
184 request_type: "x".into(),
185 session_id: None,
186 signer_id: None,
187 duration_us: d,
188 success: true,
189 error: None,
190 bytes_in: 0,
191 bytes_out: 0,
192 });
193 }
194 assert!((log.avg_duration_us() - 200.0).abs() < 0.1);
195 }
196
197 #[test]
198 fn timer_finishes() {
199 let log = RequestLog::new(100);
200 let timer = RequestTimer::start("op");
201 std::thread::sleep(std::time::Duration::from_micros(100));
202 timer.finish(&log, true, None);
203 assert_eq!(log.count(), 1);
204 assert!(log.entries()[0].duration_us >= 50);
205 }
206
207 #[test]
208 fn clear_empties() {
209 let log = RequestLog::new(100);
210 log.log(RequestLogEntry {
211 timestamp: Utc::now(),
212 request_type: "x".into(),
213 session_id: None,
214 signer_id: None,
215 duration_us: 1,
216 success: true,
217 error: None,
218 bytes_in: 0,
219 bytes_out: 0,
220 });
221 log.clear();
222 assert_eq!(log.count(), 0);
223 }
224
225 #[test]
226 fn entry_serializes() {
227 let entry = RequestLogEntry {
228 timestamp: Utc::now(),
229 request_type: "test".into(),
230 session_id: Some("s1".into()),
231 signer_id: Some("a".into()),
232 duration_us: 42,
233 success: true,
234 error: None,
235 bytes_in: 10,
236 bytes_out: 20,
237 };
238 let json = serde_json::to_string(&entry).unwrap();
239 assert!(json.contains("request_type"));
240 assert!(json.contains("s1"));
241 }
242}