use futures::{Stream, StreamExt};
use lapin::{Channel, Consumer};
use ruststream::Subscriber;
use crate::delay::DelayContext;
use crate::error::AmqpError;
use crate::message::LapinMessage;
pub struct LapinSubscriber {
_channel: Channel,
consumer: Consumer,
queue: String,
delay: Option<DelayContext>,
}
impl LapinSubscriber {
pub(crate) fn new(
channel: Channel,
consumer: Consumer,
queue: String,
delay: Option<DelayContext>,
) -> Self {
Self {
_channel: channel,
consumer,
queue,
delay,
}
}
#[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 + '_ {
let delay = self.delay.clone();
futures::stream::unfold(
(&mut self.consumer, delay),
|(consumer, delay)| async move {
let item = consumer.next().await?;
let mapped = match item {
Ok(delivery) => Ok(LapinMessage::from_delivery(delivery, delay.clone())),
Err(err) => Err(AmqpError::consume(err)),
};
Some((mapped, (consumer, delay)))
},
)
}
}