#[cfg(not(loom))]
use core::sync::atomic::{AtomicUsize, Ordering};
use core::time::Duration;
#[cfg(not(loom))]
use std::sync::{Condvar, Mutex};
#[cfg(loom)]
use loom::sync::atomic::{AtomicUsize, Ordering};
#[cfg(loom)]
use loom::sync::{Condvar, Mutex};
#[derive(Default)]
pub(super) struct Pulse {
generation: Mutex<u64>,
woken: Condvar,
}
impl Pulse {
pub(super) fn seen(&self) -> u64 {
#[expect(clippy::unwrap_used, reason = "the waiter only panics if the whole process is unwinding")]
let generation = self.generation.lock().unwrap();
*generation
}
pub(super) fn signal(&self) {
#[expect(clippy::unwrap_used, reason = "the waiter only panics if the whole process is unwinding")]
let mut generation = self.generation.lock().unwrap();
*generation = generation.wrapping_add(1);
drop(generation);
self.woken.notify_all();
}
pub(super) fn wait(&self, seen: u64, upto: Duration) {
#[expect(clippy::unwrap_used, reason = "the waiter only panics if the whole process is unwinding")]
let generation = self.generation.lock().unwrap();
if *generation != seen {
return;
}
#[cfg(not(loom))]
#[expect(clippy::unwrap_used, reason = "the waiter only panics if the whole process is unwinding")]
let (_generation, _timed_out) = self.woken.wait_timeout(generation, upto).unwrap();
#[cfg(loom)]
{
let _ = upto;
#[expect(clippy::unwrap_used, reason = "the waiter only panics if the whole process is unwinding")]
let _generation = self.woken.wait(generation).unwrap();
}
}
}
#[derive(Debug)]
pub struct Readers {
live: AtomicUsize,
peak: AtomicUsize,
}
#[cfg(not(loom))]
pub static READERS: Readers = Readers::new();
#[cfg(loom)]
loom::lazy_static! {
pub static ref READERS: Readers = Readers::new();
}
impl Readers {
#[cfg(not(loom))]
pub(super) const fn new() -> Self {
Self {
live: AtomicUsize::new(0),
peak: AtomicUsize::new(0),
}
}
#[cfg(loom)]
pub(super) fn new() -> Self {
Self {
live: AtomicUsize::new(0),
peak: AtomicUsize::new(0),
}
}
pub(super) fn started(&self) {
let live = self.live.fetch_add(1, Ordering::Relaxed).saturating_add(1);
let _raised = self.peak.fetch_max(live, Ordering::Relaxed);
}
pub(super) fn finished(&self) {
let _was = self.live.fetch_sub(1, Ordering::Relaxed);
}
#[must_use]
pub fn live(&self) -> usize {
self.live.load(Ordering::Relaxed)
}
#[must_use]
pub fn peak(&self) -> usize {
self.peak.load(Ordering::Relaxed)
}
}
#[cfg(loom)]
mod loom_models {
use core::time::Duration;
use loom::sync::Arc;
use loom::sync::atomic::AtomicUsize;
use super::{Pulse, Readers};
const CONCURRENT_READERS: usize = 2;
pub(super) fn a_pulse_wakeup_is_never_lost_under_any_interleaving() {
loom::model(|| {
let pulse = Arc::new(Pulse::default());
let seen = pulse.seen();
let announcing_reader = {
let pulse = Arc::clone(&pulse);
loom::thread::spawn(move || {
pulse.signal();
pulse.signal();
})
};
let other_reader = {
let pulse = Arc::clone(&pulse);
loom::thread::spawn(move || pulse.signal())
};
pulse.wait(seen, Duration::from_secs(0));
announcing_reader.join().unwrap();
other_reader.join().unwrap();
assert_eq!(pulse.seen(), seen.wrapping_add(3), "a concurrent signal was lost");
});
}
pub(super) fn a_readers_gauge_returns_to_exactly_zero_under_any_interleaving() {
loom::model(|| {
let readers = Arc::new(Readers {
live: AtomicUsize::new(CONCURRENT_READERS),
peak: AtomicUsize::new(CONCURRENT_READERS),
});
let mut finishers = Vec::with_capacity(CONCURRENT_READERS);
for _reader in 0..CONCURRENT_READERS {
let readers = Arc::clone(&readers);
finishers.push(loom::thread::spawn(move || readers.finished()));
}
assert!(
readers.live() <= CONCURRENT_READERS,
"a contended live read exceeded the number of started readers"
);
for finisher in finishers {
finisher.join().unwrap();
}
assert_eq!(readers.live(), 0, "a decrement was lost or double-counted");
});
}
pub(super) fn a_readers_peak_never_understates_concurrent_starts() {
loom::model(|| {
let readers = Arc::new(Readers::new());
let mut workers = Vec::with_capacity(CONCURRENT_READERS);
for _reader in 0..CONCURRENT_READERS {
let readers = Arc::clone(&readers);
workers.push(loom::thread::spawn(move || readers.started()));
}
assert!(
readers.live() <= CONCURRENT_READERS,
"a contended live read exceeded the number of starters"
);
assert!(
readers.peak() <= CONCURRENT_READERS,
"a contended peak read exceeded the number of starters"
);
for worker in workers {
worker.join().unwrap();
}
assert_eq!(readers.live(), CONCURRENT_READERS, "a concurrent increment was lost");
assert_eq!(readers.peak(), CONCURRENT_READERS, "the peak understated concurrent readers");
});
}
}
#[cfg(loom)]
pub(crate) fn run_loom_models() {
loom_models::a_pulse_wakeup_is_never_lost_under_any_interleaving();
loom_models::a_readers_gauge_returns_to_exactly_zero_under_any_interleaving();
loom_models::a_readers_peak_never_understates_concurrent_starts();
}