1use std::sync::atomic::{AtomicU64, Ordering};
2use std::sync::Arc;
3use std::time::Duration;
4
5use serde::{Deserialize, Serialize};
6
7pub struct Metrics {
24 pub pages_crawled: AtomicU64,
26 pub pages_failed: AtomicU64,
28 pub findings_generated: AtomicU64,
30 pub bytes_fetched: AtomicU64,
32 pub fetch_time_us: AtomicU64,
34 pub analysis_time_us: AtomicU64,
36 pub storage_time_us: AtomicU64,
38 pub active_connections: AtomicU64,
40 pub circuit_breaker_trips: AtomicU64,
42 pub resource_limit_hits: AtomicU64,
44 pub pages_skipped_circuit_breaker: AtomicU64,
46}
47
48impl Metrics {
49 #[must_use]
51 pub fn new() -> Self {
52 Self {
53 pages_crawled: AtomicU64::new(0),
54 pages_failed: AtomicU64::new(0),
55 findings_generated: AtomicU64::new(0),
56 bytes_fetched: AtomicU64::new(0),
57 fetch_time_us: AtomicU64::new(0),
58 analysis_time_us: AtomicU64::new(0),
59 storage_time_us: AtomicU64::new(0),
60 active_connections: AtomicU64::new(0),
61 circuit_breaker_trips: AtomicU64::new(0),
62 resource_limit_hits: AtomicU64::new(0),
63 pages_skipped_circuit_breaker: AtomicU64::new(0),
64 }
65 }
66
67 pub fn record_page_success(
69 &self,
70 bytes: u64,
71 fetch_us: u64,
72 analysis_us: u64,
73 storage_us: u64,
74 findings: u64,
75 ) {
76 self.pages_crawled.fetch_add(1, Ordering::Relaxed);
77 self.bytes_fetched.fetch_add(bytes, Ordering::Relaxed);
78 self.fetch_time_us.fetch_add(fetch_us, Ordering::Relaxed);
79 self.analysis_time_us
80 .fetch_add(analysis_us, Ordering::Relaxed);
81 self.storage_time_us
82 .fetch_add(storage_us, Ordering::Relaxed);
83 self.findings_generated
84 .fetch_add(findings, Ordering::Relaxed);
85 }
86
87 pub fn record_page_failure(&self) {
89 self.pages_failed.fetch_add(1, Ordering::Relaxed);
90 }
91
92 pub fn record_circuit_breaker_trip(&self) {
94 self.circuit_breaker_trips.fetch_add(1, Ordering::Relaxed);
95 }
96
97 pub fn record_resource_limit_hit(&self) {
99 self.resource_limit_hits.fetch_add(1, Ordering::Relaxed);
100 }
101
102 pub fn record_page_skipped_circuit_breaker(&self) {
104 self.pages_skipped_circuit_breaker
105 .fetch_add(1, Ordering::Relaxed);
106 }
107
108 pub fn inc_connections(&self) {
110 self.active_connections.fetch_add(1, Ordering::Relaxed);
111 }
112
113 pub fn dec_connections(&self) {
116 let _ =
120 self.active_connections
121 .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |current| {
122 if current > 0 {
123 Some(current - 1)
124 } else {
125 None
127 }
128 });
129 }
130
131 #[must_use]
133 pub fn pages_per_second(&self, elapsed: Duration) -> f64 {
134 let pages = self.pages_crawled.load(Ordering::Relaxed) as f64;
135 let secs = elapsed.as_secs_f64();
136 if secs <= 0.0 {
137 return 0.0;
138 }
139 pages / secs
140 }
141
142 #[must_use]
144 pub fn avg_fetch_time_ms(&self) -> f64 {
145 let pages = self.pages_crawled.load(Ordering::Relaxed);
146 if pages == 0 {
147 return 0.0;
148 }
149 let total_us = self.fetch_time_us.load(Ordering::Relaxed);
150 (total_us as f64 / pages as f64) / 1000.0
151 }
152
153 #[must_use]
155 pub fn avg_analysis_time_ms(&self) -> f64 {
156 let pages = self.pages_crawled.load(Ordering::Relaxed);
157 if pages == 0 {
158 return 0.0;
159 }
160 let total_us = self.analysis_time_us.load(Ordering::Relaxed);
161 (total_us as f64 / pages as f64) / 1000.0
162 }
163
164 #[must_use]
166 pub fn throughput_bps(&self, elapsed: Duration) -> f64 {
167 let bytes = self.bytes_fetched.load(Ordering::Relaxed) as f64;
168 let secs = elapsed.as_secs_f64();
169 if secs <= 0.0 {
170 return 0.0;
171 }
172 bytes / secs
173 }
174
175 #[must_use]
177 pub fn snapshot(&self) -> MetricsSnapshot {
178 MetricsSnapshot {
179 pages_crawled: self.pages_crawled.load(Ordering::Relaxed),
180 pages_failed: self.pages_failed.load(Ordering::Relaxed),
181 findings_generated: self.findings_generated.load(Ordering::Relaxed),
182 bytes_fetched: self.bytes_fetched.load(Ordering::Relaxed),
183 fetch_time_us: self.fetch_time_us.load(Ordering::Relaxed),
184 analysis_time_us: self.analysis_time_us.load(Ordering::Relaxed),
185 storage_time_us: self.storage_time_us.load(Ordering::Relaxed),
186 active_connections: self.active_connections.load(Ordering::Relaxed),
187 circuit_breaker_trips: self.circuit_breaker_trips.load(Ordering::Relaxed),
188 resource_limit_hits: self.resource_limit_hits.load(Ordering::Relaxed),
189 pages_skipped_circuit_breaker: self
190 .pages_skipped_circuit_breaker
191 .load(Ordering::Relaxed),
192 }
193 }
194
195 pub fn reset(&self) {
197 self.pages_crawled.store(0, Ordering::Relaxed);
198 self.pages_failed.store(0, Ordering::Relaxed);
199 self.findings_generated.store(0, Ordering::Relaxed);
200 self.bytes_fetched.store(0, Ordering::Relaxed);
201 self.fetch_time_us.store(0, Ordering::Relaxed);
202 self.analysis_time_us.store(0, Ordering::Relaxed);
203 self.storage_time_us.store(0, Ordering::Relaxed);
204 self.active_connections.store(0, Ordering::Relaxed);
205 self.circuit_breaker_trips.store(0, Ordering::Relaxed);
206 self.resource_limit_hits.store(0, Ordering::Relaxed);
207 self.pages_skipped_circuit_breaker
208 .store(0, Ordering::Relaxed);
209 }
210}
211
212impl Default for Metrics {
213 fn default() -> Self {
214 Self::new()
215 }
216}
217
218#[derive(Debug, Clone, Serialize, Deserialize)]
223pub struct MetricsSnapshot {
224 pub pages_crawled: u64,
226 pub pages_failed: u64,
228 pub findings_generated: u64,
230 pub bytes_fetched: u64,
232 pub fetch_time_us: u64,
234 pub analysis_time_us: u64,
236 pub storage_time_us: u64,
238 pub active_connections: u64,
240 pub circuit_breaker_trips: u64,
242 pub resource_limit_hits: u64,
244 pub pages_skipped_circuit_breaker: u64,
246}
247
248#[derive(Clone)]
253pub struct SharedMetrics {
254 inner: Arc<Metrics>,
255}
256
257impl SharedMetrics {
258 #[must_use]
260 pub fn new() -> Self {
261 Self {
262 inner: Arc::new(Metrics::new()),
263 }
264 }
265
266 #[must_use]
268 pub fn inner(&self) -> &Metrics {
269 &self.inner
270 }
271}
272
273impl Default for SharedMetrics {
274 fn default() -> Self {
275 Self::new()
276 }
277}
278
279#[cfg(test)]
284mod tests {
285 use super::*;
286
287 #[test]
288 fn test_metrics_record() {
289 let metrics = Metrics::new();
290 metrics.record_page_success(1024, 100, 50, 10, 3);
291 metrics.record_page_failure();
292
293 assert_eq!(metrics.pages_crawled.load(Ordering::Relaxed), 1);
294 assert_eq!(metrics.pages_failed.load(Ordering::Relaxed), 1);
295 assert_eq!(metrics.bytes_fetched.load(Ordering::Relaxed), 1024);
296 assert_eq!(metrics.findings_generated.load(Ordering::Relaxed), 3);
297 }
298
299 #[test]
300 fn test_metrics_snapshot() {
301 let metrics = Metrics::new();
302 metrics.record_page_success(1024, 100, 50, 10, 3);
303
304 let snapshot = metrics.snapshot();
305 assert_eq!(snapshot.pages_crawled, 1);
306 assert_eq!(snapshot.bytes_fetched, 1024);
307 }
308
309 #[test]
310 fn test_metrics_reset() {
311 let metrics = Metrics::new();
312 metrics.record_page_success(1024, 100, 50, 10, 3);
313 metrics.reset();
314
315 let snapshot = metrics.snapshot();
316 assert_eq!(snapshot.pages_crawled, 0);
317 }
318
319 #[test]
320 fn test_metrics_avg_times() {
321 let metrics = Metrics::new();
322 metrics.record_page_success(1024, 1000, 500, 100, 3); metrics.record_page_success(1024, 2000, 1000, 200, 5); assert!((metrics.avg_fetch_time_ms() - 1.5).abs() < 0.01);
326 assert!((metrics.avg_analysis_time_ms() - 0.75).abs() < 0.01);
327 }
328}