use std::{collections, io};
pub struct ZeroMQSource<T>
where
T: IntoIterator,
T::Item: Into<zmq::Message>,
{
socket_source: calloop::generic::Generic<calloop::generic::Fd>,
mpsc_receiver: calloop::channel::Channel<T>,
wake_ping_receiver: calloop::ping::PingSource,
wake_ping_sender: calloop::ping::Ping,
socket: zmq::Socket,
outbox: collections::VecDeque<T>,
}
impl<T> ZeroMQSource<T>
where
T: IntoIterator,
T::Item: Into<zmq::Message>,
{
pub fn from_socket(socket: zmq::Socket) -> io::Result<(Self, calloop::channel::Sender<T>)> {
let (mpsc_sender, mpsc_receiver) = calloop::channel::channel();
let (wake_ping_sender, wake_ping_receiver) = calloop::ping::make_ping()?;
let fd = socket.get_fd()?;
let socket_source =
calloop::generic::Generic::from_fd(fd, calloop::Interest::READ, calloop::Mode::Edge);
Ok((
Self {
socket,
socket_source,
mpsc_receiver,
wake_ping_receiver,
wake_ping_sender,
outbox: collections::VecDeque::new(),
},
mpsc_sender,
))
}
}
impl<T> calloop::EventSource for ZeroMQSource<T>
where
T: IntoIterator,
T::Item: Into<zmq::Message>,
{
type Event = Vec<Vec<u8>>;
type Metadata = ();
type Ret = io::Result<()>;
fn process_events<F>(
&mut self,
readiness: calloop::Readiness,
token: calloop::Token,
mut callback: F,
) -> io::Result<calloop::PostAction>
where
F: FnMut(Self::Event, &mut Self::Metadata) -> Self::Ret,
{
self.wake_ping_receiver
.process_events(readiness, token, |_, _| {})?;
let outbox = &mut self.outbox;
self.mpsc_receiver
.process_events(readiness, token, |evt, _| {
if let calloop::channel::Event::Msg(msg) = evt {
outbox.push_back(msg);
}
})?;
loop {
let events = self.socket.get_events()?;
let mut used_socket = false;
if events.contains(zmq::POLLOUT) {
if let Some(parts) = self.outbox.pop_front() {
self.socket.send_multipart(parts, 0)?;
used_socket = true;
}
}
if events.contains(zmq::POLLIN) {
let messages = self.socket.recv_multipart(0)?;
used_socket = true;
callback(messages, &mut ())?;
}
if !used_socket {
break;
}
}
Ok(calloop::PostAction::Continue)
}
fn register(
&mut self,
poll: &mut calloop::Poll,
token_factory: &mut calloop::TokenFactory,
) -> io::Result<()> {
self.socket_source.register(poll, token_factory)?;
self.mpsc_receiver.register(poll, token_factory)?;
self.wake_ping_receiver.register(poll, token_factory)?;
self.wake_ping_sender.ping();
Ok(())
}
fn reregister(
&mut self,
poll: &mut calloop::Poll,
token_factory: &mut calloop::TokenFactory,
) -> io::Result<()> {
self.socket_source.reregister(poll, token_factory)?;
self.mpsc_receiver.reregister(poll, token_factory)?;
self.wake_ping_receiver.reregister(poll, token_factory)?;
self.wake_ping_sender.ping();
Ok(())
}
fn unregister(&mut self, poll: &mut calloop::Poll) -> io::Result<()> {
self.socket_source.unregister(poll)?;
self.mpsc_receiver.unregister(poll)?;
self.wake_ping_receiver.unregister(poll)?;
Ok(())
}
}
impl<T> Drop for ZeroMQSource<T>
where
T: IntoIterator,
T::Item: Into<zmq::Message>,
{
fn drop(&mut self) {
self.socket.set_linger(0).ok();
self.socket.set_rcvtimeo(0).ok();
self.socket.set_sndtimeo(0).ok();
if let Ok(Ok(last_endpoint)) = self.socket.get_last_endpoint() {
self.socket.disconnect(&last_endpoint).ok();
}
}
}