use std::collections::HashMap;
use std::sync::Arc;
use std::time::Instant;
use chrono::{DateTime, Utc};
use serde::{Serialize, Deserialize};
use tokio::sync::RwLock;
use crate::{
QueueService, QueueMessage, QueueResult, QueueError,
};
use async_trait::async_trait;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct FifoQueueConfig {
pub enabled: bool,
pub deduplication_window_seconds: u64,
pub max_message_groups: usize,
pub enable_content_deduplication: bool,
pub content_deduplication_window_seconds: u64,
pub max_deduplicated_messages: usize,
}
impl Default for FifoQueueConfig {
fn default() -> Self {
Self {
enabled: true,
deduplication_window_seconds: 300, max_message_groups: 10000,
enable_content_deduplication: false,
content_deduplication_window_seconds: 60, max_deduplicated_messages: 100000,
}
}
}
#[derive(Debug, Clone, Default)]
pub struct MessageGroupStats {
pub message_count: u64,
pub processed_count: u64,
pub failed_count: u64,
pub avg_processing_time_ms: f64,
pub last_message_at: Option<DateTime<Utc>>,
pub created_at: DateTime<Utc>,
}
impl MessageGroupStats {
pub fn new() -> Self {
Self {
created_at: Utc::now(),
..Default::default()
}
}
pub fn success_rate(&self) -> f64 {
let total = self.processed_count + self.failed_count;
if total > 0 {
self.processed_count as f64 / total as f64 * 100.0
} else {
0.0
}
}
}
#[derive(Debug, Clone, Default)]
pub struct FifoQueueStats {
pub total_groups: u64,
pub active_groups: u64,
pub deduplicated_messages: u64,
pub out_of_order_messages: u64,
pub group_stats: HashMap<String, MessageGroupStats>,
pub deduplication_cache: HashMap<String, DateTime<Utc>>,
pub content_deduplication_cache: HashMap<String, DateTime<Utc>>,
pub last_updated: DateTime<Utc>,
}
impl FifoQueueStats {
pub fn new() -> Self {
Self {
last_updated: Utc::now(),
..Default::default()
}
}
pub fn update_group(&mut self, group_id: &str, processed: bool, processing_time_ms: u64) {
let stats = self.group_stats
.entry(group_id.to_string())
.or_default();
stats.message_count += 1;
stats.last_message_at = Some(Utc::now());
if processed {
stats.processed_count += 1;
} else {
stats.failed_count += 1;
}
let total_processed = stats.processed_count + stats.failed_count;
if total_processed > 0 {
stats.avg_processing_time_ms =
(stats.avg_processing_time_ms * (total_processed - 1) as f64 + processing_time_ms as f64)
/ total_processed as f64;
}
self.last_updated = Utc::now();
}
pub fn cleanup_expired_entries(&mut self, config: &FifoQueueConfig) {
let now = Utc::now();
self.deduplication_cache.retain(|_, timestamp| {
now.signed_duration_since(*timestamp).num_seconds() < config.deduplication_window_seconds as i64
});
self.content_deduplication_cache.retain(|_, timestamp| {
now.signed_duration_since(*timestamp).num_seconds() < config.content_deduplication_window_seconds as i64
});
if self.deduplication_cache.len() > config.max_deduplicated_messages {
let entries_to_remove = self.deduplication_cache.len() - config.max_deduplicated_messages;
let keys_to_remove: Vec<_> = self.deduplication_cache.keys().take(entries_to_remove).cloned().collect();
for key in keys_to_remove {
self.deduplication_cache.remove(&key);
}
}
if self.content_deduplication_cache.len() > config.max_deduplicated_messages {
let entries_to_remove = self.content_deduplication_cache.len() - config.max_deduplicated_messages;
let keys_to_remove: Vec<_> = self.content_deduplication_cache.keys().take(entries_to_remove).cloned().collect();
for key in keys_to_remove {
self.content_deduplication_cache.remove(&key);
}
}
if self.group_stats.len() > config.max_message_groups {
let groups_to_remove = self.group_stats.len() - config.max_message_groups;
let keys_to_remove: Vec<_> = self.group_stats.keys().take(groups_to_remove).cloned().collect();
for key in keys_to_remove {
self.group_stats.remove(&key);
}
}
}
}
#[async_trait]
pub trait FifoQueueService: Send + Sync {
async fn enqueue_fifo(&self, message: QueueMessage) -> QueueResult<String>;
async fn enqueue_fifo_batch(&self, messages: Vec<QueueMessage>) -> QueueResult<Vec<String>>;
async fn dequeue_from_group(&self, group_id: &str) -> QueueResult<Option<QueueMessage>>;
async fn get_group_stats(&self, group_id: &str) -> QueueResult<Option<MessageGroupStats>>;
async fn get_all_group_stats(&self) -> QueueResult<HashMap<String, MessageGroupStats>>;
async fn get_fifo_stats(&self) -> QueueResult<FifoQueueStats>;
async fn is_message_duplicated(&self, deduplication_id: &str) -> QueueResult<bool>;
async fn is_content_duplicated(&self, content: &str) -> QueueResult<bool>;
async fn cleanup_deduplication(&self) -> QueueResult<u64>;
}
pub struct FifoQueueServiceWrapper {
inner: Arc<dyn QueueService + Send + Sync>,
config: FifoQueueConfig,
stats: Arc<RwLock<FifoQueueStats>>,
}
impl FifoQueueServiceWrapper {
pub fn new(
queue_service: Arc<dyn QueueService + Send + Sync>,
config: FifoQueueConfig,
) -> Self {
Self {
inner: queue_service,
config,
stats: Arc::new(RwLock::new(FifoQueueStats::new())),
}
}
pub fn with_default_config(queue_service: Arc<dyn QueueService + Send + Sync>) -> Self {
Self::new(queue_service, FifoQueueConfig::default())
}
fn generate_content_hash(content: &serde_json::Value) -> String {
use std::collections::hash_map::DefaultHasher;
use std::hash::{Hash, Hasher};
let content_str = serde_json::to_string(content).unwrap_or_default();
let mut hasher = DefaultHasher::new();
content_str.hash(&mut hasher);
format!("{:x}", hasher.finish())
}
fn validate_fifo_message(&self, message: &QueueMessage) -> QueueResult<()> {
if !self.config.enabled {
return Err(QueueError::ConfigError("FIFO is disabled".to_string()));
}
if message.message_group_id.is_none() {
return Err(QueueError::ConfigError(
"FIFO messages must have message_group_id".to_string()
));
}
if message.message_deduplication_id.is_none() {
return Err(QueueError::ConfigError(
"FIFO messages must have message_deduplication_id".to_string()
));
}
let dedup_id = message.message_deduplication_id.as_ref().unwrap();
if dedup_id.is_empty() {
return Err(QueueError::ConfigError(
"message_deduplication_id cannot be empty".to_string()
));
}
if dedup_id.len() > 128 {
return Err(QueueError::ConfigError(
"message_deduplication_id cannot exceed 128 characters".to_string()
));
}
Ok(())
}
async fn generate_sequence_number(&self, group_id: &str) -> QueueResult<u64> {
let timestamp = Utc::now().timestamp_nanos_opt().unwrap_or(0) as u64;
let group_hash = {
use std::collections::hash_map::DefaultHasher;
use std::hash::{Hash, Hasher};
let mut hasher = DefaultHasher::new();
group_id.hash(&mut hasher);
hasher.finish()
};
Ok(timestamp.wrapping_add(group_hash))
}
}
#[async_trait]
impl FifoQueueService for FifoQueueServiceWrapper {
async fn enqueue_fifo(&self, mut message: QueueMessage) -> QueueResult<String> {
self.validate_fifo_message(&message)?;
if let Some(ref dedup_id) = message.message_deduplication_id {
if self.is_message_duplicated(dedup_id).await? {
return Err(QueueError::Other(
format!("Message with deduplication ID {} already exists", dedup_id)
));
}
}
if self.config.enable_content_deduplication {
let content_hash = Self::generate_content_hash(&message.payload);
if self.is_content_duplicated(&content_hash).await? {
return Err(QueueError::Other(
"Message with identical content already exists".to_string()
));
}
message.attributes.insert(
"content_hash".to_string(),
content_hash
);
}
let group_id = message.message_group_id.as_ref().unwrap().clone();
let deduplication_id = message.message_deduplication_id.clone();
let sequence_number = self.generate_sequence_number(&group_id).await?;
message.attributes.insert(
"fifo_sequence".to_string(),
sequence_number.to_string()
);
message.attributes.insert(
"fifo_group".to_string(),
group_id.clone()
);
message.attributes.insert(
"fifo_deduplication_id".to_string(),
deduplication_id.as_ref().unwrap().clone()
);
let message_id = self.inner.enqueue(message).await?;
{
let mut stats = self.stats.write().await;
stats.total_groups += 1;
if let Some(dedup_id) = deduplication_id {
stats.deduplication_cache.insert(
dedup_id.clone(),
Utc::now()
);
}
stats.update_group(&group_id, true, 0); }
Ok(message_id)
}
async fn enqueue_fifo_batch(&self, messages: Vec<QueueMessage>) -> QueueResult<Vec<String>> {
let mut message_ids = Vec::with_capacity(messages.len());
for message in &messages {
self.validate_fifo_message(message)?;
}
let mut sorted_messages = messages;
sorted_messages.sort_by(|a, b| {
let group_a = a.message_group_id.as_ref().unwrap();
let group_b = b.message_group_id.as_ref().unwrap();
group_a.cmp(group_b)
});
for message in sorted_messages {
let message_id = self.enqueue_fifo(message).await?;
message_ids.push(message_id);
}
Ok(message_ids)
}
async fn dequeue_from_group(&self, group_id: &str) -> QueueResult<Option<QueueMessage>> {
let start_time = Instant::now();
let max_attempts = 50; let mut attempts = 0;
loop {
if attempts >= max_attempts {
return Ok(None); }
attempts += 1;
let message = match self.inner.dequeue().await? {
Some(msg) => msg,
None => return Ok(None), };
if let Some(fifo_group) = message.attributes.get("fifo_group") {
if fifo_group == group_id {
let processing_time = start_time.elapsed().as_millis() as u64;
{
let mut stats = self.stats.write().await;
stats.update_group(group_id, true, processing_time);
}
return Ok(Some(message));
} else {
self.inner.enqueue(message).await?;
continue;
}
} else {
return Ok(None);
}
}
}
async fn get_group_stats(&self, group_id: &str) -> QueueResult<Option<MessageGroupStats>> {
let stats = self.stats.read().await;
Ok(stats.group_stats.get(group_id).cloned())
}
async fn get_all_group_stats(&self) -> QueueResult<HashMap<String, MessageGroupStats>> {
let stats = self.stats.read().await;
Ok(stats.group_stats.clone())
}
async fn get_fifo_stats(&self) -> QueueResult<FifoQueueStats> {
let mut stats = self.stats.write().await;
stats.cleanup_expired_entries(&self.config);
Ok(stats.clone())
}
async fn is_message_duplicated(&self, deduplication_id: &str) -> QueueResult<bool> {
let stats = self.stats.read().await;
Ok(stats.deduplication_cache.contains_key(deduplication_id))
}
async fn is_content_duplicated(&self, content: &str) -> QueueResult<bool> {
let content_hash = {
use std::collections::hash_map::DefaultHasher;
use std::hash::{Hash, Hasher};
let mut hasher = DefaultHasher::new();
content.hash(&mut hasher);
format!("{:x}", hasher.finish())
};
let stats = self.stats.read().await;
Ok(stats.content_deduplication_cache.contains_key(&content_hash))
}
async fn cleanup_deduplication(&self) -> QueueResult<u64> {
let mut stats = self.stats.write().await;
let now = Utc::now();
let initial_count = stats.deduplication_cache.len() + stats.content_deduplication_cache.len();
stats.deduplication_cache.retain(|_, timestamp| {
now.signed_duration_since(*timestamp).num_seconds() < self.config.deduplication_window_seconds as i64
});
stats.content_deduplication_cache.retain(|_, timestamp| {
now.signed_duration_since(*timestamp).num_seconds() < self.config.content_deduplication_window_seconds as i64
});
let final_count = stats.deduplication_cache.len() + stats.content_deduplication_cache.len();
Ok((initial_count - final_count) as u64)
}
}
pub mod utils {
use super::*;
pub fn validate_config(config: &FifoQueueConfig) -> Vec<String> {
let mut errors = Vec::new();
if config.deduplication_window_seconds == 0 {
errors.push("Deduplication window must be greater than 0 seconds".to_string());
}
if config.deduplication_window_seconds > 86400 * 7 { errors.push("Deduplication window cannot exceed 7 days".to_string());
}
if config.max_message_groups == 0 {
errors.push("Max message groups must be greater than 0".to_string());
}
if config.max_message_groups > 100000 {
errors.push("Max message groups cannot exceed 100000".to_string());
}
if config.content_deduplication_window_seconds == 0 && config.enable_content_deduplication {
errors.push("Content deduplication window must be greater than 0 seconds when enabled".to_string());
}
if config.max_deduplicated_messages == 0 {
errors.push("Max deduplicated messages must be greater than 0".to_string());
}
errors
}
pub fn get_recommended_config(message_volume: MessageVolume) -> FifoQueueConfig {
match message_volume {
MessageVolume::Low => FifoQueueConfig {
deduplication_window_seconds: 300, max_message_groups: 1000,
enable_content_deduplication: false,
content_deduplication_window_seconds: 60,
max_deduplicated_messages: 10000,
..Default::default()
},
MessageVolume::Medium => FifoQueueConfig {
deduplication_window_seconds: 900, max_message_groups: 10000,
enable_content_deduplication: true,
content_deduplication_window_seconds: 300, max_deduplicated_messages: 50000,
..Default::default()
},
MessageVolume::High => FifoQueueConfig {
deduplication_window_seconds: 3600, max_message_groups: 50000,
enable_content_deduplication: true,
content_deduplication_window_seconds: 1800, max_deduplicated_messages: 100000,
..Default::default()
},
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum MessageVolume {
Low,
Medium,
High,
}
}
pub use utils::MessageVolume;