Skip to main content

velesdb_core/metrics/
operational.rs

1//! Operational metrics for monitoring `VelesDB` in production.
2//!
3//! Provides thread-safe counters and gauges for:
4//! - Query throughput and errors (Prometheus-exportable)
5//! - Graph traversal statistics
6//! - Guard-rails and rate limiting metrics
7
8use std::sync::atomic::{AtomicU64, Ordering};
9use std::sync::Arc;
10
11/// Query duration histogram buckets (in seconds).
12pub const DURATION_BUCKETS: [f64; 8] = [0.001, 0.005, 0.01, 0.05, 0.1, 0.5, 1.0, 5.0];
13
14/// Traversal depth histogram buckets.
15pub const DEPTH_BUCKETS: [u64; 6] = [1, 2, 3, 5, 10, 20];
16
17/// Nodes visited histogram buckets.
18pub const NODES_BUCKETS: [u64; 7] = [10, 50, 100, 500, 1000, 5000, 10000];
19
20/// Operational metrics for `VelesDB` monitoring (EPIC-050 US-001).
21///
22/// Thread-safe counters and gauges that can be exported in Prometheus format.
23#[derive(Debug, Default)]
24pub struct OperationalMetrics {
25    /// Total queries executed
26    pub queries_total: AtomicU64,
27    /// Total query errors
28    pub query_errors: AtomicU64,
29    /// Queries rejected by guard rails (rate limiting, circuit breaker)
30    pub query_rate_limited: AtomicU64,
31    /// Vector search queries
32    pub vector_queries: AtomicU64,
33    /// Graph traversal queries
34    pub graph_queries: AtomicU64,
35    /// Hybrid queries (vector + graph)
36    pub hybrid_queries: AtomicU64,
37    /// Total documents across all collections
38    pub documents_total: AtomicU64,
39    /// Total index size in bytes
40    pub index_size_bytes: AtomicU64,
41    /// Active connections (for server)
42    pub active_connections: AtomicU64,
43}
44
45impl OperationalMetrics {
46    /// Creates a new metrics instance.
47    #[must_use]
48    pub fn new() -> Self {
49        Self::default()
50    }
51
52    /// Creates a fresh metrics instance wrapped in an `Arc` for shared
53    /// ownership across handlers.
54    ///
55    /// This is a constructor convenience — every call returns a brand-new
56    /// counter set. Callers that want a single workspace-wide instance
57    /// must hold the returned `Arc` themselves (e.g. in `AppState`) and
58    /// clone it where needed; this method does NOT back the instance with
59    /// a global cache.
60    #[must_use]
61    pub fn new_arc() -> Arc<Self> {
62        Arc::new(Self::new())
63    }
64
65    /// Deprecated alias of [`OperationalMetrics::new_arc`].
66    ///
67    /// The original name implied a global singleton, but the
68    /// implementation always returned a fresh `Arc::new(Self::new())`.
69    /// Kept as a forwarding alias to preserve the v1.13.1 public API
70    /// surface; callers should migrate to `new_arc`.
71    #[must_use]
72    #[deprecated(since = "1.13.2", note = "use `OperationalMetrics::new_arc` instead")]
73    pub fn shared() -> Arc<Self> {
74        Self::new_arc()
75    }
76
77    /// Increments the total query counter.
78    pub fn inc_queries(&self) {
79        self.queries_total.fetch_add(1, Ordering::Relaxed);
80    }
81
82    /// Increments the query error counter.
83    pub fn inc_errors(&self) {
84        self.query_errors.fetch_add(1, Ordering::Relaxed);
85    }
86
87    /// Increments the rate-limited/rejected counter.
88    pub fn inc_rate_limited(&self) {
89        self.query_rate_limited.fetch_add(1, Ordering::Relaxed);
90    }
91
92    /// Records a vector search query.
93    pub fn record_vector_query(&self) {
94        self.inc_queries();
95        self.vector_queries.fetch_add(1, Ordering::Relaxed);
96    }
97
98    /// Records a graph traversal query.
99    pub fn record_graph_query(&self) {
100        self.inc_queries();
101        self.graph_queries.fetch_add(1, Ordering::Relaxed);
102    }
103
104    /// Records a hybrid query.
105    pub fn record_hybrid_query(&self) {
106        self.inc_queries();
107        self.hybrid_queries.fetch_add(1, Ordering::Relaxed);
108    }
109
110    /// Sets the document count.
111    pub fn set_documents(&self, count: u64) {
112        self.documents_total.store(count, Ordering::Relaxed);
113    }
114
115    /// Sets the index size.
116    pub fn set_index_size(&self, bytes: u64) {
117        self.index_size_bytes.store(bytes, Ordering::Relaxed);
118    }
119
120    /// Increments active connections.
121    pub fn inc_connections(&self) {
122        self.active_connections.fetch_add(1, Ordering::Relaxed);
123    }
124
125    /// Decrements active connections.
126    ///
127    /// Uses `fetch_update` to saturate at 0, preventing underflow wrap to `u64::MAX`.
128    pub fn dec_connections(&self) {
129        self.active_connections
130            .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |x| {
131                Some(x.saturating_sub(1))
132            })
133            .ok();
134    }
135
136    /// Exports metrics in Prometheus text format.
137    #[must_use]
138    pub fn export_prometheus(&self) -> String {
139        use std::fmt::Write;
140        let mut output = String::new();
141
142        let total = self.queries_total.load(Ordering::Relaxed);
143        let errors = self.query_errors.load(Ordering::Relaxed);
144        let rate_limited = self.query_rate_limited.load(Ordering::Relaxed);
145        let success = total.saturating_sub(errors).saturating_sub(rate_limited);
146
147        Self::write_metric_header(
148            &mut output,
149            "velesdb_queries_total",
150            "counter",
151            "Total number of queries executed",
152        );
153        let _ = writeln!(
154            output,
155            "velesdb_queries_total{{status=\"success\"}} {success}"
156        );
157        let _ = writeln!(output, "velesdb_queries_total{{status=\"error\"}} {errors}");
158        let _ = writeln!(
159            output,
160            "velesdb_queries_total{{status=\"rate_limited\"}} {rate_limited}\n"
161        );
162
163        Self::write_metric_header(
164            &mut output,
165            "velesdb_queries_by_type",
166            "counter",
167            "Queries by type",
168        );
169        let _ = writeln!(
170            output,
171            "velesdb_queries_by_type{{type=\"vector\"}} {}",
172            self.vector_queries.load(Ordering::Relaxed)
173        );
174        let _ = writeln!(
175            output,
176            "velesdb_queries_by_type{{type=\"graph\"}} {}",
177            self.graph_queries.load(Ordering::Relaxed)
178        );
179        let _ = writeln!(
180            output,
181            "velesdb_queries_by_type{{type=\"hybrid\"}} {}\n",
182            self.hybrid_queries.load(Ordering::Relaxed)
183        );
184
185        Self::write_gauge(
186            &mut output,
187            "velesdb_documents_total",
188            "Total documents in database",
189            self.documents_total.load(Ordering::Relaxed),
190        );
191        Self::write_gauge(
192            &mut output,
193            "velesdb_index_size_bytes",
194            "Total index size in bytes",
195            self.index_size_bytes.load(Ordering::Relaxed),
196        );
197        Self::write_gauge(
198            &mut output,
199            "velesdb_active_connections",
200            "Current active connections",
201            self.active_connections.load(Ordering::Relaxed),
202        );
203
204        output
205    }
206
207    /// Writes a Prometheus metric header (HELP + TYPE lines).
208    fn write_metric_header(output: &mut String, name: &str, metric_type: &str, help: &str) {
209        use std::fmt::Write;
210        let _ = write!(
211            output,
212            "# HELP {name} {help}\n# TYPE {name} {metric_type}\n"
213        );
214    }
215
216    /// Writes a Prometheus gauge metric with header.
217    fn write_gauge(output: &mut String, name: &str, help: &str, value: u64) {
218        use std::fmt::Write;
219        Self::write_metric_header(output, name, "gauge", help);
220        let _ = writeln!(output, "{name} {value}\n");
221    }
222}
223
224#[cfg(test)]
225mod tests {
226    use super::*;
227
228    #[test]
229    fn test_operational_metrics_counters() {
230        let metrics = OperationalMetrics::new();
231
232        metrics.record_vector_query();
233        metrics.record_vector_query();
234        metrics.record_graph_query();
235        metrics.record_hybrid_query();
236        metrics.inc_errors();
237
238        assert_eq!(metrics.queries_total.load(Ordering::Relaxed), 4);
239        assert_eq!(metrics.vector_queries.load(Ordering::Relaxed), 2);
240        assert_eq!(metrics.graph_queries.load(Ordering::Relaxed), 1);
241        assert_eq!(metrics.hybrid_queries.load(Ordering::Relaxed), 1);
242        assert_eq!(metrics.query_errors.load(Ordering::Relaxed), 1);
243    }
244
245    #[test]
246    fn test_operational_metrics_gauges() {
247        let metrics = OperationalMetrics::new();
248
249        metrics.set_documents(1000);
250        metrics.set_index_size(1024 * 1024);
251        metrics.inc_connections();
252        metrics.inc_connections();
253        metrics.dec_connections();
254
255        assert_eq!(metrics.documents_total.load(Ordering::Relaxed), 1000);
256        assert_eq!(
257            metrics.index_size_bytes.load(Ordering::Relaxed),
258            1024 * 1024
259        );
260        assert_eq!(metrics.active_connections.load(Ordering::Relaxed), 1);
261    }
262
263    #[test]
264    fn test_operational_metrics_prometheus_export() {
265        let metrics = OperationalMetrics::new();
266        metrics.record_vector_query();
267        metrics.set_documents(100);
268
269        let output = metrics.export_prometheus();
270
271        assert!(output.contains("velesdb_queries_total"));
272        assert!(output.contains("velesdb_documents_total 100"));
273        assert!(output.contains("# TYPE"));
274        assert!(output.contains("# HELP"));
275    }
276
277    #[test]
278    fn test_operational_metrics_new_arc() {
279        let metrics = OperationalMetrics::new_arc();
280        metrics.record_vector_query();
281
282        // Clone Arc and verify shared state
283        let metrics2 = Arc::clone(&metrics);
284        metrics2.record_vector_query();
285
286        assert_eq!(metrics.queries_total.load(Ordering::Relaxed), 2);
287    }
288
289    #[test]
290    fn test_rate_limited_request_increments_status_rate_limited() {
291        let metrics = OperationalMetrics::new();
292
293        // Simulate 3 queries: 1 succeeds, 1 errors, 1 rate-limited
294        metrics.record_vector_query(); // queries_total = 1
295        metrics.record_vector_query(); // queries_total = 2
296        metrics.record_vector_query(); // queries_total = 3
297        metrics.inc_errors(); // 1 error
298        metrics.inc_rate_limited(); // 1 rate-limited
299
300        assert_eq!(metrics.queries_total.load(Ordering::Relaxed), 3);
301        assert_eq!(metrics.query_errors.load(Ordering::Relaxed), 1);
302        assert_eq!(metrics.query_rate_limited.load(Ordering::Relaxed), 1);
303
304        let output = metrics.export_prometheus();
305
306        // success = total - errors - rate_limited = 3 - 1 - 1 = 1
307        assert!(
308            output.contains("velesdb_queries_total{status=\"success\"} 1"),
309            "expected success=1 in:\n{output}"
310        );
311        assert!(
312            output.contains("velesdb_queries_total{status=\"error\"} 1"),
313            "expected error=1 in:\n{output}"
314        );
315        assert!(
316            output.contains("velesdb_queries_total{status=\"rate_limited\"} 1"),
317            "expected rate_limited=1 in:\n{output}"
318        );
319    }
320}