LastStage is a specification for exchanging events (Batch) between producers and consumers. with push-based model
Producer produce events and dispatch to subscribers by Dispatcher (no input - one/many output)
ProducerConsumer it get events from upstream and dispatch to subscribers by Dispatcher (one/many input - one/many output)
Consumer it get events from upstream and consume it (one/many input - no output)
Dispatcher first get one-many subscriber then start to dispatch events by two mode (Broadcast / RoundRobin)
Examples directory:
[Simple] (https://github.com/Rustixir/last_stage/blob/master/examples/simple.rs)
[Multi] (https://github.com/Rustixir/last_stage/blob/master/examples/multi.rs)
laststage = "1.0.0"
#[tokio::main]asyncfnmain(){let(shutdown_sender, shutdown_recv) = channel();// -----------------------------------//// Producer -> ProducerConsumer -> Consumer //// ------------------------------------// Run Consumerlet log_chan = ConsumerRunnable::new(Box::new(Log)).run(100);// Run ProducerConsumerlet filter_chan = ProducerConsumerRunnable::new(Box::new(FilterByAge),vec![log_chan],Some(DispatcherType::RoundRobin)).unwrap().run(100);// Run Producerlet _ = ProducerRunnable::new(Box::new(Prod),vec![filter_chan],None,100, shutdown_recv).unwrap().run();
tokio::time::sleep(Duration::from_secs(10)).await;}#[derive(Clone)]structProdEvent{pubfuname:String,pubage:i32}structProd;#[async_trait]implProducer<ProdEvent>forProd{asyncfninit(&mutself){}asyncfnterminate(&mutself){}asyncfnhandle_demand(&mutself,max_demand:usize) -> Vec<ProdEvent>{(0..max_demand asi32).into_iter().map(|i| {ProdEvent{funame:format!("DanyalMh-{}", i),age:(i + 30) % 35}}).collect()}}// -------------------------------------------structFilterByAge;#[async_trait]implProducerConsumer<ProdEvent,ProdEvent>forFilterByAge{asyncfninit(&mutself){}asyncfnterminate(&mutself){}asyncfnhandle_events(&mutself,events:Vec<ProdEvent>) -> Vec<ProdEvent>{
events
.into_iter().filter(|pe| pe.age > 25 && pe.age < 32).collect()}}structLog;#[async_trait]implConsumer<ProdEvent>forLog{asyncfninit(&mutself){}asyncfnterminate(&mutself){}asyncfnhandle_events(&mutself,events:Vec<ProdEvent>) -> State<ProdEvent>{
events
.into_iter().for_each(|pe| {println!("==> {} -> {}", pe.funame, pe.age)});State::Continue}}