1use std::sync::atomic::{AtomicU64, Ordering};
9use std::sync::Arc;
10
11pub const DURATION_BUCKETS: [f64; 8] = [0.001, 0.005, 0.01, 0.05, 0.1, 0.5, 1.0, 5.0];
13
14pub const DEPTH_BUCKETS: [u64; 6] = [1, 2, 3, 5, 10, 20];
16
17pub const NODES_BUCKETS: [u64; 7] = [10, 50, 100, 500, 1000, 5000, 10000];
19
20#[derive(Debug, Default)]
24pub struct OperationalMetrics {
25 pub queries_total: AtomicU64,
27 pub query_errors: AtomicU64,
29 pub query_rate_limited: AtomicU64,
31 pub vector_queries: AtomicU64,
33 pub graph_queries: AtomicU64,
35 pub hybrid_queries: AtomicU64,
37 pub documents_total: AtomicU64,
39 pub index_size_bytes: AtomicU64,
41 pub active_connections: AtomicU64,
43}
44
45impl OperationalMetrics {
46 #[must_use]
48 pub fn new() -> Self {
49 Self::default()
50 }
51
52 #[must_use]
61 pub fn new_arc() -> Arc<Self> {
62 Arc::new(Self::new())
63 }
64
65 #[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 pub fn inc_queries(&self) {
79 self.queries_total.fetch_add(1, Ordering::Relaxed);
80 }
81
82 pub fn inc_errors(&self) {
84 self.query_errors.fetch_add(1, Ordering::Relaxed);
85 }
86
87 pub fn inc_rate_limited(&self) {
89 self.query_rate_limited.fetch_add(1, Ordering::Relaxed);
90 }
91
92 pub fn record_vector_query(&self) {
94 self.inc_queries();
95 self.vector_queries.fetch_add(1, Ordering::Relaxed);
96 }
97
98 pub fn record_graph_query(&self) {
100 self.inc_queries();
101 self.graph_queries.fetch_add(1, Ordering::Relaxed);
102 }
103
104 pub fn record_hybrid_query(&self) {
106 self.inc_queries();
107 self.hybrid_queries.fetch_add(1, Ordering::Relaxed);
108 }
109
110 pub fn set_documents(&self, count: u64) {
112 self.documents_total.store(count, Ordering::Relaxed);
113 }
114
115 pub fn set_index_size(&self, bytes: u64) {
117 self.index_size_bytes.store(bytes, Ordering::Relaxed);
118 }
119
120 pub fn inc_connections(&self) {
122 self.active_connections.fetch_add(1, Ordering::Relaxed);
123 }
124
125 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 #[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 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 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 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 metrics.record_vector_query(); metrics.record_vector_query(); metrics.record_vector_query(); metrics.inc_errors(); metrics.inc_rate_limited(); 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 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}