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}