use std::time::Duration;
use crate::sync::{AtomicU32, Ordering};
#[cfg(all(not(loom), any(target_os = "linux", target_os = "android")))]
#[path = "linux.rs"]
mod imp;
#[cfg(all(not(loom), target_os = "macos"))]
#[path = "macos.rs"]
mod imp;
#[derive(Debug)]
pub(crate) struct Futex {
word: AtomicU32,
#[cfg(loom)]
lock: crate::sync::Mutex<()>,
#[cfg(loom)]
cv: crate::sync::Condvar,
}
impl Futex {
pub(crate) fn new() -> Self {
Self {
word: AtomicU32::new(0),
#[cfg(loom)]
lock: crate::sync::Mutex::new(()),
#[cfg(loom)]
cv: crate::sync::Condvar::new(),
}
}
#[inline]
pub(crate) fn load(&self, order: Ordering) -> u32 {
self.word.load(order)
}
#[inline]
pub(crate) fn bump(&self) {
self.word.fetch_add(1, Ordering::Release);
}
pub(crate) fn wait(&self, expected: u32, timeout: Option<Duration>) {
#[cfg(not(loom))]
imp::wait(&self.word, expected, timeout);
#[cfg(loom)]
{
let _ = timeout;
let guard = self
.lock
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if self.word.load(Ordering::Relaxed) == expected {
drop(
self.cv
.wait(guard)
.unwrap_or_else(std::sync::PoisonError::into_inner),
);
}
}
}
pub(crate) fn wake_one(&self) {
#[cfg(not(loom))]
imp::wake(&self.word, false);
#[cfg(loom)]
{
let _guard = self
.lock
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
self.cv.notify_one();
}
}
pub(crate) fn wake_all(&self) {
#[cfg(not(loom))]
imp::wake(&self.word, true);
#[cfg(loom)]
{
let _guard = self
.lock
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
self.cv.notify_all();
}
}
}
#[cfg(all(test, not(loom)))]
mod tests {
use super::*;
use crate::sync::Arc;
use std::thread;
use std::time::Instant;
#[test]
fn mismatched_value_returns_immediately() {
let f = Futex::new();
let start = Instant::now();
f.wait(1, None);
assert!(start.elapsed() < Duration::from_secs(1));
}
#[test]
fn timeout_is_respected() {
let f = Futex::new();
let start = Instant::now();
f.wait(0, Some(Duration::from_millis(20)));
let elapsed = start.elapsed();
assert!(elapsed < Duration::from_secs(5), "{elapsed:?}");
}
#[test]
fn wake_without_waiters_is_harmless() {
let f = Futex::new();
f.wake_one();
f.wake_all();
}
fn wait_for_change(f: &Futex) {
while f.load(Ordering::Acquire) == 0 {
f.wait(0, None);
}
}
#[test]
fn bump_and_wake_one_releases_a_sleeper() {
let f = Arc::new(Futex::new());
let sleeper = {
let f = Arc::clone(&f);
thread::spawn(move || wait_for_change(&f))
};
thread::sleep(Duration::from_millis(20));
f.bump();
f.wake_one();
sleeper.join().unwrap();
}
#[test]
fn wake_all_releases_every_sleeper() {
let f = Arc::new(Futex::new());
let sleepers: Vec<_> = (0..4)
.map(|_| {
let f = Arc::clone(&f);
thread::spawn(move || wait_for_change(&f))
})
.collect();
thread::sleep(Duration::from_millis(20));
f.bump();
f.wake_all();
for s in sleepers {
s.join().unwrap();
}
}
}