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}