use lapin::message::Delivery;
use lapin::options::{
BasicAckOptions, BasicGetOptions, BasicNackOptions, BasicPublishOptions, BasicRejectOptions,
ConfirmSelectOptions, QueueDeclareOptions,
};
use lapin::types::{AMQPValue, FieldTable, ShortString};
use lapin::{BasicProperties, Channel, Connection, ConnectionProperties};
use super::source::{MessageSource, ReceivedMessage};
use super::{message_from_wire, strip_address_prefix, Message};
use super::{retryable, MessagePublisher, TransportError};
const MESSAGE_KIND_HEADER: &str = "x-sourced-kind";
fn settle_result(context: &str, settled: bool) -> Result<(), TransportError> {
if settled {
Ok(())
} else {
Err(TransportError::retryable(format!(
"{context}: acker unavailable"
)))
}
}
pub(super) async fn connect_channel(uri: &str) -> Result<Channel, TransportError> {
let connection = Connection::connect(uri, ConnectionProperties::default())
.await
.map_err(|err| retryable("amqp connect", err))?;
connection
.create_channel()
.await
.map_err(|err| retryable("amqp channel", err))
}
pub struct RabbitPublisher {
channel: Channel,
}
impl RabbitPublisher {
pub fn new(channel: Channel) -> Self {
Self { channel }
}
pub async fn connect(uri: &str) -> Result<Self, TransportError> {
let channel = connect_channel(uri).await?;
channel
.confirm_select(ConfirmSelectOptions::default())
.await
.map_err(|err| retryable("amqp confirm_select", err))?;
Ok(Self::new(channel))
}
}
pub(super) fn message_properties(message: &Message) -> BasicProperties {
let mut headers = FieldTable::default();
headers.insert(
ShortString::from(MESSAGE_KIND_HEADER),
AMQPValue::LongString(message.kind.as_str().into()),
);
for (key, value) in &message.metadata {
headers.insert(
ShortString::from(key.as_str()),
AMQPValue::LongString(value.as_str().into()),
);
}
let mut properties = BasicProperties::default()
.with_headers(headers)
.with_content_type(ShortString::from(message.content_type.as_str()));
if let Some(id) = message.id() {
properties = properties.with_message_id(ShortString::from(id));
}
properties
}
impl MessagePublisher for RabbitPublisher {
async fn publish(&self, message: Message) -> Result<(), TransportError> {
let confirm = self
.channel
.basic_publish(
ShortString::from(""), ShortString::from(message.name()),
BasicPublishOptions::default(),
&message.payload,
message_properties(&message),
)
.await
.map_err(|err| retryable("amqp publish", err))?;
let confirmation = confirm
.await
.map_err(|err| retryable("amqp publisher confirm", err))?;
if confirmation.is_nack() {
return Err(TransportError::retryable("amqp publisher confirm: nack"));
}
Ok(())
}
}
pub struct RabbitSource {
channel: Channel,
queues: Vec<String>,
strip_prefix: Option<String>,
current: usize,
}
impl RabbitSource {
pub fn new(channel: Channel, queue: impl Into<String>) -> Self {
Self::multi(channel, vec![queue.into()], None)
}
pub(super) fn multi(
channel: Channel,
queues: Vec<String>,
strip_prefix: Option<String>,
) -> Self {
Self {
channel,
queues,
strip_prefix,
current: 0,
}
}
pub async fn connect(uri: &str, queue: &str) -> Result<Self, TransportError> {
let channel = connect_channel(uri).await?;
channel
.queue_declare(
ShortString::from(queue),
QueueDeclareOptions {
durable: true,
..Default::default()
},
FieldTable::default(),
)
.await
.map_err(|err| retryable("amqp queue_declare", err))?;
Ok(Self::new(channel, queue))
}
}
impl MessageSource for RabbitSource {
type Received = RabbitReceived;
fn transport_name(&self) -> &'static str {
"rabbitmq"
}
async fn recv(&mut self) -> Result<Option<Self::Received>, TransportError> {
for offset in 0..self.queues.len() {
let index = (self.current + offset) % self.queues.len();
let got = self
.channel
.basic_get(
ShortString::from(self.queues[index].as_str()),
BasicGetOptions::default(),
)
.await
.map_err(|err| retryable("amqp basic_get", err))?;
if let Some(get) = got {
self.current = index;
let routing_key = get.delivery.routing_key.to_string();
let name = strip_address_prefix(routing_key, self.strip_prefix.as_deref());
return Ok(Some(RabbitReceived::from_delivery_with_name(
get.delivery,
name,
)));
}
}
Ok(None)
}
}
pub struct RabbitReceived {
delivery: Delivery,
message: Message,
}
impl RabbitReceived {
pub(super) fn from_delivery_with_name(delivery: Delivery, name: String) -> Self {
let payload = delivery.data.clone();
let headers: Vec<(String, String)> = delivery
.properties
.headers()
.as_ref()
.into_iter()
.flat_map(|headers| headers.inner())
.map(|(key, value)| (key.to_string(), amqp_value_to_string(value)))
.collect();
let mut message = message_from_wire(name, payload, None, MESSAGE_KIND_HEADER, headers);
message.id = delivery
.properties
.message_id()
.as_ref()
.map(|s| s.to_string());
if let Some(content_type) = delivery.properties.content_type().as_ref() {
message.content_type = content_type.to_string();
}
Self { delivery, message }
}
}
impl ReceivedMessage for RabbitReceived {
fn message(&self) -> &Message {
&self.message
}
async fn ack(self) -> Result<(), TransportError> {
let settled = self
.delivery
.ack(BasicAckOptions::default())
.await
.map_err(|err| retryable("amqp ack", err))?;
settle_result("amqp ack", settled)
}
async fn nack(self, _reason: &str) -> Result<(), TransportError> {
let settled = self
.delivery
.nack(BasicNackOptions {
requeue: true,
..Default::default()
})
.await
.map_err(|err| retryable("amqp nack", err))?;
settle_result("amqp nack", settled)
}
async fn dead_letter(self, _reason: &str) -> Result<(), TransportError> {
let settled = self
.delivery
.reject(BasicRejectOptions { requeue: false })
.await
.map_err(|err| retryable("amqp reject", err))?;
settle_result("amqp reject", settled)
}
async fn park(self, _reason: &str) -> Result<(), TransportError> {
let settled = self
.delivery
.reject(BasicRejectOptions { requeue: false })
.await
.map_err(|err| retryable("amqp reject", err))?;
settle_result("amqp reject", settled)
}
}
fn amqp_value_to_string(value: &AMQPValue) -> String {
match value {
AMQPValue::LongString(s) => s.to_string(),
AMQPValue::ShortString(s) => s.to_string(),
other => format!("{other:?}"),
}
}