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(¤t_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
}
}