pub mod clean_log;
mod clone;
mod error;
mod fork;
mod registry;
pub mod ring_queue;
mod states;
pub use clone::CloneStream;
pub use error::{CloneStreamError, Result};
use fork::Fork;
pub use fork::ForkConfig;
use futures::Stream;
pub trait ForkStream: Stream<Item: Clone> + Sized {
fn fork(self) -> CloneStream<Self> {
CloneStream::from(Fork::new(self))
}
fn fork_with_limits(self, max_queue_size: usize, max_clone_count: usize) -> CloneStream<Self> {
let config = ForkConfig {
max_clone_count,
max_queue_size,
};
CloneStream::from(Fork::with_config(self, config))
}
}
impl<BaseStream> ForkStream for BaseStream where BaseStream: Stream<Item: Clone> {}
impl<BaseStream> From<BaseStream> for CloneStream<BaseStream>
where
BaseStream: Stream<Item: Clone>,
{
fn from(base_stream: BaseStream) -> Self {
Self::from(Fork::new(base_stream))
}
}