#![allow(clippy::new_without_default)]
use std::marker::PhantomData;
use zrx_scheduler::action::Action;
use zrx_scheduler::schedule::Subscriber;
use zrx_scheduler::Id;
use super::Stream;
mod filter;
mod filter_map;
mod join;
mod map;
mod product;
mod select;
pub use filter::Filter;
pub use filter_map::FilterMap;
pub use join::Join;
pub use map::Map;
pub use product::Product;
pub use select::Select;
pub trait Operator<'a, I, A>
where
A: Action<I>,
{
fn subscribe<S>(&self, subscriber: S) -> Stream<I, A::Output<'a>>
where
S: Into<Subscriber<'a, I, A>>;
}
impl<'a, I, A, T> Operator<'a, I, A> for Stream<I, T>
where
I: Id,
A: Action<I, Inputs = (T,)> + 'static,
{
#[inline]
fn subscribe<S>(&self, subscriber: S) -> Stream<I, A::Output<'a>>
where
S: Into<Subscriber<'a, I, A>>,
{
let id = self.workflow.with(|builder| {
builder.add([self.id], subscriber).expect("invariant")
});
Stream {
id,
workflow: self.workflow.clone(),
marker: PhantomData,
}
}
}
macro_rules! impl_operator_for_tuple {
($T1:ident $(, $T:ident)+ $(,)?) => {
impl<'a, I, A, $T1, $($T,)+> Operator<'a, I, A>
for (Stream<I, $T1>, $(Stream<I, $T>,)+)
where
I: Id,
A: Action<I, Inputs = ($T1, $($T,)+)> + 'static,
{
#[inline]
fn subscribe<S>(&self, subscriber: S) -> Stream<I, A::Output<'a>>
where
S: Into<Subscriber<'a, I, A>>,
{
#[allow(non_snake_case)]
let ($T1, $($T,)*) = self;
$(assert_eq!($T1.workflow, $T.workflow);)+
let id = $T1.workflow.with(|builder| {
builder
.add([$T1.id, $($T.id,)*], subscriber)
.expect("invariant")
});
Stream {
id,
workflow: $T1.workflow.clone(),
marker: PhantomData,
}
}
}
};
}
impl_operator_for_tuple!(T1, T2);
impl_operator_for_tuple!(T1, T2, T3);
impl_operator_for_tuple!(T1, T2, T3, T4);
impl_operator_for_tuple!(T1, T2, T3, T4, T5);
impl_operator_for_tuple!(T1, T2, T3, T4, T5, T6);
impl_operator_for_tuple!(T1, T2, T3, T4, T5, T6, T7);
impl_operator_for_tuple!(T1, T2, T3, T4, T5, T6, T7, T8);