From 90d4a0393a1c5c02723984a2c2993372c10029cd Mon Sep 17 00:00:00 2001 From: Rajas Paranjpe <52586855+ChocolateLoverRaj@users.noreply.github.com> Date: Sun, 18 Feb 2024 19:09:09 -0800 Subject: [PATCH] Add split map example --- examples/split_map.rs | 39 +++++++++++++++++++++++++++++++++++++++ 1 file changed, 39 insertions(+) create mode 100644 examples/split_map.rs diff --git a/examples/split_map.rs b/examples/split_map.rs new file mode 100644 index 0000000..04dfe1f --- /dev/null +++ b/examples/split_map.rs @@ -0,0 +1,39 @@ +use futures::StreamExt; +use split_stream_by::{Either, SplitStreamByMapExt}; +use tokio::join; + +#[derive(Debug)] +struct Request { + //... +} +#[derive(Debug)] +struct Response { + //... +} +#[derive(Debug)] +enum Message { + Request(Request), + Response(Response), +} + +#[tokio::main] +async fn main() { + let incoming_stream = futures::stream::iter([ + Message::Request(Request {}), + Message::Response(Response {}), + Message::Response(Response {}), + ]); + + let (request_stream, response_stream) = incoming_stream.split_by_map(|item| match item { + Message::Request(req) => Either::Left(req), + Message::Response(res) => Either::Right(res), + }); + + let a = tokio::spawn(async { + println!("Requests: {:?}", request_stream.collect::>().await); + }); + let b = tokio::spawn(async { + println!("Responses: {:?}", response_stream.collect::>().await); + }); + let _ = join!(a, b); +}