H3: return RecvBody along response in client, similarly to server

This commit is contained in:
stammw
2019-08-22 18:19:06 +02:00
committed by Jean-Christophe BEGUE
parent 874dafefc7
commit 13cd3cf86f
2 changed files with 27 additions and 59 deletions
+3 -2
View File
@@ -279,9 +279,10 @@ fn client(
conn.send_request_trailers(request, trailer)
.map_err(|e| format_err!("send request: {}", e))
.and_then(|response| {
.and_then(|(response, body)| {
println!("received response: {:?}", response);
let buf = Vec::with_capacity(1024 * 10); // 10K
tokio_io::io::read_to_end(response.body_reader(), buf)
tokio_io::io::read_to_end(body.into_reader(), buf)
.map_err(|e| format_err!("receive response failed: {}", e))
.and_then(|(reader, data)| {
println!("received body len = {}", data.len());
+24 -57
View File
@@ -10,7 +10,7 @@ use slog::{self, o, Logger};
use tokio_io::io::{Shutdown, WriteAll};
use crate::{
body::{Body, BodyReader, RecvBody, RecvBodyStream, SendBody},
body::{Body, RecvBody, SendBody},
connection::{ConnectionDriver, ConnectionRef},
frame::{FrameDecoder, FrameStream},
headers::DecodeHeaders,
@@ -178,10 +178,24 @@ impl SendRequest {
recv: None,
}
}
fn build_response(header: Header) -> Result<Response<()>, Error> {
let (status, headers) = header.into_response_parts()?;
let mut response = Response::builder();
response.status(status);
response.version(http::version::Version::HTTP_3);
*response
.headers_mut()
.ok_or_else(|| Error::peer("invalid response"))? = headers;
Ok(response
.body(())
.or(Err(Error::Internal("failed to build response")))?)
}
}
impl Future for SendRequest {
type Item = RecvResponse;
type Item = (Response<()>, RecvBody);
type Error = Error;
fn poll(&mut self) -> Poll<Self::Item, Self::Error> {
@@ -258,12 +272,14 @@ impl Future for SendRequest {
SendRequestState::Ready(_) => {
match mem::replace(&mut self.state, SendRequestState::Finished) {
SendRequestState::Ready(h) => {
return Ok(Async::Ready(RecvResponse::build(
h,
try_take(&mut self.recv, "Recv is none")?,
try_take(&mut self.stream_id, "stream is none")?,
self.conn.clone(),
)?));
return Ok(Async::Ready((
Self::build_response(h)?,
RecvBody::new(
try_take(&mut self.recv, "Recv is none")?,
self.conn.clone(),
try_take(&mut self.stream_id, "stream is none")?,
),
)));
}
_ => unreachable!(),
}
@@ -273,52 +289,3 @@ impl Future for SendRequest {
}
}
}
pub struct RecvResponse {
response: Response<()>,
recv: FrameStream,
stream_id: StreamId,
conn: ConnectionRef,
}
impl RecvResponse {
fn build(
header: Header,
recv: FrameStream,
stream_id: StreamId,
conn: ConnectionRef,
) -> Result<Self, Error> {
let (status, headers) = header.into_response_parts()?;
let mut response = Response::builder();
response.status(status);
response.version(http::version::Version::HTTP_3);
*response
.headers_mut()
.ok_or_else(|| Error::peer("invalid response"))? = headers;
Ok(Self {
recv,
conn,
stream_id,
response: response
.body(())
.or(Err(Error::Internal("failed to build response")))?,
})
}
pub fn response<'a>(&'a self) -> &'a Response<()> {
&self.response
}
pub fn body(self) -> RecvBody {
RecvBody::new(self.recv, self.conn.clone(), self.stream_id)
}
pub fn body_stream(self) -> RecvBodyStream {
RecvBodyStream::new(self.recv, self.conn, self.stream_id)
}
pub fn body_reader(self) -> BodyReader {
BodyReader::new(self.recv, self.conn, self.stream_id)
}
}