#[cfg(test)]
mod tests {
use super::super::*;
use crate::types::{QueueMessage, QueuePriority, MessageStatus};
use crate::redis::{RedisQueue, RedisQueueConfig, RedisQueueBuilder};
use std::collections::HashMap;
use std::time::Duration;
use tokio_test;
use uuid::Uuid;
fn create_test_config() -> RedisQueueConfig {
RedisQueueConfig {
url: "redis://localhost:6379".to_string(),
queue_name: format!("test_queue_{}", uuid::Uuid::new_v4().to_string().replace("-", "")),
key_prefix: "test".to_string(),
pool_size: 1,
health_check_interval: 1,
}
}
fn create_test_message() -> QueueMessage {
QueueMessage::builder()
.text_payload("Test message payload")
.priority(QueuePriority::Normal)
.visibility_timeout(30)
.max_receive_count(3)
.build()
}
fn create_test_message_with_priority(priority: QueuePriority) -> QueueMessage {
QueueMessage::builder()
.text_payload(format!("Test message with {:?} priority", priority))
.priority(priority)
.visibility_timeout(30)
.max_receive_count(3)
.build()
}
async fn setup_test_queue() -> QueueResult<RedisQueue> {
let config = create_test_config();
let queue = RedisQueue::new(config).await;
if let Ok(ref q) = queue {
q.purge().await?;
}
queue
}
#[tokio::test]
async fn test_redis_queue_creation() -> QueueResult<()> {
let config = create_test_config();
let queue = RedisQueue::new(config).await?;
assert_eq!(queue.backend_type(), QueueBackend::Redis);
assert!(queue.test_connection().await?);
assert!(queue.validate_config().await?);
Ok(())
}
#[tokio::test]
async fn test_enqueue_single_message() -> QueueResult<()> {
let queue = setup_test_queue().await?;
let message = create_test_message();
let message_id = queue.enqueue(message.clone()).await?;
assert_eq!(message_id, message.id);
let size = queue.size().await?;
assert_eq!(size, 1);
assert!(!queue.is_empty().await?);
Ok(())
}
#[tokio::test]
async fn test_enqueue_batch_messages() -> QueueResult<()> {
let queue = setup_test_queue().await?;
let messages = vec![
create_test_message_with_priority(QueuePriority::Low),
create_test_message_with_priority(QueuePriority::Normal),
create_test_message_with_priority(QueuePriority::High),
create_test_message_with_priority(QueuePriority::Critical),
];
let message_ids = queue.enqueue_batch(messages.clone()).await?;
assert_eq!(message_ids.len(), messages.len());
for (i, id) in message_ids.iter().enumerate() {
assert_eq!(id, &messages[i].id);
}
let size = queue.size().await?;
assert_eq!(size, 4);
Ok(())
}
#[tokio::test]
async fn test_dequeue_message_priority_order() -> QueueResult<()> {
let queue = setup_test_queue().await?;
let messages = vec![
create_test_message_with_priority(QueuePriority::Low),
create_test_message_with_priority(QueuePriority::Critical),
create_test_message_with_priority(QueuePriority::Normal),
create_test_message_with_priority(QueuePriority::High),
];
queue.enqueue_batch(messages).await?;
let mut received_priorities = Vec::new();
for _ in 0..4 {
if let Some(message) = queue.dequeue().await? {
received_priorities.push(message.priority);
}
}
assert_eq!(received_priorities, vec![
QueuePriority::Critical,
QueuePriority::High,
QueuePriority::Normal,
QueuePriority::Low,
]);
assert!(queue.is_empty().await?);
Ok(())
}
#[tokio::test]
async fn test_dequeue_batch_messages() -> QueueResult<()> {
let queue = setup_test_queue().await?;
let mut messages = Vec::new();
for i in 0..5 {
messages.push(QueueMessage::builder()
.text_payload(format!("Message {}", i))
.priority(QueuePriority::Normal)
.build());
}
queue.enqueue_batch(messages).await?;
let batch_result = queue.dequeue_batch(3).await?;
assert_eq!(batch_result.messages.len(), 3);
assert_eq!(batch_result.requested, 3);
assert_eq!(batch_result.available, 3);
assert_eq!(batch_result.total_in_queue, 2);
let batch_result2 = queue.dequeue_batch(10).await?;
assert_eq!(batch_result2.messages.len(), 2);
assert_eq!(batch_result2.total_in_queue, 0);
Ok(())
}
#[tokio::test]
async fn test_message_acknowledgment() -> QueueResult<()> {
let queue = setup_test_queue().await?;
let message = create_test_message();
queue.enqueue(message.clone()).await?;
let dequeued = queue.dequeue().await?;
assert!(dequeued.is_some());
let ack_result = queue.ack(&message.id).await?;
assert!(ack_result);
assert!(queue.is_empty().await?);
Ok(())
}
#[tokio::test]
async fn test_message_negative_acknowledgment() -> QueueResult<()> {
let queue = setup_test_queue().await?;
let message = create_test_message();
queue.enqueue(message.clone()).await?;
let dequeued = queue.dequeue().await?;
assert!(dequeued.is_some());
let nack_result = queue.nack(&message.id, Some(5)).await?;
assert!(nack_result);
tokio::time::sleep(Duration::from_millis(100)).await;
let immediate_dequeue = queue.dequeue().await?;
assert!(immediate_dequeue.is_none());
tokio::time::sleep(Duration::from_secs(6)).await;
let delayed_dequeue = queue.dequeue().await?;
assert!(delayed_dequeue.is_some());
Ok(())
}
#[tokio::test]
async fn test_message_delete() -> QueueResult<()> {
let queue = setup_test_queue().await?;
let message = create_test_message();
queue.enqueue(message.clone()).await?;
let delete_result = queue.delete(&message.id).await?;
assert!(delete_result);
assert!(queue.is_empty().await?);
Ok(())
}
#[tokio::test]
async fn test_dead_letter_queue() -> QueueResult<()> {
let queue = setup_test_queue().await?;
let message = QueueMessage::builder()
.text_payload("Test dead letter")
.max_receive_count(2) .visibility_timeout(1)
.build();
queue.enqueue(message.clone()).await?;
for _ in 0..2 {
let dequeued = queue.dequeue().await?;
assert!(dequeued.is_some());
queue.nack(&message.id, None).await?;
}
let dequeued = queue.dequeue().await?;
assert!(dequeued.is_some());
queue.nack(&message.id, None).await?;
let final_dequeue = queue.dequeue().await?;
assert!(final_dequeue.is_none());
let dead_message = queue.get_message(&message.id).await?;
assert!(dead_message.is_some());
assert_eq!(dead_message.unwrap().status, MessageStatus::DeadLettered);
Ok(())
}
#[tokio::test]
async fn test_message_visibility_timeout() -> QueueResult<()> {
let queue = setup_test_queue().await?;
let message = QueueMessage::builder()
.text_payload("Test visibility timeout")
.visibility_timeout(1) .build();
queue.enqueue(message.clone()).await?;
let dequeued = queue.dequeue().await?;
assert!(dequeued.is_some());
let immediate_dequeue = queue.dequeue().await?;
assert!(immediate_dequeue.is_none());
tokio::time::sleep(Duration::from_secs(2)).await;
let delayed_dequeue = queue.dequeue().await?;
assert!(delayed_dequeue.is_some());
Ok(())
}
#[tokio::test]
async fn test_get_message_by_id() -> QueueResult<()> {
let queue = setup_test_queue().await?;
let message = create_test_message();
queue.enqueue(message.clone()).await?;
let retrieved = queue.get_message(&message.id).await?;
assert!(retrieved.is_some());
let retrieved_message = retrieved.unwrap();
assert_eq!(retrieved_message.id, message.id);
assert_eq!(retrieved_message.payload, message.payload);
let not_found = queue.get_message("non_existent_id").await?;
assert!(not_found.is_none());
Ok(())
}
#[tokio::test]
async fn test_queue_statistics() -> QueueResult<()> {
let queue = setup_test_queue().await?;
let stats = queue.get_stats().await?;
assert_eq!(stats.visible_messages, 0);
assert_eq!(stats.invisible_messages, 0);
assert_eq!(stats.total_messages, 0);
queue.enqueue_batch(vec![
create_test_message(),
create_test_message(),
create_test_message(),
]).await?;
let stats = queue.get_stats().await?;
assert_eq!(stats.visible_messages, 3);
assert_eq!(stats.total_messages, 3);
queue.dequeue().await?;
let stats = queue.get_stats().await?;
assert_eq!(stats.visible_messages, 2);
assert_eq!(stats.invisible_messages, 1);
assert_eq!(stats.total_messages, 3);
Ok(())
}
#[tokio::test]
async fn test_queue_purge() -> QueueResult<()> {
let queue = setup_test_queue().await?;
queue.enqueue_batch(vec![
create_test_message(),
create_test_message(),
create_test_message(),
]).await?;
assert!(!queue.is_empty().await?);
let purged_count = queue.purge().await?;
assert!(queue.is_empty().await?);
assert_eq!(queue.size().await?, 0);
Ok(())
}
#[tokio::test]
async fn test_queue_health_check() -> QueueResult<()> {
let queue = setup_test_queue().await?;
let health = queue.health_check().await?;
assert_eq!(health.status, QueueHealth::Healthy);
assert_eq!(health.queue_size, 0);
assert!(health.error_rate < 0.1);
Ok(())
}
#[tokio::test]
async fn test_message_attributes() -> QueueResult<()> {
let queue = setup_test_queue().await?;
let mut attributes = HashMap::new();
attributes.insert("source".to_string(), "test".to_string());
attributes.insert("version".to_string(), "1.0".to_string());
let message = QueueMessage::builder()
.text_payload("Message with attributes")
.attributes(attributes.clone())
.message_group_id("test-group".to_string())
.message_deduplication_id("test-dedup-123".to_string())
.build();
let message_id = queue.enqueue(message.clone()).await?;
let retrieved = queue.get_message(&message_id).await?;
assert!(retrieved.is_some());
let retrieved_message = retrieved.unwrap();
assert_eq!(retrieved_message.attributes, attributes);
assert_eq!(retrieved_message.message_group_id, Some("test-group".to_string()));
assert_eq!(retrieved_message.message_deduplication_id, Some("test-dedup-123".to_string()));
Ok(())
}
#[tokio::test]
async fn test_message_builder() -> QueueResult<()> {
let message = QueueMessage::builder()
.id("test-id-123")
.text_payload("Builder test message")
.priority(QueuePriority::High)
.max_receive_count(5)
.visibility_timeout(60)
.delay(10)
.expires_in(3600)
.attribute("test_key", "test_value")
.message_group_id("group-123")
.message_deduplication_id("dedup-123")
.compress(true)
.build();
assert_eq!(message.id, "test-id-123");
assert_eq!(message.priority, QueuePriority::High);
assert_eq!(message.max_receive_count, 5);
assert_eq!(message.visibility_timeout, 60);
assert_eq!(message.delay_seconds, Some(10));
assert!(message.expires_at.is_some());
assert!(message.attributes.contains_key("test_key"));
assert_eq!(message.message_group_id, Some("group-123".to_string()));
assert_eq!(message.message_deduplication_id, Some("dedup-123".to_string()));
assert!(message.compressed);
assert!(message.validate().is_ok());
let size = message.size_bytes();
assert!(size.is_ok());
assert!(size.unwrap() > 0);
Ok(())
}
#[tokio::test]
async fn test_message_validation() -> QueueResult<()> {
let valid_message = create_test_message();
assert!(valid_message.validate().is_ok());
let mut invalid_message = create_test_message();
invalid_message.id = "".to_string();
assert!(invalid_message.validate().is_err());
let mut invalid_message = create_test_message();
invalid_message.visibility_timeout = 0;
assert!(invalid_message.validate().is_err());
let mut invalid_message = create_test_message();
invalid_message.max_receive_count = 0;
assert!(invalid_message.validate().is_err());
Ok(())
}
#[tokio::test]
async fn test_message_lifecycle() -> QueueResult<()> {
let mut message = create_test_message();
assert_eq!(message.status, MessageStatus::Pending);
assert_eq!(message.receive_count, 0);
message.mark_received();
assert_eq!(message.status, MessageStatus::Processing);
assert_eq!(message.receive_count, 1);
assert!(message.visible_at > chrono::Utc::now());
message.mark_acknowledged();
assert_eq!(message.status, MessageStatus::Acknowledged);
message.reset_for_retry(Some(5));
assert_eq!(message.status, MessageStatus::Pending);
assert!(message.visible_at > chrono::Utc::now());
message.mark_failed();
assert_eq!(message.status, MessageStatus::Failed);
message.mark_dead_lettered();
assert_eq!(message.status, MessageStatus::DeadLettered);
Ok(())
}
#[tokio::test]
async fn test_queue_priority_from_i32() -> QueueResult<()> {
assert_eq!(QueuePriority::from(1), QueuePriority::Low);
assert_eq!(QueuePriority::from(5), QueuePriority::Normal);
assert_eq!(QueuePriority::from(10), QueuePriority::High);
assert_eq!(QueuePriority::from(20), QueuePriority::Critical);
assert_eq!(QueuePriority::from(999), QueuePriority::Normal);
assert_eq!(QueuePriority::from(-1), QueuePriority::Normal);
Ok(())
}
#[tokio::test]
async fn test_queue_priority_display() -> QueueResult<()> {
assert_eq!(QueuePriority::Low.to_string(), "Low");
assert_eq!(QueuePriority::Normal.to_string(), "Normal");
assert_eq!(QueuePriority::High.to_string(), "High");
assert_eq!(QueuePriority::Critical.to_string(), "Critical");
Ok(())
}
#[tokio::test]
async fn test_redis_queue_builder() -> QueueResult<()> {
let queue = RedisQueueBuilder::new()
.url("redis://localhost:6379")
.queue_name("test_builder_queue")
.key_prefix("builder_test")
.pool_size(5)
.build()
.await?;
assert_eq!(queue.config.queue_name, "test_builder_queue");
assert_eq!(queue.config.key_prefix, "builder_test");
assert_eq!(queue.config.pool_size, 5);
assert!(queue.test_connection().await?);
queue.purge().await?;
Ok(())
}
#[tokio::test]
async fn test_large_message_handling() -> QueueResult<()> {
let queue = setup_test_queue().await?;
let large_payload = "x".repeat(100000); let large_message = QueueMessage::builder()
.text_payload(large_payload)
.build();
let message_id = queue.enqueue(large_message).await?;
assert!(!message_id.is_empty());
let retrieved = queue.get_message(&message_id).await?;
assert!(retrieved.is_some());
assert_eq!(retrieved.unwrap().payload.as_str().unwrap().len(), 100000);
Ok(())
}
#[tokio::test]
async fn test_concurrent_operations() -> QueueResult<()> {
let queue = setup_test_queue().await?;
let mut handles = Vec::new();
for i in 0..10 {
let queue_clone = queue.clone();
let handle = tokio::spawn(async move {
let message = QueueMessage::builder()
.text_payload(format!("Concurrent message {}", i))
.build();
queue_clone.enqueue(message).await
});
handles.push(handle);
}
for handle in handles {
let _ = handle.await??;
}
let size = queue.size().await?;
assert_eq!(size, 10);
let mut handles = Vec::new();
for _ in 0..5 {
let queue_clone = queue.clone();
let handle = tokio::spawn(async move {
queue_clone.dequeue().await
});
handles.push(handle);
}
let mut dequeued_count = 0;
for handle in handles {
if handle.await?.is_some() {
dequeued_count += 1;
}
}
assert_eq!(dequeued_count, 5);
Ok(())
}
#[tokio::test]
async fn test_error_handling() -> QueueResult<()> {
let queue = setup_test_queue().await?;
let ack_result = queue.ack("non_existent_id").await?;
assert!(!ack_result);
let nack_result = queue.nack("non_existent_id", None).await?;
assert!(!nack_result);
let delete_result = queue.delete("non_existent_id").await?;
assert!(!delete_result);
Ok(())
}
}