use std::collections::VecDeque;
#[doc(hidden)]
pub struct ArrayQueue<T> {
capacity: usize,
inner: shuttle::sync::Mutex<VecDeque<T>>,
}
impl<T> ArrayQueue<T> {
#[doc(hidden)]
pub fn new(capacity: usize) -> Self {
assert!(capacity > 0, "capacity must be non-zero");
Self {
capacity,
inner: shuttle::sync::Mutex::new(VecDeque::with_capacity(capacity)),
}
}
#[doc(hidden)]
pub fn force_push(&self, value: T) -> Option<T> {
let mut queue = self.inner.lock().unwrap();
let evicted = if queue.len() >= self.capacity {
queue.pop_front()
} else {
None
};
queue.push_back(value);
evicted
}
#[doc(hidden)]
pub fn pop(&self) -> Option<T> {
self.inner.lock().unwrap().pop_front()
}
#[doc(hidden)]
pub fn len(&self) -> usize {
self.inner.lock().unwrap().len()
}
#[doc(hidden)]
pub fn is_empty(&self) -> bool {
self.inner.lock().unwrap().is_empty()
}
#[doc(hidden)]
pub fn capacity(&self) -> usize {
self.capacity
}
}
#[doc(hidden)]
pub fn deadline_reached(now: std::time::Instant, deadline: std::time::Instant) -> bool {
use shuttle::rand::Rng;
now >= deadline || shuttle::rand::thread_rng().gen_bool(0.05)
}
#[doc(hidden)]
pub fn threshold_elapsed(since: std::time::Instant, threshold: std::time::Duration) -> bool {
deadline_reached(std::time::Instant::now(), since + threshold)
}
#[doc(hidden)]
pub struct Parker {
inner: std::sync::Arc<TokenState>,
}
struct TokenState {
available: shuttle::sync::Mutex<bool>,
condvar: shuttle::sync::Condvar,
}
impl Default for Parker {
fn default() -> Self {
Self {
inner: std::sync::Arc::new(TokenState {
available: shuttle::sync::Mutex::new(false),
condvar: shuttle::sync::Condvar::new(),
}),
}
}
}
impl Parker {
#[doc(hidden)]
pub fn unparker(&self) -> Unparker {
Unparker {
inner: self.inner.clone(),
}
}
#[doc(hidden)]
pub fn park_deadline(&self, _deadline: std::time::Instant) {
let mut available = self.inner.available.lock().unwrap();
while !*available {
available = self.inner.condvar.wait(available).unwrap();
}
*available = false;
}
}
#[derive(Clone)]
#[doc(hidden)]
pub struct Unparker {
inner: std::sync::Arc<TokenState>,
}
impl Unparker {
#[doc(hidden)]
pub fn unpark(&self) {
let mut available = self.inner.available.lock().unwrap();
*available = true;
self.inner.condvar.notify_one();
}
}
#[doc(hidden)]
pub struct OnceSlot<T>(shuttle::sync::Mutex<Option<T>>);
impl<T: Clone> OnceSlot<T> {
#[doc(hidden)]
pub fn new() -> Self {
Self(shuttle::sync::Mutex::new(None))
}
#[doc(hidden)]
pub fn with<R>(&self, f: impl FnOnce(Option<&T>) -> R) -> R {
let value = self.0.lock().unwrap().clone();
f(value.as_ref())
}
#[doc(hidden)]
pub fn set(&self, value: T) -> Result<(), T> {
let mut guard = self.0.lock().unwrap();
if guard.is_some() {
Err(value)
} else {
*guard = Some(value);
Ok(())
}
}
}
impl<T: Clone> Default for OnceSlot<T> {
fn default() -> Self {
Self::new()
}
}
struct GuardState<T> {
strong: usize,
weak: usize,
value: Option<T>,
}
#[doc(hidden)]
pub struct GuardArc<T>(std::sync::Arc<shuttle::sync::Mutex<GuardState<T>>>);
#[doc(hidden)]
pub struct GuardWeak<T>(std::sync::Arc<shuttle::sync::Mutex<GuardState<T>>>);
impl<T> GuardArc<T> {
#[doc(hidden)]
pub fn new(value: T) -> Self {
Self(std::sync::Arc::new(shuttle::sync::Mutex::new(GuardState {
strong: 1,
weak: 0,
value: Some(value),
})))
}
#[doc(hidden)]
pub fn downgrade(this: &Self) -> GuardWeak<T> {
this.0.lock().unwrap().weak += 1;
GuardWeak(this.0.clone())
}
#[doc(hidden)]
pub fn is_present(&self) -> bool {
self.0.lock().unwrap().value.is_some()
}
#[doc(hidden)]
pub fn take(&self) -> Option<T> {
self.0.lock().unwrap().value.take()
}
}
impl<T> Clone for GuardArc<T> {
fn clone(&self) -> Self {
self.0.lock().unwrap().strong += 1;
Self(self.0.clone())
}
}
impl<T> Drop for GuardArc<T> {
fn drop(&mut self) {
let mut state = self.0.lock().unwrap();
state.strong -= 1;
if state.strong == 0 {
state.value = None;
}
}
}
impl<T> GuardWeak<T> {
#[doc(hidden)]
pub fn upgrade(&self) -> Option<GuardArc<T>> {
let mut state = self.0.lock().unwrap();
if state.strong == 0 {
None
} else {
state.strong += 1;
Some(GuardArc(self.0.clone()))
}
}
}
impl<T> Drop for GuardWeak<T> {
fn drop(&mut self) {
self.0.lock().unwrap().weak -= 1;
}
}
#[doc(hidden)]
pub use shuttle::sync::mpsc::Sender;
#[doc(hidden)]
pub use shuttle::thread;
#[macro_export]
macro_rules! shuttle_test {
(
num_iters = $num_iters:expr, depth = $depth:expr $(, should_panic = $msg:literal)?;
$(#[$meta:meta])*
fn $name:ident() $body:block
) => {
mod $name {
#[allow(unused_imports)]
use super::*;
$(#[$meta])*
fn $name() $body
#[test]
$(#[should_panic(expected = $msg)])?
fn pct() {
::shuttle::check_pct($name, $num_iters, $depth);
}
#[test]
$(#[should_panic(expected = $msg)])?
fn determinism() {
::shuttle::check_uncontrolled_nondeterminism($name, $num_iters);
}
}
};
(
depth = $depth:expr, num_iters = $num_iters:expr $(, should_panic = $msg:literal)?;
$(#[$meta:meta])*
fn $name:ident() $body:block
) => {
$crate::shuttle_test! {
num_iters = $num_iters, depth = $depth $(, should_panic = $msg)?;
$(#[$meta])*
fn $name() $body
}
};
}
#[doc(hidden)]
pub fn channel<T>() -> (Sender<T>, RecvTimeoutReceiver<T>) {
let (tx, rx) = shuttle::sync::mpsc::channel();
(tx, RecvTimeoutReceiver { inner: rx })
}
#[doc(hidden)]
pub struct RecvTimeoutReceiver<T> {
inner: shuttle::sync::mpsc::Receiver<T>,
}
impl<T> RecvTimeoutReceiver<T> {
#[doc(hidden)]
pub fn recv_timeout(
&self,
_timeout: std::time::Duration,
) -> Result<T, std::sync::mpsc::RecvTimeoutError> {
use shuttle::rand::Rng;
use std::sync::mpsc::{RecvTimeoutError, TryRecvError};
if shuttle::rand::thread_rng().gen_bool(0.8) {
match self.inner.try_recv() {
Ok(val) => Ok(val),
Err(TryRecvError::Empty) => Err(RecvTimeoutError::Timeout),
Err(TryRecvError::Disconnected) => Err(RecvTimeoutError::Disconnected),
}
} else {
self.inner
.recv()
.map_err(|_| RecvTimeoutError::Disconnected)
}
}
}