use serde::{Serialize, Deserialize};
use std::collections::HashMap;
use chrono::{DateTime, Utc, Duration};
use uuid::Uuid;
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
#[derive(Default)]
pub enum QueuePriority {
Low = 1,
#[default]
Normal = 5,
High = 10,
Critical = 20,
}
impl std::fmt::Display for QueuePriority {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
QueuePriority::Low => write!(f, "Low"),
QueuePriority::Normal => write!(f, "Normal"),
QueuePriority::High => write!(f, "High"),
QueuePriority::Critical => write!(f, "Critical"),
}
}
}
impl From<i32> for QueuePriority {
fn from(value: i32) -> Self {
match value {
1 => Self::Low,
5 => Self::Normal,
10 => Self::High,
20 => Self::Critical,
_ => Self::Normal, }
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub enum MessageStatus {
Pending,
Processing,
Acknowledged,
Failed,
DeadLettered,
Deleted,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct QueueMessage {
pub id: String,
pub payload: serde_json::Value,
pub priority: QueuePriority,
pub receive_count: u32,
pub max_receive_count: u32,
pub enqueued_at: DateTime<Utc>,
pub created_at: DateTime<Utc>,
pub visible_at: DateTime<Utc>,
pub expires_at: Option<DateTime<Utc>>,
pub visibility_timeout: u64,
pub status: MessageStatus,
pub delay_seconds: Option<u64>,
pub attributes: HashMap<String, String>,
pub headers: HashMap<String, serde_json::Value>,
pub message_group_id: Option<String>,
pub message_deduplication_id: Option<String>,
pub routing_key: Option<String>,
pub compressed: bool,
pub original_size: Option<usize>,
}
impl QueueMessage {
pub fn builder() -> QueueMessageBuilder {
QueueMessageBuilder::new()
}
pub fn is_expired(&self) -> bool {
if let Some(expires_at) = self.expires_at {
Utc::now() > expires_at
} else {
false
}
}
pub fn is_visible(&self) -> bool {
Utc::now() >= self.visible_at
}
pub fn should_dead_letter(&self) -> bool {
self.receive_count >= self.max_receive_count
}
pub fn age_seconds(&self) -> i64 {
(Utc::now() - self.enqueued_at).num_seconds()
}
pub fn remaining_visibility(&self) -> i64 {
(self.visible_at - Utc::now()).num_seconds().max(0)
}
pub fn mark_received(&mut self) {
self.receive_count += 1;
self.status = MessageStatus::Processing;
self.visible_at = Utc::now() + Duration::seconds(self.visibility_timeout as i64);
}
pub fn mark_acknowledged(&mut self) {
self.status = MessageStatus::Acknowledged;
}
pub fn mark_failed(&mut self) {
self.status = MessageStatus::Failed;
}
pub fn mark_dead_lettered(&mut self) {
self.status = MessageStatus::DeadLettered;
}
pub fn reset_for_retry(&mut self, delay_seconds: Option<u64>) {
self.status = MessageStatus::Pending;
let delay = delay_seconds.unwrap_or(0);
self.visible_at = Utc::now() + Duration::seconds(delay as i64);
}
pub fn validate(&self) -> Result<(), String> {
if self.id.is_empty() {
return Err("Message ID cannot be empty".to_string());
}
if self.visibility_timeout == 0 {
return Err("Visibility timeout must be greater than 0".to_string());
}
if self.max_receive_count == 0 {
return Err("Max receive count must be greater than 0".to_string());
}
Ok(())
}
pub fn size_bytes(&self) -> Result<usize, serde_json::Error> {
serde_json::to_vec(self).map(|v| v.len())
}
}
pub struct QueueMessageBuilder {
message: QueueMessage,
}
impl QueueMessageBuilder {
pub fn new() -> Self {
Self {
message: QueueMessage {
id: Uuid::new_v4().to_string(),
payload: serde_json::Value::Null,
priority: QueuePriority::Normal,
receive_count: 0,
max_receive_count: crate::DEFAULT_MAX_RECEIVE_COUNT,
enqueued_at: Utc::now(),
created_at: Utc::now(),
visible_at: Utc::now(),
expires_at: None,
visibility_timeout: crate::DEFAULT_VISIBILITY_TIMEOUT,
status: MessageStatus::Pending,
delay_seconds: None,
attributes: HashMap::new(),
headers: HashMap::new(),
message_group_id: None,
message_deduplication_id: None,
routing_key: None,
compressed: false,
original_size: None,
},
}
}
pub fn id(mut self, id: impl Into<String>) -> Self {
self.message.id = id.into();
self
}
pub fn payload<T: Serialize>(mut self, payload: T) -> Result<Self, serde_json::Error> {
self.message.payload = serde_json::to_value(payload)?;
Ok(self)
}
pub fn json_payload(mut self, payload: serde_json::Value) -> Self {
self.message.payload = payload;
self
}
pub fn text_payload(mut self, text: impl Into<String>) -> Self {
self.message.payload = serde_json::Value::String(text.into());
self
}
pub fn priority(mut self, priority: QueuePriority) -> Self {
self.message.priority = priority;
self
}
pub fn max_receive_count(mut self, count: u32) -> Self {
self.message.max_receive_count = count;
self
}
pub fn visibility_timeout(mut self, timeout: u64) -> Self {
self.message.visibility_timeout = timeout;
self
}
pub fn delay(mut self, seconds: u64) -> Self {
self.message.delay_seconds = Some(seconds);
self.message.visible_at = Utc::now() + chrono::Duration::seconds(seconds as i64);
self
}
pub fn expires_at(mut self, time: DateTime<Utc>) -> Self {
self.message.expires_at = Some(time);
self
}
pub fn expires_in(mut self, seconds: u64) -> Self {
self.message.expires_at = Some(Utc::now() + chrono::Duration::seconds(seconds as i64));
self
}
pub fn attribute(mut self, key: impl Into<String>, value: impl Into<String>) -> Self {
self.message.attributes.insert(key.into(), value.into());
self
}
pub fn attributes(mut self, attrs: HashMap<String, String>) -> Self {
self.message.attributes.extend(attrs);
self
}
pub fn message_group_id(mut self, group_id: impl Into<String>) -> Self {
self.message.message_group_id = Some(group_id.into());
self
}
pub fn message_deduplication_id(mut self, dedup_id: impl Into<String>) -> Self {
self.message.message_deduplication_id = Some(dedup_id.into());
self
}
pub fn routing_key(mut self, routing_key: impl Into<String>) -> Self {
self.message.routing_key = Some(routing_key.into());
self
}
pub fn compress(mut self, enabled: bool) -> Self {
self.message.compressed = enabled;
self
}
pub fn build(self) -> QueueMessage {
self.message
}
}
impl Default for QueueMessageBuilder {
fn default() -> Self {
Self::new()
}
}
#[derive(Debug, Clone)]
pub struct BatchReceiveResult {
pub messages: Vec<QueueMessage>,
pub requested: usize,
pub available: usize,
pub total_in_queue: u64,
pub processing_time_ms: u64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub enum QueueHealth {
Healthy,
Degraded,
Unhealthy,
}
impl QueueHealth {
pub fn from_metrics(
avg_processing_time_ms: f64,
backlog_size: u64,
error_rate: f64,
) -> Self {
let processing_ok = avg_processing_time_ms < 5000.0; let backlog_ok = backlog_size < 1000;
let error_ok = error_rate < 0.05;
if processing_ok && backlog_ok && error_ok {
Self::Healthy
} else if !processing_ok || backlog_size > 5000 {
Self::Unhealthy
} else {
Self::Degraded
}
}
}
#[derive(Debug, Clone)]
pub struct QueueHealthCheck {
pub status: QueueHealth,
pub queue_size: u64,
pub avg_processing_time_ms: Option<f64>,
pub messages_per_second: Option<f64>,
pub error_rate: f64,
pub last_activity: Option<DateTime<Utc>>,
pub checked_at: DateTime<Utc>,
pub details: HashMap<String, String>,
}