use async_trait::async_trait;
use crate::{
QueueResult, QueueMessage, QueueStats, QueueConfig,
QueueHealthCheck, BatchReceiveResult, QueueBackend
};
use std::collections::HashMap;
#[async_trait]
pub trait QueueService: Send + Sync {
async fn enqueue(&self, message: QueueMessage) -> QueueResult<String>;
async fn enqueue_batch(&self, messages: Vec<QueueMessage>) -> QueueResult<Vec<String>> {
let mut ids = Vec::with_capacity(messages.len());
for message in messages {
ids.push(self.enqueue(message).await?);
}
Ok(ids)
}
async fn dequeue(&self) -> QueueResult<Option<QueueMessage>>;
async fn dequeue_batch(&self, max_messages: usize) -> QueueResult<BatchReceiveResult>;
async fn ack(&self, message_id: &str) -> QueueResult<bool>;
async fn ack_batch(&self, message_ids: Vec<String>) -> QueueResult<u64>;
async fn nack(&self, message_id: &str, delay_seconds: Option<u64>) -> QueueResult<bool>;
async fn delete(&self, message_id: &str) -> QueueResult<bool>;
async fn delete_batch(&self, messages: Vec<QueueMessage>) -> QueueResult<u64> {
let mut count = 0;
for message in messages {
if self.delete(&message.id).await? {
count += 1;
}
}
Ok(count)
}
async fn get_message(&self, message_id: &str) -> QueueResult<Option<QueueMessage>>;
async fn get_stats(&self) -> QueueResult<QueueStats>;
async fn purge(&self) -> QueueResult<u64>;
async fn size(&self) -> QueueResult<u64>;
async fn is_empty(&self) -> QueueResult<bool>;
async fn health_check(&self) -> QueueResult<QueueHealthCheck>;
async fn validate_config(&self) -> QueueResult<bool>;
async fn test_connection(&self) -> QueueResult<bool>;
fn backend_type(&self) -> QueueBackend;
}
#[async_trait]
pub trait QueueManager: Send + Sync {
async fn create_queue(&self, config: QueueConfig) -> QueueResult<bool>;
async fn delete_queue(&self, queue_name: &str) -> QueueResult<bool>;
async fn list_queues(&self) -> QueueResult<Vec<String>>;
async fn get_queue_config(&self, queue_name: &str) -> QueueResult<Option<QueueConfig>>;
async fn update_queue_config(&self, queue_name: &str, config: QueueConfig) -> QueueResult<bool>;
async fn pause_queue(&self, queue_name: &str) -> QueueResult<bool>;
async fn resume_queue(&self, queue_name: &str) -> QueueResult<bool>;
async fn get_queue_health(&self, queue_name: &str) -> QueueResult<QueueHealthCheck>;
}
#[async_trait]
pub trait MessageProcessor: Send + Sync {
async fn process_message(&self, message: QueueMessage) -> QueueResult<bool>;
async fn process_batch(&self, messages: Vec<QueueMessage>) -> QueueResult<Vec<bool>> {
let mut results = Vec::with_capacity(messages.len());
for message in messages {
results.push(self.process_message(message).await?);
}
Ok(results)
}
fn name(&self) -> &str;
fn version(&self) -> &str;
}
#[async_trait]
pub trait QueueMonitor: Send + Sync {
async fn start_monitoring(&self, queue_name: &str) -> QueueResult<()>;
async fn stop_monitoring(&self, queue_name: &str) -> QueueResult<()>;
async fn get_metrics(&self, queue_name: &str) -> QueueResult<HashMap<String, f64>>;
async fn set_alert_thresholds(&self, queue_name: &str, thresholds: HashMap<String, f64>) -> QueueResult<()>;
async fn get_alerts(&self, queue_name: &str) -> QueueResult<Vec<String>>;
}