http_streams_core/
json_nl_codec.rs1use crate::error::{StreamError, StreamErrorKind};
7use bytes::BytesMut;
8use serde::Deserialize;
9use std::marker::PhantomData;
10use tokio_util::codec::{Decoder, LinesCodec, LinesCodecError};
11
12#[derive(Debug)]
14pub struct JsonNewLineCodec<T> {
15 inner: LinesCodec,
16 _ph: PhantomData<fn() -> T>,
21}
22
23impl<T> JsonNewLineCodec<T> {
24 pub fn new_with_max_length(max_length: usize) -> Self {
26 Self {
27 inner: LinesCodec::new_with_max_length(max_length),
28 _ph: PhantomData,
29 }
30 }
31}
32
33fn frame_error(err: LinesCodecError) -> StreamError {
38 match err {
39 LinesCodecError::MaxLineLengthExceeded => StreamError::new(
40 StreamErrorKind::MaxLenReachedError,
41 None,
42 Some("Max line length reached".into()),
43 ),
44 LinesCodecError::Io(err) => StreamError::from(err),
45 }
46}
47
48fn parse<T>(line: &str) -> Option<Result<T, StreamError>>
53where
54 T: for<'de> Deserialize<'de>,
55{
56 Some(
57 serde_json::from_str(line)
58 .map_err(|err| StreamError::new(StreamErrorKind::CodecError, Some(Box::new(err)), None)),
59 )
60}
61
62impl<T> Decoder for JsonNewLineCodec<T>
63where
64 T: for<'de> Deserialize<'de>,
65{
66 type Item = Result<T, StreamError>;
67 type Error = StreamError;
68
69 fn decode(&mut self, buf: &mut BytesMut) -> Result<Option<Self::Item>, StreamError> {
70 match self.inner.decode(buf).map_err(frame_error)? {
71 Some(line) => Ok(parse(&line)),
72 None => Ok(None),
73 }
74 }
75
76 fn decode_eof(&mut self, buf: &mut BytesMut) -> Result<Option<Self::Item>, StreamError> {
77 match self.inner.decode_eof(buf).map_err(frame_error)? {
78 Some(line) => Ok(parse(&line)),
79 None => Ok(None),
80 }
81 }
82}