extern crate alloc;
use alloc::vec::Vec;
use core::num::NonZeroUsize;
use super::super::{Causal, Event, Gate, VersionVector};
use super::Ideal;
impl<T, G: Gate> Ideal<T, G> {
#[must_use = "a rejected event is returned and must be handled"]
pub fn try_insert(
&mut self,
event: Event<T, G::Dep>,
capacity: NonZeroUsize,
) -> Result<(), AtCapacity<T, G::Dep>> {
let sender = event.stamp.station_id();
if G::stale(&self.progress, sender, &event.deps) {
return Ok(());
}
if self.pending.len() >= capacity.get() {
self.prune_stale();
if self.pending.len() >= capacity.get() {
return Err(AtCapacity { event, capacity });
}
}
self.pending.push(event);
Ok(())
}
fn prune_stale(&mut self) {
let progress = &self.progress;
self.pending.retain(|event| {
let sender = event.stamp.station_id();
let stale = G::stale(progress, sender, &event.deps);
debug_assert!(
!stale || !G::deliverable(progress, sender, &event.deps),
"Gate broke stale/deliverable disjointness: a pending event was both",
);
!stale
});
}
}
impl<T> Ideal<T, Causal> {
#[must_use = "the outcome may carry an evicted event or unchanged refused input"]
pub fn try_insert_evicting_where<F>(
&mut self,
event: Event<T>,
capacity: NonZeroUsize,
mut evictable: F,
) -> Result<BoundedInsert<T>, AtCapacity<T>>
where
F: FnMut(&Event<T>) -> bool,
{
let sender = event.stamp.station_id();
if Causal::stale(&self.progress, sender, &event.deps) {
return Ok(BoundedInsert::DroppedStale);
}
if self.pending.len() < capacity.get() {
self.pending.push(event);
return Ok(BoundedInsert::Buffered);
}
self.prune_stale();
if self.pending.len() < capacity.get() {
self.pending.push(event);
return Ok(BoundedInsert::Buffered);
}
if self.pending.len() > capacity.get() {
return Err(AtCapacity { event, capacity });
}
let mut candidates: Vec<usize> = (0..self.pending.len()).collect();
candidates.sort_unstable_by_key(|&index| self.pending[index].stamp);
let Some(index) = candidates.into_iter().find(|&index| {
self.is_causally_maximal(index, &event) && evictable(&self.pending[index])
}) else {
return Err(AtCapacity { event, capacity });
};
let evicted = self.pending.swap_remove(index);
self.pending.push(event);
Ok(BoundedInsert::Evicted { evicted })
}
fn is_causally_maximal(&self, candidate: usize, arrival: &Event<T>) -> bool {
let victim = &self.pending[candidate];
!Self::depends_on(arrival, victim)
&& self
.pending
.iter()
.enumerate()
.all(|(index, event)| index == candidate || !Self::depends_on(event, victim))
}
fn depends_on(event: &Event<T>, predecessor: &Event<T>) -> bool {
let predecessor_dot = Self::dot(predecessor);
Self::dot(event) != predecessor_dot
&& event.deps.get(predecessor_dot.0) >= predecessor_dot.1
}
fn dot(event: &Event<T>) -> (u32, u64) {
let station = event.stamp.station_id();
(station, event.deps.get(station))
}
}
#[derive(Debug, Clone)]
#[must_use = "an evicted event remains the caller's responsibility"]
pub enum BoundedInsert<T> {
Buffered,
DroppedStale,
Evicted {
evicted: Event<T>,
},
}
#[derive(Debug, Clone, thiserror::Error)]
#[error("buffer at or above capacity {capacity}: the event was not buffered")]
pub struct AtCapacity<T, D = VersionVector> {
pub event: Event<T, D>,
pub capacity: NonZeroUsize,
}