From d80af4261a2c300095826509540919cc970e760e Mon Sep 17 00:00:00 2001 From: 0x676e67 Date: Tue, 6 Oct 2026 07:39:06 +0800 Subject: [PATCH] perf(body): skip waking a closed peer when a channel end drops (cherry picked from commit 663b094fe80389c393880d8d5168c73131629336) --- src/body/chan.rs | 55 ++++++++++++++++++++++++++++++++++++++++++++++-- 1 file changed, 53 insertions(+), 2 deletions(-) diff --git a/src/body/chan.rs b/src/body/chan.rs index 314a7d38e2..6e5f9f6f0d 100644 --- a/src/body/chan.rs +++ b/src/body/chan.rs @@ -207,8 +207,11 @@ impl Drop for Sender { fn drop(&mut self) { let mut state = self.shared.state.lock().panic_if_poisoned(); state.sender_open = false; + let receiver_open = state.receiver_open; drop(state); - self.shared.receiver_waker.wake(); + if receiver_open { + self.shared.receiver_waker.wake(); + } } } @@ -216,8 +219,11 @@ impl Drop for Receiver { fn drop(&mut self) { let mut state = self.shared.state.lock().panic_if_poisoned(); state.receiver_open = false; + let sender_open = state.sender_open; drop(state); - self.shared.sender_waker.wake(); + if sender_open { + self.shared.sender_waker.wake(); + } } } @@ -233,6 +239,11 @@ impl fmt::Debug for Sender { #[cfg(test)] mod tests { + use std::{ + sync::atomic::{AtomicUsize, Ordering}, + task::{Wake, Waker}, + }; + use super::*; async fn recv(rx: &mut Receiver) -> Option> { @@ -262,4 +273,44 @@ mod tests { assert!(recv(&mut rx).await.is_none()); assert_eq!(rx.take_trailers().unwrap()["x-trailer"], "value"); } + + #[test] + fn drop_wakes_only_an_open_peer() { + struct Count(AtomicUsize); + + impl Wake for Count { + fn wake(self: Arc) { + self.0.fetch_add(1, Ordering::Relaxed); + } + } + + let count = Arc::new(Count(AtomicUsize::new(0))); + let waker = Waker::from(count.clone()); + let mut cx = Context::from_waker(&waker); + let wakes = || count.0.load(Ordering::Relaxed); + + // A receiver dropped after the sender does not wake the sender's stale waker. + let (mut tx, rx) = channel(true); + assert!(tx.poll_ready(&mut cx).is_pending()); + drop(tx); + drop(rx); + assert_eq!(wakes(), 0); + + // Nor does a sender dropped after the receiver. + let (tx, mut rx) = channel(false); + assert!(rx.poll_next(&mut cx).is_pending()); + drop(rx); + drop(tx); + assert_eq!(wakes(), 0); + + // An open peer is still woken to observe the close. + let (mut tx, rx) = channel(true); + assert!(tx.poll_ready(&mut cx).is_pending()); + drop(rx); + assert_eq!(wakes(), 1); + let (tx, mut rx) = channel(false); + assert!(rx.poll_next(&mut cx).is_pending()); + drop(tx); + assert_eq!(wakes(), 2); + } }