velesdb_core/collection/graph/metrics/
mod.rs1#![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
25const BUCKET_BOUNDS_MS: [u64; 9] = [1, 5, 10, 50, 100, 500, 1000, 5000, 10000];
27
28#[derive(Debug, Default)]
32pub struct LatencyHistogram {
33 buckets: [AtomicU64; 10],
35 sum_ns: AtomicU64,
37 count: AtomicU64,
39}
40
41impl LatencyHistogram {
42 #[must_use]
44 pub fn new() -> Self {
45 Self::default()
46 }
47
48 pub fn observe(&self, duration: Duration) {
55 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 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 #[must_use]
81 pub fn count(&self) -> u64 {
82 self.count.load(Ordering::Relaxed)
83 }
84
85 #[must_use]
87 pub fn sum_ns(&self) -> u64 {
88 self.sum_ns.load(Ordering::Relaxed)
89 }
90
91 #[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 #[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 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#[derive(Debug, Default)]
144pub struct GraphMetrics {
145 nodes_total: AtomicU64,
147 node_inserts_total: AtomicU64,
148 node_deletes_total: AtomicU64,
149
150 edges_total: AtomicU64,
152 edge_inserts_total: AtomicU64,
153 edge_deletes_total: AtomicU64,
154
155 traversals_total: AtomicU64,
157 traversal_nodes_visited: AtomicU64,
158
159 pub edge_insert_latency: LatencyHistogram,
162 pub edge_delete_latency: LatencyHistogram,
164 pub traversal_latency: LatencyHistogram,
166 pub query_latency: LatencyHistogram,
168}
169
170impl GraphMetrics {
171 #[must_use]
173 pub fn new() -> Self {
174 Self::default()
175 }
176
177 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 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 #[must_use]
202 pub fn nodes_total(&self) -> u64 {
203 self.nodes_total.load(Ordering::Relaxed)
204 }
205
206 #[must_use]
208 pub fn node_inserts_total(&self) -> u64 {
209 self.node_inserts_total.load(Ordering::Relaxed)
210 }
211
212 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 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 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 #[must_use]
251 pub fn edges_total(&self) -> u64 {
252 self.edges_total.load(Ordering::Relaxed)
253 }
254
255 #[must_use]
257 pub fn edge_inserts_total(&self) -> u64 {
258 self.edge_inserts_total.load(Ordering::Relaxed)
259 }
260
261 #[must_use]
263 pub fn edge_deletes_total(&self) -> u64 {
264 self.edge_deletes_total.load(Ordering::Relaxed)
265 }
266
267 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 #[must_use]
281 pub fn traversals_total(&self) -> u64 {
282 self.traversals_total.load(Ordering::Relaxed)
283 }
284
285 #[must_use]
287 pub fn traversal_nodes_visited(&self) -> u64 {
288 self.traversal_nodes_visited.load(Ordering::Relaxed)
289 }
290
291 pub fn record_query(&self, latency: Duration) {
297 self.query_latency.observe(latency);
298 }
299
300 #[must_use]
306 pub fn to_prometheus(&self) -> String {
307 let mut output = String::with_capacity(2048);
308
309 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 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 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 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 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}