extern crate alloc;
mod capacity;
mod released;
use alloc::vec::Vec;
use super::{Causal, Event, Fifo, Gate, VersionVector};
pub use capacity::{AtCapacity, BoundedInsert};
pub use released::Released;
pub type FifoIdeal<T> = Ideal<T, Fifo>;
pub struct Ideal<T, G: Gate> {
progress: G::Progress,
pending: Vec<Event<T, G::Dep>>,
}
pub type CausalIdeal<T> = Ideal<T, Causal>;
impl<T> Ideal<T, Causal> {
#[must_use]
pub const fn new() -> Self {
Self {
progress: VersionVector::new(),
pending: Vec::new(),
}
}
#[must_use]
pub const fn delivered(&self) -> &VersionVector {
&self.progress
}
}
impl<T, G: Gate> Ideal<T, G> {
pub(crate) fn advance_progress(&mut self, sender: u32, dep: &G::Dep) {
G::advance(&mut self.progress, sender, dep);
self.pending
.retain(|event| !G::stale(&self.progress, event.stamp.station_id(), &event.deps));
}
pub fn insert(&mut self, event: Event<T, G::Dep>) {
let sender = event.stamp.station_id();
if G::stale(&self.progress, sender, &event.deps) {
return;
}
self.pending.push(event);
}
pub fn pop_ready_event(&mut self) -> Option<Released<T, G>> {
let index = self.next_deliverable()?;
Some(self.release_index(index))
}
pub fn pop_ready_event_where<F>(&mut self, mut eligible: F) -> Option<Released<T, G>>
where
F: FnMut(&Event<T, G::Dep>) -> bool,
{
let index = self.next_deliverable_where(&mut eligible)?;
if !self.is_deliverable(&self.pending[index]) {
return None;
}
Some(self.release_index(index))
}
fn release_index(&mut self, index: usize) -> Released<T, G> {
let event = self.pending.swap_remove(index);
G::advance(&mut self.progress, event.stamp.station_id(), &event.deps);
Released::new(event)
}
pub fn pop_ready(&mut self) -> Option<T> {
self.pop_ready_event().map(Released::into_payload)
}
#[must_use]
pub const fn pending_len(&self) -> usize {
self.pending.len()
}
#[must_use]
pub const fn progress(&self) -> &G::Progress {
&self.progress
}
pub fn frontier(&self) -> impl Iterator<Item = &Event<T, G::Dep>> {
self.pending
.iter()
.filter(|&event| self.is_deliverable(event))
}
fn next_deliverable(&self) -> Option<usize> {
self.pending
.iter()
.enumerate()
.filter(|&(_, event)| self.is_deliverable(event))
.min_by_key(|&(_, event)| event.stamp)
.map(|(index, _)| index)
}
fn next_deliverable_where<F>(&self, eligible: &mut F) -> Option<usize>
where
F: FnMut(&Event<T, G::Dep>) -> bool,
{
let mut candidates: Vec<usize> = (0..self.pending.len()).collect();
candidates.sort_unstable_by_key(|&index| self.pending[index].stamp);
let mut decisions = Vec::with_capacity(candidates.len());
decisions.resize(candidates.len(), None);
loop {
let mut invoked = false;
for (&index, decision) in candidates.iter().zip(&mut decisions) {
if !self.is_deliverable(&self.pending[index]) {
continue;
}
match *decision {
Some(true) => return Some(index),
Some(false) => {}
None => {
*decision = Some(eligible(&self.pending[index]));
invoked = true;
break;
}
}
}
if !invoked {
return None;
}
}
}
fn is_deliverable(&self, event: &Event<T, G::Dep>) -> bool {
G::deliverable(&self.progress, event.stamp.station_id(), &event.deps)
}
}
impl<T, G: Gate> Default for Ideal<T, G> {
fn default() -> Self {
Self {
progress: G::Progress::default(),
pending: Vec::new(),
}
}
}