extern crate alloc;
use alloc::collections::VecDeque;
use alloc::vec::Vec;
#[cfg(feature = "std")]
use std::sync::{Condvar, Mutex};
#[cfg(feature = "std")]
enum Status<A> {
Vacuus,
Perfectus(A),
}
#[cfg(feature = "std")]
pub struct Dilatum<A> {
state: Mutex<Status<A>>,
cond: Condvar,
}
#[cfg(feature = "std")]
impl<A: Clone> Dilatum<A> {
#[inline]
pub fn new() -> Self {
Dilatum {
state: Mutex::new(Status::Vacuus),
cond: Condvar::new(),
}
}
#[inline]
pub fn complete(&self, a: A) -> bool {
let mut guard = self
.state
.lock()
.expect("Dilatum::complete: state mutex poisoned — library invariant violated");
match &*guard {
Status::Perfectus(_) => false,
Status::Vacuus => {
*guard = Status::Perfectus(a);
self.cond.notify_all();
true
}
}
}
#[inline]
pub fn is_completed(&self) -> bool {
let guard = self
.state
.lock()
.expect("Dilatum::is_completed: state mutex poisoned — library invariant violated");
matches!(&*guard, Status::Perfectus(_))
}
#[inline]
pub fn try_get(&self) -> Option<A> {
let guard = self
.state
.lock()
.expect("Dilatum::try_get: state mutex poisoned — library invariant violated");
match &*guard {
Status::Perfectus(a) => Some(a.clone()),
Status::Vacuus => None,
}
}
pub fn get_blocking(&self) -> A {
let mut guard = self
.state
.lock()
.expect("Dilatum::get_blocking: state mutex poisoned — library invariant violated");
loop {
match &*guard {
Status::Perfectus(a) => return a.clone(),
Status::Vacuus => {
guard = self.cond.wait(guard).expect(
"Dilatum::get_blocking: condvar wait poisoned — library invariant violated",
);
}
}
}
}
}
#[cfg(feature = "std")]
impl<A: Clone> Default for Dilatum<A> {
fn default() -> Self {
Self::new()
}
}
#[cfg(feature = "std")]
pub struct Referentia<A> {
value: Mutex<A>,
}
#[cfg(feature = "std")]
impl<A: Clone> Referentia<A> {
#[inline]
pub fn new(a: A) -> Self {
Referentia {
value: Mutex::new(a),
}
}
#[inline]
pub fn get(&self) -> A {
let guard = self
.value
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
guard.clone()
}
#[inline]
pub fn set(&self, a: A) -> A {
let mut guard = self
.value
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
core::mem::replace(&mut *guard, a)
}
#[inline]
pub fn update<F>(&self, f: F)
where
F: FnOnce(A) -> A,
{
let mut guard = self
.value
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let new_value = f(guard.clone());
*guard = new_value;
}
#[inline]
pub fn get_and_update<F>(&self, f: F) -> A
where
F: FnOnce(A) -> A,
{
let mut guard = self
.value
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let old = guard.clone();
*guard = f(old.clone());
old
}
#[inline]
pub fn update_and_get<F>(&self, f: F) -> A
where
F: FnOnce(A) -> A,
{
let mut guard = self
.value
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let new_value = f(guard.clone());
*guard = new_value.clone();
new_value
}
#[inline]
pub fn modify<B, F>(&self, f: F) -> B
where
F: FnOnce(A) -> (A, B),
{
let mut guard = self
.value
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let (new_value, result) = f(guard.clone());
*guard = new_value;
result
}
}
#[cfg(feature = "std")]
pub struct Semaphorum {
permits: Mutex<usize>,
cond: Condvar,
}
#[cfg(feature = "std")]
impl Semaphorum {
#[inline]
pub fn new(permits: usize) -> Self {
Semaphorum {
permits: Mutex::new(permits),
cond: Condvar::new(),
}
}
#[inline]
pub fn available(&self) -> usize {
*self
.permits
.lock()
.expect("Semaphorum::available: permits mutex poisoned — library invariant violated")
}
#[inline]
pub fn try_acquire(&self) -> bool {
let mut guard = self
.permits
.lock()
.expect("Semaphorum::try_acquire: permits mutex poisoned — library invariant violated");
if *guard > 0 {
*guard -= 1;
true
} else {
false
}
}
pub fn acquire_blocking(&self) {
let mut guard = self.permits.lock().expect(
"Semaphorum::acquire_blocking: permits mutex poisoned — library invariant violated",
);
while *guard == 0 {
guard = self.cond.wait(guard).expect(
"Semaphorum::acquire_blocking: condvar wait poisoned — library invariant violated",
);
}
*guard -= 1;
}
#[inline]
pub fn try_acquire_n(&self, n: usize) -> bool {
let mut guard = self.permits.lock().expect(
"Semaphorum::try_acquire_n: permits mutex poisoned — library invariant violated",
);
if *guard >= n {
*guard -= n;
true
} else {
false
}
}
#[inline]
pub fn release(&self) {
let mut guard = self
.permits
.lock()
.expect("Semaphorum::release: permits mutex poisoned — library invariant violated");
*guard += 1;
self.cond.notify_one();
}
#[inline]
pub fn release_n(&self, n: usize) {
let mut guard = self
.permits
.lock()
.expect("Semaphorum::release_n: permits mutex poisoned — library invariant violated");
*guard += n;
self.cond.notify_all();
}
}
#[cfg(feature = "std")]
pub struct Permissum<'a> {
semaphore: &'a Semaphorum,
}
#[cfg(feature = "std")]
impl<'a> Permissum<'a> {
#[inline]
pub fn new(semaphore: &'a Semaphorum) -> Self {
semaphore.acquire_blocking();
Permissum { semaphore }
}
#[inline]
pub fn release(self) {
}
}
#[cfg(feature = "std")]
impl Drop for Permissum<'_> {
fn drop(&mut self) {
self.semaphore.release();
}
}
#[cfg(feature = "std")]
pub struct MVarSync<A> {
value: Mutex<Option<A>>,
cond: Condvar,
}
#[cfg(feature = "std")]
impl<A> MVarSync<A> {
#[inline]
pub fn new_empty() -> Self {
MVarSync {
value: Mutex::new(None),
cond: Condvar::new(),
}
}
#[inline]
pub fn new(a: A) -> Self {
MVarSync {
value: Mutex::new(Some(a)),
cond: Condvar::new(),
}
}
#[inline]
pub fn is_empty(&self) -> bool {
self.value
.lock()
.expect("MVarSync::is_empty: value mutex poisoned — library invariant violated")
.is_none()
}
#[inline]
pub fn try_take(&self) -> Option<A> {
let mut guard = self
.value
.lock()
.expect("MVarSync::try_take: value mutex poisoned — library invariant violated");
let result = guard.take();
if result.is_some() {
self.cond.notify_one();
}
result
}
pub fn take_blocking(&self) -> A {
let mut guard = self
.value
.lock()
.expect("MVarSync::take_blocking: value mutex poisoned — library invariant violated");
while guard.is_none() {
guard = self.cond.wait(guard).expect(
"MVarSync::take_blocking: condvar wait poisoned — library invariant violated",
);
}
let result = guard
.take()
.expect("MVarSync::take_blocking: value must be Some after wait loop — invariant bug");
self.cond.notify_one();
result
}
#[inline]
pub fn try_put(&self, a: A) -> bool {
let mut guard = self
.value
.lock()
.expect("MVarSync::try_put: value mutex poisoned — library invariant violated");
if guard.is_none() {
*guard = Some(a);
self.cond.notify_one();
true
} else {
false
}
}
pub fn put_blocking(&self, a: A) {
let mut guard = self
.value
.lock()
.expect("MVarSync::put_blocking: value mutex poisoned — library invariant violated");
while guard.is_some() {
guard = self.cond.wait(guard).expect(
"MVarSync::put_blocking: condvar wait poisoned — library invariant violated",
);
}
*guard = Some(a);
self.cond.notify_one();
}
pub fn read_blocking(&self) -> A
where
A: Clone,
{
let mut guard = self
.value
.lock()
.expect("MVarSync::read_blocking: value mutex poisoned — library invariant violated");
while guard.is_none() {
guard = self.cond.wait(guard).expect(
"MVarSync::read_blocking: condvar wait poisoned — library invariant violated",
);
}
guard
.clone()
.expect("MVarSync::read_blocking: value must be Some after wait loop — invariant bug")
}
pub fn swap_blocking(&self, a: A) -> A {
let mut guard = self
.value
.lock()
.expect("MVarSync::swap_blocking: value mutex poisoned — library invariant violated");
while guard.is_none() {
guard = self.cond.wait(guard).expect(
"MVarSync::swap_blocking: condvar wait poisoned — library invariant violated",
);
}
let old = guard
.take()
.expect("MVarSync::swap_blocking: value must be Some after wait loop — invariant bug");
*guard = Some(a);
self.cond.notify_one();
old
}
}
#[cfg(feature = "std")]
pub struct CaudaBackpressure<A> {
capacity: usize,
buffer: Mutex<VecDeque<A>>,
not_empty: Condvar,
not_full: Condvar,
}
#[cfg(feature = "std")]
impl<A> CaudaBackpressure<A> {
#[inline]
pub fn new(capacity: usize) -> Self {
CaudaBackpressure {
capacity,
buffer: Mutex::new(VecDeque::with_capacity(capacity)),
not_empty: Condvar::new(),
not_full: Condvar::new(),
}
}
#[inline]
pub fn capacity(&self) -> usize {
self.capacity
}
#[inline]
pub fn size(&self) -> usize {
self.buffer
.lock()
.expect("CaudaBackpressure::size: buffer mutex poisoned — library invariant violated")
.len()
}
#[inline]
pub fn is_empty(&self) -> bool {
self.buffer
.lock()
.expect(
"CaudaBackpressure::is_empty: buffer mutex poisoned — library invariant violated",
)
.is_empty()
}
#[inline]
pub fn is_full(&self) -> bool {
self.buffer
.lock()
.expect(
"CaudaBackpressure::is_full: buffer mutex poisoned — library invariant violated",
)
.len()
>= self.capacity
}
#[inline]
pub fn try_offer(&self, a: A) -> bool {
let mut guard = self.buffer.lock().expect(
"CaudaBackpressure::try_offer: buffer mutex poisoned — library invariant violated",
);
if guard.len() < self.capacity {
guard.push_back(a);
self.not_empty.notify_one();
true
} else {
false
}
}
pub fn offer_blocking(&self, a: A) {
let mut guard = self.buffer.lock().expect(
"CaudaBackpressure::offer_blocking: buffer mutex poisoned — library invariant violated",
);
while guard.len() >= self.capacity {
guard = self.not_full.wait(guard).expect(
"CaudaBackpressure::offer_blocking: not_full condvar poisoned — library invariant violated",
);
}
guard.push_back(a);
self.not_empty.notify_one();
}
#[inline]
pub fn try_take(&self) -> Option<A> {
let mut guard = self.buffer.lock().expect(
"CaudaBackpressure::try_take: buffer mutex poisoned — library invariant violated",
);
let result = guard.pop_front();
if result.is_some() {
self.not_full.notify_one();
}
result
}
pub fn take_blocking(&self) -> A {
let mut guard = self.buffer.lock().expect(
"CaudaBackpressure::take_blocking: buffer mutex poisoned — library invariant violated",
);
while guard.is_empty() {
guard = self.not_empty.wait(guard).expect(
"CaudaBackpressure::take_blocking: not_empty condvar poisoned — library invariant violated",
);
}
let result = guard.pop_front().expect(
"CaudaBackpressure::take_blocking: buffer must be non-empty after wait loop — invariant bug",
);
self.not_full.notify_one();
result
}
#[inline]
pub fn peek(&self) -> Option<A>
where
A: Clone,
{
let guard = self
.buffer
.lock()
.expect("CaudaBackpressure::peek: buffer mutex poisoned — library invariant violated");
guard.front().cloned()
}
#[inline]
pub fn drain(&self) -> Vec<A> {
let mut guard = self
.buffer
.lock()
.expect("CaudaBackpressure::drain: buffer mutex poisoned — library invariant violated");
let mut result = Vec::with_capacity(guard.len());
result.extend(guard.drain(..));
self.not_full.notify_all();
result
}
}
#[cfg(all(test, feature = "std"))]
mod tests {
use super::*;
use alloc::vec;
#[test]
fn test_referentia_new() {
let r = Referentia::new(42);
assert_eq!(r.get(), 42);
}
#[test]
fn test_referentia_set() {
let r = Referentia::new(0);
let old = r.set(42);
assert_eq!(old, 0);
assert_eq!(r.get(), 42);
}
#[test]
fn test_referentia_update() {
let r = Referentia::new(21);
r.update(|x| x * 2);
assert_eq!(r.get(), 42);
}
#[test]
fn test_referentia_modify() {
let r = Referentia::new(21);
let result = r.modify(|x| (x * 2, "done"));
assert_eq!(result, "done");
assert_eq!(r.get(), 42);
}
#[test]
fn panicking_update_does_not_poison_forever() {
let r = Referentia::new(1i32);
let _ = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
r.update(|_| -> i32 { panic!("boom") });
}));
assert_eq!(r.get(), 1);
}
#[test]
fn test_semaphorum_new() {
let sem = Semaphorum::new(3);
assert_eq!(sem.available(), 3);
}
#[test]
fn test_semaphorum_try_acquire() {
let sem = Semaphorum::new(1);
assert!(sem.try_acquire());
assert!(!sem.try_acquire());
sem.release();
assert!(sem.try_acquire());
}
#[test]
fn test_semaphorum_release_n() {
let sem = Semaphorum::new(0);
sem.release_n(5);
assert_eq!(sem.available(), 5);
}
#[test]
fn test_mvar_sync_new_empty() {
let mvar: MVarSync<i32> = MVarSync::new_empty();
assert!(mvar.is_empty());
}
#[test]
fn test_mvar_sync_new() {
let mvar = MVarSync::new(42);
assert!(!mvar.is_empty());
}
#[test]
fn test_mvar_sync_try_put_take() {
let mvar: MVarSync<i32> = MVarSync::new_empty();
assert!(mvar.try_put(42));
assert!(!mvar.try_put(43)); assert_eq!(mvar.try_take(), Some(42));
assert!(mvar.is_empty());
}
#[test]
fn test_cauda_new() {
let queue: CaudaBackpressure<i32> = CaudaBackpressure::new(10);
assert_eq!(queue.capacity(), 10);
assert!(queue.is_empty());
}
#[test]
fn test_cauda_try_offer_take() {
let queue = CaudaBackpressure::new(2);
assert!(queue.try_offer(1));
assert!(queue.try_offer(2));
assert!(!queue.try_offer(3));
assert_eq!(queue.try_take(), Some(1));
assert_eq!(queue.try_take(), Some(2));
assert_eq!(queue.try_take(), None);
}
#[test]
fn test_cauda_drain() {
let queue = CaudaBackpressure::new(10);
queue.try_offer(1);
queue.try_offer(2);
queue.try_offer(3);
let items = queue.drain();
assert_eq!(items, vec![1, 2, 3]);
assert!(queue.is_empty());
}
#[test]
fn test_dilatum_complete() {
let deferred = Dilatum::<i32>::new();
assert!(!deferred.is_completed());
assert!(deferred.complete(42));
assert!(deferred.is_completed());
assert!(!deferred.complete(43)); assert_eq!(deferred.try_get(), Some(42));
}
#[test]
fn test_dilatum_concurrent_complete_runs_exactly_once() {
use alloc::sync::Arc;
use core::sync::atomic::{AtomicU32, Ordering};
use std::thread;
let deferred = Arc::new(Dilatum::<u32>::new());
let counter = Arc::new(AtomicU32::new(0));
let handles: Vec<_> = (0..8)
.map(|_| {
let deferred = Arc::clone(&deferred);
let counter = Arc::clone(&counter);
thread::spawn(move || {
let value = counter.fetch_add(1, Ordering::SeqCst) + 1;
deferred.complete(value)
})
})
.collect();
let wins = handles
.into_iter()
.map(|h| h.join().expect("thread panicked"))
.filter(|&won| won)
.count();
assert_eq!(wins, 1, "exactly one complete() call must win the race");
assert!(deferred.is_completed());
assert!(
deferred.try_get().is_some(),
"completed deferred must never yield a torn (empty) read"
);
}
#[test]
fn test_dilatum_no_torn_read_under_race() {
use alloc::sync::Arc;
use std::thread;
let deferred = Arc::new(Dilatum::<u32>::new());
let writers: Vec<_> = (0..4)
.map(|i| {
let deferred = Arc::clone(&deferred);
thread::spawn(move || {
deferred.complete(i);
})
})
.collect();
let reader_deferred = Arc::clone(&deferred);
let reader = thread::spawn(move || {
for _ in 0..5_000 {
if reader_deferred.is_completed() {
assert!(
reader_deferred.try_get().is_some(),
"torn read: is_completed() true but try_get() returned None"
);
}
}
});
for h in writers {
h.join().expect("writer thread panicked");
}
reader.join().expect("reader thread panicked");
}
#[test]
fn test_permissum_new_consumes_exactly_one_permit() {
let sem = Semaphorum::new(2);
assert_eq!(sem.available(), 2);
let permit = Permissum::new(&sem);
assert_eq!(
sem.available(),
1,
"Permissum::new must consume exactly one permit"
);
drop(permit);
assert_eq!(
sem.available(),
2,
"releasing the permit must restore the original count exactly, no inflation"
);
}
#[test]
fn test_permissum_new_admits_exactly_requested_concurrent_holders() {
use alloc::sync::Arc;
use core::sync::atomic::{AtomicI64, AtomicUsize, Ordering};
use std::thread;
use std::time::Duration;
let sem = Arc::new(Semaphorum::new(2));
let current_holders = Arc::new(AtomicI64::new(0));
let max_holders = Arc::new(AtomicUsize::new(0));
let handles: Vec<_> = (0..5)
.map(|_| {
let sem = Arc::clone(&sem);
let current_holders = Arc::clone(¤t_holders);
let max_holders = Arc::clone(&max_holders);
thread::spawn(move || {
let _permit = Permissum::new(&sem); let now = current_holders.fetch_add(1, Ordering::SeqCst) + 1;
max_holders.fetch_max(now as usize, Ordering::SeqCst);
thread::sleep(Duration::from_millis(20));
current_holders.fetch_sub(1, Ordering::SeqCst);
})
})
.collect();
for h in handles {
h.join().expect("thread panicked");
}
assert_eq!(
max_holders.load(Ordering::SeqCst),
2,
"Semaphorum::new(2) must admit exactly 2 concurrent Permissum holders"
);
assert_eq!(
sem.available(),
2,
"all permits must be returned, no inflation"
);
}
}