mirror of
https://github.com/n0-computer/noq.git
synced 2026-09-22 03:03:39 +00:00
H3: introduce an intermediary type before any body-recv option
This commit is contained in:
committed by
Jean-Christophe BEGUE
parent
13cd3cf86f
commit
73be859e8c
@@ -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<Item = (), Error = Error> {
|
||||
println!("received request: {:?}", request);
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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<Self::Item, Self::Error> {
|
||||
@@ -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")?,
|
||||
|
||||
@@ -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<Self::Item, Self::Error> {
|
||||
@@ -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,
|
||||
|
||||
Reference in New Issue
Block a user