use std::collections::HashMap;
use std::sync::Arc;
use std::time::{Duration, Instant};
use chrono::{DateTime, Utc};
use serde::{Serialize, Deserialize};
use tokio::sync::RwLock;
use tracing::{info, warn, error, debug};
use crate::{
QueueService, QueueStats, QueueHealth, QueueHealthCheck
};
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct QueueMetrics {
pub total_messages: u64,
pub processed_messages: u64,
pub failed_messages: u64,
pub dead_letter_messages: u64,
pub visible_messages: u64,
pub invisible_messages: u64,
pub avg_processing_time_ms: f64,
pub messages_per_second: f64,
pub error_rate: f64,
pub success_rate: f64,
pub queue_depth_trend: Vec<TimestampedValue>,
pub processing_rate_trend: Vec<TimestampedValue>,
pub error_rate_trend: Vec<TimestampedValue>,
pub last_health_check: Option<DateTime<Utc>>,
pub health_status: QueueHealth,
pub uptime_percentage: f64,
pub active_alerts: Vec<String>,
pub alert_history: Vec<AlertEvent>,
pub timestamp: DateTime<Utc>,
pub collection_interval: Duration,
}
impl Default for QueueMetrics {
fn default() -> Self {
Self {
total_messages: 0,
processed_messages: 0,
failed_messages: 0,
dead_letter_messages: 0,
visible_messages: 0,
invisible_messages: 0,
avg_processing_time_ms: 0.0,
messages_per_second: 0.0,
error_rate: 0.0,
success_rate: 100.0,
queue_depth_trend: Vec::new(),
processing_rate_trend: Vec::new(),
error_rate_trend: Vec::new(),
last_health_check: None,
health_status: QueueHealth::Healthy,
uptime_percentage: 100.0,
active_alerts: Vec::new(),
alert_history: Vec::new(),
timestamp: Utc::now(),
collection_interval: Duration::from_secs(30),
}
}
}
impl QueueMetrics {
pub fn update_from_stats(&mut self, stats: &QueueStats) {
self.total_messages = stats.total_messages;
self.processed_messages = stats.total_processed;
self.failed_messages = stats.total_failed;
self.dead_letter_messages = stats.dead_letter_messages;
self.visible_messages = stats.visible_messages;
self.invisible_messages = stats.invisible_messages;
let total_attempts = self.processed_messages + self.failed_messages;
if total_attempts > 0 {
self.success_rate = (self.processed_messages as f64 / total_attempts as f64) * 100.0;
self.error_rate = (self.failed_messages as f64 / total_attempts as f64) * 100.0;
}
}
pub fn add_trend_point(&mut self, trend: &mut Vec<TimestampedValue>, value: f64) {
trend.push(TimestampedValue {
timestamp: Utc::now(),
value,
});
if trend.len() > 100 {
trend.remove(0);
}
}
pub fn update_trends(&mut self) {
let queue_depth = self.total_messages as f64;
let processing_rate = self.messages_per_second;
let error_rate = self.error_rate;
Self::add_trend_point_static(&mut self.queue_depth_trend, queue_depth);
Self::add_trend_point_static(&mut self.processing_rate_trend, processing_rate);
Self::add_trend_point_static(&mut self.error_rate_trend, error_rate);
}
fn add_trend_point_static(trend: &mut Vec<TimestampedValue>, value: f64) {
trend.push(TimestampedValue {
timestamp: Utc::now(),
value,
});
if trend.len() > 100 {
trend.remove(0);
}
}
pub fn health_score(&self) -> f64 {
let mut score = 100.0;
score -= self.error_rate * 2.0;
if self.total_messages > 1000 {
score -= (self.total_messages - 1000) as f64 / 100.0;
}
match self.health_status {
QueueHealth::Healthy => score += 0.0,
QueueHealth::Degraded => score -= 20.0,
QueueHealth::Unhealthy => score -= 50.0,
}
score.clamp(0.0, 100.0)
}
pub fn is_healthy(&self) -> bool {
self.health_score() >= 70.0
&& self.error_rate < 10.0
&& self.health_status == QueueHealth::Healthy
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct TimestampedValue {
pub timestamp: DateTime<Utc>,
pub value: f64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub enum AlertSeverity {
Info,
Warning,
Error,
Critical,
}
impl std::fmt::Display for AlertSeverity {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
AlertSeverity::Info => write!(f, "INFO"),
AlertSeverity::Warning => write!(f, "WARNING"),
AlertSeverity::Error => write!(f, "ERROR"),
AlertSeverity::Critical => write!(f, "CRITICAL"),
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct AlertEvent {
pub id: String,
pub title: String,
pub message: String,
pub severity: AlertSeverity,
pub timestamp: DateTime<Utc>,
pub resolved_at: Option<DateTime<Utc>>,
pub metadata: HashMap<String, String>,
}
impl AlertEvent {
pub fn new(
title: String,
message: String,
severity: AlertSeverity,
) -> Self {
Self {
id: uuid::Uuid::new_v4().to_string(),
title,
message,
severity,
timestamp: Utc::now(),
resolved_at: None,
metadata: HashMap::new(),
}
}
pub fn resolve(&mut self) {
self.resolved_at = Some(Utc::now());
}
pub fn is_active(&self) -> bool {
self.resolved_at.is_none()
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct AlertThresholds {
pub max_queue_size: u64,
pub max_error_rate: f64,
pub max_processing_time_ms: f64,
pub min_messages_per_second: f64,
pub max_dead_letter_messages: u64,
}
impl Default for AlertThresholds {
fn default() -> Self {
Self {
max_queue_size: 1000,
max_error_rate: 5.0, max_processing_time_ms: 5000.0, min_messages_per_second: 1.0,
max_dead_letter_messages: 10,
}
}
}
pub struct QueueMonitorService {
queue: Arc<dyn QueueService + Send + Sync>,
metrics: Arc<RwLock<QueueMetrics>>,
thresholds: AlertThresholds,
alert_callbacks: Vec<Arc<dyn AlertCallback + Send + Sync>>,
}
pub trait AlertCallback: Send + Sync {
fn on_alert(&self, alert: &AlertEvent);
fn on_alert_resolved(&self, alert: &AlertEvent);
}
impl QueueMonitorService {
pub fn new(
queue: Arc<dyn QueueService + Send + Sync>,
thresholds: AlertThresholds,
) -> Self {
Self {
queue,
metrics: Arc::new(RwLock::new(QueueMetrics::default())),
thresholds,
alert_callbacks: Vec::new(),
}
}
pub fn add_alert_callback(&mut self, callback: Arc<dyn AlertCallback + Send + Sync>) {
self.alert_callbacks.push(callback);
}
pub async fn get_metrics(&self) -> QueueMetrics {
self.metrics.read().await.clone()
}
pub async fn update_metrics(&self) -> Result<(), crate::QueueError> {
let stats = self.queue.get_stats().await?;
{
let mut metrics = self.metrics.write().await;
metrics.update_from_stats(&stats);
metrics.update_trends();
}
Ok(())
}
pub async fn health_check(&self) -> Result<QueueHealthCheck, crate::QueueError> {
let health = self.queue.health_check().await?;
{
let mut metrics = self.metrics.write().await;
metrics.health_status = health.status;
metrics.last_health_check = Some(Utc::now());
}
Ok(health)
}
async fn check_alerts_static(
metrics: &QueueMetrics,
thresholds: &AlertThresholds,
) -> Vec<AlertEvent> {
let mut alerts = Vec::new();
if metrics.total_messages > thresholds.max_queue_size {
alerts.push(AlertEvent::new(
"Queue Size Alert".to_string(),
format!("Queue size {} exceeds threshold {}", metrics.total_messages, thresholds.max_queue_size),
AlertSeverity::Warning,
));
}
if metrics.error_rate > thresholds.max_error_rate {
let severity = if metrics.error_rate > thresholds.max_error_rate * 2.0 {
AlertSeverity::Critical
} else {
AlertSeverity::Error
};
alerts.push(AlertEvent::new(
"Error Rate Alert".to_string(),
format!("Error rate {:.2}% exceeds threshold {:.2}%", metrics.error_rate * 100.0, thresholds.max_error_rate * 100.0),
severity,
));
}
if metrics.avg_processing_time_ms > thresholds.max_processing_time_ms {
alerts.push(AlertEvent::new(
"Processing Time Alert".to_string(),
format!("Average processing time {:.2}ms exceeds threshold {:.2}ms", metrics.avg_processing_time_ms, thresholds.max_processing_time_ms),
AlertSeverity::Warning,
));
}
if metrics.messages_per_second < thresholds.min_messages_per_second {
alerts.push(AlertEvent::new(
"Low Processing Rate Alert".to_string(),
format!("Processing rate {:.2} msg/sec below threshold {:.2} msg/sec", metrics.messages_per_second, thresholds.min_messages_per_second),
AlertSeverity::Warning,
));
}
if metrics.dead_letter_messages > thresholds.max_dead_letter_messages {
alerts.push(AlertEvent::new(
"Dead Letter Alert".to_string(),
format!("Dead letter messages {} exceeds threshold {}", metrics.dead_letter_messages, thresholds.max_dead_letter_messages),
AlertSeverity::Error,
));
}
alerts
}
pub async fn check_alerts(&self) -> Vec<AlertEvent> {
let metrics = self.metrics.read().await.clone();
let mut alerts = Vec::new();
if metrics.total_messages > self.thresholds.max_queue_size {
alerts.push(AlertEvent::new(
"Queue Size Alert".to_string(),
format!("Queue size exceeded threshold: {} > {}",
metrics.total_messages, self.thresholds.max_queue_size),
AlertSeverity::Warning,
));
}
if metrics.error_rate > self.thresholds.max_error_rate {
let severity = if metrics.error_rate > 20.0 {
AlertSeverity::Critical
} else if metrics.error_rate > 10.0 {
AlertSeverity::Error
} else {
AlertSeverity::Warning
};
alerts.push(AlertEvent::new(
"Error Rate Alert".to_string(),
format!("Error rate exceeded threshold: {:.1}% > {:.1}%",
metrics.error_rate, self.thresholds.max_error_rate),
severity,
));
}
if metrics.avg_processing_time_ms > self.thresholds.max_processing_time_ms {
alerts.push(AlertEvent::new(
"Processing Time Alert".to_string(),
format!("Average processing time exceeded threshold: {:.1}ms > {:.1}ms",
metrics.avg_processing_time_ms, self.thresholds.max_processing_time_ms),
AlertSeverity::Warning,
));
}
if metrics.messages_per_second < self.thresholds.min_messages_per_second {
alerts.push(AlertEvent::new(
"Low Processing Rate Alert".to_string(),
format!("Processing rate below threshold: {:.2} < {:.2}",
metrics.messages_per_second, self.thresholds.min_messages_per_second),
AlertSeverity::Info,
));
}
if metrics.dead_letter_messages > self.thresholds.max_dead_letter_messages {
alerts.push(AlertEvent::new(
"Dead Letter Queue Alert".to_string(),
format!("Dead letter messages exceeded threshold: {} > {}",
metrics.dead_letter_messages, self.thresholds.max_dead_letter_messages),
AlertSeverity::Error,
));
}
if metrics.health_status != QueueHealth::Healthy {
alerts.push(AlertEvent::new(
"Health Status Alert".to_string(),
format!("Queue health status: {:?}", metrics.health_status),
AlertSeverity::Warning,
));
}
for alert in &alerts {
for callback in &self.alert_callbacks {
callback.on_alert(alert);
}
}
alerts
}
pub async fn start_monitoring(
&self,
interval: Duration,
) -> tokio::task::JoinHandle<()> {
let queue = self.queue.clone();
let metrics = self.metrics.clone();
let thresholds = self.thresholds.clone();
let _alert_callbacks = self.alert_callbacks.clone();
tokio::spawn(async move {
let mut last_stats: Option<QueueStats> = None;
let mut start_time = Instant::now();
loop {
tokio::time::sleep(interval).await;
if let Ok(current_stats) = queue.get_stats().await {
let now = Instant::now();
let time_delta = now.duration_since(start_time).as_secs_f64();
{
let mut metrics_guard = metrics.write().await;
metrics_guard.update_from_stats(¤t_stats);
metrics_guard.update_trends();
if let Some(ref last_stats) = last_stats {
let processed_delta = current_stats.total_processed.saturating_sub(last_stats.total_processed);
let messages_per_second = processed_delta as f64 / time_delta;
metrics_guard.messages_per_second = messages_per_second;
}
}
let alerts = {
let current_metrics = metrics.read().await;
Self::check_alerts_static(¤t_metrics, &thresholds).await
};
if !alerts.is_empty() {
warn!("Generated {} alerts", alerts.len());
for alert in &alerts {
info!("ALERT: {} - {}", alert.title, alert.message);
}
}
{
let mut metrics_guard = metrics.write().await;
metrics_guard.active_alerts = alerts.iter()
.filter(|a| a.is_active())
.map(|a| format!("{}: {}", a.severity, a.title))
.collect();
}
last_stats = Some(current_stats);
start_time = now;
}
}
})
}
pub async fn generate_report(&self) -> MetricsReport {
let metrics = self.metrics.read().await.clone();
MetricsReport {
summary: MetricsSummary {
total_messages: metrics.total_messages,
processed_messages: metrics.processed_messages,
failed_messages: metrics.failed_messages,
dead_letter_messages: metrics.dead_letter_messages,
success_rate: metrics.success_rate,
error_rate: metrics.error_rate,
health_score: metrics.health_score(),
is_healthy: metrics.is_healthy(),
},
performance: PerformanceMetrics {
messages_per_second: metrics.messages_per_second,
avg_processing_time_ms: metrics.avg_processing_time_ms,
queue_depth: metrics.total_messages,
visible_messages: metrics.visible_messages,
invisible_messages: metrics.invisible_messages,
},
trends: TrendMetrics {
queue_depth_trend: metrics.queue_depth_trend.clone(),
processing_rate_trend: metrics.processing_rate_trend.clone(),
error_rate_trend: metrics.error_rate_trend.clone(),
},
health: HealthMetrics {
status: metrics.health_status,
last_check: metrics.last_health_check,
uptime_percentage: metrics.uptime_percentage,
active_alerts_count: metrics.active_alerts.len(),
},
alerts: metrics.alert_history.clone(),
generated_at: Utc::now(),
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct MetricsReport {
pub summary: MetricsSummary,
pub performance: PerformanceMetrics,
pub trends: TrendMetrics,
pub health: HealthMetrics,
pub alerts: Vec<AlertEvent>,
pub generated_at: DateTime<Utc>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct MetricsSummary {
pub total_messages: u64,
pub processed_messages: u64,
pub failed_messages: u64,
pub dead_letter_messages: u64,
pub success_rate: f64,
pub error_rate: f64,
pub health_score: f64,
pub is_healthy: bool,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PerformanceMetrics {
pub messages_per_second: f64,
pub avg_processing_time_ms: f64,
pub queue_depth: u64,
pub visible_messages: u64,
pub invisible_messages: u64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct TrendMetrics {
pub queue_depth_trend: Vec<TimestampedValue>,
pub processing_rate_trend: Vec<TimestampedValue>,
pub error_rate_trend: Vec<TimestampedValue>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct HealthMetrics {
pub status: QueueHealth,
pub last_check: Option<DateTime<Utc>>,
pub uptime_percentage: f64,
pub active_alerts_count: usize,
}
pub struct ConsoleAlertCallback;
impl AlertCallback for ConsoleAlertCallback {
fn on_alert(&self, alert: &AlertEvent) {
match alert.severity {
AlertSeverity::Info => info!("📊 ALERT: {} - {}", alert.title, alert.message),
AlertSeverity::Warning => warn!("⚠️ ALERT: {} - {}", alert.title, alert.message),
AlertSeverity::Error => error!("❌ ALERT: {} - {}", alert.title, alert.message),
AlertSeverity::Critical => error!("🚨 ALERT: {} - {}", alert.title, alert.message),
}
}
fn on_alert_resolved(&self, alert: &AlertEvent) {
info!("✅ ALERT RESOLVED: {} - {}", alert.title, alert.message);
}
}
pub struct WebhookAlertCallback {
pub webhook_url: String,
pub client: reqwest::Client,
}
impl WebhookAlertCallback {
pub fn new(webhook_url: String) -> Self {
Self {
webhook_url,
client: reqwest::Client::new(),
}
}
}
impl AlertCallback for WebhookAlertCallback {
fn on_alert(&self, alert: &AlertEvent) {
let webhook_url = self.webhook_url.clone();
let client = self.client.clone();
let alert = alert.clone();
tokio::spawn(async move {
if client
.post(&webhook_url)
.json(&alert)
.send()
.await
.is_ok()
{
debug!("Alert webhook sent successfully");
} else {
warn!("Failed to send alert webhook");
}
});
}
fn on_alert_resolved(&self, alert: &AlertEvent) {
info!("Webhook alert resolved: {}", alert.title);
}
}