Skip to main content

crawlkit_engine/
observability.rs

1use std::sync::atomic::{AtomicU64, Ordering};
2use std::sync::Arc;
3use std::time::Duration;
4
5use serde::{Deserialize, Serialize};
6
7/// Observability metrics for crawl operations.
8///
9/// Thread-safe metrics collection using atomic operations.
10/// Designed for zero-allocation in hot paths. All counters use
11/// `Relaxed` ordering for maximum throughput.
12///
13/// # Examples
14///
15/// ```rust
16/// use crawlkit_engine::Metrics;
17/// use std::time::Duration;
18///
19/// let metrics = Metrics::new();
20/// metrics.record_page_success(1024, 100, 50, 10, 3);
21/// assert_eq!(metrics.pages_crawled.load(std::sync::atomic::Ordering::Relaxed), 1);
22/// ```
23pub struct Metrics {
24    /// Total pages crawled.
25    pub pages_crawled: AtomicU64,
26    /// Total pages failed.
27    pub pages_failed: AtomicU64,
28    /// Total findings generated.
29    pub findings_generated: AtomicU64,
30    /// Total bytes fetched.
31    pub bytes_fetched: AtomicU64,
32    /// Total fetch time (microseconds).
33    pub fetch_time_us: AtomicU64,
34    /// Total analysis time (microseconds).
35    pub analysis_time_us: AtomicU64,
36    /// Total storage write time (microseconds).
37    pub storage_time_us: AtomicU64,
38    /// Active connections.
39    pub active_connections: AtomicU64,
40    /// Total circuit breaker trips (transitions to Open state).
41    pub circuit_breaker_trips: AtomicU64,
42    /// Total resource limit hits.
43    pub resource_limit_hits: AtomicU64,
44    /// Pages skipped due to circuit breaker being open.
45    pub pages_skipped_circuit_breaker: AtomicU64,
46}
47
48impl Metrics {
49    /// Create new metrics instance.
50    #[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    /// Record a successful page crawl.
68    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    /// Record a failed page crawl.
88    pub fn record_page_failure(&self) {
89        self.pages_failed.fetch_add(1, Ordering::Relaxed);
90    }
91
92    /// Record a circuit breaker trip (transition to Open state).
93    pub fn record_circuit_breaker_trip(&self) {
94        self.circuit_breaker_trips.fetch_add(1, Ordering::Relaxed);
95    }
96
97    /// Record a resource limit being hit.
98    pub fn record_resource_limit_hit(&self) {
99        self.resource_limit_hits.fetch_add(1, Ordering::Relaxed);
100    }
101
102    /// Record a page skipped because the circuit breaker was open.
103    pub fn record_page_skipped_circuit_breaker(&self) {
104        self.pages_skipped_circuit_breaker
105            .fetch_add(1, Ordering::Relaxed);
106    }
107
108    /// Increment active connections.
109    pub fn inc_connections(&self) {
110        self.active_connections.fetch_add(1, Ordering::Relaxed);
111    }
112
113    /// Decrement active connections.
114    /// Uses `fetch_update` to prevent underflow below zero.
115    pub fn dec_connections(&self) {
116        // `fetch_update` returns `Err` only if the closure returns `None` on every
117        // attempt, which in our case means the count was already zero. This is
118        // the desired behavior — we simply ignore the "already at zero" case.
119        let _ =
120            self.active_connections
121                .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |current| {
122                    if current > 0 {
123                        Some(current - 1)
124                    } else {
125                        // Already at zero; do not underflow
126                        None
127                    }
128                });
129    }
130
131    /// Get pages per second.
132    #[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    /// Get average fetch time in milliseconds.
143    #[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    /// Get average analysis time in milliseconds.
154    #[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    /// Get total throughput (bytes per second).
165    #[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    /// Get snapshot of all metrics.
176    #[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    /// Reset all metrics.
196    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/// Snapshot of metrics at a point in time.
219///
220/// Created by [`Metrics::snapshot`] for reporting and API responses.
221/// All fields are plain `u64` values for easy serialization.
222#[derive(Debug, Clone, Serialize, Deserialize)]
223pub struct MetricsSnapshot {
224    /// Total pages successfully crawled.
225    pub pages_crawled: u64,
226    /// Total pages that failed to fetch.
227    pub pages_failed: u64,
228    /// Total findings generated by analyzers.
229    pub findings_generated: u64,
230    /// Total bytes fetched.
231    pub bytes_fetched: u64,
232    /// Total fetch time in microseconds.
233    pub fetch_time_us: u64,
234    /// Total analysis time in microseconds.
235    pub analysis_time_us: u64,
236    /// Total storage write time in microseconds.
237    pub storage_time_us: u64,
238    /// Current number of active HTTP connections.
239    pub active_connections: u64,
240    /// Total circuit breaker trips.
241    pub circuit_breaker_trips: u64,
242    /// Total resource limit hits.
243    pub resource_limit_hits: u64,
244    /// Pages skipped because circuit breaker was open.
245    pub pages_skipped_circuit_breaker: u64,
246}
247
248/// Shared metrics for concurrent access.
249///
250/// Wraps [`Metrics`] in an `Arc` for sharing across tasks.
251/// Clone is cheap (reference count increment).
252#[derive(Clone)]
253pub struct SharedMetrics {
254    inner: Arc<Metrics>,
255}
256
257impl SharedMetrics {
258    /// Create shared metrics.
259    #[must_use]
260    pub fn new() -> Self {
261        Self {
262            inner: Arc::new(Metrics::new()),
263        }
264    }
265
266    /// Get inner metrics reference.
267    #[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// ---------------------------------------------------------------------------
280// Tests
281// ---------------------------------------------------------------------------
282
283#[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); // 1ms fetch, 0.5ms analysis
323        metrics.record_page_success(1024, 2000, 1000, 200, 5); // 2ms fetch, 1ms analysis
324
325        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}