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 {
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);
loop {
self.inner.readable().await.map_err(|e| Error::Io(e))?;
let mut buf = [0; CHUNK_SIZE];
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)
}
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);
let mut current_chunk = match chunks_iter.next() {
Some(v) => v,
None => return Ok(()) };
loop {
self.inner.writable().await.map_err(|e| Error::Io(e))?;
match self.inner.try_write(¤t_chunk) {
Ok(n) => {
if n != current_chunk.remaining() {
#[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))
}
}
}
pub async fn response(&self, mut response: Response) -> Result<(), Error> {
self.write_bytes(response.serialize()).await
}
pub async fn request(&self, basic_request: BasicRequest) -> Result<(), Error> {
self.write_bytes(basic_request.serialize()).await
}
}
impl AsRef<TcpStream> for Stream {
fn as_ref(&self) -> &TcpStream {
&self.inner
}
}
impl AsMut<TcpStream> for Stream {
fn as_mut(&mut self) -> &mut TcpStream {
&mut self.inner
}
}
impl Into<TcpStream> for Stream {
fn into(self) -> TcpStream {
self.inner
}
}