1use std::{
14 fmt::Write as _,
15 sync::{
16 Arc,
17 atomic::{AtomicU64, Ordering},
18 },
19 time::Instant,
20};
21
22use dashmap::DashMap;
23
24#[derive(Debug, Default)]
34pub struct MetricsCollector {
35 pub queries_total: AtomicU64,
37
38 pub queries_success: AtomicU64,
40
41 pub queries_error: AtomicU64,
43
44 pub queries_duration_us: AtomicU64,
46
47 pub db_queries_total: AtomicU64,
49
50 pub db_queries_duration_us: AtomicU64,
52
53 pub validation_errors_total: AtomicU64,
55
56 pub parse_errors_total: AtomicU64,
58
59 pub execution_errors_total: AtomicU64,
61
62 pub http_requests_total: AtomicU64,
64
65 pub http_responses_2xx: AtomicU64,
67
68 pub http_responses_4xx: AtomicU64,
70
71 pub http_responses_5xx: AtomicU64,
73
74 pub cache_hits: AtomicU64,
76
77 pub cache_misses: AtomicU64,
79
80 pub federation_entity_resolutions_total: AtomicU64,
83
84 pub federation_entity_resolutions_errors: AtomicU64,
86
87 pub federation_entity_resolution_duration_us: AtomicU64,
89
90 pub federation_subgraph_requests_total: AtomicU64,
92
93 pub federation_subgraph_requests_errors: AtomicU64,
95
96 pub federation_subgraph_request_duration_us: AtomicU64,
98
99 pub federation_mutations_total: AtomicU64,
101
102 pub federation_mutations_errors: AtomicU64,
104
105 pub federation_mutation_duration_us: AtomicU64,
107
108 pub federation_entity_cache_hits: AtomicU64,
110
111 pub federation_entity_cache_misses: AtomicU64,
113
114 pub federation_errors_total: AtomicU64,
116
117 pub operation_metrics: Arc<OperationMetricsRegistry>,
122
123 pub http_request_duration: Arc<Histogram>,
125
126 pub db_query_duration: Arc<Histogram>,
128
129 pub schema_reloads_total: AtomicU64,
131
132 pub schema_reload_errors_total: AtomicU64,
134}
135
136impl MetricsCollector {
137 #[must_use]
139 pub fn new() -> Self {
140 Self::default()
141 }
142}
143
144impl MetricsCollector {
145 pub fn record_entity_resolution(&self, duration_us: u64, success: bool) {
152 self.federation_entity_resolutions_total.fetch_add(1, Ordering::Relaxed);
153 self.federation_entity_resolution_duration_us
154 .fetch_add(duration_us, Ordering::Relaxed);
155 if !success {
156 self.federation_entity_resolutions_errors.fetch_add(1, Ordering::Relaxed);
157 self.federation_errors_total.fetch_add(1, Ordering::Relaxed);
158 }
159 }
160
161 pub fn record_subgraph_request(&self, duration_us: u64, success: bool) {
168 self.federation_subgraph_requests_total.fetch_add(1, Ordering::Relaxed);
169 self.federation_subgraph_request_duration_us
170 .fetch_add(duration_us, Ordering::Relaxed);
171 if !success {
172 self.federation_subgraph_requests_errors.fetch_add(1, Ordering::Relaxed);
173 self.federation_errors_total.fetch_add(1, Ordering::Relaxed);
174 }
175 }
176
177 pub fn record_mutation(&self, duration_us: u64, success: bool) {
184 self.federation_mutations_total.fetch_add(1, Ordering::Relaxed);
185 self.federation_mutation_duration_us.fetch_add(duration_us, Ordering::Relaxed);
186 if !success {
187 self.federation_mutations_errors.fetch_add(1, Ordering::Relaxed);
188 self.federation_errors_total.fetch_add(1, Ordering::Relaxed);
189 }
190 }
191
192 pub fn record_entity_cache_hit(&self) {
194 self.federation_entity_cache_hits.fetch_add(1, Ordering::Relaxed);
195 }
196
197 pub fn record_entity_cache_miss(&self) {
199 self.federation_entity_cache_misses.fetch_add(1, Ordering::Relaxed);
200 }
201}
202
203const HISTOGRAM_BUCKET_BOUNDS_US: [u64; 11] = [
206 1_000, 5_000, 10_000, 25_000, 50_000, 100_000, 250_000, 500_000, 1_000_000, 2_500_000,
207 5_000_000,
208];
209
210const HISTOGRAM_LE_LABELS: [&str; 11] = [
212 "0.001", "0.005", "0.01", "0.025", "0.05", "0.1", "0.25", "0.5", "1", "2.5", "5",
213];
214
215#[derive(Debug)]
217pub struct OperationMetrics {
218 count: AtomicU64,
219 duration_us: AtomicU64,
220 error_count: AtomicU64,
221 bucket_counts: [AtomicU64; 11],
222}
223
224impl OperationMetrics {
225 fn new() -> Self {
226 Self {
227 count: AtomicU64::new(0),
228 duration_us: AtomicU64::new(0),
229 error_count: AtomicU64::new(0),
230 bucket_counts: std::array::from_fn(|_| AtomicU64::new(0)),
231 }
232 }
233
234 fn record(&self, duration_us: u64, is_error: bool) {
235 self.count.fetch_add(1, Ordering::Relaxed);
236 self.duration_us.fetch_add(duration_us, Ordering::Relaxed);
237 if is_error {
238 self.error_count.fetch_add(1, Ordering::Relaxed);
239 }
240 for (i, &bound) in HISTOGRAM_BUCKET_BOUNDS_US.iter().enumerate() {
243 if duration_us <= bound {
244 self.bucket_counts[i].fetch_add(1, Ordering::Relaxed);
245 return;
246 }
247 }
248 }
250}
251
252#[derive(Debug)]
257pub struct OperationMetricsRegistry {
258 operations: DashMap<String, OperationMetrics>,
259 max_operations: usize,
260 overflow: OperationMetrics,
261}
262
263impl OperationMetricsRegistry {
264 #[must_use]
266 pub fn new(max_operations: usize) -> Self {
267 Self {
268 operations: DashMap::new(),
269 max_operations,
270 overflow: OperationMetrics::new(),
271 }
272 }
273
274 pub fn record(&self, name: &str, duration_us: u64, is_error: bool) {
276 let canonical = if name.is_empty() {
277 "__anonymous__"
278 } else {
279 name
280 };
281
282 if let Some(entry) = self.operations.get(canonical) {
284 entry.record(duration_us, is_error);
285 return;
286 }
287
288 if self.operations.len() >= self.max_operations {
290 self.overflow.record(duration_us, is_error);
291 return;
292 }
293
294 self.operations
296 .entry(canonical.to_owned())
297 .or_insert_with(OperationMetrics::new)
298 .record(duration_us, is_error);
299 }
300
301 #[must_use]
303 pub fn to_prometheus_format(&self) -> String {
304 let mut out = String::new();
305
306 let mut entries: Vec<(String, u64, u64, u64, [u64; 11])> = self
308 .operations
309 .iter()
310 .map(|e| {
311 let buckets: [u64; 11] =
312 std::array::from_fn(|i| e.value().bucket_counts[i].load(Ordering::Relaxed));
313 (
314 e.key().clone(),
315 e.value().count.load(Ordering::Relaxed),
316 e.value().duration_us.load(Ordering::Relaxed),
317 e.value().error_count.load(Ordering::Relaxed),
318 buckets,
319 )
320 })
321 .collect();
322
323 let overflow_count = self.overflow.count.load(Ordering::Relaxed);
325 if overflow_count > 0 {
326 let buckets: [u64; 11] =
327 std::array::from_fn(|i| self.overflow.bucket_counts[i].load(Ordering::Relaxed));
328 entries.push((
329 "__overflow__".to_owned(),
330 overflow_count,
331 self.overflow.duration_us.load(Ordering::Relaxed),
332 self.overflow.error_count.load(Ordering::Relaxed),
333 buckets,
334 ));
335 }
336
337 if entries.is_empty() {
338 return out;
339 }
340
341 entries.sort_by(|a, b| a.0.cmp(&b.0));
342
343 out.push_str(
345 "\n# HELP fraiseql_query_duration_seconds Per-operation query duration histogram\n\
346 # TYPE fraiseql_query_duration_seconds histogram\n",
347 );
348 for (name, count, duration_us, _, buckets) in &entries {
349 let mut cumulative: u64 = 0;
350 for (i, &bucket_count) in buckets.iter().enumerate() {
351 cumulative += bucket_count;
352 let _ = writeln!(
353 out,
354 "fraiseql_query_duration_seconds_bucket{{operation=\"{name}\",le=\"{}\"}} \
355 {cumulative}",
356 HISTOGRAM_LE_LABELS[i],
357 );
358 }
359 let _ = writeln!(
360 out,
361 "fraiseql_query_duration_seconds_bucket{{operation=\"{name}\",le=\"+Inf\"}} \
362 {count}",
363 );
364 #[allow(clippy::cast_precision_loss)]
365 let sum_secs = *duration_us as f64 / 1_000_000.0;
367 let _ = writeln!(
368 out,
369 "fraiseql_query_duration_seconds_sum{{operation=\"{name}\"}} {sum_secs:.6}",
370 );
371 let _ = writeln!(
372 out,
373 "fraiseql_query_duration_seconds_count{{operation=\"{name}\"}} {count}",
374 );
375 }
376
377 out.push_str(
379 "\n# HELP fraiseql_query_errors_total Per-operation query error count\n\
380 # TYPE fraiseql_query_errors_total counter\n",
381 );
382 for (name, _, _, error_count, _) in &entries {
383 let _ =
384 writeln!(out, "fraiseql_query_errors_total{{operation=\"{name}\"}} {error_count}",);
385 }
386
387 out
388 }
389}
390
391impl Default for OperationMetricsRegistry {
392 fn default() -> Self {
393 Self::new(500)
394 }
395}
396
397#[derive(Debug)]
399pub struct Histogram {
400 count: AtomicU64,
401 sum_us: AtomicU64,
402 bucket_counts: [AtomicU64; 11],
403}
404
405impl Histogram {
406 #[must_use]
408 pub fn new() -> Self {
409 Self {
410 count: AtomicU64::new(0),
411 sum_us: AtomicU64::new(0),
412 bucket_counts: std::array::from_fn(|_| AtomicU64::new(0)),
413 }
414 }
415
416 pub fn observe_us(&self, duration_us: u64) {
418 self.count.fetch_add(1, Ordering::Relaxed);
419 self.sum_us.fetch_add(duration_us, Ordering::Relaxed);
420 for (i, &bound) in HISTOGRAM_BUCKET_BOUNDS_US.iter().enumerate() {
421 if duration_us <= bound {
422 self.bucket_counts[i].fetch_add(1, Ordering::Relaxed);
423 return;
424 }
425 }
426 }
427
428 #[allow(clippy::cast_precision_loss)] #[must_use]
431 pub fn to_prometheus_lines(&self, name: &str, help: &str) -> String {
432 let mut out = String::new();
433 let _ = writeln!(out, "\n# HELP {name} {help}");
434 let _ = writeln!(out, "# TYPE {name} histogram");
435 let count = self.count.load(Ordering::Relaxed);
436 let sum_us = self.sum_us.load(Ordering::Relaxed);
437 let mut cumulative = 0u64;
438 for (i, le) in HISTOGRAM_LE_LABELS.iter().enumerate() {
439 cumulative += self.bucket_counts[i].load(Ordering::Relaxed);
440 let _ = writeln!(out, "{name}_bucket{{le=\"{le}\"}} {cumulative}");
441 }
442 let _ = writeln!(out, "{name}_bucket{{le=\"+Inf\"}} {count}");
443 let sum_secs = sum_us as f64 / 1_000_000.0;
444 let _ = writeln!(out, "{name}_sum {sum_secs:.6}");
445 let _ = writeln!(out, "{name}_count {count}");
446 out
447 }
448}
449
450impl Default for Histogram {
451 fn default() -> Self {
452 Self::new()
453 }
454}
455
456impl Histogram {
457 #[allow(clippy::cast_precision_loss)] #[allow(clippy::cast_possible_truncation)] #[allow(clippy::cast_sign_loss)] #[must_use]
465 pub fn estimate_quantile_us(&self, quantile: f64) -> u64 {
466 let total = self.count.load(Ordering::Relaxed);
467 if total == 0 {
468 return 0;
469 }
470
471 let target = (total as f64 * quantile) as u64;
472 let mut cumulative = 0u64;
473
474 for (i, &bound) in HISTOGRAM_BUCKET_BOUNDS_US.iter().enumerate() {
475 let bucket_count = self.bucket_counts[i].load(Ordering::Relaxed);
476 cumulative += bucket_count;
477 if cumulative >= target {
478 return bound;
481 }
482 }
483
484 HISTOGRAM_BUCKET_BOUNDS_US[HISTOGRAM_BUCKET_BOUNDS_US.len() - 1]
486 }
487}
488
489pub struct TimingGuard {
491 start: Instant,
492 duration_atomic: Arc<AtomicU64>,
493}
494
495impl TimingGuard {
496 pub fn new(duration_atomic: Arc<AtomicU64>) -> Self {
498 Self {
499 start: Instant::now(),
500 duration_atomic,
501 }
502 }
503
504 pub fn record(self) {
506 #[allow(clippy::cast_possible_truncation)]
507 let duration_us = self.start.elapsed().as_micros() as u64;
509 self.duration_atomic.fetch_add(duration_us, Ordering::Relaxed);
510 }
511}
512
513#[derive(Debug)]
515pub struct PrometheusMetrics {
516 pub queries_total: u64,
518 pub queries_success: u64,
520 pub queries_error: u64,
522 pub queries_avg_duration_ms: f64,
524 pub db_queries_total: u64,
526 pub db_queries_avg_duration_ms: f64,
528 pub validation_errors_total: u64,
530 pub parse_errors_total: u64,
532 pub execution_errors_total: u64,
534 pub http_requests_total: u64,
536 pub http_responses_2xx: u64,
538 pub http_responses_4xx: u64,
540 pub http_responses_5xx: u64,
542 pub cache_hits: u64,
544 pub cache_misses: u64,
546 pub cache_hit_ratio: f64,
548}
549
550impl PrometheusMetrics {
551 #[must_use]
553 pub fn to_prometheus_format(&self) -> String {
554 format!(
555 r"# HELP fraiseql_graphql_queries_total Total GraphQL queries executed
556# TYPE fraiseql_graphql_queries_total counter
557fraiseql_graphql_queries_total {}
558
559# HELP fraiseql_graphql_queries_success Total successful GraphQL queries
560# TYPE fraiseql_graphql_queries_success counter
561fraiseql_graphql_queries_success {}
562
563# HELP fraiseql_graphql_queries_error Total failed GraphQL queries
564# TYPE fraiseql_graphql_queries_error counter
565fraiseql_graphql_queries_error {}
566
567# HELP fraiseql_graphql_query_duration_ms Average query execution time in milliseconds
568# TYPE fraiseql_graphql_query_duration_ms gauge
569fraiseql_graphql_query_duration_ms {}
570
571# HELP fraiseql_database_queries_total Total database queries executed
572# TYPE fraiseql_database_queries_total counter
573fraiseql_database_queries_total {}
574
575# HELP fraiseql_database_query_duration_ms Average database query time in milliseconds
576# TYPE fraiseql_database_query_duration_ms gauge
577fraiseql_database_query_duration_ms {}
578
579# HELP fraiseql_validation_errors_total Total validation errors
580# TYPE fraiseql_validation_errors_total counter
581fraiseql_validation_errors_total {}
582
583# HELP fraiseql_parse_errors_total Total parse errors
584# TYPE fraiseql_parse_errors_total counter
585fraiseql_parse_errors_total {}
586
587# HELP fraiseql_execution_errors_total Total execution errors
588# TYPE fraiseql_execution_errors_total counter
589fraiseql_execution_errors_total {}
590
591# HELP fraiseql_http_requests_total Total HTTP requests
592# TYPE fraiseql_http_requests_total counter
593fraiseql_http_requests_total {}
594
595# HELP fraiseql_http_responses_2xx Total 2xx HTTP responses
596# TYPE fraiseql_http_responses_2xx counter
597fraiseql_http_responses_2xx {}
598
599# HELP fraiseql_http_responses_4xx Total 4xx HTTP responses
600# TYPE fraiseql_http_responses_4xx counter
601fraiseql_http_responses_4xx {}
602
603# HELP fraiseql_http_responses_5xx Total 5xx HTTP responses
604# TYPE fraiseql_http_responses_5xx counter
605fraiseql_http_responses_5xx {}
606
607# HELP fraiseql_cache_hits Total cache hits
608# TYPE fraiseql_cache_hits counter
609fraiseql_cache_hits {}
610
611# HELP fraiseql_cache_misses Total cache misses
612# TYPE fraiseql_cache_misses counter
613fraiseql_cache_misses {}
614
615# HELP fraiseql_cache_hit_ratio Cache hit ratio (0-1)
616# TYPE fraiseql_cache_hit_ratio gauge
617fraiseql_cache_hit_ratio {:.3}
618",
619 self.queries_total,
620 self.queries_success,
621 self.queries_error,
622 self.queries_avg_duration_ms,
623 self.db_queries_total,
624 self.db_queries_avg_duration_ms,
625 self.validation_errors_total,
626 self.parse_errors_total,
627 self.execution_errors_total,
628 self.http_requests_total,
629 self.http_responses_2xx,
630 self.http_responses_4xx,
631 self.http_responses_5xx,
632 self.cache_hits,
633 self.cache_misses,
634 self.cache_hit_ratio,
635 )
636 }
637}
638
639impl From<&MetricsCollector> for PrometheusMetrics {
640 fn from(collector: &MetricsCollector) -> Self {
641 let queries_total = collector.queries_total.load(Ordering::Relaxed);
642 let queries_success = collector.queries_success.load(Ordering::Relaxed);
643 let queries_error = collector.queries_error.load(Ordering::Relaxed);
644 let queries_duration_us = collector.queries_duration_us.load(Ordering::Relaxed);
645
646 let db_queries_total = collector.db_queries_total.load(Ordering::Relaxed);
647 let db_queries_duration_us = collector.db_queries_duration_us.load(Ordering::Relaxed);
648
649 let cache_hits = collector.cache_hits.load(Ordering::Relaxed);
650 let cache_misses = collector.cache_misses.load(Ordering::Relaxed);
651 let cache_total = cache_hits + cache_misses;
652
653 Self {
654 queries_total,
655 queries_success,
656 queries_error,
657 #[allow(clippy::cast_precision_loss)] queries_avg_duration_ms: if queries_total > 0 {
659 (queries_duration_us as f64 / queries_total as f64) / 1000.0
660 } else {
661 0.0
662 },
663 db_queries_total,
664 #[allow(clippy::cast_precision_loss)] db_queries_avg_duration_ms: if db_queries_total > 0 {
666 (db_queries_duration_us as f64 / db_queries_total as f64) / 1000.0
667 } else {
668 0.0
669 },
670 validation_errors_total: collector.validation_errors_total.load(Ordering::Relaxed),
671 parse_errors_total: collector.parse_errors_total.load(Ordering::Relaxed),
672 execution_errors_total: collector.execution_errors_total.load(Ordering::Relaxed),
673 http_requests_total: collector.http_requests_total.load(Ordering::Relaxed),
674 http_responses_2xx: collector.http_responses_2xx.load(Ordering::Relaxed),
675 http_responses_4xx: collector.http_responses_4xx.load(Ordering::Relaxed),
676 http_responses_5xx: collector.http_responses_5xx.load(Ordering::Relaxed),
677 cache_hits,
678 cache_misses,
679 #[allow(clippy::cast_precision_loss)] cache_hit_ratio: if cache_total > 0 {
681 cache_hits as f64 / cache_total as f64
682 } else {
683 0.0
684 },
685 }
686 }
687}