cataclysm/
stream.rs

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
use tokio::net::TcpStream;
use bytes::Buf;
use crate::{Error, http::{Response, BasicRequest}};

const CHUNK_SIZE: usize = 4_096;

pub struct Stream {
    inner: TcpStream
}

impl Stream {
    /// Generates a new stream
    pub fn new(stream: TcpStream) -> Stream {
        Stream{inner: stream}
    }

    pub async fn try_read_response(&self) -> Result<Response, Error> {
        let mut response_bytes = Vec::with_capacity(CHUNK_SIZE);
        // First we read
        loop {
            self.inner.readable().await.map_err(|e| Error::Io(e))?;
            
            // being stored in the async task.
            let mut buf = [0; CHUNK_SIZE];

            // Try to read data, this may still fail with `WouldBlock`
            // if the readiness event is a false positive.
            match self.inner.try_read(&mut buf) {
                Ok(0) => {
                    break
                },
                Ok(n) => {
                    response_bytes.extend_from_slice(&buf[0..n]);
                },
                Err(ref e) if e.kind() == std::io::ErrorKind::WouldBlock => {
                    if response_bytes.is_empty() {
                        continue
                    } else {
                        break
                    }
                }
                Err(e) => return Err(Error::Io(e))
            }
        }

        Response::parse(response_bytes)
    }

    /// Writes bytes through the tcp connection
    pub async fn write_bytes<A: AsRef<[u8]>>(&self, bytes: A) -> Result<(), Error> {
        let bytes_ref: &[u8] = bytes.as_ref();
        let mut chunks_iter = bytes_ref.chunks(CHUNK_SIZE);
        #[cfg(feature = "full_log")]
        log::trace!("writting {} chunks of maximum {} bytes each", chunks_iter.len(), CHUNK_SIZE);
        // We check the first chunk
        let mut current_chunk = match chunks_iter.next() {
            Some(v) => v,
            None => return Ok(()) // Zero length response
        };
        loop {
            // Wait for the socket to be writable
            self.inner.writable().await.map_err(|e| Error::Io(e))?;
    
            // Try to write data, this may still fail with `WouldBlock`
            // if the readiness event is a false positive.        
            match self.inner.try_write(&current_chunk) {
                Ok(n) => {
                    if n != current_chunk.remaining() {
                        // There are some bytes still to be written in this chunk
                        #[cfg(feature = "full_log")]
                        log::debug!("incomplete chunk, trying to serve remaining bytes ({}/{})", current_chunk.len(), CHUNK_SIZE);
                        current_chunk.advance(n);
                        continue;
                    } else {
                        current_chunk = match chunks_iter.next() {
                            Some(v) => v,
                            None => return Ok(())
                        }
                    }
                }
                Err(ref e) if e.kind() == tokio::io::ErrorKind::WouldBlock => {
                    continue;
                }
                Err(e) => break Err(Error::Io(e))
            }
        }
    }

    /// Allows to send a response through the stream
    pub async fn response(&self, mut response: Response) -> Result<(), Error> {
        self.write_bytes(response.serialize()).await
    }

    /// Allows to send a basic request through the stream
    pub async fn request(&self, basic_request: BasicRequest) -> Result<(), Error> {
        self.write_bytes(basic_request.serialize()).await
    }
}

// Reference access to the inner structure
impl AsRef<TcpStream> for Stream {
    fn as_ref(&self) -> &TcpStream {
        &self.inner
    }
}

// Mutable reference access to the inner structure
impl AsMut<TcpStream> for Stream {
    fn as_mut(&mut self) -> &mut TcpStream {
        &mut self.inner
    }
}

// Conversion to inner type
impl Into<TcpStream> for Stream {
    fn into(self) -> TcpStream {
        self.inner
    }
}