use std::{
collections::VecDeque,
sync::{
Mutex, PoisonError,
atomic::{AtomicU32, Ordering},
},
};
#[cfg(test)]
mod tests;
#[derive(Debug)]
pub struct SampleCursor {
next: AtomicU32,
end: u32,
local_retry: Mutex<VecDeque<(u32, u32)>>,
requeued: Mutex<VecDeque<(u32, u32)>>,
}
impl SampleCursor {
#[must_use]
pub const fn new(start: u32, end: u32) -> Self {
Self {
next: AtomicU32::new(start),
end,
local_retry: Mutex::new(VecDeque::new()),
requeued: Mutex::new(VecDeque::new()),
}
}
#[must_use]
pub const fn end(&self) -> u32 {
self.end
}
pub fn claim(&self, want: u32) -> Option<(u32, u32)> {
if want == 0 {
return None;
}
let start = self.next.fetch_add(want, Ordering::Relaxed);
if start >= self.end {
return None;
}
Some((start, want.min(self.end - start)))
}
pub fn claim_local(&self, want: u32) -> Option<(u32, u32)> {
let retried = self
.local_retry
.lock()
.unwrap_or_else(PoisonError::into_inner)
.pop_front();
if retried.is_some() {
return retried;
}
self.claim(want)
}
pub fn claim_local_bounded(&self, want: u32) -> Option<(u32, u32)> {
if want == 0 {
return None;
}
pop_bounded(&self.local_retry, want).or_else(|| self.claim(want))
}
pub fn claim_any(&self, want: u32) -> Option<(u32, u32)> {
if want == 0 {
return None;
}
pop_bounded(&self.requeued, want).or_else(|| self.claim(want))
}
pub fn shared_pool_exhausted(&self) -> bool {
self.next.load(Ordering::Relaxed) >= self.end
}
pub fn fully_claimed(&self) -> bool {
self.shared_pool_exhausted()
&& self
.requeued
.lock()
.unwrap_or_else(PoisonError::into_inner)
.is_empty()
&& self
.local_retry
.lock()
.unwrap_or_else(PoisonError::into_inner)
.is_empty()
}
pub fn return_to_local(&self, start: u32, count: u32) {
if count == 0 {
return;
}
self.local_retry
.lock()
.unwrap_or_else(PoisonError::into_inner)
.push_back((start, count));
}
pub fn requeue(&self, start: u32, count: u32) {
if count == 0 {
return;
}
self.requeued
.lock()
.unwrap_or_else(PoisonError::into_inner)
.push_back((start, count));
}
}
fn pop_bounded(pile: &Mutex<VecDeque<(u32, u32)>>, want: u32) -> Option<(u32, u32)> {
let mut pile = pile.lock().unwrap_or_else(PoisonError::into_inner);
let head = pile.pop_front();
let claimed = match head {
Some((start, count)) if count > want => {
pile.push_front((start + want, count - want));
Some((start, want))
}
other => other,
};
drop(pile);
claimed
}