#![cfg(loom)]
use loom::sync::Arc;
use loom::thread;
use parkring::{BlockingQueue, LockFreeQueue, PopError, TryPopError, TryPushError};
fn model(f: impl Fn() + Sync + Send + 'static) {
model_with_bound(3, f);
}
fn model_with_bound(max_preemptions: usize, f: impl Fn() + Sync + Send + 'static) {
let mut builder = loom::model::Builder::new();
builder.max_branches = std::env::var("LOOM_MAX_BRANCHES")
.ok()
.and_then(|v| v.parse().ok())
.unwrap_or(20_000);
let from_env = std::env::var("LOOM_MAX_PREEMPTIONS")
.ok()
.and_then(|v| v.parse().ok());
builder.preemption_bound =
Some(from_env.map_or(max_preemptions, |n: usize| n.min(max_preemptions)));
builder.check(f);
}
#[test]
fn consumer_is_not_lost_when_parking_races_a_push() {
model(|| {
let q = Arc::new(LockFreeQueue::new(2));
let consumer = {
let q = q.clone();
thread::spawn(move || q.pop())
};
q.push(7).unwrap();
assert_eq!(consumer.join().unwrap(), Ok(7));
});
}
#[test]
fn producer_is_not_lost_when_parking_races_a_pop() {
model(|| {
let q = Arc::new(LockFreeQueue::new(2));
q.push(0).unwrap();
q.push(1).unwrap();
let producer = {
let q = q.clone();
thread::spawn(move || q.push(2))
};
assert_eq!(q.pop(), Ok(0));
producer.join().unwrap().unwrap();
assert_eq!(q.try_pop(), Ok(1));
assert_eq!(q.try_pop(), Ok(2));
});
}
#[test]
fn close_wakes_a_parked_consumer() {
model(|| {
let q = Arc::new(LockFreeQueue::<u32>::new(2));
let consumer = {
let q = q.clone();
thread::spawn(move || q.pop())
};
q.close();
assert_eq!(consumer.join().unwrap(), Err(PopError));
});
}
#[test]
fn close_wakes_a_parked_producer() {
model(|| {
let q = Arc::new(LockFreeQueue::new(2));
q.push(0).unwrap();
q.push(1).unwrap();
let producer = {
let q = q.clone();
thread::spawn(move || q.push(2))
};
q.close();
assert_eq!(producer.join().unwrap().unwrap_err().into_inner(), 2);
});
}
#[test]
fn items_pushed_before_close_are_drained() {
model(|| {
let q = Arc::new(LockFreeQueue::new(2));
let consumer = {
let q = q.clone();
thread::spawn(move || (q.pop(), q.pop()))
};
q.push(1).unwrap();
q.close();
assert_eq!(consumer.join().unwrap(), (Ok(1), Err(PopError)));
});
}
#[test]
fn close_and_push_are_linearizable() {
model(|| {
let q = Arc::new(LockFreeQueue::new(2));
let producer = {
let q = q.clone();
thread::spawn(move || q.try_push(1))
};
q.close();
let popped = q.try_pop();
let pushed = producer.join().unwrap();
match pushed {
Ok(()) => {
assert_ne!(popped, Err(TryPopError::Closed), "item stranded");
let later = q.try_pop();
assert!(
popped == Ok(1) || later == Ok(1),
"successful push was never delivered"
);
}
Err(e) => {
assert_eq!(e, TryPushError::Closed(1));
assert_eq!(popped, Err(TryPopError::Closed));
}
}
});
}
#[test]
fn notify_one_does_not_strand_a_second_waiter() {
model_with_bound(1, || {
let q = Arc::new(LockFreeQueue::new(2));
let consumers: Vec<_> = (0..2)
.map(|_| {
let q = q.clone();
thread::spawn(move || q.pop().unwrap())
})
.collect();
q.push(1).unwrap();
q.push(2).unwrap();
let mut got: Vec<_> = consumers.into_iter().map(|c| c.join().unwrap()).collect();
got.sort_unstable();
assert_eq!(got, vec![1, 2]);
});
}
#[test]
fn concurrent_try_push_never_fails_spuriously() {
model(|| {
let q = Arc::new(LockFreeQueue::new(2));
let other = {
let q = q.clone();
thread::spawn(move || q.try_push(1))
};
assert_eq!(q.try_push(2), Ok(()));
assert_eq!(other.join().unwrap(), Ok(()));
assert_eq!(q.len(), 2);
});
}
#[test]
fn concurrent_try_pop_never_fails_spuriously() {
model(|| {
let q = Arc::new(LockFreeQueue::new(2));
q.push(1).unwrap();
q.push(2).unwrap();
let other = {
let q = q.clone();
thread::spawn(move || q.try_pop())
};
let mine = q.try_pop().unwrap();
let theirs = other.join().unwrap().unwrap();
assert_ne!(mine, theirs);
});
}
#[test]
fn fifo_across_laps() {
model(|| {
let q = Arc::new(LockFreeQueue::new(2));
let producer = {
let q = q.clone();
thread::spawn(move || {
for i in 0..4 {
q.push(i).unwrap();
}
})
};
for i in 0..4 {
assert_eq!(q.pop(), Ok(i));
}
producer.join().unwrap();
});
}
#[test]
fn drop_on_another_thread_sees_published_items() {
model(|| {
let q = Arc::new(LockFreeQueue::new(2));
let item = std::sync::Arc::new(());
let producer = {
let (q, item) = (q.clone(), item.clone());
thread::spawn(move || q.push(item).unwrap())
};
drop(q);
producer.join().unwrap();
assert_eq!(std::sync::Arc::strong_count(&item), 1);
});
}
#[test]
fn close_after_one_of_two_parked_consumers_is_served() {
model_with_bound(1, || {
let q = Arc::new(LockFreeQueue::new(2));
let consumers: Vec<_> = (0..2)
.map(|_| {
let q = q.clone();
thread::spawn(move || q.pop())
})
.collect();
q.push(1).unwrap();
q.close();
let got: Vec<_> = consumers.into_iter().map(|c| c.join().unwrap()).collect();
assert!(
got == [Ok(1), Err(PopError)] || got == [Err(PopError), Ok(1)],
"{got:?}"
);
});
}
#[test]
fn timed_pop_racing_a_push() {
model(|| {
let q = Arc::new(LockFreeQueue::new(2));
let consumer = {
let q = q.clone();
thread::spawn(move || q.pop_timeout(std::time::Duration::from_secs(3600)))
};
q.push(9).unwrap();
assert_eq!(consumer.join().unwrap(), Ok(9));
});
}
#[test]
fn blocking_queue_close_wakes_consumer_and_drains() {
model(|| {
let q = Arc::new(BlockingQueue::new(1));
let consumer = {
let q = q.clone();
thread::spawn(move || (q.pop(), q.pop()))
};
q.push(1).unwrap();
q.close();
assert_eq!(consumer.join().unwrap(), (Ok(1), Err(PopError)));
});
}