Skip to main content

rustfs_audit/
observability.rs

1//  Copyright 2024 RustFS Team
2//
3//  Licensed under the Apache License, Version 2.0 (the "License");
4//  you may not use this file except in compliance with the License.
5//  You may obtain a copy of the License at
6//
7//      http://www.apache.org/licenses/LICENSE-2.0
8//
9//  Unless required by applicable law or agreed to in writing, software
10//  distributed under the License is distributed on an "AS IS" BASIS,
11//  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12//  See the License for the specific language governing permissions and
13//  limitations under the License.
14
15//! Observability and metrics for the audit system
16//!
17//! This module provides comprehensive observability features including:
18//! - Performance metrics (EPS, latency)
19//! - Target health monitoring  
20//! - Configuration change tracking
21//! - Error rate monitoring
22//! - Queue depth monitoring
23
24use metrics::{counter, describe_counter, describe_gauge, describe_histogram, gauge, histogram};
25use std::sync::atomic::{AtomicU64, Ordering};
26use std::sync::{Arc, OnceLock};
27use std::time::{Duration, Instant};
28use tokio::sync::RwLock;
29use tracing::info;
30
31const RUSTFS_AUDIT_METRICS_NAMESPACE: &str = "rustfs.audit.";
32const LOG_COMPONENT_AUDIT: &str = "audit";
33const LOG_SUBSYSTEM_OBSERVABILITY: &str = "observability";
34const EVENT_AUDIT_OBSERVABILITY_STATE: &str = "audit_observability_state";
35
36const M_AUDIT_EVENTS_TOTAL: &str = const_str::concat!(RUSTFS_AUDIT_METRICS_NAMESPACE, "events.total");
37const M_AUDIT_EVENTS_FAILED: &str = const_str::concat!(RUSTFS_AUDIT_METRICS_NAMESPACE, "events.failed");
38const M_AUDIT_DISPATCH_NS: &str = const_str::concat!(RUSTFS_AUDIT_METRICS_NAMESPACE, "dispatch.ns");
39const M_AUDIT_EPS: &str = const_str::concat!(RUSTFS_AUDIT_METRICS_NAMESPACE, "eps");
40const M_AUDIT_TARGET_OPS: &str = const_str::concat!(RUSTFS_AUDIT_METRICS_NAMESPACE, "target.ops");
41const M_AUDIT_CONFIG_RELOADS: &str = const_str::concat!(RUSTFS_AUDIT_METRICS_NAMESPACE, "config.reloads");
42const M_AUDIT_SYSTEM_STARTS: &str = const_str::concat!(RUSTFS_AUDIT_METRICS_NAMESPACE, "system.starts");
43
44const L_RESULT: &str = "result";
45const L_STATUS: &str = "status";
46
47const V_SUCCESS: &str = "success";
48const V_FAILURE: &str = "failure";
49
50/// One-time registration of indicator meta information
51/// This function ensures that metric descriptors are registered only once.
52pub fn init_observability_metrics() {
53    static METRICS_DESC_INIT: OnceLock<()> = OnceLock::new();
54    METRICS_DESC_INIT.get_or_init(|| {
55        // Event/Time-consuming
56        describe_counter!(M_AUDIT_EVENTS_TOTAL, "Total audit events (labeled by result).");
57        describe_counter!(M_AUDIT_EVENTS_FAILED, "Total failed audit events.");
58        describe_histogram!(M_AUDIT_DISPATCH_NS, "Dispatch time per event (ns).");
59        describe_gauge!(M_AUDIT_EPS, "Events per second since last reset.");
60
61        // Target operation/system event
62        describe_counter!(M_AUDIT_TARGET_OPS, "Total target operations (labeled by status).");
63        describe_counter!(M_AUDIT_CONFIG_RELOADS, "Total configuration reloads.");
64        describe_counter!(M_AUDIT_SYSTEM_STARTS, "Total system starts.");
65    });
66}
67
68/// Metrics collector for audit system observability
69#[derive(Debug)]
70pub struct AuditMetrics {
71    // Performance metrics
72    total_events_processed: AtomicU64,
73    total_events_failed: AtomicU64,
74    total_dispatch_time_ns: AtomicU64,
75
76    // Target metrics
77    target_success_count: AtomicU64,
78    target_failure_count: AtomicU64,
79
80    // System metrics
81    config_reload_count: AtomicU64,
82    system_start_count: AtomicU64,
83
84    // Performance tracking
85    last_reset_time: Arc<RwLock<Instant>>,
86}
87
88impl Default for AuditMetrics {
89    fn default() -> Self {
90        Self::new()
91    }
92}
93
94impl AuditMetrics {
95    /// Creates a new metrics collector
96    pub fn new() -> Self {
97        init_observability_metrics();
98        Self {
99            total_events_processed: AtomicU64::new(0),
100            total_events_failed: AtomicU64::new(0),
101            total_dispatch_time_ns: AtomicU64::new(0),
102            target_success_count: AtomicU64::new(0),
103            target_failure_count: AtomicU64::new(0),
104            config_reload_count: AtomicU64::new(0),
105            system_start_count: AtomicU64::new(0),
106            last_reset_time: Arc::new(RwLock::new(Instant::now())),
107        }
108    }
109
110    // Suggestion: Call this auxiliary function in the existing "Successful Event Recording" method body to complete the instrumentation
111    #[inline]
112    fn emit_event_success_metrics(&self, dispatch_time: Duration) {
113        // count + histogram
114        counter!(M_AUDIT_EVENTS_TOTAL, L_RESULT => V_SUCCESS).increment(1);
115        histogram!(M_AUDIT_DISPATCH_NS).record(dispatch_time.as_nanos() as f64);
116    }
117
118    // Suggestion: Call this auxiliary function in the existing "Failure Event Recording" method body to complete the instrumentation
119    #[inline]
120    fn emit_event_failure_metrics(&self, dispatch_time: Duration) {
121        counter!(M_AUDIT_EVENTS_TOTAL, L_RESULT => V_FAILURE).increment(1);
122        counter!(M_AUDIT_EVENTS_FAILED).increment(1);
123        histogram!(M_AUDIT_DISPATCH_NS).record(dispatch_time.as_nanos() as f64);
124    }
125
126    /// Records a successful event dispatch
127    pub fn record_event_success(&self, dispatch_time: Duration) {
128        self.total_events_processed.fetch_add(1, Ordering::Relaxed);
129        self.total_dispatch_time_ns
130            .fetch_add(dispatch_time.as_nanos() as u64, Ordering::Relaxed);
131        self.emit_event_success_metrics(dispatch_time);
132    }
133
134    /// Records a failed event dispatch
135    pub fn record_event_failure(&self, dispatch_time: Duration) {
136        self.total_events_failed.fetch_add(1, Ordering::Relaxed);
137        self.total_dispatch_time_ns
138            .fetch_add(dispatch_time.as_nanos() as u64, Ordering::Relaxed);
139        self.emit_event_failure_metrics(dispatch_time);
140    }
141
142    /// Records a successful target operation
143    pub fn record_target_success(&self) {
144        self.target_success_count.fetch_add(1, Ordering::Relaxed);
145        counter!(M_AUDIT_TARGET_OPS, L_STATUS => V_SUCCESS).increment(1);
146    }
147
148    /// Records a failed target operation
149    pub fn record_target_failure(&self) {
150        self.target_failure_count.fetch_add(1, Ordering::Relaxed);
151        counter!(M_AUDIT_TARGET_OPS, L_STATUS => V_FAILURE).increment(1);
152    }
153
154    /// Records a configuration reload
155    pub fn record_config_reload(&self) {
156        self.config_reload_count.fetch_add(1, Ordering::Relaxed);
157        counter!(M_AUDIT_CONFIG_RELOADS).increment(1);
158        info!(
159            event = EVENT_AUDIT_OBSERVABILITY_STATE,
160            component = LOG_COMPONENT_AUDIT,
161            subsystem = LOG_SUBSYSTEM_OBSERVABILITY,
162            state = "config_reloaded",
163            "audit observability state"
164        );
165    }
166
167    /// Records a system start
168    pub fn record_system_start(&self) {
169        self.system_start_count.fetch_add(1, Ordering::Relaxed);
170        counter!(M_AUDIT_SYSTEM_STARTS).increment(1);
171        info!(
172            event = EVENT_AUDIT_OBSERVABILITY_STATE,
173            component = LOG_COMPONENT_AUDIT,
174            subsystem = LOG_SUBSYSTEM_OBSERVABILITY,
175            state = "system_started",
176            "audit observability state"
177        );
178    }
179
180    /// Gets the current events per second (EPS)
181    pub async fn get_events_per_second(&self) -> f64 {
182        let reset_time = *self.last_reset_time.read().await;
183        let elapsed = reset_time.elapsed();
184        let total_events = self.total_events_processed.load(Ordering::Relaxed) + self.total_events_failed.load(Ordering::Relaxed);
185
186        let eps = if elapsed.as_secs_f64() > 0.0 {
187            total_events as f64 / elapsed.as_secs_f64()
188        } else {
189            0.0
190        };
191        // EPS is reported in gauge
192        gauge!(M_AUDIT_EPS).set(eps);
193        eps
194    }
195
196    /// Gets the average dispatch latency in milliseconds
197    pub fn get_average_latency_ms(&self) -> f64 {
198        let total_events = self.total_events_processed.load(Ordering::Relaxed) + self.total_events_failed.load(Ordering::Relaxed);
199        let total_time_ns = self.total_dispatch_time_ns.load(Ordering::Relaxed);
200
201        if total_events > 0 {
202            (total_time_ns as f64 / total_events as f64) / 1_000_000.0 // Convert ns to ms
203        } else {
204            0.0
205        }
206    }
207
208    /// Gets the error rate as a percentage
209    pub fn get_error_rate(&self) -> f64 {
210        let total_events = self.total_events_processed.load(Ordering::Relaxed) + self.total_events_failed.load(Ordering::Relaxed);
211        let failed_events = self.total_events_failed.load(Ordering::Relaxed);
212
213        if total_events > 0 {
214            (failed_events as f64 / total_events as f64) * 100.0
215        } else {
216            0.0
217        }
218    }
219
220    /// Gets target success rate as a percentage
221    pub fn get_target_success_rate(&self) -> f64 {
222        let total_ops = self.target_success_count.load(Ordering::Relaxed) + self.target_failure_count.load(Ordering::Relaxed);
223        let success_ops = self.target_success_count.load(Ordering::Relaxed);
224
225        if total_ops > 0 {
226            (success_ops as f64 / total_ops as f64) * 100.0
227        } else {
228            100.0 // No operations = 100% success rate
229        }
230    }
231
232    /// Resets all metrics and timing
233    pub async fn reset(&self) {
234        self.total_events_processed.store(0, Ordering::Relaxed);
235        self.total_events_failed.store(0, Ordering::Relaxed);
236        self.total_dispatch_time_ns.store(0, Ordering::Relaxed);
237        self.target_success_count.store(0, Ordering::Relaxed);
238        self.target_failure_count.store(0, Ordering::Relaxed);
239        self.config_reload_count.store(0, Ordering::Relaxed);
240        self.system_start_count.store(0, Ordering::Relaxed);
241
242        let mut reset_time = self.last_reset_time.write().await;
243        *reset_time = Instant::now();
244
245        // Reset EPS to zero after reset
246        gauge!(M_AUDIT_EPS).set(0.0);
247        info!(
248            event = EVENT_AUDIT_OBSERVABILITY_STATE,
249            component = LOG_COMPONENT_AUDIT,
250            subsystem = LOG_SUBSYSTEM_OBSERVABILITY,
251            state = "metrics_reset",
252            "audit observability state"
253        );
254    }
255
256    /// Generates a comprehensive metrics report
257    pub async fn generate_report(&self) -> AuditMetricsReport {
258        AuditMetricsReport {
259            events_per_second: self.get_events_per_second().await,
260            average_latency_ms: self.get_average_latency_ms(),
261            error_rate_percent: self.get_error_rate(),
262            target_success_rate_percent: self.get_target_success_rate(),
263            total_events_processed: self.total_events_processed.load(Ordering::Relaxed),
264            total_events_failed: self.total_events_failed.load(Ordering::Relaxed),
265            config_reload_count: self.config_reload_count.load(Ordering::Relaxed),
266            system_start_count: self.system_start_count.load(Ordering::Relaxed),
267        }
268    }
269
270    /// Validates performance requirements
271    pub async fn validate_performance_requirements(&self) -> PerformanceValidation {
272        let eps = self.get_events_per_second().await;
273        let avg_latency_ms = self.get_average_latency_ms();
274        let error_rate = self.get_error_rate();
275
276        let mut validation = PerformanceValidation {
277            meets_eps_requirement: eps >= 3000.0,
278            meets_latency_requirement: avg_latency_ms <= 30.0,
279            meets_error_rate_requirement: error_rate <= 1.0, // Less than 1% error rate
280            current_eps: eps,
281            current_latency_ms: avg_latency_ms,
282            current_error_rate: error_rate,
283            recommendations: Vec::new(),
284        };
285
286        // Generate recommendations
287        if !validation.meets_eps_requirement {
288            validation.recommendations.push(format!(
289                "EPS ({eps:.0}) is below requirement (3000). Consider optimizing target dispatch or adding more target instances."
290            ));
291        }
292
293        if !validation.meets_latency_requirement {
294            validation.recommendations.push(format!(
295                "Average latency ({avg_latency_ms:.2}ms) exceeds requirement (30ms). Consider optimizing target responses or increasing timeout values."
296            ));
297        }
298
299        if !validation.meets_error_rate_requirement {
300            validation.recommendations.push(format!(
301                "Error rate ({error_rate:.2}%) exceeds recommendation (1%). Check target connectivity and configuration."
302            ));
303        }
304
305        if validation.meets_eps_requirement && validation.meets_latency_requirement && validation.meets_error_rate_requirement {
306            validation
307                .recommendations
308                .push("All performance requirements are met.".to_string());
309        }
310
311        validation
312    }
313}
314
315/// Comprehensive metrics report
316#[derive(Debug, Clone)]
317pub struct AuditMetricsReport {
318    pub events_per_second: f64,
319    pub average_latency_ms: f64,
320    pub error_rate_percent: f64,
321    pub target_success_rate_percent: f64,
322    pub total_events_processed: u64,
323    pub total_events_failed: u64,
324    pub config_reload_count: u64,
325    pub system_start_count: u64,
326}
327
328impl AuditMetricsReport {
329    /// Formats the report as a human-readable string
330    pub fn format(&self) -> String {
331        format!(
332            "Audit System Metrics Report:\n\
333             Events per Second: {:.2}\n\
334             Average Latency: {:.2}ms\n\
335             Error Rate: {:.2}%\n\
336             Target Success Rate: {:.2}%\n\
337             Total Events Processed: {}\n\
338             Total Events Failed: {}\n\
339             Configuration Reloads: {}\n\
340             System Starts: {}",
341            self.events_per_second,
342            self.average_latency_ms,
343            self.error_rate_percent,
344            self.target_success_rate_percent,
345            self.total_events_processed,
346            self.total_events_failed,
347            self.config_reload_count,
348            self.system_start_count
349        )
350    }
351}
352
353/// Performance validation results
354#[derive(Debug, Clone)]
355pub struct PerformanceValidation {
356    pub meets_eps_requirement: bool,
357    pub meets_latency_requirement: bool,
358    pub meets_error_rate_requirement: bool,
359    pub current_eps: f64,
360    pub current_latency_ms: f64,
361    pub current_error_rate: f64,
362    pub recommendations: Vec<String>,
363}
364
365impl PerformanceValidation {
366    /// Checks if all performance requirements are met
367    pub fn all_requirements_met(&self) -> bool {
368        self.meets_eps_requirement && self.meets_latency_requirement && self.meets_error_rate_requirement
369    }
370
371    /// Formats the validation as a human-readable string
372    pub fn format(&self) -> String {
373        let status = if self.all_requirements_met() { "✅ PASS" } else { "❌ FAIL" };
374
375        let mut result = format!(
376            "Performance Requirements Validation: {}\n\
377             EPS Requirement (≥3000): {} ({:.2})\n\
378             Latency Requirement (≤30ms): {} ({:.2}ms)\n\
379             Error Rate Requirement (≤1%): {} ({:.2}%)\n\
380             \nRecommendations:",
381            status,
382            if self.meets_eps_requirement { "✅" } else { "❌" },
383            self.current_eps,
384            if self.meets_latency_requirement { "✅" } else { "❌" },
385            self.current_latency_ms,
386            if self.meets_error_rate_requirement { "✅" } else { "❌" },
387            self.current_error_rate
388        );
389
390        for rec in &self.recommendations {
391            result.push_str(&format!("\n• {rec}"));
392        }
393
394        result
395    }
396}
397
398/// Global metrics instance
399static GLOBAL_METRICS: OnceLock<Arc<AuditMetrics>> = OnceLock::new();
400
401/// Get or initialize the global metrics instance
402pub fn global_metrics() -> Arc<AuditMetrics> {
403    GLOBAL_METRICS.get_or_init(|| Arc::new(AuditMetrics::new())).clone()
404}
405
406/// Record a successful audit event dispatch
407pub fn record_audit_success(dispatch_time: Duration) {
408    global_metrics().record_event_success(dispatch_time);
409}
410
411/// Record a failed audit event dispatch
412pub fn record_audit_failure(dispatch_time: Duration) {
413    global_metrics().record_event_failure(dispatch_time);
414}
415
416/// Record a successful target operation
417pub fn record_target_success() {
418    global_metrics().record_target_success();
419}
420
421/// Record a failed target operation
422pub fn record_target_failure() {
423    global_metrics().record_target_failure();
424}
425
426/// Record a configuration reload
427pub fn record_config_reload() {
428    global_metrics().record_config_reload();
429}
430
431/// Record a system start
432pub fn record_system_start() {
433    global_metrics().record_system_start();
434}
435
436/// Get the current metrics report
437pub async fn get_metrics_report() -> AuditMetricsReport {
438    global_metrics().generate_report().await
439}
440
441/// Validate performance requirements
442pub async fn validate_performance() -> PerformanceValidation {
443    global_metrics().validate_performance_requirements().await
444}
445
446/// Reset all metrics
447pub async fn reset_metrics() {
448    global_metrics().reset().await;
449}