use crate::exchanges::DeclareExchange;
use anyhow::Result;
use core::pin::Pin;
use futures::{
future::Future,
stream::Stream,
task::{Context, Poll},
};
use lapin::{message::Delivery, Channel, Consumer};
use std::sync::{Arc, Weak};
use tokio::sync::mpsc::{Receiver, Sender};
use super::comms::*;
type ConsumerCreator = Box<dyn RabbitDispatcher<Object = Consumer>>;
pub type CreatorResult<T> = Pin<Box<dyn Future<Output = Result<T>> + Send>>;
pub type Creator<T> =
Pin<Box<dyn Fn(Arc<Channel>, Option<Arc<DeclareExchange>>) -> CreatorResult<T> + Send + Sync>>;
pub type ChannelReceiver = Receiver<Weak<Channel>>;
type NextFuture = Pin<Box<dyn Future<Output = (Delivery, Data)> + Send>>;
struct Data {
consumer: Consumer,
channel: Weak<Channel>,
creator: ConsumerCreator,
channel_sender: ChannelSender,
channel_receiver: Receiver<Weak<Channel>>,
channel_requester: Arc<Sender<CommsMsg>>,
}
enum State {
Idle(Data),
Next {
next: NextFuture,
},
}
pub struct ConsumerWrapper {
state: Option<State>,
}
impl ConsumerWrapper {
pub(crate) async fn new(creator: ConsumerCreator) -> Self {
log::trace!("Getting channel requester");
let channel_requester = Comms::get_channel_comms();
let (channel_sender, mut channel_receiver) = Comms::create_channel_channel();
log::trace!("Creating the consumer using the creator");
let (consumer, channel) = creator
.start_dispatch(
None,
&channel_sender,
&mut channel_receiver,
&channel_requester,
)
.await;
log::trace!("Consumer wrapper created");
Self {
state: Some(State::Idle(Data {
consumer,
channel,
creator,
channel_sender,
channel_receiver,
channel_requester,
})),
}
}
async fn next_item(mut data: Data) -> (Delivery, Data) {
loop {
use futures::stream::StreamExt;
log::trace!("Polling consumer");
match data.consumer.next().await {
Some(Ok(delivery)) => {
log::trace!("Got delivery");
return (delivery, data);
}
Some(Err(err)) => {
log::error!("Failed to consume a message: {}", err);
}
None => {
log::error!("Consumer has finished for some reason!");
}
}
log::warn!("Consumer is broken, waiting for a new connection");
let (consumer, channel) = data
.creator
.start_dispatch(
Some(data.channel.clone()),
&data.channel_sender,
&mut data.channel_receiver,
&data.channel_requester,
)
.await;
data.consumer = consumer;
data.channel = channel;
}
}
}
impl Stream for ConsumerWrapper {
type Item = Delivery;
fn poll_next(self: Pin<&mut Self>, cx: &mut Context) -> Poll<Option<Self::Item>> {
log::trace!("Poll next");
let this = Pin::into_inner(self);
loop {
match this.state.take() {
Some(State::Idle(data)) => {
this.state = Some(State::Next {
next: Box::pin(Self::next_item(data)),
});
}
Some(State::Next { mut next }) => {
let action = next.as_mut();
return match Future::poll(action, cx) {
Poll::Pending => {
this.state = Some(State::Next { next });
log::trace!("Pending");
Poll::Pending
}
Poll::Ready((delivery, data)) => {
this.state = Some(State::Idle(data));
log::trace!("Ready");
Poll::Ready(Some(delivery))
}
};
}
None => unreachable!(),
}
}
}
}