Skip to main content

ntex_codec/
lib.rs

1//! Traits for encoding and decoding frames.
2//!
3//! A codec turns a byte stream into frames and back. [`Decoder`] splits
4//! incoming bytes into frames, [`Encoder`] serializes outgoing frames. The
5//! `ntex-io` crate drives them: `Io::recv()` and `IoRef::decode()` run the
6//! decoder on the read buffer, `Io::send()` and `IoRef::encode()` run the
7//! encoder on the write buffer, and the `ntex-dispatcher` crate uses both to
8//! connect a stream to a service.
9//!
10//! Both traits take `&self`, since a codec is shared between the reading and
11//! writing sides. A codec that keeps state between calls, such as a parser
12//! position or negotiated settings, has to use interior mutability, e.g.
13//! [`Cell`](std::cell::Cell) or [`RefCell`](std::cell::RefCell).
14//!
15//! # Example
16//!
17//! A codec for newline-terminated lines:
18//!
19//! ```
20//! use std::io;
21//!
22//! use ntex_bytes::{BytePages, Bytes, BytesMut};
23//! use ntex_codec::{Decoder, Encoder};
24//!
25//! struct LineCodec;
26//!
27//! impl Decoder for LineCodec {
28//!     type Item = Bytes;
29//!     type Error = io::Error;
30//!
31//!     fn decode(&self, src: &mut BytesMut) -> Result<Option<Bytes>, io::Error> {
32//!         match src.iter().position(|b| *b == b'\n') {
33//!             // consume the line and its terminator
34//!             Some(n) => Ok(Some(src.split_to(n + 1).slice(..n))),
35//!             // incomplete line, wait for more input
36//!             None => Ok(None),
37//!         }
38//!     }
39//!
40//!     fn decode_eof(&self, src: &mut BytesMut) -> Result<Option<Bytes>, io::Error> {
41//!         match self.decode(src)? {
42//!             Some(line) => Ok(Some(line)),
43//!             // the last line has no terminator
44//!             None if !src.is_empty() => Ok(Some(src.split_to(src.len()))),
45//!             None => Ok(None),
46//!         }
47//!     }
48//! }
49//!
50//! impl Encoder for LineCodec {
51//!     type Item = Bytes;
52//!     type Error = io::Error;
53//!
54//!     fn encode(&self, item: Bytes, dst: &mut BytePages) -> Result<(), io::Error> {
55//!         if item.contains(&b'\n') {
56//!             return Err(io::Error::new(io::ErrorKind::InvalidInput, "newline in line"));
57//!         }
58//!         dst.append(item);
59//!         dst.extend_from_slice(b"\n");
60//!         Ok(())
61//!     }
62//! }
63//!
64//! let mut src = BytesMut::from(&b"one\ntwo"[..]);
65//! assert_eq!(LineCodec.decode(&mut src).unwrap().unwrap(), "one");
66//! assert!(LineCodec.decode(&mut src).unwrap().is_none());
67//! assert_eq!(LineCodec.decode_eof(&mut src).unwrap().unwrap(), "two");
68//! assert!(LineCodec.decode_eof(&mut src).unwrap().is_none());
69//! ```
70
71use std::{fmt, io, rc::Rc};
72
73use ntex_bytes::{BytePages, Bytes, BytesMut};
74
75/// Serializes frames into bytes.
76pub trait Encoder {
77    /// The type of frames consumed by the encoder.
78    type Item;
79
80    /// The type of encoding errors.
81    type Error: fmt::Debug;
82
83    /// Encodes a frame and appends it to `dst`.
84    ///
85    /// `dst` is the write buffer, a list of pages. [`BytePages::append`] adds
86    /// a `Bytes` value as a page of its own without copying it, while
87    /// [`BytePages::extend_from_slice`] and the [`BufMut`](ntex_bytes::BufMut)
88    /// methods copy into the pages.
89    ///
90    /// Output written to `dst` is not rolled back when this returns an
91    /// error, it is sent like any other output. Validate the frame before
92    /// writing it, so that an error does not leave a partial frame behind.
93    fn encode(&self, item: Self::Item, dst: &mut BytePages) -> Result<(), Self::Error>;
94}
95
96/// Splits a byte stream into frames.
97pub trait Decoder {
98    /// The type of decoded frames.
99    type Item: fmt::Debug;
100
101    /// The type of unrecoverable frame decoding errors.
102    ///
103    /// If an individual message is ill-formed but can be ignored without
104    /// interfering with the processing of future messages, it may be more
105    /// useful to report the failure as an `Item`.
106    type Error: fmt::Debug;
107
108    /// Attempts to decode a frame from the buffered input.
109    ///
110    /// Returns `Ok(Some(item))` after removing exactly one frame's bytes from
111    /// the front of `src`, and `Ok(None)` if `src` does not hold a complete
112    /// frame yet. On `None` the partial frame must stay in `src`, it is
113    /// passed in again together with the input that arrives next.
114    ///
115    /// This is called repeatedly while it returns frames, and may be called
116    /// with an empty buffer.
117    fn decode(&self, src: &mut BytesMut) -> Result<Option<Self::Item>, Self::Error>;
118
119    /// Attempts to decode a frame once the transport reached a clean EOF.
120    ///
121    /// The peer closed its write half and no further input will arrive, so
122    /// this is the decoder's chance to produce a frame from whatever is left
123    /// in `src`, or to report a truncated one as an error. Once the transport
124    /// is at EOF it is used instead of [`decode`](Self::decode) for every
125    /// decode attempt, it may be called again after it returned `None`, and
126    /// with an empty buffer.
127    ///
128    /// The default implementation calls [`decode`](Self::decode).
129    fn decode_eof(&self, src: &mut BytesMut) -> Result<Option<Self::Item>, Self::Error> {
130        self.decode(src)
131    }
132}
133
134impl<T> Encoder for Rc<T>
135where
136    T: Encoder,
137{
138    type Item = T::Item;
139    type Error = T::Error;
140
141    fn encode(&self, item: Self::Item, dst: &mut BytePages) -> Result<(), Self::Error> {
142        (**self).encode(item, dst)
143    }
144}
145
146impl<T> Decoder for Rc<T>
147where
148    T: Decoder,
149{
150    type Item = T::Item;
151    type Error = T::Error;
152
153    fn decode(&self, src: &mut BytesMut) -> Result<Option<Self::Item>, Self::Error> {
154        (**self).decode(src)
155    }
156
157    fn decode_eof(&self, src: &mut BytesMut) -> Result<Option<Self::Item>, Self::Error> {
158        (**self).decode_eof(src)
159    }
160}
161
162/// Passes bytes through unchanged.
163///
164/// Decoding returns everything buffered as one frame, and returns `None` only
165/// when the buffer is empty. Encoding appends the `Bytes` value to the write
166/// buffer as is, without copying it.
167#[derive(Debug, Copy, Clone)]
168pub struct BytesCodec;
169
170impl Encoder for BytesCodec {
171    type Item = Bytes;
172    type Error = io::Error;
173
174    #[inline]
175    fn encode(&self, item: Bytes, dst: &mut BytePages) -> Result<(), Self::Error> {
176        dst.append(item);
177        Ok(())
178    }
179}
180
181impl Decoder for BytesCodec {
182    type Item = Bytes;
183    type Error = io::Error;
184
185    fn decode(&self, src: &mut BytesMut) -> Result<Option<Self::Item>, Self::Error> {
186        if src.is_empty() {
187            Ok(None)
188        } else {
189            Ok(Some(src.split_to(src.len())))
190        }
191    }
192}
193
194#[cfg(test)]
195mod tests {
196    use std::cell::Cell;
197
198    use super::*;
199
200    /// Decodes one-byte frames and counts `decode_eof` calls.
201    #[derive(Default)]
202    struct ByteCodec {
203        eof: Cell<usize>,
204    }
205
206    impl Decoder for ByteCodec {
207        type Item = u8;
208        type Error = io::Error;
209
210        fn decode(&self, src: &mut BytesMut) -> Result<Option<u8>, io::Error> {
211            if src.is_empty() {
212                Ok(None)
213            } else {
214                Ok(Some(src.split_to(1)[0]))
215            }
216        }
217
218        fn decode_eof(&self, src: &mut BytesMut) -> Result<Option<u8>, io::Error> {
219            self.eof.set(self.eof.get() + 1);
220            self.decode(src)
221        }
222    }
223
224    /// Uses the default `decode_eof`.
225    struct DefaultEof;
226
227    impl Decoder for DefaultEof {
228        type Item = usize;
229        type Error = io::Error;
230
231        fn decode(&self, src: &mut BytesMut) -> Result<Option<usize>, io::Error> {
232            let len = src.len();
233            src.clear();
234            Ok(Some(len))
235        }
236    }
237
238    #[test]
239    fn bytes_codec() {
240        let codec = BytesCodec;
241        let codec2 = codec;
242        assert_eq!(format!("{:?}", codec2.clone()), "BytesCodec");
243
244        let mut src = BytesMut::new();
245        assert!(codec.decode(&mut src).unwrap().is_none());
246        assert!(codec.decode_eof(&mut src).unwrap().is_none());
247
248        src.extend_from_slice(b"hello");
249        assert_eq!(codec.decode(&mut src).unwrap().unwrap(), "hello");
250        assert!(src.is_empty());
251        src.extend_from_slice(b"eof");
252        assert_eq!(codec.decode_eof(&mut src).unwrap().unwrap(), "eof");
253        assert!(src.is_empty());
254
255        let mut dst = BytePages::default();
256        codec
257            .encode(Bytes::from_static(b"hello "), &mut dst)
258            .unwrap();
259        codec
260            .encode(Bytes::from_static(b"world"), &mut dst)
261            .unwrap();
262        codec.encode(Bytes::new(), &mut dst).unwrap();
263        assert_eq!(dst.len(), 11);
264        assert_eq!(dst.freeze(), "hello world");
265    }
266
267    #[test]
268    fn default_decode_eof() {
269        let mut src = BytesMut::from(&b"abc"[..]);
270        assert_eq!(DefaultEof.decode_eof(&mut src).unwrap(), Some(3));
271        assert_eq!(DefaultEof.decode_eof(&mut src).unwrap(), Some(0));
272    }
273
274    #[test]
275    fn rc_codec() {
276        let codec = Rc::new(ByteCodec::default());
277        let mut src = BytesMut::from(&b"ab"[..]);
278        assert_eq!(codec.decode(&mut src).unwrap(), Some(b'a'));
279        assert_eq!(codec.eof.get(), 0);
280        // forwards to the inner `decode_eof`, not the default one
281        assert_eq!(codec.decode_eof(&mut src).unwrap(), Some(b'b'));
282        assert_eq!(codec.decode_eof(&mut src).unwrap(), None);
283        assert_eq!(codec.eof.get(), 2);
284
285        let codec = Rc::new(BytesCodec);
286        let mut dst = BytePages::default();
287        codec.encode(Bytes::from_static(b"rc"), &mut dst).unwrap();
288        assert_eq!(dst.freeze(), "rc");
289    }
290}