use std::collections::HashSet;
use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
use std::sync::Mutex;
pub(super) struct QueueGate {
pub(super) depth: AtomicUsize,
pub(super) cap: usize,
pub(super) rejected: AtomicU64,
}
impl QueueGate {
pub(super) fn new(cap: usize) -> Self {
QueueGate {
depth: AtomicUsize::new(0),
cap,
rejected: AtomicU64::new(0),
}
}
pub(super) fn try_reserve(&self) -> Result<(), usize> {
let mut current = self.depth.load(Ordering::Acquire);
loop {
if current >= self.cap {
self.rejected.fetch_add(1, Ordering::Relaxed);
return Err(current);
}
match self.depth.compare_exchange_weak(
current,
current + 1,
Ordering::AcqRel,
Ordering::Acquire,
) {
Ok(_) => return Ok(()),
Err(actual) => current = actual,
}
}
}
pub(super) fn release(&self) {
let previous = self.depth.fetch_sub(1, Ordering::AcqRel);
debug_assert!(previous > 0, "queue depth underflow");
}
pub(super) fn depth(&self) -> usize {
self.depth.load(Ordering::Relaxed)
}
pub(super) fn rejected(&self) -> u64 {
self.rejected.load(Ordering::Relaxed)
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord)]
pub(super) struct AbortId(pub(super) u64);
#[derive(Default)]
pub(super) struct AbortInbox {
pub(super) pending: Mutex<HashSet<AbortId>>,
pub(super) next_id: AtomicU64,
pub(super) aborted: AtomicU64,
}
impl AbortInbox {
pub(super) fn next_id(&self) -> AbortId {
AbortId(self.next_id.fetch_add(1, Ordering::Relaxed))
}
pub(super) fn enqueue(&self, id: AbortId) {
self.pending
.lock()
.unwrap_or_else(|p| p.into_inner())
.insert(id);
}
pub(super) fn drain(&self) -> HashSet<AbortId> {
std::mem::take(&mut *self.pending.lock().unwrap_or_else(|p| p.into_inner()))
}
pub(super) fn aborted(&self) -> u64 {
self.aborted.load(Ordering::Relaxed)
}
}