| { | ||
| "git": { | ||
| "sha1": "4abe9d732eb01f7b092a571c3dcc4fbd266f4067" | ||
| "sha1": "d87569164fb61145e79e7ffe0b25783569cc8f93" | ||
| }, | ||
| "path_in_vcs": "tokio" | ||
| } |
+30
-30
@@ -136,5 +136,5 @@ # This file is automatically @generated by Cargo. | ||
| name = "cc" | ||
| version = "1.2.61" | ||
| version = "1.2.62" | ||
| source = "registry+https://github.com/rust-lang/crates.io-index" | ||
| checksum = "d16d90359e986641506914ba71350897565610e87ce0ad9e6f28569db3dd5c6d" | ||
| checksum = "a1dce859f0832a7d088c4f1119888ab94ef4b5d6795d1ce05afb7fe159d79f98" | ||
| dependencies = [ | ||
@@ -452,5 +452,5 @@ "find-msvc-tools", | ||
| name = "js-sys" | ||
| version = "0.3.97" | ||
| version = "0.3.98" | ||
| source = "registry+https://github.com/rust-lang/crates.io-index" | ||
| checksum = "a1840c94c045fbcf8ba2812c95db44499f7c64910a912551aaaa541decebcacf" | ||
| checksum = "67df7112613f8bfd9150013a0314e196f4800d3201ae742489d999db2f979f08" | ||
| dependencies = [ | ||
@@ -701,5 +701,5 @@ "cfg-if", | ||
| name = "pin-project" | ||
| version = "1.1.11" | ||
| version = "1.1.12" | ||
| source = "registry+https://github.com/rust-lang/crates.io-index" | ||
| checksum = "f1749c7ed4bcaf4c3d0a3efc28538844fb29bcdd7d2b67b2be7e20ba861ff517" | ||
| checksum = "cbf0d9e68100b3a7989b4901972f265cd542e560a3a8a724e1e20322f4d06ce9" | ||
| dependencies = [ | ||
@@ -711,5 +711,5 @@ "pin-project-internal", | ||
| name = "pin-project-internal" | ||
| version = "1.1.11" | ||
| version = "1.1.12" | ||
| source = "registry+https://github.com/rust-lang/crates.io-index" | ||
| checksum = "d9b20ed30f105399776b9c883e68e536ef602a16ae6f596d2c473591d6ad64c6" | ||
| checksum = "a990e22f43e84855daf260dded30524ef4a9021cc7541c26540500a50b624389" | ||
| dependencies = [ | ||
@@ -1092,5 +1092,5 @@ "proc-macro2", | ||
| name = "tokio" | ||
| version = "1.52.1" | ||
| version = "1.52.2" | ||
| source = "registry+https://github.com/rust-lang/crates.io-index" | ||
| checksum = "b67dee974fe86fd92cc45b7a95fdd2f99a36a6d7b0d431a231178d3d670bbcc6" | ||
| checksum = "110a78583f19d5cdb2c5ccf321d1290344e71313c6c37d43520d386027d18386" | ||
| dependencies = [ | ||
@@ -1102,3 +1102,3 @@ "pin-project-lite", | ||
| name = "tokio" | ||
| version = "1.52.2" | ||
| version = "1.52.3" | ||
| dependencies = [ | ||
@@ -1155,3 +1155,3 @@ "async-stream", | ||
| "pin-project-lite", | ||
| "tokio 1.52.1", | ||
| "tokio 1.52.2", | ||
| ] | ||
@@ -1166,3 +1166,3 @@ | ||
| "futures-core", | ||
| "tokio 1.52.1", | ||
| "tokio 1.52.2", | ||
| "tokio-stream", | ||
@@ -1182,3 +1182,3 @@ ] | ||
| "pin-project-lite", | ||
| "tokio 1.52.1", | ||
| "tokio 1.52.2", | ||
| ] | ||
@@ -1326,5 +1326,5 @@ | ||
| name = "wasm-bindgen" | ||
| version = "0.2.120" | ||
| version = "0.2.121" | ||
| source = "registry+https://github.com/rust-lang/crates.io-index" | ||
| checksum = "df52b6d9b87e0c74c9edfa1eb2d9bf85e5d63515474513aa50fa181b3c4f5db1" | ||
| checksum = "49ace1d07c165b0864824eee619580c4689389afa9dc9ed3a4c75040d82e6790" | ||
| dependencies = [ | ||
@@ -1340,5 +1340,5 @@ "cfg-if", | ||
| name = "wasm-bindgen-futures" | ||
| version = "0.4.70" | ||
| version = "0.4.71" | ||
| source = "registry+https://github.com/rust-lang/crates.io-index" | ||
| checksum = "af934872acec734c2d80e6617bbb5ff4f12b052dd8e6332b0817bce889516084" | ||
| checksum = "96492d0d3ffba25305a7dc88720d250b1401d7edca02cc3bcd50633b424673b8" | ||
| dependencies = [ | ||
@@ -1351,5 +1351,5 @@ "js-sys", | ||
| name = "wasm-bindgen-macro" | ||
| version = "0.2.120" | ||
| version = "0.2.121" | ||
| source = "registry+https://github.com/rust-lang/crates.io-index" | ||
| checksum = "78b1041f495fb322e64aca85f5756b2172e35cd459376e67f2a6c9dffcedb103" | ||
| checksum = "8e68e6f4afd367a562002c05637acb8578ff2dea1943df76afb9e83d177c8578" | ||
| dependencies = [ | ||
@@ -1362,5 +1362,5 @@ "quote", | ||
| name = "wasm-bindgen-macro-support" | ||
| version = "0.2.120" | ||
| version = "0.2.121" | ||
| source = "registry+https://github.com/rust-lang/crates.io-index" | ||
| checksum = "9dcd0ff20416988a18ac686d4d4d0f6aae9ebf08a389ff5d29012b05af2a1b41" | ||
| checksum = "d95a9ec35c64b2a7cb35d3fead40c4238d0940c86d107136999567a4703259f2" | ||
| dependencies = [ | ||
@@ -1376,5 +1376,5 @@ "bumpalo", | ||
| name = "wasm-bindgen-shared" | ||
| version = "0.2.120" | ||
| version = "0.2.121" | ||
| source = "registry+https://github.com/rust-lang/crates.io-index" | ||
| checksum = "49757b3c82ebf16c57d69365a142940b384176c24df52a087fb748e2085359ea" | ||
| checksum = "c4e0100b01e9f0d03189a92b96772a1fb998639d981193d7dbab487302513441" | ||
| dependencies = [ | ||
@@ -1386,5 +1386,5 @@ "unicode-ident", | ||
| name = "wasm-bindgen-test" | ||
| version = "0.3.70" | ||
| version = "0.3.71" | ||
| source = "registry+https://github.com/rust-lang/crates.io-index" | ||
| checksum = "29826f9d9ecaa314c480d376b276d1c790e6cb6a4681fab8532da69cbabf977d" | ||
| checksum = "af5ec93229ad9ccd0a545a516dec76dc276613f278f6a91aa6b463d5b33d42d0" | ||
| dependencies = [ | ||
@@ -1409,5 +1409,5 @@ "async-trait", | ||
| name = "wasm-bindgen-test-macro" | ||
| version = "0.3.70" | ||
| version = "0.3.71" | ||
| source = "registry+https://github.com/rust-lang/crates.io-index" | ||
| checksum = "c610311887f9e6599a546d278d12d69dfd3a3e92639b2129e4b11ad6cf1961d6" | ||
| checksum = "3c81b9fef827e575e0e54431736d1baa0d700315d8c62cfef1f61fa3aad0cbeb" | ||
| dependencies = [ | ||
@@ -1421,5 +1421,5 @@ "proc-macro2", | ||
| name = "wasm-bindgen-test-shared" | ||
| version = "0.2.120" | ||
| version = "0.2.121" | ||
| source = "registry+https://github.com/rust-lang/crates.io-index" | ||
| checksum = "60238e5b4b1b295701d6f9a66d2a126fe19990348f5fb9dae3b623a370119d94" | ||
| checksum = "4f4d8ae7ad5440360e9799dfd42857d126454a88441ddf72d288ef83fa47f527" | ||
@@ -1426,0 +1426,0 @@ [[package]] |
+1
-1
@@ -16,3 +16,3 @@ # THIS FILE IS AUTOMATICALLY GENERATED BY CARGO | ||
| name = "tokio" | ||
| version = "1.52.2" | ||
| version = "1.52.3" | ||
| authors = ["Tokio Contributors <team@tokio.rs>"] | ||
@@ -19,0 +19,0 @@ build = false |
+1
-1
@@ -63,3 +63,3 @@ *[TokioConf 2026 program and tickets are now available!](https://tokioconf.com)* | ||
| [dependencies] | ||
| tokio = { version = "1.52.2", features = ["full"] } | ||
| tokio = { version = "1.52.3", features = ["full"] } | ||
| ``` | ||
@@ -66,0 +66,0 @@ Then, on your main.rs: |
@@ -223,7 +223,2 @@ use crate::loom::cell::UnsafeCell; | ||
| pub(crate) unsafe fn is_closed(&self) -> bool { | ||
| let ready_bits = self.header.ready_slots.load(Acquire); | ||
| is_tx_closed(ready_bits) | ||
| } | ||
| /// Resets the block to a blank state. This enables reusing blocks in the | ||
@@ -230,0 +225,0 @@ /// channel. |
@@ -439,3 +439,5 @@ use crate::loom::cell::UnsafeCell; | ||
| // If close() was called, an empty queue should report Disconnected. | ||
| TryPopResult::Empty if rx_fields.rx_closed => { | ||
| TryPopResult::Empty | ||
| if rx_fields.rx_closed && self.inner.semaphore.is_idle() => | ||
| { | ||
| return Err(TryRecvError::Disconnected) | ||
@@ -442,0 +444,0 @@ } |
+50
-12
@@ -242,11 +242,2 @@ //! A concurrent, lock-free, FIFO list. | ||
| } | ||
| pub(crate) fn is_closed(&self) -> bool { | ||
| let tail = self.block_tail.load(Acquire); | ||
| unsafe { | ||
| let tail_block = &*tail; | ||
| tail_block.is_closed() | ||
| } | ||
| } | ||
| } | ||
@@ -275,7 +266,54 @@ | ||
| // Guaranteed to return true if `slot_index` is the fake message sent on channel close. | ||
| // Guaranteed to return false if `slot_index` is a fully sent message. | ||
| // | ||
| // For messages that are partially sent, may return either true or false. | ||
| fn is_maybe_closed(&self, tx: &Tx<T>, slot_index: usize) -> bool { | ||
| let start_index = block::start_index(slot_index); | ||
| let tail = tx.block_tail.load(Acquire); | ||
| // SAFETY: Only the receiver frees blocks, so since we are the receiver, this will not be | ||
| // freed right now. | ||
| let tail_ref = unsafe { &*tail }; | ||
| if tail_ref.is_at_index(start_index) { | ||
| return !tail_ref.has_value(slot_index); | ||
| } | ||
| // This method is optimized for checking whether the last value is present, so most of the | ||
| // time it is in `block_tail`. However, this isn't always the case since it's possible | ||
| // that the list was grown with an empty block, in which case `block_tail` points one block | ||
| // too far. To handle this case, we walk the list from the head. | ||
| let mut block_ptr = Some(self.head); | ||
| while let Some(block) = block_ptr { | ||
| // SAFETY: Only the receiver frees blocks, so since we are the receiver, this will not | ||
| // be freed right now. | ||
| let block_ref = unsafe { block.as_ref() }; | ||
| if block_ref.is_at_index(start_index) { | ||
| return !block_ref.has_value(slot_index); | ||
| } | ||
| block_ptr = block_ref.load_next(Acquire); | ||
| } | ||
| true | ||
| } | ||
| pub(crate) fn len(&self, tx: &Tx<T>) -> usize { | ||
| // When all the senders are dropped, there will be a last block in the tail position, | ||
| // but it will be closed | ||
| let tail_position = tx.tail_position.load(Acquire); | ||
| tail_position - self.index - (tx.is_closed() as usize) | ||
| let mut len = tail_position.wrapping_sub(self.index); | ||
| debug_assert!(0 <= len as isize); | ||
| if len == 0 { | ||
| return 0; | ||
| } | ||
| // There are messages present in the queue. However, it's possible that the last message is | ||
| // a fake "closed" message that we do not wish to count. To avoid counting it, we do not | ||
| // count the last message if the ready bit is unset. | ||
| // | ||
| // Note that it is also possible for the ready bit to be unset on a normal message, but | ||
| // this happens only if that message is currently being sent *right now* in parallel on | ||
| // another thread. That is okay because it is optional to count messages that are currently | ||
| // being sent. | ||
| if self.is_maybe_closed(tx, tail_position.wrapping_sub(1)) { | ||
| len -= 1; | ||
| } | ||
| len | ||
| } | ||
@@ -282,0 +320,0 @@ |
@@ -140,8 +140,8 @@ #![cfg_attr(not(feature = "sync"), allow(dead_code, unreachable_pub))] | ||
| #[cfg(all(target_pointer_width = "64", not(loom)))] | ||
| const BLOCK_CAP: usize = 32; | ||
| pub(crate) const BLOCK_CAP: usize = 32; | ||
| #[cfg(all(not(target_pointer_width = "64"), not(loom)))] | ||
| const BLOCK_CAP: usize = 16; | ||
| pub(crate) const BLOCK_CAP: usize = 16; | ||
| #[cfg(loom)] | ||
| const BLOCK_CAP: usize = 2; | ||
| pub(crate) const BLOCK_CAP: usize = 2; |
+11
-1
@@ -268,3 +268,3 @@ use crate::sync::batch_semaphore::{Semaphore, TryAcquireError}; | ||
| /// | ||
| /// Panics if `max_reads` is more than `u32::MAX >> 3`. | ||
| /// Panics if `max_reads` is `0` or is bigger than `u32::MAX >> 3`. | ||
| #[track_caller] | ||
@@ -275,2 +275,3 @@ pub fn with_max_readers(value: T, max_reads: u32) -> RwLock<T> | ||
| { | ||
| assert_ne!(max_reads, 0, "a RwLock may not be created with 0 readers"); | ||
| assert!( | ||
@@ -371,2 +372,6 @@ max_reads <= MAX_READS, | ||
| /// ``` | ||
| /// | ||
| /// # Panics | ||
| /// | ||
| /// Panics if `max_reads` is `0` or is bigger than `u32::MAX >> 3`. | ||
| #[cfg(not(all(loom, test)))] | ||
@@ -377,2 +382,3 @@ pub const fn const_with_max_readers(value: T, max_reads: u32) -> RwLock<T> | ||
| { | ||
| assert!(max_reads != 0, "a RwLock may not be created with 0 readers"); | ||
| assert!(max_reads <= MAX_READS); | ||
@@ -780,2 +786,3 @@ | ||
| let acquire_fut = async { | ||
| debug_assert_ne!(self.mr, 0); | ||
| self.s.acquire(self.mr as usize).await.unwrap_or_else(|_| { | ||
@@ -919,2 +926,3 @@ // The semaphore was closed. but, we never explicitly close it, and we have a | ||
| let acquire_fut = async { | ||
| debug_assert_ne!(self.mr, 0); | ||
| self.s.acquire(self.mr as usize).await.unwrap_or_else(|_| { | ||
@@ -984,2 +992,3 @@ // The semaphore was closed. but, we never explicitly close it, and we have a | ||
| pub fn try_write(&self) -> Result<RwLockWriteGuard<'_, T>, TryLockError> { | ||
| debug_assert_ne!(self.mr, 0); | ||
| match self.s.try_acquire(self.mr as usize) { | ||
@@ -1043,2 +1052,3 @@ Ok(permit) => permit, | ||
| pub fn try_write_owned(self: Arc<Self>) -> Result<OwnedRwLockWriteGuard<T>, TryLockError> { | ||
| debug_assert_ne!(self.mr, 0); | ||
| match self.s.try_acquire(self.mr as usize) { | ||
@@ -1045,0 +1055,0 @@ Ok(permit) => permit, |
@@ -1,2 +0,2 @@ | ||
| use crate::sync::mpsc; | ||
| use crate::sync::mpsc::{self, BLOCK_CAP}; | ||
@@ -225,1 +225,54 @@ use loom::future::block_on; | ||
| } | ||
| #[test] | ||
| fn is_empty_during_close() { | ||
| loom::model(|| { | ||
| let (tx, rx) = mpsc::channel::<()>(1); | ||
| let th1 = thread::spawn(move || { | ||
| assert!(rx.is_empty()); | ||
| }); | ||
| drop(tx); | ||
| th1.join().unwrap(); | ||
| }); | ||
| } | ||
| fn len_during_close_helper(n: usize) { | ||
| loom::model(move || { | ||
| let (tx, rx) = mpsc::channel::<()>(n + 1); | ||
| for _ in 0..n { | ||
| tx.try_send(()).unwrap(); | ||
| } | ||
| let th1 = thread::spawn(move || { | ||
| assert_eq!(rx.len(), n); | ||
| }); | ||
| drop(tx); | ||
| th1.join().unwrap(); | ||
| }); | ||
| } | ||
| #[test] | ||
| fn len_during_close_0() { | ||
| len_during_close_helper(0); | ||
| } | ||
| #[test] | ||
| fn len_during_close_1() { | ||
| len_during_close_helper(1); | ||
| } | ||
| #[test] | ||
| fn len_during_close_block_cap() { | ||
| len_during_close_helper(BLOCK_CAP); | ||
| } | ||
| #[test] | ||
| fn len_during_close_block_cap_plus_1() { | ||
| len_during_close_helper(BLOCK_CAP + 1); | ||
| } |
+98
-0
@@ -791,2 +791,87 @@ #![allow(clippy::redundant_clone)] | ||
| #[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] | ||
@@ -1004,2 +1089,15 @@ async fn dropping_rx_closes_channel() { | ||
| #[test] | ||
| fn try_recv_after_receiver_close_with_permit() { | ||
| let (tx, mut rx) = mpsc::channel::<()>(5); | ||
| let permit = tx.try_reserve().unwrap(); | ||
| assert_eq!(Err(TryRecvError::Empty), rx.try_recv()); | ||
| rx.close(); | ||
| assert_eq!(Err(TryRecvError::Empty), rx.try_recv()); | ||
| drop(permit); | ||
| assert_eq!(Err(TryRecvError::Disconnected), rx.try_recv()); | ||
| } | ||
| #[test] | ||
| fn try_recv_close_while_empty_bounded() { | ||
@@ -1006,0 +1104,0 @@ let (tx, mut rx) = mpsc::channel::<()>(5); |
+12
-0
@@ -81,2 +81,14 @@ #![warn(rust_2018_idioms)] | ||
| #[test] | ||
| #[should_panic(expected = "a RwLock may not be created with 0 readers")] | ||
| fn zero_max_readers() { | ||
| RwLock::with_max_readers(100, 0); | ||
| } | ||
| #[test] | ||
| #[should_panic(expected = "a RwLock may not be created with 0 readers")] | ||
| fn zero_max_readers_const() { | ||
| RwLock::const_with_max_readers(100, 0); | ||
| } | ||
| // When there is an active exclusive owner, subsequent exclusive access should not be possible | ||
@@ -83,0 +95,0 @@ #[test] |
Sorry, the diff of this file is not supported yet
Sorry, the diff of this file is too big to display
Sorry, the diff of this file is too big to display