use std::sync::OnceLock;
use std::time::Duration;
use parking_lot::{Condvar, Mutex};
use tokio::sync::watch;
pub struct Wake {
generation: Mutex<u64>,
changed: Condvar,
published: watch::Sender<u64>,
}
impl Wake {
pub fn new() -> Self {
Self {
generation: Mutex::new(0),
changed: Condvar::new(),
published: watch::Sender::new(0),
}
}
pub fn bump(&self) {
let mut generation = self.generation.lock();
*generation = generation.wrapping_add(1);
self.changed.notify_all();
self.published.send_replace(*generation);
}
pub fn generation(&self) -> u64 {
*self.generation.lock()
}
pub fn wait(&self, seen: u64) -> u64 {
let mut generation = self.generation.lock();
while *generation == seen {
self.changed.wait(&mut generation);
}
*generation
}
pub fn subscribe(&self) -> watch::Receiver<u64> {
self.published.subscribe()
}
pub fn wait_until(&self, seen: u64, timeout: Duration) -> u64 {
let mut generation = self.generation.lock();
if *generation == seen {
self.changed.wait_for(&mut generation, timeout);
}
*generation
}
}
impl Default for Wake {
fn default() -> Self {
Self::new()
}
}
pub fn engine_changed() -> &'static Wake {
static WAKE: OnceLock<Wake> = OnceLock::new();
WAKE.get_or_init(Wake::new)
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
#[test]
fn a_waiter_sleeps_until_something_moves() {
let wake = Arc::new(Wake::new());
let seen = wake.generation();
let woken = Arc::new(AtomicBool::new(false));
let waiter = {
let (wake, woken) = (Arc::clone(&wake), Arc::clone(&woken));
std::thread::spawn(move || {
wake.wait(seen);
woken.store(true, Ordering::Relaxed);
})
};
std::thread::sleep(Duration::from_millis(50));
assert!(
!woken.load(Ordering::Relaxed),
"woke with nothing to wake for"
);
wake.bump();
waiter.join().unwrap();
}
#[test]
fn a_bump_before_the_wait_is_not_slept_through() {
let wake = Wake::new();
let seen = wake.generation();
wake.bump();
assert_ne!(wake.wait(seen), seen);
}
}