#[cfg(test)]
mod tests {
use super::*;
use crate::{
QueueStats, QueueHealth, QueueHealthCheck,
redis::RedisQueueBuilder,
types::{QueueMessage, QueuePriority, MessageStatus}
};
use std::collections::HashMap;
use std::time::Duration;
use chrono::{Utc, DateTime};
async fn create_test_queue() -> Result<RedisQueue, Box<dyn std::error::Error>> {
let queue = RedisQueueBuilder::new()
.url("redis://localhost:6379")
.queue_name("monitoring_test_queue")
.key_prefix("test_monitoring")
.build()
.await?;
queue.purge().await?;
Ok(queue)
}
fn create_test_monitor(queue: RedisQueue) -> QueueMonitor {
let thresholds = AlertThresholds {
max_queue_size: 100,
max_error_rate: 5.0,
max_processing_time_ms: 1000.0,
min_messages_per_second: 1.0,
max_dead_letter_messages: 5,
};
QueueMonitor::new(Arc::new(queue), thresholds)
}
#[tokio::test]
async fn test_queue_metrics_creation() -> Result<(), Box<dyn std::error::Error>> {
let metrics = QueueMetrics::new();
assert_eq!(metrics.total_messages, 0);
assert_eq!(metrics.processed_messages, 0);
assert_eq!(metrics.failed_messages, 0);
assert_eq!(metrics.success_rate, 100.0);
assert_eq!(metrics.error_rate, 0.0);
assert_eq!(metrics.health_score(), 100.0);
assert!(metrics.is_healthy());
Ok(())
}
#[tokio::test]
async fn test_queue_metrics_update_from_stats() -> Result<(), Box<dyn std::error::Error>> {
let mut metrics = QueueMetrics::new();
let stats = QueueStats {
total_messages: 150,
visible_messages: 100,
invisible_messages: 50,
dead_letter_messages: 5,
total_processed: 140,
total_failed: 10,
..Default::default()
};
metrics.update_from_stats(&stats);
assert_eq!(metrics.total_messages, 150);
assert_eq!(metrics.processed_messages, 140);
assert_eq!(metrics.failed_messages, 10);
assert_eq!(metrics.success_rate, 93.3); assert_eq!(metrics.error_rate, 6.7);
Ok(())
}
#[tokio::test]
async fn test_queue_metrics_trends() -> Result<(), Box<dyn std::error::Error>> {
let mut metrics = QueueMetrics::new();
metrics.add_trend_point(&mut metrics.queue_depth_trend, 10.0);
metrics.add_trend_point(&mut metrics.queue_depth_trend, 20.0);
metrics.add_trend_point(&mut metrics.queue_depth_trend, 15.0);
assert_eq!(metrics.queue_depth_trend.len(), 3);
assert_eq!(metrics.queue_depth_trend[0].value, 10.0);
assert_eq!(metrics.queue_depth_trend[1].value, 20.0);
assert_eq!(metrics.queue_depth_trend[2].value, 15.0);
for i in 0..110 {
metrics.add_trend_point(&mut metrics.queue_depth_trend, i as f64);
}
assert_eq!(metrics.queue_depth_trend.len(), 100);
assert_eq!(metrics.queue_depth_trend[0].value, 10.0); assert_eq!(metrics.queue_depth_trend[99].value, 109.0);
Ok(())
}
#[tokio::test]
async fn test_queue_metrics_health_score() -> Result<(), Box<dyn std::error::Error>> {
let mut metrics = QueueMetrics::new();
assert_eq!(metrics.health_score(), 100.0);
metrics.error_rate = 5.0; assert_eq!(metrics.health_score(), 90.0);
metrics.total_messages = 1500; assert_eq!(metrics.health_score(), 80.0);
metrics.health_status = QueueHealth::Unhealthy;
assert_eq!(metrics.health_score(), 30.0);
assert!(!metrics.is_healthy());
Ok(())
}
#[tokio::test]
async fn test_alert_event_creation() -> Result<(), Box<dyn std::error::Error>> {
let alert = AlertEvent::new(
"Test Alert".to_string(),
"This is a test alert".to_string(),
AlertSeverity::Warning,
);
assert!(!alert.id.is_empty());
assert_eq!(alert.title, "Test Alert");
assert_eq!(alert.message, "This is a test alert");
assert_eq!(alert.severity, AlertSeverity::Warning);
assert!(alert.is_active());
assert!(alert.resolved_at.is_none());
let mut resolved_alert = alert;
resolved_alert.resolve();
assert!(resolved_alert.is_resolved());
Ok(())
}
#[tokio::test]
async fn test_alert_thresholds_default() -> Result<(), Box<dyn std::error::Error>> {
let thresholds = AlertThresholds::default();
assert_eq!(thresholds.max_queue_size, 1000);
assert_eq!(thresholds.max_error_rate, 5.0);
assert_eq!(thresholds.max_processing_time_ms, 5000.0);
assert_eq!(thresholds.min_messages_per_second, 1.0);
assert_eq!(thresholds.max_dead_letter_messages, 10);
Ok(())
}
#[tokio::test]
async fn test_queue_monitor_creation() -> Result<(), Box<dyn std::error::Error>> {
let queue = create_test_queue().await?;
let monitor = create_test_monitor(queue);
Ok(())
}
#[tokio::test]
async fn test_queue_monitor_metrics_update() -> Result<(), Box<dyn std::error::Error>> {
let queue = create_test_queue().await?;
let monitor = create_test_monitor(queue);
let message = QueueMessage::builder()
.payload("Test message")
.build();
queue.enqueue(message).await?;
monitor.update_metrics().await?;
let metrics = monitor.get_metrics().await;
assert!(metrics.total_messages >= 1);
Ok(())
}
#[tokio::test]
async fn test_queue_monitor_health_check() -> Result<(), Box<dyn std::error::Error>> {
let queue = create_test_queue().await?;
let monitor = create_test_monitor(queue);
let health = monitor.health_check().await?;
assert!(matches!(health.status, QueueHealth::Healthy | QueueHealth::Degraded));
Ok(())
}
#[tokio::test]
async fn test_queue_monitor_alerts_no_issues() -> Result<(), Box<dyn std::error::Error>> {
let queue = create_test_queue().await?;
let monitor = create_test_monitor(queue);
{
let mut metrics = monitor.metrics.write().await;
metrics.total_messages = 50;
metrics.processed_messages = 48;
metrics.failed_messages = 2;
metrics.messages_per_second = 5.0;
metrics.health_status = QueueHealth::Healthy;
}
let alerts = monitor.check_alerts().await;
assert!(alerts.is_empty());
Ok(())
}
#[tokio::test]
async fn test_queue_monitor_alerts_with_issues() -> Result<(), Box<dyn std::error::Error>> {
let queue = create_test_queue().await?;
let monitor = create_test_monitor(queue);
{
let mut metrics = monitor.metrics.write().await;
metrics.total_messages = 2000; metrics.processed_messages = 1600;
metrics.failed_messages = 400; metrics.messages_per_second = 0.5; metrics.health_status = QueueHealth::Degraded;
}
let alerts = monitor.check_alerts().await;
assert!(!alerts.is_empty());
let alert_titles: Vec<String> = alerts.iter().map(|a| a.title.clone()).collect();
assert!(alert_titles.contains(&"Queue Size Alert".to_string()));
assert!(alert_titles.contains(&"Error Rate Alert".to_string()));
assert!(alert_titles.contains(&"Low Processing Rate Alert".to_string()));
assert!(alert_titles.contains(&"Health Status Alert".to_string()));
Ok(())
}
#[tokio::test]
async fn test_queue_monitor_severity_levels() -> Result<(), Box<dyn std::error::Error>> {
let queue = create_test_queue().await?;
let monitor = create_test_monitor(queue);
{
let mut metrics = monitor.metrics.write().await;
metrics.total_messages = 100;
metrics.processed_messages = 50;
metrics.failed_messages = 50; }
let alerts = monitor.check_alerts().await;
let critical_alerts: Vec<_> = alerts.iter()
.filter(|a| a.severity == AlertSeverity::Critical)
.collect();
assert!(!critical_alerts.is_empty());
assert!(critical_alerts[0].title.contains("Error Rate Alert"));
Ok(())
}
#[tokio::test]
async fn test_console_alert_callback() -> Result<(), Box<dyn std::error::Error>> {
let callback = ConsoleAlertCallback;
let alert = AlertEvent::new(
"Test Alert".to_string(),
"Test message".to_string(),
AlertSeverity::Info,
);
callback.on_alert(&alert);
callback.on_alert_resolved(&alert);
Ok(())
}
#[tokio::test]
async fn test_webhook_alert_callback() -> Result<(), Box<dyn std::error::Error>> {
let callback = WebhookAlertCallback::new(
"https://example.com/webhook".to_string()
);
let alert = AlertEvent::new(
"Test Webhook Alert".to_string(),
"Test webhook message".to_string(),
AlertSeverity::Warning,
);
callback.on_alert(&alert);
callback.on_alert_resolved(&alert);
Ok(())
}
#[tokio::test]
async fn test_metrics_report_generation() -> Result<(), Box<dyn std::error::Error>> {
let queue = create_test_queue().await?;
let monitor = create_test_monitor(queue);
{
let mut metrics = monitor.metrics.write().await;
metrics.total_messages = 500;
metrics.processed_messages = 450;
metrics.failed_messages = 50;
metrics.messages_per_second = 10.0;
metrics.avg_processing_time_ms = 150.0;
metrics.health_status = QueueHealth::Healthy;
metrics.last_health_check = Some(Utc::now());
}
let report = monitor.generate_report().await;
assert_eq!(report.summary.total_messages, 500);
assert_eq!(report.summary.processed_messages, 450);
assert_eq!(report.summary.failed_messages, 50);
assert_eq!(report.performance.messages_per_second, 10.0);
assert_eq!(report.health.status, QueueHealth::Healthy);
Ok(())
}
#[tokio::test]
async fn test_timestamped_value() -> Result<(), Box<dyn std::error::Error>> {
let timestamp = Utc::now();
let value = TimestampedValue {
timestamp,
value: 42.5,
};
assert_eq!(value.value, 42.5);
assert!(value.timestamp <= Utc::now());
Ok(())
}
#[tokio::test]
async fn test_metrics_report_serialization() -> Result<(), Box<dyn std::error::Error>> {
let queue = create_test_queue().await?;
let monitor = create_test_monitor(queue);
let report = monitor.generate_report().await;
let json = serde_json::to_string(&report)?;
let deserialized: MetricsReport = serde_json::from_str(&json)?;
assert_eq!(deserialized.summary.total_messages, report.summary.total_messages);
assert_eq!(deserialized.performance.messages_per_second, report.performance.messages_per_second);
Ok(())
}
#[tokio::test]
async fn test_monitoring_integration() -> Result<(), Box<dyn std::error::Error>> {
let queue = create_test_queue().await?;
let monitor = create_test_monitor(queue);
for i in 0..5 {
let message = QueueMessage::builder()
.id(format!("test-msg-{}", i))
.payload(format!("Test message {}", i))
.priority(QueuePriority::Normal)
.build();
queue.enqueue(message).await?;
}
monitor.update_metrics().await?;
let metrics_after = monitor.get_metrics().await;
assert!(metrics_after.total_messages >= 5);
let health = monitor.health_check().await?;
assert!(matches!(health.status, QueueHealth::Healthy | QueueHealth::Degraded));
let alerts = monitor.check_alerts().await;
let report = monitor.generate_report().await;
assert_eq!(report.summary.total_messages, metrics_after.total_messages);
Ok(())
}
#[tokio::test]
async fn test_monitoring_with_alert_callback() -> Result<(), Box<dyn std::error::Error>> {
let queue = create_test_queue().await?;
let mut monitor = create_test_monitor(queue);
let alert_log = Arc::new(std::sync::Mutex::new(Vec::new()));
let alert_log_clone = alert_log.clone();
struct TestCallback {
alert_log: Arc<std::sync::Mutex<Vec<AlertEvent>>>,
}
impl backbone_queue::monitoring::AlertCallback for TestCallback {
fn on_alert(&self, alert: &AlertEvent) {
let mut log = self.alert_log.lock().unwrap();
log.push(alert.clone());
}
fn on_alert_resolved(&self, alert: &AlertEvent) {
let mut log = self.alert_log.lock().unwrap();
log.push(alert.clone());
}
}
monitor.add_alert_callback(Box::new(TestCallback {
alert_log: alert_log_clone,
}));
{
let mut metrics = monitor.metrics.write().await;
metrics.total_messages = 2000; }
let alerts = monitor.check_alerts().await;
let log = alert_log.lock().unwrap();
assert!(!log.is_empty());
assert_eq!(log.len(), alerts.len());
Ok(())
}
#[tokio::test]
async fn test_monitoring_continuous_updates() -> Result<(), Box<dyn std::error::Error>> {
let queue = create_test_queue().await?;
let monitor = create_test_monitor(queue);
let handle = monitor.start_monitoring(Duration::from_millis(100));
tokio::time::sleep(Duration::from_millis(250)).await;
for i in 0..3 {
let message = QueueMessage::builder()
.payload(format!("Continuous test {}", i))
.build();
queue.enqueue(message).await?;
}
tokio::time::sleep(Duration::from_millis(200)).await;
let metrics = monitor.get_metrics().await;
assert!(metrics.total_messages >= 3);
handle.abort();
Ok(())
}
}