use crate::{
alias,
service::event::{InternalEvent, InternalEventSender},
};
use futures::{
io::{BufReader, BufWriter, ReadHalf, WriteHalf},
AsyncReadExt, AsyncWriteExt, StreamExt,
};
use libp2p::{swarm::NegotiatedSubstream, PeerId};
use log::*;
use tokio::sync::mpsc;
use tokio_stream::wrappers::UnboundedReceiverStream;
const MSG_BUFFER_LEN: usize = 32768;
pub type GossipSender = mpsc::UnboundedSender<Vec<u8>>;
pub type GossipReceiver = UnboundedReceiverStream<Vec<u8>>;
pub fn channel() -> (GossipSender, GossipReceiver) {
let (sender, receiver) = mpsc::unbounded_channel();
(sender, UnboundedReceiverStream::new(receiver))
}
pub fn start_incoming_processor(
peer_id: PeerId,
mut reader: BufReader<ReadHalf<Box<NegotiatedSubstream>>>,
incoming_tx: GossipSender,
internal_event_sender: InternalEventSender,
) {
tokio::spawn(async move {
let mut msg_buf = vec![0u8; MSG_BUFFER_LEN];
loop {
if let Some(len) = (&mut reader).read(&mut msg_buf).await.ok().filter(|len| *len > 0) {
if incoming_tx.send(msg_buf[..len].to_vec()).is_err() {
debug!("gossip-in: receiver dropped locally.");
break;
}
} else {
debug!("gossip-in: stream closed remotely.");
internal_event_sender
.send(InternalEvent::ProtocolDropped { peer_id })
.expect("The service must not shutdown as long as there are gossip tasks running.");
break;
}
}
debug!("gossip-in: exiting gossip-in processor for {}.", alias!(peer_id));
});
}
pub fn start_outgoing_processor(
peer_id: PeerId,
mut writer: BufWriter<WriteHalf<Box<NegotiatedSubstream>>>,
outgoing_rx: GossipReceiver,
internal_event_sender: InternalEventSender,
) {
tokio::spawn(async move {
let mut outgoing_gossip_receiver = outgoing_rx.fuse();
while let Some(message) = outgoing_gossip_receiver.next().await {
if message.is_empty() {
debug!("gossip-out: received shutdown message.");
internal_event_sender
.send(InternalEvent::ProtocolDropped { peer_id })
.expect("The service must not shutdown as long as there are gossip tasks running.");
break;
}
if (&mut writer).write_all(&message).await.is_err() || (&mut writer).flush().await.is_err() {
debug!("gossip-out: stream closed remotely");
break;
}
}
debug!("gossip-out: exiting gossip-out processor for {}.", alias!(peer_id));
});
}