Skip to content

Commit 0a9b4ec

Browse files
committed
refactor: streamline ReqBody by removing intermediate channels
1 parent bae42d4 commit 0a9b4ec

8 files changed

Lines changed: 278 additions & 251 deletions

File tree

‎Cargo.toml‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -51,6 +51,7 @@ thiserror = "2"
5151

5252
arc-swap = "1.7"
5353
once_cell = "1.21"
54+
triomphe = "0.1"
5455

5556

5657
matchit = "0.8"

‎crates/http/Cargo.toml‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -31,6 +31,7 @@ futures.workspace = true
3131
trait-variant.workspace = true
3232

3333
thiserror.workspace = true
34+
triomphe.workspace = true
3435

3536
[dev-dependencies]
3637
indoc = "2.0.5"

‎crates/http/src/connection/http_connection.rs‎

Lines changed: 13 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -37,18 +37,18 @@ where
3737
R: Debug,
3838
W: Debug,
3939
{
40-
framed_read: FramedRead<R, RequestDecoder>,
40+
framed_read: Option<FramedRead<R, RequestDecoder>>,
4141
message_writer: MessageWriter<W>,
4242
}
4343

4444
impl<R, W> HttpConnection<R, W>
4545
where
46-
R: AsyncRead + Unpin + Debug,
46+
R: AsyncRead + Unpin + Send + Debug,
4747
W: AsyncWrite + Unpin + Debug,
4848
{
4949
pub fn new(reader: R, writer: W) -> Self {
5050
Self {
51-
framed_read: FramedRead::with_capacity(reader, RequestDecoder::new(), 8 * 1024),
51+
framed_read: Some(FramedRead::with_capacity(reader, RequestDecoder::new(), 8 * 1024)),
5252
message_writer: MessageWriter::with_capacity(writer, 8 * 1024),
5353
}
5454
}
@@ -60,7 +60,9 @@ where
6060
<H::RespBody as Body>::Error: Display,
6161
{
6262
loop {
63-
match self.framed_read.next().await {
63+
let framed_read = self.framed_read.as_mut().expect("framed reader must be available while processing requests");
64+
65+
match framed_read.next().await {
6466
Some(Ok(Message::Header((header, payload_size)))) => {
6567
self.do_process(header, payload_size, handler).await?;
6668
}
@@ -112,19 +114,14 @@ where
112114
}
113115
}
114116

115-
let (req_body, maybe_body_sender) = ReqBody::create_req_body(&mut self.framed_read, payload_size);
117+
let framed_read = self.framed_read.take().expect("framed reader must exist when creating request body");
118+
let (req_body, req_body_state) = ReqBody::create_req_body(framed_read, payload_size);
116119
let request = header.body(req_body);
117120

118-
let response_result = match maybe_body_sender {
119-
None => handler.call(request).await,
120-
Some(mut body_sender) => {
121-
let (handler_result, body_send_result) = tokio::join!(handler.call(request), body_sender.start());
121+
let response_result = handler.call(request).await;
122122

123-
// check if body sender has error
124-
body_send_result?;
125-
handler_result
126-
}
127-
};
123+
let framed_read = req_body_state.finish().await?;
124+
self.framed_read = Some(framed_read);
128125

129126
self.send_response(response_result).await
130127
}
@@ -152,15 +149,8 @@ where
152149
{
153150
let (header_parts, mut body) = response.into_parts();
154151

155-
let payload_size = {
156-
let size_hint = body.size_hint();
157-
match size_hint.exact() {
158-
Some(0) => PayloadSize::Empty,
159-
Some(length) => PayloadSize::Length(length),
160-
None => PayloadSize::Chunked,
161-
}
162-
};
163-
152+
let payload_size: PayloadSize = body.size_hint().into();
153+
164154
let header = Message::<_, T::Data>::Header((ResponseHead::from_parts(header_parts, ()), payload_size));
165155

166156
self.message_writer.write(header)?;

‎crates/http/src/protocol/body/body_channel.rs‎

Lines changed: 0 additions & 191 deletions
This file was deleted.

‎crates/http/src/protocol/body/mod.rs‎

Lines changed: 12 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -13,10 +13,13 @@
1313
//! The body handling system consists of two main components:
1414
//!
1515
//! - [`ReqBody`]: The consumer side that implements `http_body::Body` trait
16-
//! - [`ReqBodySender`]: The producer side that reads from the raw payload stream
16+
//! - [`ReqBodyState`]: Connection-side guard that owns the streaming state
1717
//!
18-
//! These components communicate through channels to enable concurrent processing while
19-
//! maintaining backpressure.
18+
//! The connection hands a [`ReqBody`] to the request handler and keeps the
19+
//! associated [`ReqBodyState`]. Once the handler finishes, the connection uses
20+
//! the state to finish draining any unread body data and to reclaim ownership of
21+
//! the underlying decoder. This design removes the previous channel-based
22+
//! indirection and lets the handler poll the decoder directly.
2023
//!
2124
//! # Design Goals
2225
//!
@@ -28,27 +31,19 @@
2831
//! - Ensure complete body consumption even if handler abandons reading
2932
//! - Maintain proper connection state for keep-alive support
3033
//!
31-
//! 3. **Concurrent Processing**
32-
//! - Allow request handling to proceed while body streams
33-
//! - Support cancellation and cleanup in error cases
34+
//! 3. **Graceful Cancellation**
35+
//! - Ensure complete body consumption even if the handler drops the body
36+
//! without reading it
37+
//! - Support cleanup in error cases without spawning helper tasks
3438
//!
3539
//! 4. **Clean Abstractions**
3640
//! - Hide channel complexity from consumers
3741
//! - Provide standard http_body::Body interface
3842
//!
3943
//! # Implementation Details
4044
//!
41-
//! The body handling implementation uses:
42-
//!
43-
//! - MPSC channel for signaling between consumer and producer
44-
//! - Oneshot channels for individual chunk transfers
45-
//! - EOF tracking to ensure complete body consumption
46-
//! - Automatic cleanup of unread data
47-
//!
48-
//! See individual component documentation for more details.
49-
50-
//mod req_body_2;
51-
mod body_channel;
5245
mod req_body;
5346

5447
pub use req_body::ReqBody;
48+
#[allow(unused_imports)]
49+
pub(crate) use req_body::ReqBodyState;

0 commit comments

Comments
 (0)