Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,10 @@ All notable changes to this project will be documented in this file.

* Add bounded MPMC `reserve` and `try_reserve` methods returning a borrowed `Permit`, so callers can wait for capacity before constructing a value; sends and reservations receive capacity in wait-queue order, and unused permits release it.

### Bug fixes

* Prevent deadlocks in `singleflight::Group::work` and `try_work` when a duplicate key's destructor calls back into the same group.

### Notable changes

* Executor waker operations, including cloning, waking, and dropping, are expected not to panic; recovery from panicking waker callbacks is no longer supported.
Expand Down
23 changes: 12 additions & 11 deletions asyncband/src/singleflight/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -74,17 +74,18 @@ where
fn get_or_insert(&self, key: K) -> Arc<Entry<K, V>> {
let hash = self.hasher.hash_one(&key);
let mut entries = self.entries.lock();
entries
.entry(hash, |entry| entry.key.eq(&key), |entry| entry.hash)
.or_insert_with(|| {
Arc::new(Entry {
hash,
key,
cell: OnceCell::new(),
})
})
.into_mut()
.clone()
// Drop duplicate keys after unlocking: their destructors may reenter the group.
if let Some(entry) = entries.find(hash, |entry| entry.key.eq(&key)) {
return entry.clone();
}

let entry = Arc::new(Entry {
hash,
key,
cell: OnceCell::new(),
});
entries.insert_unique(hash, entry.clone(), |entry| entry.hash);
entry
}

fn remove<Q>(&self, key: &Q)
Expand Down
89 changes: 89 additions & 0 deletions tests-integration/tests/singleflight_test.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,10 +15,18 @@
// specific language governing permissions and limitations
// under the License.

use std::borrow::Borrow;
use std::cell::Cell;
use std::hash::Hash;
use std::hash::Hasher;
use std::pin::pin;
use std::rc::Rc;
use std::sync::atomic::AtomicUsize;
use std::sync::atomic::Ordering;
use std::task::Poll;

use asyncband::singleflight::Group;
use tests_integration::assert_completes_without_deadlock;
use tests_integration::poll_once;

#[tokio::test]
Expand Down Expand Up @@ -164,3 +172,84 @@ async fn concurrent_try_work_is_coalesced() {
}
assert_eq!(counter.load(Ordering::SeqCst), 1);
}

#[test]
fn duplicate_key_destructor_can_forget_another_key() {
struct Key {
value: usize,
on_drop: Option<Box<dyn FnOnce()>>,
}

impl Borrow<usize> for Key {
fn borrow(&self) -> &usize {
&self.value
}
}

impl PartialEq for Key {
fn eq(&self, other: &Self) -> bool {
self.value == other.value
}
}

impl Eq for Key {}

impl Hash for Key {
fn hash<H: Hasher>(&self, state: &mut H) {
self.value.hash(state);
}
}

impl Drop for Key {
fn drop(&mut self) {
if let Some(on_drop) = self.on_drop.take() {
on_drop();
}
}
}

assert_completes_without_deadlock(|| {
for fallible in [false, true] {
let group = Rc::new(Group::new());
let (release_tx, release_rx) = tokio::sync::oneshot::channel();
let mut leader = pin!(group.work(
Key {
value: 7,
on_drop: None,
},
async || {
release_rx.await.unwrap();
"leader"
}
));
assert!(poll_once(leader.as_mut()).is_pending());

let key_dropped = Rc::new(Cell::new(false));
let dropped = key_dropped.clone();
let weak = Rc::downgrade(&group);
let key = Key {
value: 7,
on_drop: Some(Box::new(move || {
weak.upgrade().unwrap().forget(&42_usize);
dropped.set(true);
})),
};
let mut duplicate = pin!(async {
if fallible {
group
.try_work(key, async || Ok::<_, ()>("duplicate"))
.await
.unwrap()
} else {
group.work(key, async || "duplicate").await
}
});
assert!(poll_once(duplicate.as_mut()).is_pending());
assert!(key_dropped.get());

release_tx.send(()).unwrap();
assert_eq!(poll_once(leader.as_mut()), Poll::Ready("leader"));
assert_eq!(poll_once(duplicate.as_mut()), Poll::Ready("leader"));
}
});
}
Loading