hypershell_tokio_components/providers/
line.rs1use core::marker::PhantomData;
2
3use cgp::extra::handler::{Handler, HandlerComponent};
4use cgp::prelude::*;
5use futures::Stream;
6use hypershell_components::dsl::StreamToLines;
7use tokio::io::AsyncRead;
8use tokio_util::codec::{FramedRead, LinesCodec, LinesCodecError};
9
10#[cgp_new_provider]
11impl<Context, Input> Handler<Context, StreamToLines, Input> for HandleStreamToLines
12where
13 Context: HasAsyncErrorType,
14 Input: Send + AsyncRead + Unpin + 'static,
15{
16 type Output = Box<dyn Stream<Item = Result<String, LinesCodecError>> + Send>;
17
18 async fn handle(
19 _context: &Context,
20 _tag: PhantomData<StreamToLines>,
21 input: Input,
22 ) -> Result<Box<dyn Stream<Item = Result<String, LinesCodecError>> + Send>, Context::Error>
23 {
24 let stream = FramedRead::new(input, LinesCodec::new());
25
26 Ok(Box::new(stream))
27 }
28}