use std::marker::PhantomData;
use zrx_scheduler::action::context::Binding;
use zrx_scheduler::action::options::Interest;
use zrx_scheduler::action::{Action, Context, Options};
use zrx_scheduler::schedule::Subscriber;
use zrx_scheduler::step::{IntoSteps, Scope};
use zrx_scheduler::{Id, Key, Value};
use crate::stream::barrier::{Barrier, Barriers};
use crate::stream::Stream;
use super::Operator;
#[derive(Debug)]
pub struct Select<I, T> {
barriers: Barriers<I>,
marker: PhantomData<T>,
}
impl<I, T> Stream<I, T>
where
I: Id + Value,
T: Value,
{
#[inline]
pub fn select<B>(&self, iter: B) -> Stream<I, Vec<(Key<I>, T)>>
where
B: IntoIterator<Item = (Key<I>, Barrier<I>)>,
{
let options = Options::default().interest(Interest::Enter);
let barriers = Barriers::from_iter(iter);
self.subscribe(
Subscriber::new(Select { barriers, marker: PhantomData })
.with_options(options),
)
}
}
impl<I, T> Action<I> for Select<I, T>
where
I: Id + Value,
T: Value,
{
type Inputs = (T,);
type Output<'a> = Vec<(Key<I>, T)>;
fn execute(&mut self, ctx: Context<I, Self>) -> impl IntoSteps<I, Self> {
let Binding {
events,
scopes,
inputs,
mut output,
..
} = ctx.bind();
for event in events {
self.barriers.handle(&event);
}
for scope in scopes {
if inputs.contains_key(scope.key()) {
self.barriers.notify(scope.key());
}
}
self.barriers.drain().map(move |advance| {
let new_key = advance.scope().clone();
output.insert(
new_key.clone(),
advance
.into_iter()
.cloned()
.map(|key| {
let value = inputs.get(&key).expect("invariant");
(key, value.clone())
})
.collect(),
);
Scope::from(new_key).done()
})
}
}