use core::sync::atomic::{AtomicU32, Ordering};
use crate::Futex;
use crate::futex::TimeoutError;
pub struct ResetEvent {
state: AtomicU32,
}
impl Default for ResetEvent {
fn default() -> Self {
Self::new()
}
}
const UNSET: u32 = 0;
const WAITING: u32 = 1;
const IS_SET: u32 = 2;
impl ResetEvent {
pub const fn new() -> Self {
Self {
state: AtomicU32::new(UNSET),
}
}
#[inline]
pub fn is_set(&self) -> bool {
self.state.load(Ordering::Acquire) == IS_SET
}
pub fn wait(&self) {
match self.wait_inner(None) {
Ok(()) => {}
Err(TimeoutError::Timeout) => unreachable!(), }
}
pub fn timed_wait(&self, timeout_ns: u64) -> Result<(), TimeoutError> {
self.wait_inner(Some(timeout_ns))
}
#[inline]
fn wait_inner(&self, timeout: Option<u64>) -> Result<(), TimeoutError> {
if !self.is_set() {
return self.wait_until_set(timeout);
}
Ok(())
}
#[cold]
fn wait_until_set(&self, timeout: Option<u64>) -> Result<(), TimeoutError> {
let mut state = self.state.load(Ordering::Acquire);
if state == UNSET {
state = match self.state.compare_exchange(
state,
WAITING,
Ordering::Acquire,
Ordering::Acquire,
) {
Ok(_) => WAITING,
Err(s) => s,
};
}
if state == WAITING {
let mut futex_deadline = Futex::Deadline::init(timeout);
loop {
let wait_result = futex_deadline.wait(&self.state, WAITING);
state = self.state.load(Ordering::Acquire);
if state != WAITING {
break;
}
wait_result?;
}
}
debug_assert!(state == IS_SET);
Ok(())
}
pub fn set(&self) {
if self.state.load(Ordering::Relaxed) == IS_SET {
return;
}
if self.state.swap(IS_SET, Ordering::Release) == WAITING {
Futex::wake(&self.state, u32::MAX);
}
}
pub fn reset(&self) {
self.state.store(UNSET, Ordering::Relaxed);
}
}