use std::{
io,
task::{Context, Poll},
};
use futures::{
channel::{mpsc, oneshot},
SinkExt, StreamExt,
};
pub(crate) fn new<Req, Res>(capacity: usize) -> (Sender<Req, Res>, Receiver<Req, Res>) {
let (sender, receiver) = mpsc::channel(capacity);
(
Sender {
inner: futures::lock::Mutex::new(sender),
},
Receiver { inner: receiver },
)
}
pub(crate) struct Sender<Req, Res> {
inner: futures::lock::Mutex<mpsc::Sender<(Req, oneshot::Sender<Res>)>>,
}
impl<Req, Res> Sender<Req, Res> {
pub(crate) async fn send(&self, req: Req) -> io::Result<Res> {
let (sender, receiver) = oneshot::channel();
self.inner
.lock()
.await
.send((req, sender))
.await
.map_err(io::Error::other)?;
let res = receiver.await.map_err(io::Error::other)?;
Ok(res)
}
}
pub(crate) struct Receiver<Req, Res> {
inner: mpsc::Receiver<(Req, oneshot::Sender<Res>)>,
}
impl<Req, Res> Receiver<Req, Res> {
pub(crate) fn poll_next_unpin(
&mut self,
cx: &mut Context<'_>,
) -> Poll<Option<(Req, oneshot::Sender<Res>)>> {
self.inner.poll_next_unpin(cx)
}
}