From 7d5b12c50947929326bdfaadb78155ee6593f209 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?O=C4=9Fuz=20Bilgener?= Date: Wed, 20 Jan 2021 17:12:51 -0500 Subject: [PATCH] sync: fix panic in broadcast::Receiver drop (#3434) --- tokio/src/sync/broadcast.rs | 2 +- tokio/src/sync/tests/loom_broadcast.rs | 27 ++++++++++++++++++++++++++ 2 files changed, 28 insertions(+), 1 deletion(-) diff --git a/tokio/src/sync/broadcast.rs b/tokio/src/sync/broadcast.rs index 2cf2b1c32..58ea481ca 100644 --- a/tokio/src/sync/broadcast.rs +++ b/tokio/src/sync/broadcast.rs @@ -929,7 +929,7 @@ impl Drop for Receiver { drop(tail); - while self.next != until { + while self.next < until { match self.recv_ref(None) { Ok(_) => {} // The channel is closed diff --git a/tokio/src/sync/tests/loom_broadcast.rs b/tokio/src/sync/tests/loom_broadcast.rs index 4b1f034f4..039b01bf4 100644 --- a/tokio/src/sync/tests/loom_broadcast.rs +++ b/tokio/src/sync/tests/loom_broadcast.rs @@ -178,3 +178,30 @@ fn drop_rx() { assert_ok!(th2.join()); }); } + +#[test] +fn drop_multiple_rx_with_overflow() { + loom::model(move || { + // It is essential to have multiple senders and receivers in this test case. + let (tx, mut rx) = broadcast::channel(1); + let _rx2 = tx.subscribe(); + + let _ = tx.send(()); + let tx2 = tx.clone(); + let th1 = thread::spawn(move || { + block_on(async { + for _ in 0..100 { + let _ = tx2.send(()); + } + }); + }); + let _ = tx.send(()); + + let th2 = thread::spawn(move || { + block_on(async { while let Ok(_) = rx.recv().await {} }); + }); + + assert_ok!(th1.join()); + assert_ok!(th2.join()); + }); +}