use futures::{Stream, StreamExt};
use lapin::{Channel, Consumer};
use ruststream::Subscriber;
use crate::error::AmqpError;
use crate::message::LapinMessage;
pub struct LapinSubscriber {
_channel: Channel,
consumer: Consumer,
queue: String,
}
impl LapinSubscriber {
pub(crate) fn new(channel: Channel, consumer: Consumer, queue: String) -> Self {
Self {
_channel: channel,
consumer,
queue,
}
}
#[must_use]
pub fn queue(&self) -> &str {
&self.queue
}
}
impl std::fmt::Debug for LapinSubscriber {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("LapinSubscriber")
.field("queue", &self.queue)
.finish_non_exhaustive()
}
}
impl Subscriber for LapinSubscriber {
type Message = LapinMessage;
type Error = AmqpError;
fn stream(&mut self) -> impl Stream<Item = Result<Self::Message, Self::Error>> + Send + '_ {
futures::stream::unfold(&mut self.consumer, |consumer| async move {
let item = consumer.next().await?;
let mapped = match item {
Ok(delivery) => Ok(LapinMessage::from_delivery(delivery)),
Err(err) => Err(AmqpError::consume(err)),
};
Some((mapped, consumer))
})
}
}