pub mod redis;
pub mod sqs;
pub mod rabbitmq_simple;
pub mod traits;
pub mod types;
pub mod monitoring;
pub mod compression;
pub mod fifo;
pub mod queue_manager;
pub mod deduplication;
pub mod processor;
#[cfg(test)]
mod rabbitmq_tests;
pub use traits::*;
pub use types::*;
pub use redis::*;
pub use sqs::*;
pub use rabbitmq_simple::*;
pub use monitoring::{
QueueMetrics, QueueMonitorService, AlertEvent, AlertSeverity, AlertThresholds,
TimestampedValue, AlertCallback, ConsoleAlertCallback,
WebhookAlertCallback, MetricsReport
};
pub use compression::{
MessageCompressor, CompressionConfig, CompressionAlgorithm, CompressionStats,
CompressedMessageBuilder, utils as compression_utils
};
pub use fifo::{
FifoQueueService, FifoQueueServiceWrapper, FifoQueueConfig, FifoQueueStats,
MessageGroupStats, utils as fifo_utils, MessageVolume
};
pub use queue_manager::{
QueueManager, QueueConfig as QueueManagerConfig, QueueAdminService, MaintenanceAction, MaintenanceResult
};
pub use types::QueueHealthCheck;
pub use deduplication::{
MessageDeduplicator, DeduplicationConfig, DeduplicationStrategy, DeduplicationCache,
DeduplicationEntry, ProcessingStatus, ExactlyOnceRecord, ProcessingResult,
DeduplicationStats, DeduplicationCacheBackend, ExactlyOnceStorage,
MemoryDeduplicationCache, MemoryExactlyOnceStorage
};
pub use processor::{
MessageProcessor, ProcessingOutcome, ProcessedMessage, BatchProcessingResult,
ProcessingContext, BatchContext, BatchConfig, BatchTimeoutPolicy, RetryConfig,
RetryPolicy, RetryHandler, ProcessorStats, BatchingProcessor, SimpleMessageProcessor
};
pub const VERSION: &str = env!("CARGO_PKG_VERSION");
pub const DEFAULT_QUEUE_NAME: &str = "default";
pub const MAX_MESSAGE_SIZE: usize = 256 * 1024;
pub const DEFAULT_VISIBILITY_TIMEOUT: u64 = 30;
pub const DEFAULT_MAX_RECEIVE_COUNT: u32 = 5;
#[derive(thiserror::Error, Debug)]
pub enum QueueError {
#[error("Redis connection error: {0}")]
RedisConnection(String),
#[error("Redis operation error: {0}")]
RedisOperation(String),
#[error("AWS SQS error: {0}")]
SqsError(String),
#[error("Message serialization error: {0}")]
Serialization(String),
#[error("Message deserialization error: {0}")]
Deserialization(String),
#[error("Message too large: {size} bytes (max: {max} bytes)")]
MessageTooLarge { size: usize, max: usize },
#[error("Invalid queue name: {0}")]
InvalidQueueName(String),
#[error("Invalid message ID: {0}")]
InvalidMessageId(String),
#[error("Queue not found: {0}")]
QueueNotFound(String),
#[error("Resource not found: {0}")]
NotFound(String),
#[error("Configuration error: {0}")]
ConfigError(String),
#[error("Network error: {0}")]
NetworkError(String),
#[error("Connection error: {0}")]
ConnectionError(String),
#[error("Publish error: {0}")]
PublishError(String),
#[error("Consumer error: {0}")]
ConsumerError(String),
#[error("Queue error: {0}")]
Other(String),
}
pub type QueueResult<T> = Result<T, QueueError>;
#[derive(Debug, Clone)]
pub struct QueueConfig {
pub queue_name: String,
pub visibility_timeout: u64,
pub message_retention_period: Option<u64>,
pub max_receive_count: u32,
pub dead_letter_queue: Option<String>,
pub enable_compression: bool,
pub compression_threshold: usize,
pub default_priority: QueuePriority,
pub enable_batch_operations: bool,
pub batch_size: usize,
}
impl Default for QueueConfig {
fn default() -> Self {
Self {
queue_name: DEFAULT_QUEUE_NAME.to_string(),
visibility_timeout: DEFAULT_VISIBILITY_TIMEOUT,
message_retention_period: None,
max_receive_count: DEFAULT_MAX_RECEIVE_COUNT,
dead_letter_queue: None,
enable_compression: false,
compression_threshold: 1024,
default_priority: QueuePriority::Normal,
enable_batch_operations: true,
batch_size: 10,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum QueueBackend {
Redis,
Sqs,
RabbitMQ,
}
impl QueueBackend {
pub fn name(&self) -> &'static str {
match self {
Self::Redis => "redis",
Self::Sqs => "sqs",
Self::RabbitMQ => "rabbitmq",
}
}
}
#[derive(Debug, Clone)]
#[derive(Default)]
pub struct QueueStats {
pub total_messages: u64,
pub visible_messages: u64,
pub invisible_messages: u64,
pub delayed_messages: u64,
pub dead_letter_messages: u64,
pub avg_processing_time_ms: Option<f64>,
pub messages_per_second: Option<f64>,
pub total_processed: u64,
pub total_failed: u64,
pub queue_age_seconds: Option<u64>,
pub backend_stats: std::collections::HashMap<String, serde_json::Value>,
}
impl QueueStats {
pub fn success_rate(&self) -> f64 {
let total = self.total_processed + self.total_failed;
if total > 0 {
self.total_processed as f64 / total as f64
} else {
0.0
}
}
pub fn failure_rate(&self) -> f64 {
1.0 - self.success_rate()
}
}