use std::collections::VecDeque;
use tokio::sync::mpsc;
use tracing::warn;
use uuid::Uuid;
use super::message::Envelope;
pub struct Mailbox {
receiver: mpsc::Receiver<Envelope>,
capacity: usize,
retry_queue: VecDeque<RetryEntry>,
dead_letters: VecDeque<DeadLetter>,
inflight_retry_count: Option<u32>,
max_retries: u32,
dlq_capacity: usize,
}
struct RetryEntry {
envelope: Envelope,
retry_count: u32,
}
pub struct DeadLetter {
pub msg_id: Uuid,
pub trace_id: Uuid,
pub retry_count: u32,
pub last_error: String,
pub envelope: Option<Envelope>,
}
impl std::fmt::Debug for DeadLetter {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("DeadLetter")
.field("msg_id", &self.msg_id)
.field("retry_count", &self.retry_count)
.field("last_error", &self.last_error)
.finish()
}
}
#[derive(Clone)]
pub struct MailboxSender {
sender: mpsc::Sender<Envelope>,
}
impl Mailbox {
pub fn new(capacity: usize) -> (Self, MailboxSender) {
let (sender, receiver) = mpsc::channel(capacity);
(
Self {
receiver,
capacity,
retry_queue: VecDeque::new(),
dead_letters: VecDeque::new(),
inflight_retry_count: None,
max_retries: 3,
dlq_capacity: 100,
},
MailboxSender { sender },
)
}
pub fn default_capacity() -> (Self, MailboxSender) {
Self::new(256)
}
pub fn set_max_retries(&mut self, max: u32) {
self.max_retries = max;
}
pub async fn recv(&mut self) -> Option<Envelope> {
if let Some(entry) = self.retry_queue.pop_front() {
self.inflight_retry_count = Some(entry.retry_count);
Some(entry.envelope)
} else {
let envelope = self.receiver.recv().await?;
self.inflight_retry_count = Some(0);
Some(envelope)
}
}
pub fn commit(&mut self) {
self.inflight_retry_count = None;
}
pub fn nack(&mut self, envelope: Envelope, error: &str) {
let prev_count = self.inflight_retry_count.take().unwrap_or(0);
let retry_count = prev_count + 1;
if retry_count >= self.max_retries {
warn!(
msg_id = %envelope.id,
retries = retry_count,
error = %error,
"message exceeded max retries, moving to DLQ"
);
let msg_id = envelope.id;
let trace_id = envelope.trace_id;
self.push_dead_letter(DeadLetter {
msg_id,
trace_id,
retry_count,
last_error: error.to_string(),
envelope: Some(envelope),
});
} else {
warn!(
msg_id = %envelope.id,
retry = retry_count,
error = %error,
"nack: message will be retried"
);
self.retry_queue.push_back(RetryEntry {
envelope,
retry_count,
});
}
}
pub fn record_failure(&mut self, msg_id: Uuid, trace_id: Uuid, error: &str) {
let retry_count = self.inflight_retry_count.take().unwrap_or(0) + 1;
warn!(
msg_id = %msg_id,
retries = retry_count,
error = %error,
"message processing failed, recorded to DLQ"
);
self.push_dead_letter(DeadLetter {
msg_id,
trace_id,
retry_count,
last_error: error.to_string(),
envelope: None,
});
}
fn push_dead_letter(&mut self, dl: DeadLetter) {
if self.dead_letters.len() >= self.dlq_capacity {
self.dead_letters.pop_front();
}
self.dead_letters.push_back(dl);
}
pub fn dead_letter_count(&self) -> usize {
self.dead_letters.len()
}
pub fn retry_queue_len(&self) -> usize {
self.retry_queue.len()
}
pub fn drain_dead_letters(&mut self) -> Vec<DeadLetter> {
self.dead_letters.drain(..).collect()
}
pub fn capacity(&self) -> usize {
self.capacity
}
}
impl MailboxSender {
pub async fn send(&self, envelope: Envelope) -> Result<(), MailboxSendError> {
self.sender
.send(envelope)
.await
.map_err(|_| MailboxSendError::ActorStopped)
}
pub fn try_send(&self, envelope: Envelope) -> Result<(), MailboxSendError> {
self.sender.try_send(envelope).map_err(|e| match e {
mpsc::error::TrySendError::Full(_) => {
warn!("mailbox full, message dropped");
MailboxSendError::MailboxFull
}
mpsc::error::TrySendError::Closed(_) => MailboxSendError::ActorStopped,
})
}
pub fn is_closed(&self) -> bool {
self.sender.is_closed()
}
}
#[derive(Debug, thiserror::Error)]
pub enum MailboxSendError {
#[error("actor has stopped, mailbox closed")]
ActorStopped,
#[error("mailbox is full")]
MailboxFull,
}