use bytes::Buf; use bytes::BytesMut; use futures_core::ready; use std::io::Error as IoError; use std::io::ErrorKind as IoErrorKind; use std::io::IoSlice; use std::pin::Pin; use std::sync::{Arc, Mutex}; use std::task::{Context, Poll, Waker}; use tokio::io::{AsyncRead, AsyncWrite, ReadBuf};
type IoResult<T> = Result<T, IoError>;
const CLOSED_ERROR_MSG: &str = "simplex has been closed";
#[derive(Debug)] struct Inner { /// `poll_write` will return [`Poll::Pending`] if the backpressure boundary is reached
backpressure_boundary: usize,
/// either [`Sender`] or [`Receiver`] is closed
is_closed: bool,
/// Waker used to wake the [`Receiver`]
receiver_waker: Option<Waker>,
/// Waker used to wake the [`Sender`]
sender_waker: Option<Waker>,
/// Buffer used to read and write data
buf: BytesMut,
}
/// Receiver of the simplex channel. /// /// # Cancellation safety /// /// The `Receiver` is cancel safe. If it is used as the event in a /// [`tokio::select!`](macro@tokio::select) statement and some other branch /// completes first, it is guaranteed that no bytes were received on this /// channel. /// /// You can still read the remaining data from the buffer /// even if the write half has been dropped. /// See [`Sender::poll_shutdown`] and [`Sender::drop`] for more details. #[derive(Debug)] pubstruct Receiver {
inner: Arc<Mutex<Inner>>,
}
impl Drop for Receiver { /// This also wakes up the [`Sender`]. fn drop(&mutself) { let maybe_waker = { letmut inner = self.inner.lock().unwrap();
inner.close_receiver()
};
let waker = inner.take_sender_waker();
drop(inner); // unlock before waking up iflet Some(waker) = waker {
waker.wake();
}
Poll::Ready(Ok(()))
}
}
/// Sender of the simplex channel. /// /// # Cancellation safety /// /// The `Sender` is cancel safe. If it is used as the event in a /// [`tokio::select!`](macro@tokio::select) statement and some other branch /// completes first, it is guaranteed that no bytes were sent on this /// channel. /// /// # Shutdown /// /// See [`Sender::poll_shutdown`]. #[derive(Debug)] pubstruct Sender {
inner: Arc<Mutex<Inner>>,
}
impl Drop for Sender { /// This also wakes up the [`Receiver`]. fn drop(&mutself) { let maybe_waker = { letmut inner = self.inner.lock().unwrap();
inner.close_sender()
};
impl AsyncWrite for Sender { /// # Errors /// /// This method will return [`IoErrorKind::BrokenPipe`] /// if the channel has been closed. fn poll_write(self: Pin<&mutSelf>, cx: &mut Context<'_>, buf: &[u8]) -> Poll<IoResult<usize>> { let coop = ready!(poll_proceed(cx));
letmut inner = self.inner.lock().unwrap();
if inner.is_closed() { return Poll::Ready(Err(IoError::new(IoErrorKind::BrokenPipe, CLOSED_ERROR_MSG)));
}
let free = inner
.backpressure_boundary
.checked_sub(inner.buf.len())
.expect("backpressure boundary overflow"); let to_write = buf.len().min(free); if to_write == 0 { if buf.is_empty() { return Poll::Ready(Ok(0));
}
let old_waker = inner.register_sender_waker(cx.waker()); let waker = inner.take_receiver_waker();
// unlock before waking up and dropping old waker
drop(inner);
drop(old_waker); iflet Some(waker) = waker {
waker.wake();
}
return Poll::Pending;
}
// this is to avoid starving other tasks
coop.made_progress();
inner.buf.extend_from_slice(&buf[..to_write]);
let waker = inner.take_receiver_waker();
drop(inner); // unlock before waking up iflet Some(waker) = waker {
waker.wake();
}
Poll::Ready(Ok(to_write))
}
/// # Errors /// /// This method will return [`IoErrorKind::BrokenPipe`] /// if the channel has been closed. fn poll_flush(self: Pin<&mutSelf>, _cx: &mut Context<'_>) -> Poll<IoResult<()>> { let inner = self.inner.lock().unwrap(); if inner.is_closed() {
Poll::Ready(Err(IoError::new(IoErrorKind::BrokenPipe, CLOSED_ERROR_MSG)))
} else {
Poll::Ready(Ok(()))
}
}
/// After returns [`Poll::Ready`], all the following call to /// [`Sender::poll_write`] and [`Sender::poll_flush`] /// will return error. /// /// The [`Receiver`] can still be used to read remaining data /// until all bytes have been consumed. fn poll_shutdown(self: Pin<&mutSelf>, _cx: &mut Context<'_>) -> Poll<IoResult<()>> { let maybe_waker = { letmut inner = self.inner.lock().unwrap();
inner.close_sender()
};
let free = inner
.backpressure_boundary
.checked_sub(inner.buf.len())
.expect("backpressure boundary overflow"); if free == 0 { let old_waker = inner.register_sender_waker(cx.waker()); let maybe_waker = inner.take_receiver_waker();
// unlock before waking up and dropping old waker
drop(inner);
drop(old_waker); iflet Some(waker) = maybe_waker {
waker.wake();
}
return Poll::Pending;
}
// this is to avoid starving other tasks
coop.made_progress();
letmut rem = free; for buf in bufs { if rem == 0 { break;
}
let to_write = buf.len().min(rem); if to_write == 0 {
assert_ne!(rem, 0);
assert_eq!(buf.len(), 0); continue;
}
inner.buf.extend_from_slice(&buf[..to_write]);
rem -= to_write;
}
let waker = inner.take_receiver_waker();
drop(inner); // unlock before waking up iflet Some(waker) = waker {
waker.wake();
}
Poll::Ready(Ok(free - rem))
}
}
/// Create a simplex channel. /// /// The `capacity` parameter specifies the maximum number of bytes that can be /// stored in the channel without making the [`Sender::poll_write`] /// return [`Poll::Pending`]. /// /// # Panics /// /// This function will panic if `capacity` is zero. pubfn new(capacity: usize) -> (Sender, Receiver) {
assert_ne!(capacity, 0, "capacity must be greater than zero");
let inner = Arc::new(Mutex::new(Inner::with_capacity(capacity))); let tx = Sender {
inner: Arc::clone(&inner),
}; let rx = Receiver { inner };
(tx, rx)
}
Die Informationen auf dieser Webseite wurden
nach bestem Wissen sorgfältig zusammengestellt. Es wird jedoch weder Vollständigkeit, noch Richtigkeit,
noch Qualität der bereit gestellten Informationen zugesichert.
Bemerkung:
Die farbliche Syntaxdarstellung und die Messung sind noch experimentell.