use zrx_scheduler::action::options::Event;
use zrx_scheduler::{Id, Key};
use zrx_store::stash::Items;
use zrx_store::{Stash, Store};
use super::advance::Advance;
use super::lifecycle::Lifecycle;
use super::Barrier;
mod drain;
pub use drain::Drain;
#[derive(Clone, Debug)]
pub struct Barriers<I> {
inner: Stash<Key<I>, Barrier<I>>,
scopes: Stash<Key<I>, Items>,
lifecycle: Lifecycle,
fulfilled: Items,
}
impl<I> Barriers<I>
where
I: Id,
{
#[must_use]
pub fn new() -> Self {
Self {
inner: Stash::new(),
scopes: Stash::new(),
lifecycle: Lifecycle::default(),
fulfilled: Items::new(),
}
}
pub fn insert<S, B>(&mut self, scope: S, barrier: B)
where
S: Into<Key<I>>,
B: Into<Barrier<I>>,
{
let b = self.inner.insert(scope.into(), barrier.into());
let barrier = &mut self.inner[b];
for (s, (scope, items)) in self.scopes.slots_mut() {
if barrier.contains(scope) {
barrier.items.insert(s);
items.insert(b);
}
}
if self.inner[b].is_complete(&self.lifecycle) {
self.fulfilled.insert(b);
}
}
pub fn remove(&mut self, scope: &Key<I>) -> Option<Barrier<I>> {
let b = self.inner.get(scope)?;
let (_, mut barrier) = self.inner.remove(b)?;
barrier.items.clear();
self.fulfilled.remove(b);
Some(barrier)
}
pub fn handle(&mut self, event: &Event<I>) {
match event {
Event::Insert(scope) => {
let s = self.scopes.get(scope).unwrap_or_else(|| {
self.scopes.insert(scope.clone(), Items::new())
});
if !self.lifecycle.submit(s) {
return;
}
for (b, (_, barrier)) in self.inner.slots_mut() {
if barrier.contains(scope) {
self.scopes[s].insert(b);
if barrier.insert(s)
&& barrier.is_complete(&self.lifecycle)
{
self.fulfilled.insert(b);
} else {
self.fulfilled.remove(b);
}
}
}
}
Event::Remove(scope) => {
let Some(s) = self.scopes.get(scope) else {
return;
};
if !self.lifecycle.withdraw(s) {
return;
}
for b in &self.scopes[s] {
if self.inner[b].remove(s)
&& self.inner[b].is_complete(&self.lifecycle)
{
self.fulfilled.insert(b);
}
}
self.scopes.remove(s);
}
}
}
pub fn notify(&mut self, scope: &Key<I>) {
let Some(s) = self.scopes.get(scope) else {
return;
};
self.lifecycle.complete(s);
for b in &self.scopes[s] {
if self.inner[b].is_complete(&self.lifecycle) {
self.fulfilled.insert(b);
}
}
}
}
#[allow(clippy::must_use_candidate)]
impl<I> Barriers<I>
where
I: Id,
{
#[inline]
pub fn len(&self) -> usize {
self.inner.len()
}
#[inline]
pub fn is_empty(&self) -> bool {
self.inner.is_empty()
}
}
impl<S, I> FromIterator<(S, Barrier<I>)> for Barriers<I>
where
S: Into<Key<I>>,
I: Id,
{
#[inline]
fn from_iter<T>(iter: T) -> Self
where
T: IntoIterator<Item = (S, Barrier<I>)>,
{
let mut barriers = Barriers::new();
for (scope, barrier) in iter {
barriers.insert(scope.into(), barrier);
}
barriers
}
}
impl<I> Default for Barriers<I>
where
I: Id,
{
#[inline]
fn default() -> Self {
Self::new()
}
}