use crate::error::StorageResult;
use async_trait::async_trait;
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use uuid::Uuid;
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
pub enum MessagePriority {
Normal,
High,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct MessageRecord {
pub id: Uuid,
pub from: Vec<u8>,
pub payload: Vec<u8>,
pub priority: MessagePriority,
pub created_at: DateTime<Utc>,
pub status: MessageStatus,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub enum MessageStatus {
Queued,
Inflight,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct MailboxStats {
pub queued_messages: u64,
pub inflight_messages: u64,
pub queued_by_priority: std::collections::HashMap<MessagePriority, u64>,
}
#[async_trait]
pub trait Mailbox: Send + Sync {
async fn enqueue(
&self,
from: Vec<u8>,
payload: Vec<u8>,
priority: MessagePriority,
) -> StorageResult<Uuid>;
async fn dequeue(&self) -> StorageResult<Vec<MessageRecord>>;
async fn ack(&self, message_id: Uuid) -> StorageResult<()>;
async fn status(&self) -> StorageResult<MailboxStats>;
}