From f085b6211b8ebb6aba21f1f1f91e7b8b243aa815 Mon Sep 17 00:00:00 2001 From: Alice Ryhl Date: Thu, 7 May 2026 09:32:14 +0200 Subject: [PATCH] sync: notify receivers in mpsc `OwnedPermit::release()` method (#8075) --- tokio/src/sync/mpsc/bounded.rs | 19 ++------ tokio/tests/sync_mpsc.rs | 85 ++++++++++++++++++++++++++++++++++ 2 files changed, 88 insertions(+), 16 deletions(-) diff --git a/tokio/src/sync/mpsc/bounded.rs b/tokio/src/sync/mpsc/bounded.rs index 06eeffc3f..f7b0081c6 100644 --- a/tokio/src/sync/mpsc/bounded.rs +++ b/tokio/src/sync/mpsc/bounded.rs @@ -1844,14 +1844,12 @@ impl OwnedPermit { /// /// [`Sender`]: Sender pub fn release(mut self) -> Sender { - use chan::Semaphore; - let chan = self.chan.take().unwrap_or_else(|| { unreachable!("OwnedPermit channel is only taken when the permit is moved") }); // Add the permit back to the semaphore - chan.semaphore().add_permit(); + drop(Permit { chan: &chan }); Sender { chan } } @@ -1910,21 +1908,10 @@ impl OwnedPermit { impl Drop for OwnedPermit { fn drop(&mut self) { - use chan::Semaphore; - // Are we still holding onto the sender? if let Some(chan) = self.chan.take() { - let semaphore = chan.semaphore(); - - // Add the permit back to the semaphore - semaphore.add_permit(); - - // If this `OwnedPermit` is holding the last sender for this - // channel, wake the receiver so that it can be notified that the - // channel is closed. - if semaphore.is_closed() && semaphore.is_idle() { - chan.wake_rx(); - } + // Reuse Drop impl of non-owned Permit. + drop(Permit { chan: &chan }); } // Otherwise, do nothing. diff --git a/tokio/tests/sync_mpsc.rs b/tokio/tests/sync_mpsc.rs index 27157288e..93804581b 100644 --- a/tokio/tests/sync_mpsc.rs +++ b/tokio/tests/sync_mpsc.rs @@ -788,6 +788,91 @@ async fn drop_permit_iterator_releases_permits() { } } +#[test] +fn dropping_last_permit_wakes_closed_receiver() { + let (tx, mut rx) = mpsc::channel::<()>(100); + + let permit = tx.try_reserve().unwrap(); + rx.close(); + + let mut recv = tokio_test::task::spawn(rx.recv()); + assert_pending!(recv.poll()); + drop(permit); + assert!(recv.is_woken()); + assert_ready!(recv.poll()); +} + +#[test] +fn dropping_last_owned_permit_wakes_closed_receiver() { + let (tx, mut rx) = mpsc::channel::<()>(100); + + let permit = tx.try_reserve_owned().unwrap(); + rx.close(); + + let mut recv = tokio_test::task::spawn(rx.recv()); + assert_pending!(recv.poll()); + drop(permit); + assert!(recv.is_woken()); + assert_ready!(recv.poll()); +} + +#[test] +fn dropping_last_permit_iterator_wakes_closed_receiver() { + let (tx, mut rx) = mpsc::channel::<()>(100); + + let permits = tx.try_reserve_many(1).unwrap(); + rx.close(); + + let mut recv = tokio_test::task::spawn(rx.recv()); + assert_pending!(recv.poll()); + drop(permits); + assert!(recv.is_woken()); + assert_ready!(recv.poll()); +} + +#[test] +fn sending_last_permit_wakes_closed_receiver() { + let (tx, mut rx) = mpsc::channel::<()>(100); + + let permit = tx.try_reserve().unwrap(); + rx.close(); + + let mut recv = tokio_test::task::spawn(rx.recv()); + assert_pending!(recv.poll()); + permit.send(()); + assert!(recv.is_woken()); + assert_ready!(recv.poll()); +} + +#[test] +fn sending_last_owned_permit_wakes_closed_receiver() { + let (tx, mut rx) = mpsc::channel::<()>(100); + + let permit = tx.try_reserve_owned().unwrap(); + rx.close(); + + let mut recv = tokio_test::task::spawn(rx.recv()); + assert_pending!(recv.poll()); + permit.send(()); + assert!(recv.is_woken()); + assert_ready!(recv.poll()); +} + +#[test] +fn releasing_last_owned_permit_wakes_closed_receiver() { + let (tx, mut rx) = mpsc::channel::<()>(100); + + let permit = tx.try_reserve_owned().unwrap(); + rx.close(); + + let mut recv = tokio_test::task::spawn(rx.recv()); + assert_pending!(recv.poll()); + let inert_sender = permit.release(); + assert!(recv.is_woken()); + assert_ready!(recv.poll()); + drop(inert_sender); +} + #[maybe_tokio_test] async fn dropping_rx_closes_channel() { let (tx, rx) = mpsc::channel(100);