#[cfg(test)]
mod tests {
use super::super::*;
use crate::types::{QueueMessage, QueuePriority, MessageStatus};
use crate::sqs::{SqsQueue, SqsQueueConfig, SqsQueueBuilder};
use std::collections::HashMap;
use std::time::Duration;
use aws_sdk_sqs::primitives::Blob;
use uuid::Uuid;
fn create_test_config() -> SqsQueueConfig {
SqsQueueConfig {
queue_url: "https://sqs.us-east-1.amazonaws.com/123456789012/test-queue".to_string(),
region: "us-east-1".to_string(),
max_receive_count: 3,
visibility_timeout: 30,
wait_time_seconds: 5,
max_messages: 10,
message_retention_period: 345600, dead_letter_queue_url: None,
fifo_queue: false,
content_based_deduplication: false,
}
}
fn create_test_message() -> QueueMessage {
QueueMessage::builder()
.text_payload("Test SQS 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 SQS message with {:?} priority", priority))
.priority(priority)
.visibility_timeout(30)
.max_receive_count(3)
.build()
}
fn create_mock_sqs_client() -> aws_sdk_sqs::Client {
let config = aws_config::defaults(aws_config::BehaviorVersion::latest())
.region(aws_sdk_sqs::config::Region::new("us-east-1"))
.load();
aws_sdk_sqs::Client::new(&config)
}
async fn create_test_sqs_queue() -> QueueResult<SqsQueue> {
let mut config = create_test_config();
config.queue_url = format!("https://sqs.us-east-1.amazonaws.com/123456789012/test-queue-{}",
Uuid::new_v4().to_string().replace("-", ""));
SqsQueue::new(config).await
}
#[tokio::test]
async fn test_sqs_queue_creation() -> QueueResult<()> {
let config = create_test_config();
let queue = SqsQueue::new(config).await?;
assert_eq!(queue.backend_type(), QueueBackend::Sqs);
assert!(queue.validate_config().await?);
Ok(())
}
#[tokio::test]
async fn test_sqs_message_serialization() -> QueueResult<()> {
let queue = create_test_sqs_queue().await?;
let message = create_test_message();
let (body, attributes) = queue.convert_to_sqs_message_data(&message)?;
assert!(!body.is_empty());
assert!(body.len() > 100);
assert!(attributes.contains_key("Priority"));
assert!(attributes.contains_key("ReceiveCount"));
assert!(attributes.contains_key("MaxReceiveCount"));
assert!(attributes.contains_key("VisibilityTimeout"));
assert!(attributes.contains_key("EnqueuedAt"));
assert!(attributes.contains_key("ExpiresAt"));
if let Some(priority_attr) = attributes.get("Priority") {
assert_eq!(priority_attr.data_type(), Some("Number"));
assert_eq!(priority_attr.string_value(), Some("5")); }
Ok(())
}
#[tokio::test]
async fn test_sqs_message_deserialization() -> QueueResult<()> {
let queue = create_test_sqs_queue().await?;
let original_message = create_test_message_with_priority(QueuePriority::High);
let (body, attributes) = queue.convert_to_sqs_message_data(&original_message)?;
let sqs_message = aws_sdk_sqs::types::Message::builder()
.message_id(Uuid::new_v4().to_string())
.receipt_handle("test-receipt-handle".to_string())
.body(body)
.set_message_attributes(Some(attributes))
.build()
.map_err(|e| QueueError::SqsError(format!("Failed to build SQS message: {}", e)))?;
let converted_message = queue.convert_from_sqs_message(&sqs_message)?;
assert_eq!(converted_message.priority, QueuePriority::High);
assert_eq!(converted_message.payload, original_message.payload);
assert_eq!(converted_message.max_receive_count, original_message.max_receive_count);
assert_eq!(converted_message.visibility_timeout, original_message.visibility_timeout);
assert!(!converted_message.id.is_empty());
Ok(())
}
#[tokio::test]
async fn test_sqs_priority_attribute_conversion() -> QueueResult<()> {
let queue = create_test_sqs_queue().await?;
let priorities = vec![
(QueuePriority::Low, 1),
(QueuePriority::Normal, 5),
(QueuePriority::High, 10),
(QueuePriority::Critical, 20),
];
for (priority, expected_value) in priorities {
let message = QueueMessage::builder()
.text_payload(format!("Test message for {:?}", priority))
.priority(priority)
.build();
let (_body, attributes) = queue.convert_to_sqs_message_data(&message)?;
if let Some(priority_attr) = attributes.get("Priority") {
assert_eq!(priority_attr.string_value(), Some(&expected_value.to_string()));
assert_eq!(priority_attr.data_type(), Some("Number"));
} else {
return Err(QueueError::SqsError("Priority attribute not found".to_string()));
}
}
Ok(())
}
#[tokio::test]
async fn test_sqs_message_with_attributes() -> QueueResult<()> {
let queue = create_test_sqs_queue().await?;
let mut attributes = HashMap::new();
attributes.insert("source".to_string(), "test".to_string());
attributes.insert("version".to_string(), "1.0".to_string());
attributes.insert("user_id".to_string(), "user-123".to_string());
let message = QueueMessage::builder()
.text_payload("Message with custom attributes")
.attributes(attributes.clone())
.message_group_id("test-group".to_string())
.message_deduplication_id("test-dedup-123".to_string())
.compress(true)
.build();
let (body, sqs_attributes) = queue.convert_to_sqs_message_data(&message)?;
assert!(sqs_attributes.contains_key("source"));
assert!(sqs_attributes.contains_key("version"));
assert!(sqs_attributes.contains_key("user_id"));
assert!(sqs_attributes.contains_key("MessageGroupId"));
assert!(sqs_attributes.contains_key("MessageDeduplicationId"));
assert!(sqs_attributes.contains_key("Compressed"));
if let Some(source_attr) = sqs_attributes.get("source") {
assert_eq!(source_attr.string_value(), Some("test"));
}
Ok(())
}
#[tokio::test]
async fn test_sqs_queue_builder() -> QueueResult<()> {
let config = SqsQueueBuilder::new()
.queue_url("https://sqs.us-east-1.amazonaws.com/123456789012/builder-test-queue")
.region("us-west-2")
.max_receive_count(5)
.visibility_timeout(60)
.wait_time_seconds(10)
.max_messages(20)
.fifo_queue(true)
.content_based_deduplication(true)
.dead_letter_queue_url("https://sqs.us-east-1.amazonaws.com/123456789012/dead-letter-queue")
.build();
let queue = SqsQueue::new(config).await?;
assert_eq!(queue.config.queue_url, "https://sqs.us-east-1.amazonaws.com/123456789012/builder-test-queue");
assert_eq!(queue.config.region, "us-west-2");
assert_eq!(queue.config.max_receive_count, 5);
assert_eq!(queue.config.visibility_timeout, 60);
assert_eq!(queue.config.wait_time_seconds, 10);
assert_eq!(queue.config.max_messages, 20);
assert!(queue.config.fifo_queue);
assert!(queue.config.content_based_deduplication);
assert_eq!(queue.config.dead_letter_queue_url, Some("https://sqs.us-east-1.amazonaws.com/123456789012/dead-letter-queue".to_string()));
Ok(())
}
#[tokio::test]
async fn test_sqs_fifo_queue_configuration() -> QueueResult<()> {
let mut config = create_test_config();
config.fifo_queue = true;
config.content_based_deduplication = true;
config.queue_url = "https://sqs.us-east-1.amazonaws.com/123456789012/test-queue.fifo".to_string();
let queue = SqsQueue::new(config).await?;
assert!(queue.config.fifo_queue);
assert!(queue.config.content_based_deduplication);
assert!(queue.config.queue_url.ends_with(".fifo"));
let fifo_message = QueueMessage::builder()
.text_payload("FIFO test message")
.message_group_id("group-123")
.message_deduplication_id("dedup-456")
.build();
let (body, attributes) = queue.convert_to_sqs_message_data(&fifo_message)?;
assert!(attributes.contains_key("MessageGroupId"));
assert!(attributes.contains_key("MessageDeduplicationId"));
if let Some(group_id_attr) = attributes.get("MessageGroupId") {
assert_eq!(group_id_attr.string_value(), Some("group-123"));
}
Ok(())
}
#[tokio::test]
async fn test_sqs_dead_letter_queue_configuration() -> QueueResult<()> {
let dead_letter_url = "https://sqs.us-east-1.amazonaws.com/123456789012/test-dead-letter-queue".to_string();
let mut config = create_test_config();
config.dead_letter_queue_url = Some(dead_letter_url.clone());
let queue = SqsQueue::new(config).await?;
assert_eq!(queue.config.dead_letter_queue_url, Some(dead_letter_url));
let message = QueueMessage::builder()
.text_payload("Test dead letter message")
.max_receive_count(1) .receive_count(2) .build();
assert!(message.should_dead_letter());
Ok(())
}
#[tokio::test]
async fn test_sqs_message_expiration() -> QueueResult<()> {
let queue = create_test_sqs_queue().await?;
let message = QueueMessage::builder()
.text_payload("Message with expiration")
.expires_in(3600) .build();
assert!(!message.is_expired());
let (body, attributes) = queue.convert_to_sqs_message_data(&message)?;
assert!(attributes.contains_key("ExpiresAt"));
let expired_message = QueueMessage::builder()
.text_payload("Expired message")
.expires_at(chrono::Utc::now() - chrono::Duration::hours(1))
.build();
assert!(expired_message.is_expired());
Ok(())
}
#[tokio::test]
async fn test_sqs_large_message_handling() -> QueueResult<()> {
let queue = create_test_sqs_queue().await?;
let large_payload = "x".repeat(100000); let large_message = QueueMessage::builder()
.text_payload(large_payload.clone())
.build();
let (body, attributes) = queue.convert_to_sqs_message_data(&large_message)?;
assert!(body.len() > large_payload.len()); assert!(attributes.contains_key("OriginalSize"));
if let Some(size_attr) = attributes.get("OriginalSize") {
assert_eq!(size_attr.string_value(), Some(&large_payload.len().to_string()));
}
Ok(())
}
#[tokio::test]
async fn test_sqs_message_visibility_calculation() -> QueueResult<()> {
let message = QueueMessage::builder()
.text_payload("Visibility test message")
.visibility_timeout(120) .build();
assert!(message.is_visible());
let mut received_message = message;
received_message.mark_received();
assert!(!received_message.is_visible());
let now = chrono::Utc::now();
received_message.visible_at = now - chrono::Duration::seconds(1);
assert!(received_message.is_visible());
Ok(())
}
#[tokio::test]
async fn test_sqs_message_receive_count_tracking() -> QueueResult<()> {
let mut message = QueueMessage::builder()
.text_payload("Receive count test")
.max_receive_count(3)
.build();
assert_eq!(message.receive_count, 0);
assert_eq!(message.status, MessageStatus::Pending);
assert!(!message.should_dead_letter());
message.mark_received();
assert_eq!(message.receive_count, 1);
assert_eq!(message.status, MessageStatus::Processing);
assert!(!message.should_dead_letter());
message.reset_for_retry(None);
assert_eq!(message.receive_count, 1); assert_eq!(message.status, MessageStatus::Pending);
assert!(!message.should_dead_letter());
message.mark_received();
assert_eq!(message.receive_count, 2);
assert!(!message.should_dead_letter());
message.mark_received();
assert_eq!(message.receive_count, 3);
assert!(message.should_dead_letter());
}
#[tokio::test]
async fn test_sqs_error_handling() -> QueueResult<()> {
let queue = create_test_sqs_queue().await?;
let invalid_message = QueueMessage::builder()
.text_payload("")
.build();
let result = queue.convert_to_sqs_message_data(&invalid_message);
assert!(result.is_ok());
let mut message = create_test_message();
message.attributes.insert("invalid".to_string(), "{invalid json}".to_string());
let result = queue.convert_to_sqs_message_data(&message);
assert!(result.is_ok());
Ok(())
}
#[tokio::test]
async fn test_sqs_batch_message_processing() -> QueueResult<()> {
let queue = create_test_sqs_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 mut sqs_messages = Vec::new();
for message in &messages {
let (body, attributes) = queue.convert_to_sqs_message_data(message)?;
let sqs_message = aws_sdk_sqs::types::Message::builder()
.message_id(Uuid::new_v4().to_string())
.receipt_handle(format!("receipt-{}", Uuid::new_v4().to_string()))
.body(body)
.set_message_attributes(Some(attributes))
.build()
.map_err(|e| QueueError::SqsError(format!("Failed to build SQS message: {}", e)))?;
sqs_messages.push(sqs_message);
}
assert_eq!(sqs_messages.len(), messages.len());
let mut converted_messages = Vec::new();
for sqs_message in &sqs_messages {
let queue_message = queue.convert_from_sqs_message(sqs_message)?;
converted_messages.push(queue_message);
}
assert_eq!(converted_messages.len(), messages.len());
let priorities: Vec<_> = converted_messages.iter().map(|m| m.priority).collect();
assert_eq!(priorities, vec![
QueuePriority::Low,
QueuePriority::Normal,
QueuePriority::High,
QueuePriority::Critical,
]);
Ok(())
}
#[tokio::test]
async fn test_sqs_message_size_calculation() -> QueueResult<()> {
let queue = create_test_sqs_queue().await;
let message = create_test_message();
let size = message.size_bytes();
assert!(size.is_ok());
assert!(size.unwrap() > 100);
let compressed_message = QueueMessage::builder()
.text_payload("Compressed message")
.compress(true)
.original_size(Some(1000))
.build();
assert!(compressed_message.compressed);
assert_eq!(compressed_message.original_size, Some(1000));
Ok(())
}
#[tokio::test]
async fn test_sqs_configuration_validation() -> QueueResult<()> {
let valid_config = create_test_config();
let queue = SqsQueue::new(valid_config).await?;
assert!(queue.validate_config().await?);
let mut invalid_fifo_config = create_test_config();
invalid_fifo_config.fifo_queue = true;
invalid_fifo_config.queue_url = "https://sqs.us-east-1.amazonaws.com/123456789012/regular-queue".to_string();
let invalid_queue = SqsQueue::new(invalid_fifo_config).await?;
let result = invalid_queue.validate_config().await;
assert!(result.is_ok());
let mut invalid_visibility_config = create_test_config();
invalid_visibility_config.visibility_timeout = 0;
let invalid_queue = SqsQueue::new(invalid_visibility_config).await?;
let result = invalid_queue.validate_config().await;
assert!(result.is_ok());
Ok(())
}
#[tokio::test]
async fn test_sqs_queue_health_metrics() -> QueueResult<()> {
let queue = create_test_sqs_queue().await?;
let health = queue.health_check().await?;
assert!(matches!(health.status, QueueHealth::Healthy | QueueHealth::Degraded | QueueHealth::Unhealthy));
assert!(health.error_rate >= 0.0);
assert!(health.error_rate <= 1.0);
let connection_ok = queue.test_connection().await?;
Ok(())
}
}