mod builder;
mod chain;
mod handoff;
mod split;
#[cfg(test)]
mod tests;
pub use builder::{
Assemble, ChainBuilder, ChainFactory, FilterPart, FlatMapPart, InspectPart, MapPart, Root,
RoutedSplit, SinkedChain, SplitBuilder, TryMapPart, chain, chain_owned,
};
pub use chain::{Emitter, Filter, FlatMap, Inspect, Map, StageLifecycle, TryMap, TypedChain};
pub use handoff::{ChunkConfig, SinkHandoff};
pub use split::{Sink, SinkCtx, SplitEmitter, SplitTerminal};
use crate::deser::RecFamily;
use crate::error::FatalError;
use crate::record::{Flow, Record};
use crate::source::PayloadBatch;
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
#[non_exhaustive]
pub enum BlockReason {
Capacity,
NotReady,
}
#[derive(Debug)]
#[non_exhaustive]
pub enum PushOutcome {
Done,
Blocked {
resume_at: usize,
reason: BlockReason,
},
Fatal(FatalError),
}
pub trait RunnableChain: Send {
fn push_batch<'buf>(&mut self, batch: &mut dyn PayloadBatch<'buf>, from: usize) -> PushOutcome;
fn flush(&mut self) -> PushOutcome;
fn abandon_batch(&mut self) {}
}
pub trait Collector<T> {
fn push(&mut self, rec: Record<T>) -> Flow;
}
pub trait CollectorFor<F: RecFamily> {
fn push_rec<'buf>(&mut self, rec: Record<F::Rec<'buf>>) -> Flow;
}
impl<F, C> CollectorFor<F> for C
where
F: RecFamily,
C: for<'buf> Collector<<F as RecFamily>::Rec<'buf>>,
{
#[inline(always)]
fn push_rec<'buf>(&mut self, rec: Record<F::Rec<'buf>>) -> Flow {
self.push(rec)
}
}