Skip to main content

ntex_io/
framed.rs

1use std::{fmt, io};
2
3use ntex_codec::{Decoder, Encoder};
4use ntex_util::future::Either;
5
6use crate::IoBoxed;
7
8/// A codec-driven interface to an underlying I/O stream.
9///
10/// Uses the `Encoder` and `Decoder` traits to encode and decode frames.
11pub struct Framed<U> {
12    io: IoBoxed,
13    codec: U,
14}
15
16impl<U> Framed<U> {
17    /// Wraps an I/O object with a codec.
18    pub fn new<Io>(io: Io, codec: U) -> Framed<U>
19    where
20        IoBoxed: From<Io>,
21    {
22        Framed {
23            codec,
24            io: IoBoxed::from(io),
25        }
26    }
27
28    #[inline]
29    /// Returns a reference to the underlying I/O stream wrapped by `Framed`.
30    pub fn get_io(&self) -> &IoBoxed {
31        &self.io
32    }
33
34    #[inline]
35    /// Returns a reference to the underlying codec.
36    pub fn get_codec(&self) -> &U {
37        &self.codec
38    }
39
40    #[inline]
41    /// Returns the underlying I/O object and codec.
42    pub fn into_inner(self) -> (IoBoxed, U) {
43        (self.io, self.codec)
44    }
45}
46
47impl<U> Framed<U>
48where
49    U: Decoder + Encoder,
50{
51    /// Flushes encoded data to the transport.
52    ///
53    /// If `full` is `true`, waits until all buffered data has been written.
54    pub async fn flush(&self, full: bool) -> io::Result<()> {
55        self.io.flush(full).await
56    }
57
58    /// Gracefully shuts down the I/O stream.
59    pub async fn shutdown(&self) -> io::Result<()> {
60        self.io.shutdown().await
61    }
62}
63
64impl<U> Framed<U>
65where
66    U: Decoder,
67{
68    #[inline]
69    /// Reads and decodes the next item.
70    ///
71    /// Returns `Ok(None)` when the connection can no longer produce another
72    /// complete item and no error was recorded. Codec errors are returned in
73    /// `Either::Left`; dispatcher timeouts and connection errors from the
74    /// transport, a filter, or shutdown are returned in `Either::Right`.
75    pub async fn recv(&self) -> Result<Option<U::Item>, Either<U::Error, io::Error>> {
76        self.io.recv(&self.codec).await
77    }
78}
79
80impl<U> Framed<U>
81where
82    U: Encoder,
83{
84    #[inline]
85    /// Encodes an item and fully flushes it to the transport.
86    ///
87    /// Codec errors are returned in `Either::Left`; connection errors from the
88    /// transport, a filter, or shutdown, including an expired write timeout,
89    /// are returned in `Either::Right`.
90    pub async fn send(
91        &self,
92        item: <U as Encoder>::Item,
93    ) -> Result<(), Either<U::Error, io::Error>> {
94        self.io.send(item, &self.codec).await
95    }
96}
97
98impl<U> fmt::Debug for Framed<U>
99where
100    U: fmt::Debug,
101{
102    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
103        f.debug_struct("Framed")
104            .field("codec", &self.codec)
105            .finish()
106    }
107}
108
109#[cfg(test)]
110mod tests {
111    use ntex_bytes::Bytes;
112    use ntex_codec::BytesCodec;
113
114    use super::*;
115    use crate::{Io, testing::IoTest};
116
117    #[ntex::test]
118    async fn framed() {
119        let (client, server) = IoTest::create();
120        client.remote_buffer_cap(1024);
121        client.write(b"chunk-0");
122
123        let server = Framed::new(Io::from(server), BytesCodec);
124        server.get_codec();
125        server.get_io();
126        assert!(format!("{server:?}").contains("Framed"));
127
128        let item = server.recv().await.unwrap().unwrap();
129        assert_eq!(item, b"chunk-0".as_ref());
130
131        let data = Bytes::from_static(b"chunk-1");
132        server.send(data).await.unwrap();
133        server.flush(true).await.unwrap();
134        assert_eq!(client.read_any(), b"chunk-1".as_ref());
135
136        server.shutdown().await.unwrap();
137        assert!(client.is_closed());
138
139        let (io, _codec) = server.into_inner();
140        assert!(io.is_closed());
141    }
142}