mirror of
https://github.com/tokio-rs/tokio.git
synced 2026-08-17 00:00:11 +02:00
sync: expand oneshot docs and TryRecvError (#1874)
`oneshot::Receiver::try_recv` does not provide any information as to the reason **why** receiving failed. The two cases are that the channel is empty or that the channel closed. `TryRecvError` is changed to be an enum of those two cases. This is backwards compatible as `TryRecvError` was an opaque struct. This also expands on `oneshot` API documentation, adding details and examples. Closes #1872
This commit is contained in:
+225
-26
@@ -41,7 +41,13 @@ pub mod error {
|
||||
|
||||
/// Error returned by the `try_recv` function on `Receiver`.
|
||||
#[derive(Debug)]
|
||||
pub struct TryRecvError(pub(super) ());
|
||||
pub enum TryRecvError {
|
||||
/// The send half of the channel has not yet sent a value.
|
||||
Empty,
|
||||
|
||||
/// The send half of the channel was dropped without sending a value.
|
||||
Closed,
|
||||
}
|
||||
|
||||
// ===== impl RecvError =====
|
||||
|
||||
@@ -57,7 +63,10 @@ pub mod error {
|
||||
|
||||
impl fmt::Display for TryRecvError {
|
||||
fn fmt(&self, fmt: &mut fmt::Formatter<'_>) -> fmt::Result {
|
||||
write!(fmt, "channel closed")
|
||||
match self {
|
||||
TryRecvError::Empty => write!(fmt, "channel empty"),
|
||||
TryRecvError::Closed => write!(fmt, "channel closed"),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -132,15 +141,43 @@ pub fn channel<T>() -> (Sender<T>, Receiver<T>) {
|
||||
}
|
||||
|
||||
impl<T> Sender<T> {
|
||||
/// Completes this oneshot with a successful result.
|
||||
/// Attempts to send a value on this channel, returning it back if it could
|
||||
/// not be sent.
|
||||
///
|
||||
/// The function consumes `self` and notifies the `Receiver` handle that a
|
||||
/// value is ready to be received.
|
||||
/// The function consumes `self` as only one value may ever be sent on a
|
||||
/// one-shot channel.
|
||||
///
|
||||
/// If the value is successfully enqueued for the remote end to receive,
|
||||
/// then `Ok(())` is returned. If the receiving end was dropped before this
|
||||
/// function was called, however, then `Err` is returned with the value
|
||||
/// provided.
|
||||
/// A successful send occurs when it is determined that the other end of the
|
||||
/// channel has not hung up already. An unsuccessful send would be one where
|
||||
/// the corresponding receiver has already been deallocated. Note that a
|
||||
/// return value of `Err` means that the data will never be received, but
|
||||
/// a return value of `Ok` does *not* mean that the data will be received.
|
||||
/// It is possible for the corresponding receiver to hang up immediately
|
||||
/// after this function returns `Ok`.
|
||||
///
|
||||
/// # Examples
|
||||
///
|
||||
/// Send a value to another task
|
||||
///
|
||||
/// ```
|
||||
/// use tokio::sync::oneshot;
|
||||
///
|
||||
/// #[tokio::main]
|
||||
/// async fn main() {
|
||||
/// let (tx, rx) = oneshot::channel();
|
||||
///
|
||||
/// tokio::spawn(async move {
|
||||
/// if let Err(_) = tx.send(3) {
|
||||
/// println!("the receiver dropped");
|
||||
/// }
|
||||
/// });
|
||||
///
|
||||
/// match rx.await {
|
||||
/// Ok(v) => println!("got = {:?}", v),
|
||||
/// Err(_) => println!("the sender dropped"),
|
||||
/// }
|
||||
/// }
|
||||
/// ```
|
||||
pub fn send(mut self, t: T) -> Result<(), T> {
|
||||
let inner = self.inner.take().unwrap();
|
||||
|
||||
@@ -200,16 +237,25 @@ impl<T> Sender<T> {
|
||||
Pending
|
||||
}
|
||||
|
||||
/// Wait for the associated [`Receiver`] handle to drop.
|
||||
/// Wait for the associated [`Receiver`] handle to close.
|
||||
///
|
||||
/// A [`Receiver`] is closed by either calling [`close`] explicitly or the
|
||||
/// [`Receiver`] value is dropped.
|
||||
///
|
||||
/// This function is useful when paired with `select!` to abort a
|
||||
/// computation when the receiver is no longer interested in the result.
|
||||
///
|
||||
/// # Return
|
||||
///
|
||||
/// Returns a `Future` which must be awaited on.
|
||||
///
|
||||
/// [`Receiver`]: struct.Receiver.html
|
||||
/// [`Receiver`]: Receiver
|
||||
/// [`close`]: Receiver::close
|
||||
///
|
||||
/// # Examples
|
||||
///
|
||||
/// Basic usage
|
||||
///
|
||||
/// ```
|
||||
/// use tokio::sync::oneshot;
|
||||
///
|
||||
@@ -225,19 +271,72 @@ impl<T> Sender<T> {
|
||||
/// println!("the receiver dropped");
|
||||
/// }
|
||||
/// ```
|
||||
///
|
||||
/// Paired with select
|
||||
///
|
||||
/// ```
|
||||
/// use tokio::sync::oneshot;
|
||||
/// use tokio::time::{self, Duration};
|
||||
///
|
||||
/// use futures::{select, FutureExt};
|
||||
///
|
||||
/// async fn compute() -> String {
|
||||
/// // Complex computation returning a `String`
|
||||
/// # "hello".to_string()
|
||||
/// }
|
||||
///
|
||||
/// #[tokio::main]
|
||||
/// async fn main() {
|
||||
/// let (mut tx, rx) = oneshot::channel();
|
||||
///
|
||||
/// tokio::spawn(async move {
|
||||
/// select! {
|
||||
/// _ = tx.closed().fuse() => {
|
||||
/// // The receiver dropped, no need to do any further work
|
||||
/// }
|
||||
/// value = compute().fuse() => {
|
||||
/// tx.send(value).unwrap()
|
||||
/// }
|
||||
/// }
|
||||
/// });
|
||||
///
|
||||
/// // Wait for up to 10 seconds
|
||||
/// let _ = time::timeout(Duration::from_secs(10), rx).await;
|
||||
/// }
|
||||
/// ```
|
||||
pub async fn closed(&mut self) {
|
||||
use crate::future::poll_fn;
|
||||
|
||||
poll_fn(|cx| self.poll_closed(cx)).await
|
||||
}
|
||||
|
||||
/// Check if the associated [`Receiver`] handle has been dropped.
|
||||
/// Returns `true` if the associated [`Receiver`] handle has been dropped.
|
||||
///
|
||||
/// Unlike [`poll_closed`], this function does not register a task for
|
||||
/// wakeup upon close.
|
||||
/// A [`Receiver`] is closed by either calling [`close`] explicitly or the
|
||||
/// [`Receiver`] value is dropped.
|
||||
///
|
||||
/// [`Receiver`]: struct.Receiver.html
|
||||
/// [`poll_closed`]: struct.Sender.html#method.poll_closed
|
||||
/// If `true` is returned, a call to `send` will always result in an error.
|
||||
///
|
||||
/// [`Receiver`]: Receiver
|
||||
/// [`close`]: Receiver::close
|
||||
///
|
||||
/// # Examples
|
||||
///
|
||||
/// ```
|
||||
/// use tokio::sync::oneshot;
|
||||
///
|
||||
/// #[tokio::main]
|
||||
/// async fn main() {
|
||||
/// let (tx, rx) = oneshot::channel();
|
||||
///
|
||||
/// assert!(!tx.is_closed());
|
||||
///
|
||||
/// drop(rx);
|
||||
///
|
||||
/// assert!(tx.is_closed());
|
||||
/// assert!(tx.send("never received").is_err());
|
||||
/// }
|
||||
/// ```
|
||||
pub fn is_closed(&self) -> bool {
|
||||
let inner = self.inner.as_ref().unwrap();
|
||||
|
||||
@@ -262,22 +361,122 @@ impl<T> Receiver<T> {
|
||||
/// receive a value if one was sent **before** the call to `close`
|
||||
/// completed.
|
||||
///
|
||||
/// [`Sender`]: struct.Sender.html
|
||||
/// This function is useful to perform a graceful shutdown and ensure that a
|
||||
/// value will not be sent into the channel and never received.
|
||||
///
|
||||
/// [`Sender`]: Sender
|
||||
///
|
||||
/// # Examples
|
||||
///
|
||||
/// Prevent a value from being sent
|
||||
///
|
||||
/// ```
|
||||
/// use tokio::sync::oneshot;
|
||||
/// use tokio::sync::oneshot::error::TryRecvError;
|
||||
///
|
||||
/// #[tokio::main]
|
||||
/// async fn main() {
|
||||
/// let (tx, mut rx) = oneshot::channel();
|
||||
///
|
||||
/// assert!(!tx.is_closed());
|
||||
///
|
||||
/// rx.close();
|
||||
///
|
||||
/// assert!(tx.is_closed());
|
||||
/// assert!(tx.send("never received").is_err());
|
||||
///
|
||||
/// match rx.try_recv() {
|
||||
/// Err(TryRecvError::Closed) => {}
|
||||
/// _ => unreachable!(),
|
||||
/// }
|
||||
/// }
|
||||
/// ```
|
||||
///
|
||||
/// Receive a value sent **before** calling `close`
|
||||
///
|
||||
/// ```
|
||||
/// use tokio::sync::oneshot;
|
||||
///
|
||||
/// #[tokio::main]
|
||||
/// async fn main() {
|
||||
/// let (tx, mut rx) = oneshot::channel();
|
||||
///
|
||||
/// assert!(tx.send("will receive").is_ok());
|
||||
///
|
||||
/// rx.close();
|
||||
///
|
||||
/// let msg = rx.try_recv().unwrap();
|
||||
/// assert_eq!(msg, "will receive");
|
||||
/// }
|
||||
/// ```
|
||||
pub fn close(&mut self) {
|
||||
let inner = self.inner.as_ref().unwrap();
|
||||
inner.close();
|
||||
}
|
||||
|
||||
/// Attempts to receive a value outside of the context of a task.
|
||||
/// Attempts to receive a value.
|
||||
///
|
||||
/// Does not register a task if no value has been sent.
|
||||
/// If a pending value exists in the channel, it is returned. If no value
|
||||
/// has been sent, the current task **will not** be registered for
|
||||
/// future notification.
|
||||
///
|
||||
/// A return value of `None` must be considered immediately stale (out of
|
||||
/// date) unless [`close`] has been called first.
|
||||
/// This function is useful to call from outside the context of an
|
||||
/// asynchronous task.
|
||||
///
|
||||
/// Returns an error if the sender was dropped.
|
||||
/// # Return
|
||||
///
|
||||
/// [`close`]: #method.close
|
||||
/// - `Ok(T)` if a value is pending in the channel.
|
||||
/// - `Err(TryRecvError::Empty)` if no value has been sent yet.
|
||||
/// - `Err(TryRecvError::Closed)` if the sender has dropped without sending
|
||||
/// a value.
|
||||
///
|
||||
/// # Examples
|
||||
///
|
||||
/// `try_recv` before a value is sent, then after.
|
||||
///
|
||||
/// ```
|
||||
/// use tokio::sync::oneshot;
|
||||
/// use tokio::sync::oneshot::error::TryRecvError;
|
||||
///
|
||||
/// #[tokio::main]
|
||||
/// async fn main() {
|
||||
/// let (tx, mut rx) = oneshot::channel();
|
||||
///
|
||||
/// match rx.try_recv() {
|
||||
/// // The channel is currently empty
|
||||
/// Err(TryRecvError::Empty) => {}
|
||||
/// _ => unreachable!(),
|
||||
/// }
|
||||
///
|
||||
/// // Send a value
|
||||
/// tx.send("hello").unwrap();
|
||||
///
|
||||
/// match rx.try_recv() {
|
||||
/// Ok(value) => assert_eq!(value, "hello"),
|
||||
/// _ => unreachable!(),
|
||||
/// }
|
||||
/// }
|
||||
/// ```
|
||||
///
|
||||
/// `try_recv` when the sender dropped before sending a value
|
||||
///
|
||||
/// ```
|
||||
/// use tokio::sync::oneshot;
|
||||
/// use tokio::sync::oneshot::error::TryRecvError;
|
||||
///
|
||||
/// #[tokio::main]
|
||||
/// async fn main() {
|
||||
/// let (tx, mut rx) = oneshot::channel::<()>();
|
||||
///
|
||||
/// drop(tx);
|
||||
///
|
||||
/// match rx.try_recv() {
|
||||
/// // The channel will never receive a value.
|
||||
/// Err(TryRecvError::Closed) => {}
|
||||
/// _ => unreachable!(),
|
||||
/// }
|
||||
/// }
|
||||
/// ```
|
||||
pub fn try_recv(&mut self) -> Result<T, TryRecvError> {
|
||||
let result = if let Some(inner) = self.inner.as_ref() {
|
||||
let state = State::load(&inner.state, Acquire);
|
||||
@@ -285,13 +484,13 @@ impl<T> Receiver<T> {
|
||||
if state.is_complete() {
|
||||
match unsafe { inner.consume_value() } {
|
||||
Some(value) => Ok(value),
|
||||
None => Err(TryRecvError(())),
|
||||
None => Err(TryRecvError::Closed),
|
||||
}
|
||||
} else if state.is_closed() {
|
||||
Err(TryRecvError(()))
|
||||
Err(TryRecvError::Closed)
|
||||
} else {
|
||||
// Not ready, this does not clear `inner`
|
||||
return Err(TryRecvError(()));
|
||||
return Err(TryRecvError::Empty);
|
||||
}
|
||||
} else {
|
||||
panic!("called after complete");
|
||||
|
||||
Reference in New Issue
Block a user