sync: notify receivers in mpsc OwnedPermit::release() method (#8075)

This commit is contained in:
Alice Ryhl
2026-05-07 09:32:14 +02:00
committed by GitHub
parent 30d25ccb8b
commit f085b6211b
2 changed files with 88 additions and 16 deletions
+3 -16
View File
@@ -1844,14 +1844,12 @@ impl<T> OwnedPermit<T> {
///
/// [`Sender`]: Sender
pub fn release(mut self) -> Sender<T> {
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<T> OwnedPermit<T> {
impl<T> Drop for OwnedPermit<T> {
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.
+85
View File
@@ -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);