diff --git a/quinn-h3/examples/h3.rs b/quinn-h3/examples/h3.rs index 928b949d2..b7b5d7596 100644 --- a/quinn-h3/examples/h3.rs +++ b/quinn-h3/examples/h3.rs @@ -17,7 +17,7 @@ use url::Url; use quinn::ConnectionDriver as QuicDriver; use quinn_h3::{ self, - body::RecvBody, + body::Receiver, client::Builder as ClientBuilder, connection::ConnectionDriver, server::{Builder as ServerBuilder, IncomingRequest, Sender}, @@ -191,7 +191,7 @@ fn handle_connection( fn handle_request( request: Request<()>, - body: RecvBody, + body: Receiver, sender: Sender, ) -> impl Future { println!("received request: {:?}", request); diff --git a/quinn-h3/src/body.rs b/quinn-h3/src/body.rs index 012e6cffb..66018a57a 100644 --- a/quinn-h3/src/body.rs +++ b/quinn-h3/src/body.rs @@ -97,6 +97,34 @@ impl Future for SendBody { } } +pub struct Receiver { + recv: FrameStream, + conn: ConnectionRef, + stream_id: StreamId, +} + +impl Receiver { + pub(crate) fn new(recv: FrameStream, conn: ConnectionRef, stream_id: StreamId) -> Self { + Self { + conn, + stream_id, + recv, + } + } + + pub fn into_future(self) -> RecvBody { + RecvBody::new(self.recv, self.conn, self.stream_id) + } + + pub fn into_reader(self) -> BodyReader { + BodyReader::new(self.recv, self.conn, self.stream_id) + } + + pub fn into_stream(self) -> RecvBodyStream { + RecvBodyStream::new(self.recv, self.conn, self.stream_id) + } +} + pub struct RecvBody { state: RecvBodyState, capacity: usize, diff --git a/quinn-h3/src/client.rs b/quinn-h3/src/client.rs index 53724ee9d..95163eaf1 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, RecvBody, SendBody}, + body::{Body, Receiver, SendBody}, connection::{ConnectionDriver, ConnectionRef}, frame::{FrameDecoder, FrameStream}, headers::DecodeHeaders, @@ -195,7 +195,7 @@ impl SendRequest { } impl Future for SendRequest { - type Item = (Response<()>, RecvBody); + type Item = (Response<()>, Receiver); type Error = Error; fn poll(&mut self) -> Poll { @@ -274,7 +274,7 @@ impl Future for SendRequest { SendRequestState::Ready(h) => { return Ok(Async::Ready(( Self::build_response(h)?, - RecvBody::new( + Receiver::new( try_take(&mut self.recv, "Recv is none")?, self.conn.clone(), try_take(&mut self.stream_id, "stream is none")?, diff --git a/quinn-h3/src/server.rs b/quinn-h3/src/server.rs index 6f17aab08..500b7efd6 100644 --- a/quinn-h3/src/server.rs +++ b/quinn-h3/src/server.rs @@ -10,7 +10,7 @@ use slog::{self, o, Logger}; use tokio_io::io::{Shutdown, WriteAll}; use crate::{ - body::{Body, RecvBody, SendBody}, + body::{Body, Receiver, SendBody}, connection::{ConnectionDriver, ConnectionRef}, frame::{FrameDecoder, FrameStream}, headers::DecodeHeaders, @@ -177,7 +177,7 @@ impl RecvRequest { } impl Future for RecvRequest { - type Item = (Request<()>, RecvBody, Sender); + type Item = (Request<()>, Receiver, Sender); type Error = Error; fn poll(&mut self) -> Poll { @@ -200,7 +200,7 @@ impl Future for RecvRequest { let (recv, send) = try_take(&mut self.streams, "Recv request invalid state")?; return Ok(Async::Ready(( Self::build_request(header)?, - RecvBody::new(recv, self.conn.clone(), self.stream_id), + Receiver::new(recv, self.conn.clone(), self.stream_id), Sender { send, stream_id: self.stream_id,