Skip to main content

spaniel/protocol/codec/
mod.rs

1use bytes::Bytes;
2use bytes::BytesMut;
3use futures::Async;
4use futures::AsyncSink;
5use futures::Poll;
6use futures::Sink;
7use futures::StartSend;
8use futures::Stream;
9use protocol::frames::Frame;
10use protocol::frames::FrameHead;
11use tokio_io::codec::length_delimited::{self, Framed};
12use tokio_io::AsyncRead;
13use tokio_io::AsyncWrite;
14
15mod buffer;
16pub mod reader;
17pub mod writer;
18
19pub struct FrameCodec<T>
20where
21    T: AsyncRead + AsyncWrite,
22{
23    inner: Framed<T, Bytes>,
24}
25
26impl<T> FrameCodec<T>
27where
28    T: AsyncRead + AsyncWrite,
29{
30    pub fn new(conn: T) -> Self {
31        Self {
32            inner: length_delimited::Builder::new()
33                .big_endian()
34                .length_adjustment(-4)
35                .length_field_offset(0)
36                .length_field_length(4)
37                .max_frame_length(::std::u32::MAX as usize)
38                .new_framed(conn),
39        }
40    }
41}
42
43impl<T> Stream for FrameCodec<T>
44where
45    T: AsyncRead + AsyncWrite,
46{
47    type Item = Frame;
48    type Error = (); // TODO err
49
50    fn poll(&mut self) -> Poll<Option<Self::Item>, Self::Error> {
51        let frame = match try_ready!(self.inner.poll().map_err(|_err| {
52            // TODO err
53            ()
54        })) {
55            Some(buf) => Some(Frame::decode_from(buf).expect("deserialization")),
56            None => None,
57        };
58        Ok(Async::Ready(frame))
59    }
60}
61
62impl<T> Sink for FrameCodec<T>
63where
64    T: AsyncRead + AsyncWrite,
65{
66    // TODO err
67    type SinkItem = Option<Frame>;
68    type SinkError = ();
69
70    fn start_send(&mut self, item: Self::SinkItem) -> StartSend<Self::SinkItem, Self::SinkError> {
71        match item {
72            None => {
73                return Ok(AsyncSink::Ready);
74            }
75            Some(frame) => {
76                // TODO buffer provider
77                let size = FrameHead::encoded_len() + frame.encoded_len();
78                let mut buf = BytesMut::with_capacity(size);
79                frame.encode_into(&mut buf).expect("serialization");
80
81                match self.inner.start_send(buf.freeze()) {
82                    Ok(AsyncSink::NotReady(_)) => Ok(AsyncSink::NotReady(Some(frame))),
83                    Ok(AsyncSink::Ready) => Ok(AsyncSink::Ready),
84                    Err(_err) => Err(()),
85                }
86            }
87        }
88    }
89
90    fn poll_complete(&mut self) -> Poll<(), Self::SinkError> {
91        self.inner.poll_complete().map_err(|_| ()) // TODO err
92    }
93
94    fn close(&mut self) -> Poll<(), Self::SinkError> {
95        self.poll_complete()
96    }
97}