diff --git a/tokio/src/sync/oneshot.rs b/tokio/src/sync/oneshot.rs index ed3801c81..67094645e 100644 --- a/tokio/src/sync/oneshot.rs +++ b/tokio/src/sync/oneshot.rs @@ -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() -> (Sender, Receiver) { } impl Sender { - /// 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 Sender { 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 Sender { /// 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 Receiver { /// 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 { let result = if let Some(inner) = self.inner.as_ref() { let state = State::load(&inner.state, Acquire); @@ -285,13 +484,13 @@ impl Receiver { 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");