#![cfg(loom)]
use loom::sync::atomic::{AtomicUsize, Ordering};
use loom::sync::{Arc, Mutex};
use loom::thread;
const CAPACITY: usize = 2;
const WATERMARK: usize = 0;
const NOTHING: usize = usize::MAX;
struct RingModel {
cursor: AtomicUsize,
trackers: Mutex<Vec<Arc<AtomicUsize>>>,
}
impl RingModel {
fn new() -> Self {
RingModel {
cursor: AtomicUsize::new(NOTHING),
trackers: Mutex::new(Vec::new()),
}
}
fn register(&self) -> Arc<AtomicUsize> {
let mut trackers = self.trackers.lock().unwrap();
let head = self.cursor.load(Ordering::Acquire);
let start = if head == NOTHING { 0 } else { head + 1 };
let tracker = Arc::new(AtomicUsize::new(start));
trackers.push(tracker.clone());
tracker
}
fn slowest(&self) -> Option<usize> {
let trackers = self.trackers.lock().unwrap();
trackers.iter().map(|t| t.load(Ordering::Acquire)).min()
}
fn would_clobber(&self, seq: usize) -> bool {
let trackers = self.trackers.lock().unwrap();
trackers
.iter()
.map(|t| t.load(Ordering::Acquire))
.any(|t| seq >= t + CAPACITY)
}
}
struct PublisherModel {
ring: Arc<RingModel>,
seq: usize,
cached_slowest: usize,
}
impl PublisherModel {
fn new(ring: Arc<RingModel>) -> Self {
PublisherModel {
ring,
seq: 0,
cached_slowest: 0,
}
}
fn has_room(&mut self) -> bool {
let effective = CAPACITY - WATERMARK;
if self.seq >= self.cached_slowest + effective {
if let Some(slowest) = self.ring.slowest() {
self.cached_slowest = slowest;
if self.seq >= slowest + effective {
return false;
}
}
}
true
}
fn try_publish(&mut self) -> bool {
if !self.has_room() {
return false;
}
assert!(
!self.ring.would_clobber(self.seq),
"published seq {} over a slot a registered subscriber had not read \
(cached_slowest = {})",
self.seq,
self.cached_slowest
);
self.ring.cursor.store(self.seq, Ordering::Release);
self.seq += 1;
true
}
}
#[test]
fn subscriber_joining_a_live_ring_is_never_lapped() {
loom::model(|| {
let ring = Arc::new(RingModel::new());
let producer_ring = ring.clone();
let producer = thread::spawn(move || {
let mut p = PublisherModel::new(producer_ring);
for _ in 0..(CAPACITY + 1) {
if !p.try_publish() {
break;
}
}
});
let joiner_ring = ring.clone();
let joiner = thread::spawn(move || {
let tracker = joiner_ring.register();
let mine = tracker.load(Ordering::Relaxed);
let head = joiner_ring.cursor.load(Ordering::Acquire);
if head != NOTHING && head >= mine {
tracker.store(mine + 1, Ordering::Release);
}
});
producer.join().unwrap();
joiner.join().unwrap();
});
}
#[test]
fn an_idle_subscriber_holds_the_publisher() {
loom::model(|| {
let ring = Arc::new(RingModel::new());
let _idle = ring.register();
let active = ring.register();
let producer_ring = ring.clone();
let producer = thread::spawn(move || {
let mut p = PublisherModel::new(producer_ring);
for _ in 0..(CAPACITY + 1) {
if !p.try_publish() {
break;
}
}
});
let consumer_ring = ring.clone();
let consumer = thread::spawn(move || {
let mine = active.load(Ordering::Relaxed);
let head = consumer_ring.cursor.load(Ordering::Acquire);
if head != NOTHING && head >= mine {
active.store(mine + 1, Ordering::Release);
}
});
producer.join().unwrap();
consumer.join().unwrap();
});
}