Skip to main content

tokio/sync/mpsc/
unbounded.rs

1use crate::loom::sync::{atomic::AtomicUsize, Arc};
2use crate::sync::mpsc::chan;
3use crate::sync::mpsc::error::{SendError, TryRecvError};
4
5use std::fmt;
6use std::task::{Context, Poll};
7
8/// Send values to the associated `UnboundedReceiver`.
9///
10/// Instances are created by the [`unbounded_channel`] function.
11pub struct UnboundedSender<T> {
12    chan: chan::Tx<T, Semaphore>,
13}
14
15/// An unbounded sender that does not prevent the channel from being closed.
16///
17/// If all [`UnboundedSender`] instances of a channel were dropped and only
18/// `WeakUnboundedSender` instances remain, the channel is closed.
19///
20/// In order to send messages, the `WeakUnboundedSender` needs to be upgraded using
21/// [`WeakUnboundedSender::upgrade`], which returns `Option<UnboundedSender>`. It returns `None`
22/// if all `UnboundedSender`s have been dropped, and otherwise it returns an `UnboundedSender`.
23///
24/// [`UnboundedSender`]: UnboundedSender
25/// [`WeakUnboundedSender::upgrade`]: WeakUnboundedSender::upgrade
26///
27/// # Examples
28///
29/// ```
30/// use tokio::sync::mpsc::unbounded_channel;
31///
32/// # #[tokio::main(flavor = "current_thread")]
33/// # async fn main() {
34/// let (tx, _rx) = unbounded_channel::<i32>();
35/// let tx_weak = tx.downgrade();
36///
37/// // Upgrading will succeed because `tx` still exists.
38/// assert!(tx_weak.upgrade().is_some());
39///
40/// // If we drop `tx`, then it will fail.
41/// drop(tx);
42/// assert!(tx_weak.clone().upgrade().is_none());
43/// # }
44/// ```
45pub struct WeakUnboundedSender<T> {
46    chan: Arc<chan::Chan<T, Semaphore>>,
47}
48
49impl<T> Clone for UnboundedSender<T> {
50    fn clone(&self) -> Self {
51        UnboundedSender {
52            chan: self.chan.clone(),
53        }
54    }
55}
56
57impl<T> fmt::Debug for UnboundedSender<T> {
58    fn fmt(&self, fmt: &mut fmt::Formatter<'_>) -> fmt::Result {
59        fmt.debug_struct("UnboundedSender")
60            .field("chan", &self.chan)
61            .finish()
62    }
63}
64
65/// Receive values from the associated `UnboundedSender`.
66///
67/// Instances are created by the [`unbounded_channel`] function.
68///
69/// This receiver can be turned into a `Stream` using [`UnboundedReceiverStream`].
70///
71/// [`UnboundedReceiverStream`]: https://docs.rs/tokio-stream/0.1/tokio_stream/wrappers/struct.UnboundedReceiverStream.html
72pub struct UnboundedReceiver<T> {
73    /// The channel receiver
74    chan: chan::Rx<T, Semaphore>,
75}
76
77impl<T> fmt::Debug for UnboundedReceiver<T> {
78    fn fmt(&self, fmt: &mut fmt::Formatter<'_>) -> fmt::Result {
79        fmt.debug_struct("UnboundedReceiver")
80            .field("chan", &self.chan)
81            .finish()
82    }
83}
84
85/// Creates an unbounded mpsc channel for communicating between asynchronous
86/// tasks without backpressure.
87///
88/// A `send` on this channel will always succeed as long as the receive half has
89/// not been closed. If the receiver falls behind, messages will be arbitrarily
90/// buffered.
91///
92/// **Note** that the amount of available system memory is an implicit bound to
93/// the channel. Using an `unbounded` channel has the ability of causing the
94/// process to run out of memory. In this case, the process will be aborted.
95pub fn unbounded_channel<T>() -> (UnboundedSender<T>, UnboundedReceiver<T>) {
96    let (tx, rx) = chan::channel(Semaphore(AtomicUsize::new(0)));
97
98    let tx = UnboundedSender::new(tx);
99    let rx = UnboundedReceiver::new(rx);
100
101    (tx, rx)
102}
103
104#[cfg(all(test, not(loom)))]
105pub(crate) fn unbounded_channel_from_index<T>(
106    start_index: usize,
107) -> (UnboundedSender<T>, UnboundedReceiver<T>) {
108    let (tx, rx) = chan::channel_from_index(start_index, Semaphore(AtomicUsize::new(0)));
109
110    let tx = UnboundedSender::new(tx);
111    let rx = UnboundedReceiver::new(rx);
112
113    (tx, rx)
114}
115
116/// No capacity
117#[derive(Debug)]
118pub(crate) struct Semaphore(pub(crate) AtomicUsize);
119
120impl<T> UnboundedReceiver<T> {
121    pub(crate) fn new(chan: chan::Rx<T, Semaphore>) -> UnboundedReceiver<T> {
122        UnboundedReceiver { chan }
123    }
124
125    /// Receives the next value for this receiver.
126    ///
127    /// This method returns `None` if the channel has been closed and there are
128    /// no remaining messages in the channel's buffer. This indicates that no
129    /// further values can ever be received from this `Receiver`. The channel is
130    /// closed when all senders have been dropped, or when [`close`] is called.
131    ///
132    /// If there are no messages in the channel's buffer, but the channel has
133    /// not yet been closed, this method will sleep until a message is sent or
134    /// the channel is closed.
135    ///
136    /// # Cancel safety
137    ///
138    /// This method is cancel safe. If `recv` is used as a branch in
139    /// [`tokio::select!`](crate::select) and another branch completes first,
140    /// it is guaranteed that no messages were received on this
141    /// channel.
142    ///
143    /// [`close`]: Self::close
144    ///
145    /// # Examples
146    ///
147    /// ```
148    /// use tokio::sync::mpsc;
149    ///
150    /// # #[tokio::main(flavor = "current_thread")]
151    /// # async fn main() {
152    /// let (tx, mut rx) = mpsc::unbounded_channel();
153    ///
154    /// tokio::spawn(async move {
155    ///     tx.send("hello").unwrap();
156    /// });
157    ///
158    /// assert_eq!(Some("hello"), rx.recv().await);
159    /// assert_eq!(None, rx.recv().await);
160    /// # }
161    /// ```
162    ///
163    /// Values are buffered:
164    ///
165    /// ```
166    /// use tokio::sync::mpsc;
167    ///
168    /// # #[tokio::main(flavor = "current_thread")]
169    /// # async fn main() {
170    /// let (tx, mut rx) = mpsc::unbounded_channel();
171    ///
172    /// tx.send("hello").unwrap();
173    /// tx.send("world").unwrap();
174    ///
175    /// assert_eq!(Some("hello"), rx.recv().await);
176    /// assert_eq!(Some("world"), rx.recv().await);
177    /// # }
178    /// ```
179    pub async fn recv(&mut self) -> Option<T> {
180        use std::future::poll_fn;
181
182        poll_fn(|cx| self.poll_recv(cx)).await
183    }
184
185    /// Receives the next values for this receiver and extends `buffer`.
186    ///
187    /// This method extends `buffer` by no more than a fixed number of values
188    /// as specified by `limit`. If `limit` is zero, the function returns
189    /// immediately with `0`. The return value is the number of values added to
190    /// `buffer`.
191    ///
192    /// For `limit > 0`, if there are no messages in the channel's queue,
193    /// but the channel has not yet been closed, this method will sleep
194    /// until a message is sent or the channel is closed.
195    ///
196    /// For non-zero values of `limit`, this method will never return `0` unless
197    /// the channel has been closed and there are no remaining messages in the
198    /// channel's queue. This indicates that no further values can ever be
199    /// received from this `Receiver`. The channel is closed when all senders
200    /// have been dropped, or when [`close`] is called.
201    ///
202    /// The capacity of `buffer` is increased as needed.
203    ///
204    /// # Cancel safety
205    ///
206    /// This method is cancel safe. If `recv_many` is used as a branch in
207    /// [`tokio::select!`](crate::select) and another branch completes first,
208    /// it is guaranteed that no messages were received on this
209    /// channel.
210    ///
211    /// [`close`]: Self::close
212    ///
213    /// # Examples
214    ///
215    /// ```
216    /// use tokio::sync::mpsc;
217    ///
218    /// # #[tokio::main(flavor = "current_thread")]
219    /// # async fn main() {
220    /// let mut buffer: Vec<&str> = Vec::with_capacity(2);
221    /// let limit = 2;
222    /// let (tx, mut rx) = mpsc::unbounded_channel();
223    /// let tx2 = tx.clone();
224    /// tx2.send("first").unwrap();
225    /// tx2.send("second").unwrap();
226    /// tx2.send("third").unwrap();
227    ///
228    /// // Call `recv_many` to receive up to `limit` (2) values.
229    /// assert_eq!(2, rx.recv_many(&mut buffer, limit).await);
230    /// assert_eq!(vec!["first", "second"], buffer);
231    ///
232    /// // If the buffer is full, the next call to `recv_many`
233    /// // reserves additional capacity.
234    /// assert_eq!(1, rx.recv_many(&mut buffer, limit).await);
235    ///
236    /// tokio::spawn(async move {
237    ///     tx.send("fourth").unwrap();
238    /// });
239    ///
240    /// // 'tx' is dropped, but `recv_many`
241    /// // is guaranteed not to return 0 as the channel
242    /// // is not yet closed.
243    /// assert_eq!(1, rx.recv_many(&mut buffer, limit).await);
244    /// assert_eq!(vec!["first", "second", "third", "fourth"], buffer);
245    ///
246    /// // Once the last sender is dropped, the channel is
247    /// // closed and `recv_many` returns 0, capacity unchanged.
248    /// drop(tx2);
249    /// assert_eq!(0, rx.recv_many(&mut buffer, limit).await);
250    /// assert_eq!(vec!["first", "second", "third", "fourth"], buffer);
251    /// # }
252    /// ```
253    pub async fn recv_many(&mut self, buffer: &mut Vec<T>, limit: usize) -> usize {
254        use std::future::poll_fn;
255        poll_fn(|cx| self.chan.recv_many(cx, buffer, limit)).await
256    }
257
258    /// Tries to receive the next value for this receiver.
259    ///
260    /// This method returns the [`Empty`] error if the channel is currently
261    /// empty, but there are still outstanding [senders] or [permits].
262    ///
263    /// This method returns the [`Disconnected`] error if the channel is
264    /// currently empty, and there are no outstanding [senders] or [permits].
265    ///
266    /// Unlike the [`poll_recv`] method, this method will never return an
267    /// [`Empty`] error spuriously.
268    ///
269    /// [`Empty`]: crate::sync::mpsc::error::TryRecvError::Empty
270    /// [`Disconnected`]: crate::sync::mpsc::error::TryRecvError::Disconnected
271    /// [`poll_recv`]: Self::poll_recv
272    /// [senders]: crate::sync::mpsc::Sender
273    /// [permits]: crate::sync::mpsc::Permit
274    ///
275    /// # Examples
276    ///
277    /// ```
278    /// use tokio::sync::mpsc;
279    /// use tokio::sync::mpsc::error::TryRecvError;
280    ///
281    /// # #[tokio::main(flavor = "current_thread")]
282    /// # async fn main() {
283    /// let (tx, mut rx) = mpsc::unbounded_channel();
284    ///
285    /// tx.send("hello").unwrap();
286    ///
287    /// assert_eq!(Ok("hello"), rx.try_recv());
288    /// assert_eq!(Err(TryRecvError::Empty), rx.try_recv());
289    ///
290    /// tx.send("hello").unwrap();
291    /// // Drop the last sender, closing the channel.
292    /// drop(tx);
293    ///
294    /// assert_eq!(Ok("hello"), rx.try_recv());
295    /// assert_eq!(Err(TryRecvError::Disconnected), rx.try_recv());
296    /// # }
297    /// ```
298    pub fn try_recv(&mut self) -> Result<T, TryRecvError> {
299        self.chan.try_recv()
300    }
301
302    /// Blocking receive to call outside of asynchronous contexts.
303    ///
304    /// # Panics
305    ///
306    /// This function panics if called within an asynchronous execution
307    /// context.
308    ///
309    /// # Examples
310    ///
311    /// ```
312    /// # #[cfg(not(target_family = "wasm"))]
313    /// # {
314    /// use std::thread;
315    /// use tokio::sync::mpsc;
316    ///
317    /// #[tokio::main]
318    /// async fn main() {
319    ///     let (tx, mut rx) = mpsc::unbounded_channel::<u8>();
320    ///
321    ///     let sync_code = thread::spawn(move || {
322    ///         assert_eq!(Some(10), rx.blocking_recv());
323    ///     });
324    ///
325    ///     let _ = tx.send(10);
326    ///     sync_code.join().unwrap();
327    /// }
328    /// # }
329    /// ```
330    #[track_caller]
331    #[cfg(feature = "sync")]
332    #[cfg_attr(docsrs, doc(alias = "recv_blocking"))]
333    pub fn blocking_recv(&mut self) -> Option<T> {
334        crate::future::block_on(self.recv())
335    }
336
337    /// Variant of [`Self::recv_many`] for blocking contexts.
338    ///
339    /// The same conditions as in [`Self::blocking_recv`] apply.
340    #[track_caller]
341    #[cfg(feature = "sync")]
342    #[cfg_attr(docsrs, doc(alias = "recv_many_blocking"))]
343    pub fn blocking_recv_many(&mut self, buffer: &mut Vec<T>, limit: usize) -> usize {
344        crate::future::block_on(self.recv_many(buffer, limit))
345    }
346
347    /// Closes the receiving half of a channel, without dropping it.
348    ///
349    /// This prevents any further messages from being sent on the channel while
350    /// still enabling the receiver to drain messages that are buffered.
351    ///
352    /// To guarantee that no messages are dropped, after calling `close()`,
353    /// `recv()` must be called until `None` is returned.
354    pub fn close(&mut self) {
355        self.chan.close();
356    }
357
358    /// Checks if a channel is closed.
359    ///
360    /// This method returns `true` if the channel has been closed. The channel is closed
361    /// when all [`UnboundedSender`] have been dropped, or when [`UnboundedReceiver::close`] is called.
362    ///
363    /// [`UnboundedSender`]: crate::sync::mpsc::UnboundedSender
364    /// [`UnboundedReceiver::close`]: crate::sync::mpsc::UnboundedReceiver::close
365    ///
366    /// # Examples
367    /// ```
368    /// use tokio::sync::mpsc;
369    ///
370    /// # #[tokio::main(flavor = "current_thread")]
371    /// # async fn main() {
372    /// let (_tx, mut rx) = mpsc::unbounded_channel::<()>();
373    /// assert!(!rx.is_closed());
374    ///
375    /// rx.close();
376    ///
377    /// assert!(rx.is_closed());
378    /// # }
379    /// ```
380    pub fn is_closed(&self) -> bool {
381        self.chan.is_closed()
382    }
383
384    /// Checks if a channel is empty.
385    ///
386    /// This method returns `true` if the channel has no messages.
387    ///
388    /// # Examples
389    /// ```
390    /// use tokio::sync::mpsc;
391    ///
392    /// # #[tokio::main(flavor = "current_thread")]
393    /// # async fn main() {
394    /// let (tx, rx) = mpsc::unbounded_channel();
395    /// assert!(rx.is_empty());
396    ///
397    /// tx.send(0).unwrap();
398    /// assert!(!rx.is_empty());
399    /// # }
400    ///
401    /// ```
402    pub fn is_empty(&self) -> bool {
403        self.chan.is_empty()
404    }
405
406    /// Returns the number of messages in the channel.
407    ///
408    /// # Examples
409    /// ```
410    /// use tokio::sync::mpsc;
411    ///
412    /// # #[tokio::main(flavor = "current_thread")]
413    /// # async fn main() {
414    /// let (tx, rx) = mpsc::unbounded_channel();
415    /// assert_eq!(0, rx.len());
416    ///
417    /// tx.send(0).unwrap();
418    /// assert_eq!(1, rx.len());
419    /// # }
420    /// ```
421    pub fn len(&self) -> usize {
422        self.chan.len()
423    }
424
425    /// Polls to receive the next message on this channel.
426    ///
427    /// This method returns:
428    ///
429    ///  * `Poll::Pending` if no messages are available but the channel is not
430    ///    closed, or if a spurious failure happens.
431    ///  * `Poll::Ready(Some(message))` if a message is available.
432    ///  * `Poll::Ready(None)` if the channel has been closed and all messages
433    ///    sent before it was closed have been received.
434    ///
435    /// When the method returns `Poll::Pending`, the `Waker` in the provided
436    /// `Context` is scheduled to receive a wakeup when a message is sent on any
437    /// receiver, or when the channel is closed.  Note that on multiple calls to
438    /// `poll_recv` or `poll_recv_many`, only the `Waker` from the `Context`
439    /// passed to the most recent call is scheduled to receive a wakeup.
440    ///
441    /// If this method returns `Poll::Pending` due to a spurious failure, then
442    /// the `Waker` will be notified when the situation causing the spurious
443    /// failure has been resolved. Note that receiving such a wakeup does not
444    /// guarantee that the next call will succeed — it could fail with another
445    /// spurious failure.
446    pub fn poll_recv(&mut self, cx: &mut Context<'_>) -> Poll<Option<T>> {
447        self.chan.recv(cx)
448    }
449
450    /// Polls to receive multiple messages on this channel, extending the provided buffer.
451    ///
452    /// This method returns:
453    /// * `Poll::Pending` if no messages are available but the channel is not closed, or if a
454    ///   spurious failure happens.
455    /// * `Poll::Ready(count)` where `count` is the number of messages successfully received and
456    ///   stored in `buffer`. This can be less than, or equal to, `limit`.
457    /// * `Poll::Ready(0)` if `limit` is set to zero or when the channel is closed.
458    ///
459    /// When the method returns `Poll::Pending`, the `Waker` in the provided
460    /// `Context` is scheduled to receive a wakeup when a message is sent on any
461    /// receiver, or when the channel is closed.  Note that on multiple calls to
462    /// `poll_recv` or `poll_recv_many`, only the `Waker` from the `Context`
463    /// passed to the most recent call is scheduled to receive a wakeup.
464    ///
465    /// Note that this method does not guarantee that exactly `limit` messages
466    /// are received. Rather, if at least one message is available, it returns
467    /// as many messages as it can up to the given limit. This method returns
468    /// zero only if the channel is closed (or if `limit` is zero).
469    ///
470    /// # Examples
471    ///
472    /// ```
473    /// # #[cfg(not(target_family = "wasm"))]
474    /// # {
475    /// use std::task::{Context, Poll};
476    /// use std::pin::Pin;
477    /// use tokio::sync::mpsc;
478    /// use futures::Future;
479    ///
480    /// struct MyReceiverFuture<'a> {
481    ///     receiver: mpsc::UnboundedReceiver<i32>,
482    ///     buffer: &'a mut Vec<i32>,
483    ///     limit: usize,
484    /// }
485    ///
486    /// impl<'a> Future for MyReceiverFuture<'a> {
487    ///     type Output = usize; // Number of messages received
488    ///
489    ///     fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
490    ///         let MyReceiverFuture { receiver, buffer, limit } = &mut *self;
491    ///
492    ///         // Now `receiver` and `buffer` are mutable references, and `limit` is copied
493    ///         match receiver.poll_recv_many(cx, *buffer, *limit) {
494    ///             Poll::Pending => Poll::Pending,
495    ///             Poll::Ready(count) => Poll::Ready(count),
496    ///         }
497    ///     }
498    /// }
499    ///
500    /// # #[tokio::main(flavor = "current_thread")]
501    /// # async fn main() {
502    /// let (tx, rx) = mpsc::unbounded_channel::<i32>();
503    /// let mut buffer = Vec::new();
504    ///
505    /// let my_receiver_future = MyReceiverFuture {
506    ///     receiver: rx,
507    ///     buffer: &mut buffer,
508    ///     limit: 3,
509    /// };
510    ///
511    /// for i in 0..10 {
512    ///     tx.send(i).expect("Unable to send integer");
513    /// }
514    ///
515    /// let count = my_receiver_future.await;
516    /// assert_eq!(count, 3);
517    /// assert_eq!(buffer, vec![0,1,2])
518    /// # }
519    /// # }
520    /// ```
521    pub fn poll_recv_many(
522        &mut self,
523        cx: &mut Context<'_>,
524        buffer: &mut Vec<T>,
525        limit: usize,
526    ) -> Poll<usize> {
527        self.chan.recv_many(cx, buffer, limit)
528    }
529
530    /// Returns the number of [`UnboundedSender`] handles.
531    pub fn sender_strong_count(&self) -> usize {
532        self.chan.sender_strong_count()
533    }
534
535    /// Returns the number of [`WeakUnboundedSender`] handles.
536    pub fn sender_weak_count(&self) -> usize {
537        self.chan.sender_weak_count()
538    }
539}
540
541impl<T> UnboundedSender<T> {
542    pub(crate) fn new(chan: chan::Tx<T, Semaphore>) -> UnboundedSender<T> {
543        UnboundedSender { chan }
544    }
545
546    /// Attempts to send a message on this `UnboundedSender` without blocking.
547    ///
548    /// This method is not marked as `async` because sending a message to an unbounded channel
549    /// never requires any form of waiting. This is due to the channel's infinite capacity,
550    /// allowing the `send` operation to complete immediately. As a result, the `send` method can be
551    /// used in both synchronous and asynchronous code without issues.
552    ///
553    /// If the receive half of the channel is closed, either due to [`close`]
554    /// being called or the [`UnboundedReceiver`] having been dropped, this
555    /// function returns an error. The error includes the value passed to `send`.
556    ///
557    /// [`close`]: UnboundedReceiver::close
558    /// [`UnboundedReceiver`]: UnboundedReceiver
559    pub fn send(&self, message: T) -> Result<(), SendError<T>> {
560        if !self.inc_num_messages() {
561            return Err(SendError(message));
562        }
563
564        self.chan.send(message);
565        Ok(())
566    }
567
568    fn inc_num_messages(&self) -> bool {
569        use std::process;
570        use std::sync::atomic::Ordering::{AcqRel, Acquire};
571
572        let mut curr = self.chan.semaphore().0.load(Acquire);
573
574        loop {
575            if curr & 1 == 1 {
576                return false;
577            }
578
579            if curr == usize::MAX ^ 1 {
580                // Overflowed the ref count. There is no safe way to recover, so
581                // abort the process. In practice, this should never happen.
582                process::abort()
583            }
584
585            match self
586                .chan
587                .semaphore()
588                .0
589                .compare_exchange(curr, curr + 2, AcqRel, Acquire)
590            {
591                Ok(_) => return true,
592                Err(actual) => {
593                    curr = actual;
594                }
595            }
596        }
597    }
598
599    /// Completes when the receiver has dropped.
600    ///
601    /// This allows the producers to get notified when interest in the produced
602    /// values is canceled and immediately stop doing work.
603    ///
604    /// # Cancel safety
605    ///
606    /// This method is cancel safe. Once the channel is closed, it stays closed
607    /// forever and all future calls to `closed` will return immediately.
608    ///
609    /// # Examples
610    ///
611    /// ```
612    /// use tokio::sync::mpsc;
613    ///
614    /// # #[tokio::main(flavor = "current_thread")]
615    /// # async fn main() {
616    /// let (tx1, rx) = mpsc::unbounded_channel::<()>();
617    /// let tx2 = tx1.clone();
618    /// let tx3 = tx1.clone();
619    /// let tx4 = tx1.clone();
620    /// let tx5 = tx1.clone();
621    /// tokio::spawn(async move {
622    ///     drop(rx);
623    /// });
624    ///
625    /// futures::join!(
626    ///     tx1.closed(),
627    ///     tx2.closed(),
628    ///     tx3.closed(),
629    ///     tx4.closed(),
630    ///     tx5.closed()
631    /// );
632    /// println!("Receiver dropped");
633    /// # }
634    /// ```
635    pub async fn closed(&self) {
636        self.chan.closed().await;
637    }
638
639    /// Checks if the channel has been closed. This happens when the
640    /// [`UnboundedReceiver`] is dropped, or when the
641    /// [`UnboundedReceiver::close`] method is called.
642    ///
643    /// [`UnboundedReceiver`]: crate::sync::mpsc::UnboundedReceiver
644    /// [`UnboundedReceiver::close`]: crate::sync::mpsc::UnboundedReceiver::close
645    ///
646    /// ```
647    /// let (tx, rx) = tokio::sync::mpsc::unbounded_channel::<()>();
648    /// assert!(!tx.is_closed());
649    ///
650    /// let tx2 = tx.clone();
651    /// assert!(!tx2.is_closed());
652    ///
653    /// drop(rx);
654    /// assert!(tx.is_closed());
655    /// assert!(tx2.is_closed());
656    /// ```
657    pub fn is_closed(&self) -> bool {
658        self.chan.is_closed()
659    }
660
661    /// Returns `true` if senders belong to the same channel.
662    ///
663    /// # Examples
664    ///
665    /// ```
666    /// let (tx, rx) = tokio::sync::mpsc::unbounded_channel::<()>();
667    /// let  tx2 = tx.clone();
668    /// assert!(tx.same_channel(&tx2));
669    ///
670    /// let (tx3, rx3) = tokio::sync::mpsc::unbounded_channel::<()>();
671    /// assert!(!tx3.same_channel(&tx2));
672    /// ```
673    pub fn same_channel(&self, other: &Self) -> bool {
674        self.chan.same_channel(&other.chan)
675    }
676
677    /// Converts the `UnboundedSender` to a [`WeakUnboundedSender`] that does not count
678    /// towards RAII semantics, i.e. if all `UnboundedSender` instances of the
679    /// channel were dropped and only `WeakUnboundedSender` instances remain,
680    /// the channel is closed.
681    #[must_use = "Downgrade creates a WeakSender without destroying the original non-weak sender."]
682    pub fn downgrade(&self) -> WeakUnboundedSender<T> {
683        WeakUnboundedSender {
684            chan: self.chan.downgrade(),
685        }
686    }
687
688    /// Returns the number of [`UnboundedSender`] handles.
689    pub fn strong_count(&self) -> usize {
690        self.chan.strong_count()
691    }
692
693    /// Returns the number of [`WeakUnboundedSender`] handles.
694    pub fn weak_count(&self) -> usize {
695        self.chan.weak_count()
696    }
697}
698
699impl<T> Clone for WeakUnboundedSender<T> {
700    fn clone(&self) -> Self {
701        self.chan.increment_weak_count();
702
703        WeakUnboundedSender {
704            chan: self.chan.clone(),
705        }
706    }
707}
708
709impl<T> Drop for WeakUnboundedSender<T> {
710    fn drop(&mut self) {
711        self.chan.decrement_weak_count();
712    }
713}
714
715impl<T> WeakUnboundedSender<T> {
716    /// Tries to convert a `WeakUnboundedSender` into an [`UnboundedSender`].
717    /// This will return `Some` if there are other `Sender` instances alive and
718    /// the channel wasn't previously dropped, otherwise `None` is returned.
719    pub fn upgrade(&self) -> Option<UnboundedSender<T>> {
720        chan::Tx::upgrade(self.chan.clone()).map(UnboundedSender::new)
721    }
722
723    /// Returns the number of [`UnboundedSender`] handles.
724    pub fn strong_count(&self) -> usize {
725        self.chan.strong_count()
726    }
727
728    /// Returns the number of [`WeakUnboundedSender`] handles.
729    pub fn weak_count(&self) -> usize {
730        self.chan.weak_count()
731    }
732}
733
734impl<T> fmt::Debug for WeakUnboundedSender<T> {
735    fn fmt(&self, fmt: &mut fmt::Formatter<'_>) -> fmt::Result {
736        fmt.debug_struct("WeakUnboundedSender").finish()
737    }
738}