Skip to main content

hypershell_tokio_components/providers/
line.rs

1use 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}