use std::{
collections::VecDeque,
sync::{Condvar, Mutex, MutexGuard, PoisonError},
};
#[cfg(test)]
mod tests;
#[derive(Debug)]
#[must_use = "hand the ticket back through ItemQueue::complete or ItemQueue::fail"]
pub struct ItemTicket<T> {
item: T,
failed_on: Vec<usize>,
}
impl<T> ItemTicket<T> {
#[must_use]
pub const fn item(&self) -> &T {
&self.item
}
#[must_use]
pub const fn previous_failures(&self) -> usize {
self.failed_on.len()
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum FailOutcome<T> {
Requeued,
Abandoned(T),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct QueueCounts {
pub pending: usize,
pub in_flight: usize,
pub completed: usize,
pub abandoned: usize,
}
#[derive(Debug)]
struct Entry<T> {
item: T,
failed_on: Vec<usize>,
}
#[derive(Debug)]
struct QueueState<T> {
fresh: VecDeque<Entry<T>>,
retry: VecDeque<Entry<T>>,
live: Vec<bool>,
in_flight: usize,
completed: usize,
abandoned: usize,
abandoned_unclaimed: Vec<T>,
closed: bool,
}
impl<T> QueueState<T> {
fn is_live(&self, lane: usize) -> bool {
self.live.get(lane).copied().unwrap_or(false)
}
fn nobody_left_for(&self, entry: &Entry<T>) -> bool {
self.live
.iter()
.enumerate()
.all(|(lane, &live)| !live || entry.failed_on.contains(&lane))
}
fn take_for(&mut self, lane: usize) -> Option<Entry<T>> {
let retried = self
.retry
.iter()
.position(|entry| !entry.failed_on.contains(&lane))
.and_then(|index| self.retry.remove(index));
retried.or_else(|| self.fresh.pop_front())
}
fn retire(&mut self, lane: usize) -> Vec<T> {
if let Some(live) = self.live.get_mut(lane) {
*live = false;
}
let mut dropped = Vec::new();
let retry = std::mem::take(&mut self.retry);
for entry in retry {
if self.nobody_left_for(&entry) {
dropped.push(entry.item);
} else {
self.retry.push_back(entry);
}
}
if !self.live.iter().any(|&live| live) {
dropped.extend(self.fresh.drain(..).map(|entry| entry.item));
}
self.abandoned += dropped.len();
dropped
}
}
#[derive(Debug)]
pub struct ItemQueue<T> {
state: Mutex<QueueState<T>>,
wake: Condvar,
}
impl<T> ItemQueue<T> {
pub fn new(items: impl IntoIterator<Item = T>, lane_count: usize) -> Self {
Self {
state: Mutex::new(QueueState {
fresh: items
.into_iter()
.map(|item| Entry {
item,
failed_on: Vec::new(),
})
.collect(),
retry: VecDeque::new(),
live: vec![true; lane_count],
in_flight: 0,
completed: 0,
abandoned: 0,
abandoned_unclaimed: Vec::new(),
closed: false,
}),
wake: Condvar::new(),
}
}
fn lock(&self) -> MutexGuard<'_, QueueState<T>> {
self.state.lock().unwrap_or_else(PoisonError::into_inner)
}
pub fn claim(&self, lane: usize) -> Option<ItemTicket<T>> {
let mut state = self.lock();
loop {
if state.closed || !state.is_live(lane) {
return None;
}
if let Some(entry) = state.take_for(lane) {
state.in_flight += 1;
return Some(ItemTicket {
item: entry.item,
failed_on: entry.failed_on,
});
}
if state.in_flight == 0 {
let dropped = state.retire(lane);
state.abandoned_unclaimed.extend(dropped);
drop(state);
self.wake.notify_all();
return None;
}
state = self
.wake
.wait(state)
.unwrap_or_else(PoisonError::into_inner);
}
}
pub fn complete(&self, ticket: ItemTicket<T>) -> T {
let mut state = self.lock();
state.in_flight -= 1;
state.completed += 1;
drop(state);
self.wake.notify_all();
ticket.item
}
#[must_use = "an abandoned item must be reported"]
pub fn fail(&self, lane: usize, ticket: ItemTicket<T>) -> FailOutcome<T> {
let mut entry = Entry {
item: ticket.item,
failed_on: ticket.failed_on,
};
if !entry.failed_on.contains(&lane) {
entry.failed_on.push(lane);
}
let mut state = self.lock();
state.in_flight -= 1;
let outcome = if state.nobody_left_for(&entry) {
state.abandoned += 1;
FailOutcome::Abandoned(entry.item)
} else {
state.retry.push_back(entry);
FailOutcome::Requeued
};
drop(state);
self.wake.notify_all();
outcome
}
#[must_use = "the returned items were abandoned and must be reported"]
pub fn retire(&self, lane: usize) -> Vec<T> {
let dropped = self.lock().retire(lane);
self.wake.notify_all();
dropped
}
pub fn take_abandoned(&self) -> Vec<T> {
std::mem::take(&mut self.lock().abandoned_unclaimed)
}
#[must_use = "the returned items were never processed"]
pub fn close(&self) -> Vec<T> {
let mut state = self.lock();
state.closed = true;
let mut left = std::mem::take(&mut state.abandoned_unclaimed);
left.extend(state.retry.drain(..).map(|entry| entry.item));
left.extend(state.fresh.drain(..).map(|entry| entry.item));
drop(state);
self.wake.notify_all();
left
}
#[must_use]
pub fn counts(&self) -> QueueCounts {
let state = self.lock();
let counts = QueueCounts {
pending: state.fresh.len() + state.retry.len(),
in_flight: state.in_flight,
completed: state.completed,
abandoned: state.abandoned,
};
drop(state);
counts
}
}