fix(threading): make auto-reset regression deterministic
This commit is contained in:
@@ -913,34 +913,50 @@ impl WaitHandleAsyncFactory {
|
||||
mod tests {
|
||||
use super::*;
|
||||
use std::sync::atomic::{AtomicBool, Ordering};
|
||||
use std::sync::{Barrier, mpsc};
|
||||
|
||||
#[test]
|
||||
fn auto_reset_releases_exactly_one_waiter_per_signal() {
|
||||
let event = Arc::new(ManagedAutoResetEvent::new_with_constructor().unwrap());
|
||||
let first = Arc::clone(&event);
|
||||
let second = Arc::clone(&event);
|
||||
let first_done = Arc::new(AtomicBool::new(false));
|
||||
let second_done = Arc::new(AtomicBool::new(false));
|
||||
let first_flag = Arc::clone(&first_done);
|
||||
let second_flag = Arc::clone(&second_done);
|
||||
let a = std::thread::spawn(move || {
|
||||
first.wait_one_with_method().unwrap();
|
||||
first_flag.store(true, Ordering::SeqCst);
|
||||
});
|
||||
let b = std::thread::spawn(move || {
|
||||
second.wait_one_with_method().unwrap();
|
||||
second_flag.store(true, Ordering::SeqCst);
|
||||
});
|
||||
std::thread::sleep(Duration::from_millis(20));
|
||||
event.set().unwrap();
|
||||
std::thread::sleep(Duration::from_millis(20));
|
||||
assert_ne!(
|
||||
first_done.load(Ordering::SeqCst),
|
||||
second_done.load(Ordering::SeqCst)
|
||||
);
|
||||
event.set().unwrap();
|
||||
a.join().unwrap();
|
||||
b.join().unwrap();
|
||||
for _ in 0..32 {
|
||||
let event = Arc::new(ManagedAutoResetEvent::new_with_constructor().unwrap());
|
||||
let ready = Arc::new(Barrier::new(3));
|
||||
let (completed, completions) = mpsc::channel();
|
||||
let mut waiters = Vec::new();
|
||||
for id in 0..2 {
|
||||
let event = Arc::clone(&event);
|
||||
let ready = Arc::clone(&ready);
|
||||
let completed = completed.clone();
|
||||
waiters.push(std::thread::spawn(move || {
|
||||
ready.wait();
|
||||
event.wait_one_with_method().unwrap();
|
||||
completed.send(id).unwrap();
|
||||
}));
|
||||
}
|
||||
drop(completed);
|
||||
|
||||
// Release both contenders together. The signal is deliberately
|
||||
// allowed to race their actual wait calls: an auto-reset event
|
||||
// must retain one signal even when no waiter has blocked yet.
|
||||
ready.wait();
|
||||
event.set().unwrap();
|
||||
let first = completions
|
||||
.recv_timeout(Duration::from_secs(2))
|
||||
.expect("one waiter must consume the first signal");
|
||||
assert_eq!(
|
||||
completions.recv_timeout(Duration::from_millis(10)),
|
||||
Err(mpsc::RecvTimeoutError::Timeout),
|
||||
"one auto-reset signal released both waiters"
|
||||
);
|
||||
|
||||
event.set().unwrap();
|
||||
let second = completions
|
||||
.recv_timeout(Duration::from_secs(2))
|
||||
.expect("the second waiter must consume the second signal");
|
||||
assert_ne!(first, second);
|
||||
for waiter in waiters {
|
||||
waiter.join().unwrap();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
||||
Reference in New Issue
Block a user