1use 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
50pub fn init_observability_metrics() {
53 static METRICS_DESC_INIT: OnceLock<()> = OnceLock::new();
54 METRICS_DESC_INIT.get_or_init(|| {
55 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 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#[derive(Debug)]
70pub struct AuditMetrics {
71 total_events_processed: AtomicU64,
73 total_events_failed: AtomicU64,
74 total_dispatch_time_ns: AtomicU64,
75
76 target_success_count: AtomicU64,
78 target_failure_count: AtomicU64,
79
80 config_reload_count: AtomicU64,
82 system_start_count: AtomicU64,
83
84 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 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 #[inline]
112 fn emit_event_success_metrics(&self, dispatch_time: Duration) {
113 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 #[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 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 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 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 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 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 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 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 gauge!(M_AUDIT_EPS).set(eps);
193 eps
194 }
195
196 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 } else {
204 0.0
205 }
206 }
207
208 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 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 }
230 }
231
232 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 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 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 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, current_eps: eps,
281 current_latency_ms: avg_latency_ms,
282 current_error_rate: error_rate,
283 recommendations: Vec::new(),
284 };
285
286 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#[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 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#[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 pub fn all_requirements_met(&self) -> bool {
368 self.meets_eps_requirement && self.meets_latency_requirement && self.meets_error_rate_requirement
369 }
370
371 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
398static GLOBAL_METRICS: OnceLock<Arc<AuditMetrics>> = OnceLock::new();
400
401pub fn global_metrics() -> Arc<AuditMetrics> {
403 GLOBAL_METRICS.get_or_init(|| Arc::new(AuditMetrics::new())).clone()
404}
405
406pub fn record_audit_success(dispatch_time: Duration) {
408 global_metrics().record_event_success(dispatch_time);
409}
410
411pub fn record_audit_failure(dispatch_time: Duration) {
413 global_metrics().record_event_failure(dispatch_time);
414}
415
416pub fn record_target_success() {
418 global_metrics().record_target_success();
419}
420
421pub fn record_target_failure() {
423 global_metrics().record_target_failure();
424}
425
426pub fn record_config_reload() {
428 global_metrics().record_config_reload();
429}
430
431pub fn record_system_start() {
433 global_metrics().record_system_start();
434}
435
436pub async fn get_metrics_report() -> AuditMetricsReport {
438 global_metrics().generate_report().await
439}
440
441pub async fn validate_performance() -> PerformanceValidation {
443 global_metrics().validate_performance_requirements().await
444}
445
446pub async fn reset_metrics() {
448 global_metrics().reset().await;
449}