use std::time::Duration;
use async_trait::async_trait;
use lapin::{
Acker,
options::{BasicAckOptions, BasicRejectOptions},
};
use queuey_core::{Delivery, Envelope, Result};
use tracing::warn;
use crate::{
error::{RabbitMqError, amqp},
publisher::{Hold, Publisher},
};
pub struct RabbitMqDelivery {
envelope: Envelope,
acker: Acker,
publisher: Publisher,
}
impl std::fmt::Debug for RabbitMqDelivery {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("RabbitMqDelivery")
.field("job_id", &self.envelope.job_id)
.field("job_type", &self.envelope.job_type)
.field("queue", &self.envelope.queue)
.field("attempt", &self.envelope.attempt)
.finish_non_exhaustive()
}
}
impl RabbitMqDelivery {
pub(crate) fn new(envelope: Envelope, acker: Acker, publisher: Publisher) -> Self {
Self {
envelope,
acker,
publisher,
}
}
fn settled(&self, settled: bool, operation: &'static str) -> Result<()> {
if settled {
return Ok(());
}
Err(RabbitMqError::AlreadySettled {
operation,
job_id: self.envelope.job_id.to_string(),
queue: self.envelope.queue.clone(),
}
.into_core())
}
async fn ack_original(&self) -> Result<()> {
let acked = self
.acker
.ack(BasicAckOptions::default())
.await
.map_err(amqp)?;
self.settled(acked, "ack")
}
async fn reject_original(&self) -> Result<()> {
let rejected = self
.acker
.reject(BasicRejectOptions { requeue: false })
.await
.map_err(amqp)?;
self.settled(rejected, "reject")
}
}
#[async_trait]
impl Delivery for RabbitMqDelivery {
fn envelope(&self) -> &Envelope {
&self.envelope
}
async fn ack(self: Box<Self>) -> Result<()> {
self.ack_original().await
}
async fn dead_letter(self: Box<Self>, reason: &str) -> Result<()> {
if !self.publisher.options().declare_dead_letter_queues {
warn!(
job_id = %self.envelope.job_id,
queue = %self.envelope.queue,
reason,
"dead-letter queues are disabled; rejecting the delivery instead of publishing \
to the dead-letter queue (the broker's own dead-letter policy applies, if any, \
otherwise the job is dropped)"
);
return self.reject_original().await;
}
self.publisher
.publish_dead_letter(&self.envelope, reason)
.await?;
self.ack_original().await
}
async fn retry(self: Box<Self>, next: Envelope, delay: Duration) -> Result<()> {
self.publisher
.publish_held(&next, delay, Hold::Retry)
.await?;
self.ack_original().await
}
async fn defer(self: Box<Self>, next: Envelope, delay: Duration) -> Result<()> {
self.publisher
.publish_held(&next, delay, Hold::Deferral)
.await?;
self.ack_original().await
}
}