use std::time::Duration;
use bytes::Bytes;
use lapin::Acker;
use lapin::message::Delivery;
use lapin::options::{BasicAckOptions, BasicNackOptions, BasicRejectOptions};
use ruststream::{AckError, Headers, IncomingMessage, Partitioned};
use crate::convert;
use crate::delay::DelayContext;
pub const PARTITION_KEY_HEADER: &str = "amqp-partition-key";
#[derive(Debug)]
pub struct LapinMessage {
payload: Bytes,
headers: Headers,
exchange: String,
routing_key: String,
redelivered: bool,
delivery_tag: u64,
acker: Option<Acker>,
delay: Option<DelayContext>,
}
impl LapinMessage {
pub(crate) fn from_delivery(delivery: Delivery, delay: Option<DelayContext>) -> Self {
let headers = convert::headers_from_properties(&delivery.properties);
Self {
payload: Bytes::from(delivery.data),
headers,
exchange: delivery.exchange.to_string(),
routing_key: delivery.routing_key.to_string(),
redelivered: delivery.redelivered,
delivery_tag: delivery.delivery_tag,
acker: Some(delivery.acker),
delay,
}
}
pub(crate) fn from_delivery_no_ack(delivery: Delivery) -> Self {
let mut msg = Self::from_delivery(delivery, None);
msg.acker = None;
msg
}
#[must_use]
pub fn exchange(&self) -> &str {
&self.exchange
}
#[must_use]
pub fn routing_key(&self) -> &str {
&self.routing_key
}
#[must_use]
pub fn redelivered(&self) -> bool {
self.redelivered
}
#[must_use]
pub fn delivery_tag(&self) -> u64 {
self.delivery_tag
}
async fn settle<F, Fut>(mut self, op: F, what: &'static str) -> Result<(), AckError>
where
F: FnOnce(Acker) -> Fut,
Fut: Future<Output = lapin::Result<bool>>,
{
let Some(acker) = self.acker.take() else {
return Ok(());
};
match op(acker).await {
Ok(true) => Ok(()),
Ok(false) => Err(AckError::Broker(
format!("{what} was not sent: the delivery channel is closed or errored").into(),
)),
Err(err) => Err(AckError::Broker(Box::new(err))),
}
}
}
impl IncomingMessage for LapinMessage {
fn payload(&self) -> &[u8] {
&self.payload
}
fn headers(&self) -> &Headers {
&self.headers
}
async fn ack(self) -> Result<(), AckError> {
self.settle(
|acker| async move { acker.ack(BasicAckOptions::default()).await },
"basic.ack",
)
.await
}
async fn nack(self, requeue: bool) -> Result<(), AckError> {
if requeue {
self.settle(
|acker| async move {
acker
.nack(BasicNackOptions {
multiple: false,
requeue: true,
})
.await
},
"basic.nack",
)
.await
} else {
self.settle(
|acker| async move { acker.reject(BasicRejectOptions { requeue: false }).await },
"basic.reject",
)
.await
}
}
fn partition_key(&self) -> Option<&[u8]> {
self.headers.get(PARTITION_KEY_HEADER)
}
fn supports_nack_after(&self) -> bool {
self.delay.is_some()
}
async fn nack_after(mut self, delay: Duration) -> Result<(), AckError> {
let Some(context) = self.delay.take() else {
return Err(AckError::Unsupported);
};
context
.republish(&self.payload, &self.headers, delay)
.await
.map_err(|err| AckError::Broker(Box::new(err)))?;
self.settle(
|acker| async move { acker.ack(BasicAckOptions::default()).await },
"basic.ack",
)
.await
}
}
impl Partitioned for LapinMessage {
fn partition_key(&self) -> Option<&[u8]> {
self.headers.get(PARTITION_KEY_HEADER)
}
}