diff --git a/quinn-h3/examples/h3.rs b/quinn-h3/examples/h3.rs index 7d7f0fd31..928b949d2 100644 --- a/quinn-h3/examples/h3.rs +++ b/quinn-h3/examples/h3.rs @@ -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()); diff --git a/quinn-h3/src/client.rs b/quinn-h3/src/client.rs index 1e3737963..53724ee9d 100644 --- a/quinn-h3/src/client.rs +++ b/quinn-h3/src/client.rs @@ -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, 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 { @@ -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 { - 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) - } -}