use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::OnceLock;
#[derive(Debug, Default)]
pub struct TraversalMetrics {
pub nodes_visited_total: AtomicU64,
pub max_depth_reached: AtomicU64,
pub edges_scanned_total: AtomicU64,
pub traversal_count: AtomicU64,
}
impl TraversalMetrics {
#[must_use]
pub fn new() -> Self {
Self::default()
}
pub fn record_traversal(&self, nodes_visited: u64, depth: u64, edges_scanned: u64) {
self.traversal_count.fetch_add(1, Ordering::Relaxed);
self.nodes_visited_total
.fetch_add(nodes_visited, Ordering::Relaxed);
self.edges_scanned_total
.fetch_add(edges_scanned, Ordering::Relaxed);
self.max_depth_reached.fetch_max(depth, Ordering::Relaxed);
}
#[must_use]
pub fn export_prometheus(&self) -> String {
use std::fmt::Write;
let mut output = String::new();
let _ = writeln!(
output,
"# HELP velesdb_traversal_nodes_visited_total Total nodes visited in traversals"
);
let _ = writeln!(
output,
"# TYPE velesdb_traversal_nodes_visited_total counter"
);
let _ = writeln!(
output,
"velesdb_traversal_nodes_visited_total {}",
self.nodes_visited_total.load(Ordering::Relaxed)
);
let _ = writeln!(output);
let _ = writeln!(
output,
"# HELP velesdb_traversal_max_depth Maximum traversal depth reached"
);
let _ = writeln!(output, "# TYPE velesdb_traversal_max_depth gauge");
let _ = writeln!(
output,
"velesdb_traversal_max_depth {}",
self.max_depth_reached.load(Ordering::Relaxed)
);
let _ = writeln!(output);
let _ = writeln!(
output,
"# HELP velesdb_traversal_edges_scanned_total Total edges scanned"
);
let _ = writeln!(
output,
"# TYPE velesdb_traversal_edges_scanned_total counter"
);
let _ = writeln!(
output,
"velesdb_traversal_edges_scanned_total {}",
self.edges_scanned_total.load(Ordering::Relaxed)
);
output
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
#[non_exhaustive]
pub enum LimitType {
Timeout,
Depth,
Cardinality,
Memory,
}
impl LimitType {
#[must_use]
pub fn as_str(self) -> &'static str {
match self {
Self::Timeout => "timeout",
Self::Depth => "depth",
Self::Cardinality => "cardinality",
Self::Memory => "memory",
}
}
}
#[derive(Debug, Default)]
pub struct GuardRailsMetrics {
pub timeout_exceeded: AtomicU64,
pub depth_exceeded: AtomicU64,
pub cardinality_exceeded: AtomicU64,
pub memory_exceeded: AtomicU64,
pub rate_limit_allowed: AtomicU64,
pub rate_limit_rejected: AtomicU64,
pub like_guardrail_rejected: AtomicU64,
pub parser_depth_limit_rejected: AtomicU64,
pub invalid_offset_read_errors: AtomicU64,
pub cache_collision_fallback_total: AtomicU64,
pub wal_replay_corrupt_entries: AtomicU64,
}
impl GuardRailsMetrics {
#[must_use]
pub fn new() -> Self {
Self::default()
}
pub fn record_limit_exceeded(&self, limit_type: LimitType) {
match limit_type {
LimitType::Timeout => self.timeout_exceeded.fetch_add(1, Ordering::Relaxed),
LimitType::Depth => self.depth_exceeded.fetch_add(1, Ordering::Relaxed),
LimitType::Cardinality => self.cardinality_exceeded.fetch_add(1, Ordering::Relaxed),
LimitType::Memory => self.memory_exceeded.fetch_add(1, Ordering::Relaxed),
};
}
pub fn record_rate_limit(&self, allowed: bool) {
if allowed {
self.rate_limit_allowed.fetch_add(1, Ordering::Relaxed);
} else {
self.rate_limit_rejected.fetch_add(1, Ordering::Relaxed);
}
}
pub fn record_like_guardrail_rejected(&self) {
self.like_guardrail_rejected.fetch_add(1, Ordering::Relaxed);
}
pub fn record_parser_depth_limit_rejected(&self) {
self.parser_depth_limit_rejected
.fetch_add(1, Ordering::Relaxed);
}
pub fn record_invalid_offset_read_error(&self) {
self.invalid_offset_read_errors
.fetch_add(1, Ordering::Relaxed);
}
pub fn record_cache_collision_fallback(&self) {
self.cache_collision_fallback_total
.fetch_add(1, Ordering::Relaxed);
}
pub fn record_wal_replay_corrupt_entry(&self) {
self.wal_replay_corrupt_entries
.fetch_add(1, Ordering::Relaxed);
}
#[must_use]
pub fn export_prometheus(&self) -> String {
let mut output = String::new();
self.write_limits_metrics(&mut output);
self.write_rate_limit_metrics(&mut output);
self.write_specialized_metrics(&mut output);
self.write_storage_metrics(&mut output);
output
}
fn write_limits_metrics(&self, output: &mut String) {
use std::fmt::Write;
let _ = writeln!(
output,
"# HELP velesdb_limits_exceeded_total Guard-rail limits exceeded"
);
let _ = writeln!(output, "# TYPE velesdb_limits_exceeded_total counter");
let _ = writeln!(
output,
"velesdb_limits_exceeded_total{{limit_type=\"timeout\"}} {}",
self.timeout_exceeded.load(Ordering::Relaxed)
);
let _ = writeln!(
output,
"velesdb_limits_exceeded_total{{limit_type=\"depth\"}} {}",
self.depth_exceeded.load(Ordering::Relaxed)
);
let _ = writeln!(
output,
"velesdb_limits_exceeded_total{{limit_type=\"cardinality\"}} {}",
self.cardinality_exceeded.load(Ordering::Relaxed)
);
let _ = writeln!(
output,
"velesdb_limits_exceeded_total{{limit_type=\"memory\"}} {}",
self.memory_exceeded.load(Ordering::Relaxed)
);
let _ = writeln!(output);
}
fn write_rate_limit_metrics(&self, output: &mut String) {
use std::fmt::Write;
let _ = writeln!(
output,
"# HELP velesdb_rate_limit_requests_total Rate limit decisions"
);
let _ = writeln!(output, "# TYPE velesdb_rate_limit_requests_total counter");
let _ = writeln!(
output,
"velesdb_rate_limit_requests_total{{decision=\"allowed\"}} {}",
self.rate_limit_allowed.load(Ordering::Relaxed)
);
let _ = writeln!(
output,
"velesdb_rate_limit_requests_total{{decision=\"rejected\"}} {}",
self.rate_limit_rejected.load(Ordering::Relaxed)
);
let _ = writeln!(output);
}
fn write_specialized_metrics(&self, output: &mut String) {
use std::fmt::Write;
let _ = writeln!(
output,
"# HELP velesdb_like_guardrail_rejected_total LIKE/ILIKE guardrail rejections"
);
let _ = writeln!(
output,
"# TYPE velesdb_like_guardrail_rejected_total counter"
);
let _ = writeln!(
output,
"velesdb_like_guardrail_rejected_total {}",
self.like_guardrail_rejected.load(Ordering::Relaxed)
);
let _ = writeln!(output);
let _ = writeln!(
output,
"# HELP velesdb_parser_depth_limit_rejected_total Parser depth-limit rejections"
);
let _ = writeln!(
output,
"# TYPE velesdb_parser_depth_limit_rejected_total counter"
);
let _ = writeln!(
output,
"velesdb_parser_depth_limit_rejected_total {}",
self.parser_depth_limit_rejected.load(Ordering::Relaxed)
);
let _ = writeln!(output);
}
fn write_storage_metrics(&self, output: &mut String) {
use std::fmt::Write;
let _ = writeln!(
output,
"# HELP velesdb_invalid_offset_read_errors_total Invalid storage offset read errors"
);
let _ = writeln!(
output,
"# TYPE velesdb_invalid_offset_read_errors_total counter"
);
let _ = writeln!(
output,
"velesdb_invalid_offset_read_errors_total {}",
self.invalid_offset_read_errors.load(Ordering::Relaxed)
);
let _ = writeln!(output);
let _ = writeln!(
output,
"# HELP velesdb_cache_collision_fallback_total Query-cache collision fallbacks"
);
let _ = writeln!(
output,
"# TYPE velesdb_cache_collision_fallback_total counter"
);
let _ = writeln!(
output,
"velesdb_cache_collision_fallback_total {}",
self.cache_collision_fallback_total.load(Ordering::Relaxed)
);
let _ = writeln!(output);
let _ = writeln!(
output,
"# HELP velesdb_wal_replay_corrupt_entries_total Mid-stream corrupt WAL entries skipped during replay"
);
let _ = writeln!(
output,
"# TYPE velesdb_wal_replay_corrupt_entries_total counter"
);
let _ = writeln!(
output,
"velesdb_wal_replay_corrupt_entries_total {}",
self.wal_replay_corrupt_entries.load(Ordering::Relaxed)
);
}
}
static GUARDRAILS_METRICS: OnceLock<GuardRailsMetrics> = OnceLock::new();
#[must_use]
pub fn global_guardrails_metrics() -> &'static GuardRailsMetrics {
GUARDRAILS_METRICS.get_or_init(GuardRailsMetrics::default)
}
#[cfg(test)]
#[path = "guardrails_tests.rs"]
mod tests;