mirror of
https://github.com/tokio-rs/tokio.git
synced 2026-08-28 00:00:11 +02:00
sync: Add is_closed method to mpsc senders (#2726)
Co-authored-by: Alice Ryhl <[email protected]>
This commit is contained in:
co-authored by
Alice Ryhl
parent
99d4061203
commit
078d0a2ebc
@@ -523,6 +523,28 @@ impl<T> Sender<T> {
|
|||||||
enter_handle.block_on(self.send(value)).unwrap()
|
enter_handle.block_on(self.send(value)).unwrap()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Checks if the channel has been closed. This happens when the
|
||||||
|
/// [`Receiver`] is dropped, or when the [`Receiver::close`] method is
|
||||||
|
/// called.
|
||||||
|
///
|
||||||
|
/// [`Receiver`]: crate::sync::mpsc::Receiver
|
||||||
|
/// [`Receiver::close`]: crate::sync::mpsc::Receiver::close
|
||||||
|
///
|
||||||
|
/// ```
|
||||||
|
/// let (tx, rx) = tokio::sync::mpsc::channel::<()>(42);
|
||||||
|
/// assert!(!tx.is_closed());
|
||||||
|
///
|
||||||
|
/// let tx2 = tx.clone();
|
||||||
|
/// assert!(!tx2.is_closed());
|
||||||
|
///
|
||||||
|
/// drop(rx);
|
||||||
|
/// assert!(tx.is_closed());
|
||||||
|
/// assert!(tx2.is_closed());
|
||||||
|
/// ```
|
||||||
|
pub fn is_closed(&self) -> bool {
|
||||||
|
self.chan.is_closed()
|
||||||
|
}
|
||||||
|
|
||||||
/// Wait for channel capacity. Once capacity to send one message is
|
/// Wait for channel capacity. Once capacity to send one message is
|
||||||
/// available, it is reserved for the caller.
|
/// available, it is reserved for the caller.
|
||||||
///
|
///
|
||||||
|
|||||||
@@ -143,6 +143,10 @@ impl<T, S> Tx<T, S> {
|
|||||||
}
|
}
|
||||||
|
|
||||||
impl<T, S: Semaphore> Tx<T, S> {
|
impl<T, S: Semaphore> Tx<T, S> {
|
||||||
|
pub(crate) fn is_closed(&self) -> bool {
|
||||||
|
self.inner.semaphore.is_closed()
|
||||||
|
}
|
||||||
|
|
||||||
pub(crate) async fn closed(&mut self) {
|
pub(crate) async fn closed(&mut self) {
|
||||||
use std::future::Future;
|
use std::future::Future;
|
||||||
use std::pin::Pin;
|
use std::pin::Pin;
|
||||||
|
|||||||
@@ -245,4 +245,25 @@ impl<T> UnboundedSender<T> {
|
|||||||
pub async fn closed(&mut self) {
|
pub async fn closed(&mut self) {
|
||||||
self.chan.closed().await
|
self.chan.closed().await
|
||||||
}
|
}
|
||||||
|
/// Checks if the channel has been closed. This happens when the
|
||||||
|
/// [`UnboundedReceiver`] is dropped, or when the
|
||||||
|
/// [`UnboundedReceiver::close`] method is called.
|
||||||
|
///
|
||||||
|
/// [`UnboundedReceiver`]: crate::sync::mpsc::UnboundedReceiver
|
||||||
|
/// [`UnboundedReceiver::close`]: crate::sync::mpsc::UnboundedReceiver::close
|
||||||
|
///
|
||||||
|
/// ```
|
||||||
|
/// let (tx, rx) = tokio::sync::mpsc::unbounded_channel::<()>();
|
||||||
|
/// assert!(!tx.is_closed());
|
||||||
|
///
|
||||||
|
/// let tx2 = tx.clone();
|
||||||
|
/// assert!(!tx2.is_closed());
|
||||||
|
///
|
||||||
|
/// drop(rx);
|
||||||
|
/// assert!(tx.is_closed());
|
||||||
|
/// assert!(tx2.is_closed());
|
||||||
|
/// ```
|
||||||
|
pub fn is_closed(&self) -> bool {
|
||||||
|
self.chan.is_closed()
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user