Skip to main content

velesdb_core/collection/graph/metrics/
mod.rs

1//! Performance metrics for graph operations (EPIC-019 US-006).
2//!
3//! Provides low-overhead, thread-safe metrics for monitoring:
4//! - Operation counters (inserts, deletes, traversals)
5//! - Latency histograms
6//! - Memory usage estimates
7//!
8//! Metrics use atomic operations with relaxed ordering for minimal overhead (~1-5ns per op).
9
10// Reason: Numeric casts in metrics are intentional:
11// - All casts are for histogram bucketing and latency calculations
12// - f64/u64 conversions for computing percentiles and averages
13// - Values bounded by practical limits (bucket counts, durations)
14// - Precision loss acceptable for metrics (approximate by design)
15#![allow(clippy::cast_precision_loss)]
16#![allow(clippy::cast_possible_truncation)]
17
18#[cfg(test)]
19mod tests;
20
21use std::fmt::Write;
22use std::sync::atomic::{AtomicU64, Ordering};
23use std::time::Duration;
24
25/// Latency histogram buckets (milliseconds).
26const BUCKET_BOUNDS_MS: [u64; 9] = [1, 5, 10, 50, 100, 500, 1000, 5000, 10000];
27
28/// Simple latency histogram with fixed buckets.
29///
30/// Buckets: <1ms, <5ms, <10ms, <50ms, <100ms, <500ms, <1s, <5s, <10s, ≥10s
31#[derive(Debug, Default)]
32pub struct LatencyHistogram {
33    /// Bucket counts [<1ms, <5ms, <10ms, <50ms, <100ms, <500ms, <1s, <5s, <10s, ≥10s]
34    buckets: [AtomicU64; 10],
35    /// Sum of all observed durations in nanoseconds
36    sum_ns: AtomicU64,
37    /// Total number of observations
38    count: AtomicU64,
39}
40
41impl LatencyHistogram {
42    /// Creates a new empty histogram.
43    #[must_use]
44    pub fn new() -> Self {
45        Self::default()
46    }
47
48    /// Records a duration observation.
49    ///
50    /// # Note
51    ///
52    /// For extremely large durations (> 584 years), nanoseconds are capped at u64::MAX
53    /// to prevent truncation. This is acceptable since such durations indicate a bug.
54    pub fn observe(&self, duration: Duration) {
55        // Cap at u64::MAX for durations > 584 years (u128 -> u64 truncation protection)
56        let ns_u128 = duration.as_nanos();
57        let ns = if ns_u128 > u128::from(u64::MAX) {
58            u64::MAX
59        } else {
60            ns_u128 as u64
61        };
62        self.sum_ns.fetch_add(ns, Ordering::Relaxed);
63        self.count.fetch_add(1, Ordering::Relaxed);
64
65        // Same protection for milliseconds (though less likely to overflow)
66        let ms_u128 = duration.as_millis();
67        let ms = if ms_u128 > u128::from(u64::MAX) {
68            u64::MAX
69        } else {
70            ms_u128 as u64
71        };
72        let bucket_idx = BUCKET_BOUNDS_MS
73            .iter()
74            .position(|&bound| ms < bound)
75            .unwrap_or(9);
76        self.buckets[bucket_idx].fetch_add(1, Ordering::Relaxed);
77    }
78
79    /// Returns the total count of observations.
80    #[must_use]
81    pub fn count(&self) -> u64 {
82        self.count.load(Ordering::Relaxed)
83    }
84
85    /// Returns the sum of all durations in nanoseconds.
86    #[must_use]
87    pub fn sum_ns(&self) -> u64 {
88        self.sum_ns.load(Ordering::Relaxed)
89    }
90
91    /// Returns the average duration in nanoseconds.
92    #[must_use]
93    pub fn avg_ns(&self) -> f64 {
94        let count = self.count();
95        if count == 0 {
96            0.0
97        } else {
98            self.sum_ns() as f64 / count as f64
99        }
100    }
101
102    /// Returns bucket counts as an array.
103    #[must_use]
104    pub fn bucket_counts(&self) -> [u64; 10] {
105        let mut counts = [0u64; 10];
106        for (i, bucket) in self.buckets.iter().enumerate() {
107            counts[i] = bucket.load(Ordering::Relaxed);
108        }
109        counts
110    }
111
112    /// Resets all counters to zero.
113    pub fn reset(&self) {
114        self.sum_ns.store(0, Ordering::Relaxed);
115        self.count.store(0, Ordering::Relaxed);
116        for bucket in &self.buckets {
117            bucket.store(0, Ordering::Relaxed);
118        }
119    }
120}
121
122/// Graph-specific performance metrics.
123///
124/// Thread-safe counters and histograms for monitoring graph operations.
125///
126/// # Example
127///
128/// ```rust,ignore
129/// use velesdb_core::collection::graph::GraphMetrics;
130/// use std::time::Instant;
131///
132/// let metrics = GraphMetrics::new();
133///
134/// // Record an edge insertion
135/// let start = Instant::now();
136/// // ... perform insertion ...
137/// metrics.record_edge_insert(start.elapsed());
138///
139/// // Get statistics
140/// println!("Total edges inserted: {}", metrics.edge_inserts_total());
141/// println!("Avg insert latency: {:.2}µs", metrics.edge_insert_latency.avg_ns() / 1000.0);
142/// ```
143#[derive(Debug, Default)]
144pub struct GraphMetrics {
145    // Node counters
146    nodes_total: AtomicU64,
147    node_inserts_total: AtomicU64,
148    node_deletes_total: AtomicU64,
149
150    // Edge counters
151    edges_total: AtomicU64,
152    edge_inserts_total: AtomicU64,
153    edge_deletes_total: AtomicU64,
154
155    // Traversal counters
156    traversals_total: AtomicU64,
157    traversal_nodes_visited: AtomicU64,
158
159    // Latency histograms
160    /// Edge insertion latency histogram
161    pub edge_insert_latency: LatencyHistogram,
162    /// Edge deletion latency histogram
163    pub edge_delete_latency: LatencyHistogram,
164    /// Traversal latency histogram
165    pub traversal_latency: LatencyHistogram,
166    /// Query latency histogram
167    pub query_latency: LatencyHistogram,
168}
169
170impl GraphMetrics {
171    /// Creates a new metrics instance with all counters at zero.
172    #[must_use]
173    pub fn new() -> Self {
174        Self::default()
175    }
176
177    // =========================================================================
178    // Node metrics
179    // =========================================================================
180
181    /// Records a node insertion.
182    pub fn record_node_insert(&self) {
183        self.node_inserts_total.fetch_add(1, Ordering::Relaxed);
184        self.nodes_total.fetch_add(1, Ordering::Relaxed);
185    }
186
187    /// Records a node deletion.
188    ///
189    /// Uses saturating subtraction to prevent underflow if called
190    /// more times than `record_node_insert`.
191    pub fn record_node_delete(&self) {
192        self.node_deletes_total.fetch_add(1, Ordering::Relaxed);
193        self.nodes_total
194            .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |x| {
195                Some(x.saturating_sub(1))
196            })
197            .ok();
198    }
199
200    /// Returns total node count.
201    #[must_use]
202    pub fn nodes_total(&self) -> u64 {
203        self.nodes_total.load(Ordering::Relaxed)
204    }
205
206    /// Returns total node insertions.
207    #[must_use]
208    pub fn node_inserts_total(&self) -> u64 {
209        self.node_inserts_total.load(Ordering::Relaxed)
210    }
211
212    // =========================================================================
213    // Edge metrics
214    // =========================================================================
215
216    /// Records an edge insertion with latency.
217    pub fn record_edge_insert(&self, latency: Duration) {
218        self.edge_inserts_total.fetch_add(1, Ordering::Relaxed);
219        self.edges_total.fetch_add(1, Ordering::Relaxed);
220        self.edge_insert_latency.observe(latency);
221    }
222
223    /// Records a batch edge insertion.
224    ///
225    /// Bumps the insert/edge counters by `count` and observes the batch
226    /// `latency` once, avoiding a per-edge `Instant::now()` on the bulk path.
227    pub fn record_edge_inserts_batch(&self, count: u64, latency: Duration) {
228        if count == 0 {
229            return;
230        }
231        self.edge_inserts_total.fetch_add(count, Ordering::Relaxed);
232        self.edges_total.fetch_add(count, Ordering::Relaxed);
233        self.edge_insert_latency.observe(latency);
234    }
235
236    /// Records an edge deletion with latency.
237    ///
238    /// Uses saturating subtraction to prevent underflow.
239    pub fn record_edge_delete(&self, latency: Duration) {
240        self.edge_deletes_total.fetch_add(1, Ordering::Relaxed);
241        self.edges_total
242            .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |x| {
243                Some(x.saturating_sub(1))
244            })
245            .ok();
246        self.edge_delete_latency.observe(latency);
247    }
248
249    /// Returns total edge count.
250    #[must_use]
251    pub fn edges_total(&self) -> u64 {
252        self.edges_total.load(Ordering::Relaxed)
253    }
254
255    /// Returns total edge insertions.
256    #[must_use]
257    pub fn edge_inserts_total(&self) -> u64 {
258        self.edge_inserts_total.load(Ordering::Relaxed)
259    }
260
261    /// Returns total edge deletions.
262    #[must_use]
263    pub fn edge_deletes_total(&self) -> u64 {
264        self.edge_deletes_total.load(Ordering::Relaxed)
265    }
266
267    // =========================================================================
268    // Traversal metrics
269    // =========================================================================
270
271    /// Records a traversal with latency and nodes visited.
272    pub fn record_traversal(&self, latency: Duration, nodes_visited: u64) {
273        self.traversals_total.fetch_add(1, Ordering::Relaxed);
274        self.traversal_nodes_visited
275            .fetch_add(nodes_visited, Ordering::Relaxed);
276        self.traversal_latency.observe(latency);
277    }
278
279    /// Returns total traversal count.
280    #[must_use]
281    pub fn traversals_total(&self) -> u64 {
282        self.traversals_total.load(Ordering::Relaxed)
283    }
284
285    /// Returns total nodes visited across all traversals.
286    #[must_use]
287    pub fn traversal_nodes_visited(&self) -> u64 {
288        self.traversal_nodes_visited.load(Ordering::Relaxed)
289    }
290
291    // =========================================================================
292    // Query metrics
293    // =========================================================================
294
295    /// Records a query latency.
296    pub fn record_query(&self, latency: Duration) {
297        self.query_latency.observe(latency);
298    }
299
300    // =========================================================================
301    // Export
302    // =========================================================================
303
304    /// Exports metrics in Prometheus text format.
305    #[must_use]
306    pub fn to_prometheus(&self) -> String {
307        let mut output = String::with_capacity(2048);
308
309        // Node metrics
310        output.push_str("# HELP velesdb_graph_nodes_total Current number of nodes\n");
311        output.push_str("# TYPE velesdb_graph_nodes_total gauge\n");
312        let _ = writeln!(output, "velesdb_graph_nodes_total {}\n", self.nodes_total());
313
314        output.push_str("# HELP velesdb_graph_node_inserts_total Total node insertions\n");
315        output.push_str("# TYPE velesdb_graph_node_inserts_total counter\n");
316        let _ = writeln!(
317            output,
318            "velesdb_graph_node_inserts_total {}\n",
319            self.node_inserts_total()
320        );
321
322        // Edge metrics
323        output.push_str("# HELP velesdb_graph_edges_total Current number of edges\n");
324        output.push_str("# TYPE velesdb_graph_edges_total gauge\n");
325        let _ = writeln!(output, "velesdb_graph_edges_total {}\n", self.edges_total());
326
327        output.push_str("# HELP velesdb_graph_edge_inserts_total Total edge insertions\n");
328        output.push_str("# TYPE velesdb_graph_edge_inserts_total counter\n");
329        let _ = writeln!(
330            output,
331            "velesdb_graph_edge_inserts_total {}\n",
332            self.edge_inserts_total()
333        );
334
335        // Latency histograms
336        Self::append_histogram_prometheus(&mut output, "edge_insert", &self.edge_insert_latency);
337        Self::append_histogram_prometheus(&mut output, "traversal", &self.traversal_latency);
338
339        // Traversal metrics
340        output.push_str("# HELP velesdb_graph_traversals_total Total traversals executed\n");
341        output.push_str("# TYPE velesdb_graph_traversals_total counter\n");
342        let _ = writeln!(
343            output,
344            "velesdb_graph_traversals_total {}\n",
345            self.traversals_total()
346        );
347
348        output
349    }
350
351    fn append_histogram_prometheus(output: &mut String, name: &str, histogram: &LatencyHistogram) {
352        let bucket_bounds = [
353            "0.001", "0.005", "0.01", "0.05", "0.1", "0.5", "1", "5", "10", "+Inf",
354        ];
355        let counts = histogram.bucket_counts();
356        let mut cumulative = 0u64;
357
358        let _ = writeln!(
359            output,
360            "# HELP velesdb_graph_{}_duration_seconds {} latency histogram",
361            name,
362            name.replace('_', " ")
363        );
364        let _ = writeln!(
365            output,
366            "# TYPE velesdb_graph_{name}_duration_seconds histogram"
367        );
368
369        for (i, &bound) in bucket_bounds.iter().enumerate() {
370            cumulative += counts[i];
371            let _ = writeln!(
372                output,
373                "velesdb_graph_{name}_duration_seconds_bucket{{le=\"{bound}\"}} {cumulative}",
374            );
375        }
376
377        let _ = writeln!(
378            output,
379            "velesdb_graph_{}_duration_seconds_sum {}",
380            name,
381            histogram.sum_ns() as f64 / 1_000_000_000.0
382        );
383        let _ = writeln!(
384            output,
385            "velesdb_graph_{}_duration_seconds_count {}\n",
386            name,
387            histogram.count()
388        );
389    }
390
391    /// Resets all metrics to zero.
392    pub fn reset(&self) {
393        self.nodes_total.store(0, Ordering::Relaxed);
394        self.node_inserts_total.store(0, Ordering::Relaxed);
395        self.node_deletes_total.store(0, Ordering::Relaxed);
396        self.edges_total.store(0, Ordering::Relaxed);
397        self.edge_inserts_total.store(0, Ordering::Relaxed);
398        self.edge_deletes_total.store(0, Ordering::Relaxed);
399        self.traversals_total.store(0, Ordering::Relaxed);
400        self.traversal_nodes_visited.store(0, Ordering::Relaxed);
401        self.edge_insert_latency.reset();
402        self.edge_delete_latency.reset();
403        self.traversal_latency.reset();
404        self.query_latency.reset();
405    }
406}