Skip to content

Commit ccf8746

Browse files
committed
vey-proxy: support h2 informational response
1 parent bff1ad1 commit ccf8746

13 files changed

Lines changed: 481 additions & 144 deletions

File tree

lib/vey-h2/src/client/mod.rs

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,7 @@
1+
/*
2+
* SPDX-License-Identifier: Apache-2.0
3+
* Copyright 2026 VEY-OSS developers.
4+
*/
5+
6+
mod response;
7+
pub use response::H2ResponseHeaderReceiver;

lib/vey-h2/src/client/response.rs

Lines changed: 47 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,47 @@
1+
/*
2+
* SPDX-License-Identifier: Apache-2.0
3+
* Copyright 2026 VEY-OSS developers.
4+
*/
5+
6+
use std::future::poll_fn;
7+
use std::pin::Pin;
8+
use std::task::{Context, Poll, ready};
9+
10+
use h2::client::ResponseFuture;
11+
use h2::{Error, RecvStream};
12+
use http::Response;
13+
14+
pub struct H2ResponseHeaderReceiver {
15+
rsp_fut: ResponseFuture,
16+
rsp_body: Option<RecvStream>,
17+
}
18+
19+
impl H2ResponseHeaderReceiver {
20+
pub fn new(rsp_fut: ResponseFuture) -> Self {
21+
Self {
22+
rsp_fut,
23+
rsp_body: None,
24+
}
25+
}
26+
27+
fn poll_recv(&mut self, cx: &mut Context<'_>) -> Poll<Result<Response<()>, Error>> {
28+
if let Some(r) = ready!(self.rsp_fut.poll_informational(cx)) {
29+
return Poll::Ready(r);
30+
}
31+
32+
let r = ready!(Pin::new(&mut self.rsp_fut).poll(cx))?;
33+
let (header, body) = r.into_parts();
34+
35+
self.rsp_body = Some(body);
36+
Poll::Ready(Ok(Response::from_parts(header, ())))
37+
}
38+
39+
pub async fn recv_header(&mut self) -> Result<Response<()>, Error> {
40+
poll_fn(move |cx| self.poll_recv(cx)).await
41+
}
42+
43+
#[inline]
44+
pub fn take_body(&mut self) -> Option<RecvStream> {
45+
self.rsp_body.take()
46+
}
47+
}

lib/vey-h2/src/ext/request.rs

Lines changed: 19 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -7,14 +7,15 @@ use std::io::Write;
77

88
use bytes::BufMut;
99
use http::uri::Authority;
10-
use http::{HeaderMap, Method, Request, Uri};
10+
use http::{HeaderMap, Method, Request, Uri, header};
1111

1212
use vey_http::server::HttpAdaptedRequest;
1313

1414
pub trait RequestExt {
1515
fn serialize_for_adapter(&self) -> Vec<u8>;
1616
fn adapt_to(self, other: &HttpAdaptedRequest) -> Self;
1717
fn clone_header(&self) -> Request<()>;
18+
fn expect_100_continue(&self) -> bool;
1819
}
1920

2021
impl<T> RequestExt for Request<T> {
@@ -30,7 +31,7 @@ impl<T> RequestExt for Request<T> {
3031
let _ = write!(buf, "{method} / HTTP/1.1\r\n");
3132
}
3233
for (name, value) in self.headers() {
33-
if matches!(name, &http::header::TE) {
34+
if matches!(name, &header::TE) {
3435
// skip hop-by-hop headers
3536
continue;
3637
}
@@ -39,7 +40,7 @@ impl<T> RequestExt for Request<T> {
3940
buf.put_slice(value.as_bytes());
4041
buf.put_slice(b"\r\n");
4142
}
42-
if !self.headers().contains_key(http::header::HOST)
43+
if !self.headers().contains_key(header::HOST)
4344
&& let Some(host) = uri.host()
4445
{
4546
buf.put_slice(b"Host: ");
@@ -53,19 +54,19 @@ impl<T> RequestExt for Request<T> {
5354
fn adapt_to(self, other: &HttpAdaptedRequest) -> Self {
5455
let mut headers = HeaderMap::from(&other.headers);
5556
// add hop-by-hop headers
56-
if let Some(v) = self.headers().get(http::header::TE) {
57-
headers.insert(http::header::TE, v.into());
57+
if let Some(v) = self.headers().get(header::TE) {
58+
headers.insert(header::TE, v.into());
5859
}
5960
let (mut parts, body) = self.into_parts();
6061
parts.method = other.method.clone();
6162
let mut uri_parts = other.uri.clone().into_parts();
6263
uri_parts.scheme = parts.uri.scheme().cloned();
6364
uri_parts.authority = parts.uri.authority().cloned();
64-
if let Some(host) = headers.remove(http::header::HOST) {
65+
if let Some(host) = headers.remove(header::HOST) {
6566
// we should always remove the Host header to be compatible with Google,
6667
// but let's keep the same as client behaviour here
67-
if parts.headers.contains_key(http::header::HOST) {
68-
headers.insert(http::header::HOST, host.clone());
68+
if parts.headers.contains_key(header::HOST) {
69+
headers.insert(header::HOST, host.clone());
6970
}
7071
if uri_parts.authority.is_none()
7172
&& let Ok(authority) = Authority::from_maybe_shared(host.clone())
@@ -90,4 +91,14 @@ impl<T> RequestExt for Request<T> {
9091
parts.headers = self.headers().clone();
9192
Request::from_parts(parts, ())
9293
}
94+
95+
fn expect_100_continue(&self) -> bool {
96+
for v in self.headers().get_all(header::EXPECT) {
97+
if v.as_bytes() == b"100-continue" {
98+
return true;
99+
}
100+
}
101+
102+
false
103+
}
93104
}

lib/vey-h2/src/lib.rs

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -8,3 +8,6 @@ pub use body::*;
88

99
mod ext;
1010
pub use ext::{RequestExt, ResponseExt};
11+
12+
mod client;
13+
pub use client::H2ResponseHeaderReceiver;

lib/vey-icap-client/src/reqmod/h2/bidirectional.rs

Lines changed: 47 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -8,11 +8,12 @@ use std::time::Duration;
88

99
use bytes::Bytes;
1010
use h2::client::SendRequest;
11-
use http::Request;
11+
use h2::server::SendResponse;
12+
use http::{Request, Response, StatusCode};
1213

1314
use vey_h2::{
14-
H2StreamFromChunkedTransfer, H2StreamFromChunkedTransferError, H2StreamToChunkedTransfer,
15-
H2StreamToChunkedTransferError, RequestExt,
15+
H2ResponseHeaderReceiver, H2StreamFromChunkedTransfer, H2StreamFromChunkedTransferError,
16+
H2StreamToChunkedTransfer, H2StreamToChunkedTransferError, RequestExt,
1617
};
1718
use vey_http::server::HttpAdaptedRequest;
1819
use vey_io_ext::{IdleCheck, LimitedBufReadExt, StreamCopyConfig};
@@ -110,7 +111,9 @@ impl<I: IdleCheck> BidirectionalRecvHttpRequest<'_, I> {
110111
state: &mut ReqmodAdaptationRunState,
111112
mut clt_body_transfer: &mut H2StreamToChunkedTransfer<'_, IcapClientWriter>,
112113
orig_http_request: Request<()>,
113-
mut ups_send_request: SendRequest<Bytes>,
114+
mut ups_send_req: SendRequest<Bytes>,
115+
clt_send_rsp: &mut SendResponse<Bytes>,
116+
mut allow_continue: bool,
114117
) -> Result<ReqmodAdaptationEndState, H2ReqmodAdaptationError> {
115118
let http_req = HttpAdaptedRequest::parse(
116119
self.icap_reader,
@@ -120,7 +123,7 @@ impl<I: IdleCheck> BidirectionalRecvHttpRequest<'_, I> {
120123
.await?;
121124

122125
let final_req = orig_http_request.adapt_to(&http_req);
123-
let (mut ups_recv_rsp, mut ups_send_stream) = ups_send_request
126+
let (rsp_fut, mut ups_send_stream) = ups_send_req
124127
.send_request(final_req, false)
125128
.map_err(H2ReqmodAdaptationError::HttpUpstreamSendHeadFailed)?;
126129
state.mark_ups_send_header();
@@ -133,22 +136,51 @@ impl<I: IdleCheck> BidirectionalRecvHttpRequest<'_, I> {
133136
self.http_trailer_max_size,
134137
);
135138

139+
let mut ups_recv_rsp = H2ResponseHeaderReceiver::new(rsp_fut);
140+
136141
let mut idle_interval = self.idle_checker.interval_timer();
137142
let mut idle_count = 0;
138143

139144
loop {
140145
tokio::select! {
141-
r = &mut ups_recv_rsp => {
142-
return match r {
146+
r = ups_recv_rsp.recv_header() => {
147+
match r {
143148
Ok(ups_rsp) => {
144-
state.mark_ups_recv_header();
145-
if ups_body_transfer.finished() {
146-
self.icap_read_finished = true;
149+
match ups_rsp.status() {
150+
StatusCode::CONTINUE => {
151+
if allow_continue {
152+
clt_send_rsp
153+
.send_informational(ups_rsp)
154+
.map_err(H2ReqmodAdaptationError::HttpClientSendResponseFailed)?;
155+
allow_continue = false;
156+
} else {
157+
return Err(H2ReqmodAdaptationError::InvalidUpstreamContinueResponse);
158+
}
159+
}
160+
StatusCode::EARLY_HINTS => {
161+
clt_send_rsp
162+
.send_informational(ups_rsp)
163+
.map_err(H2ReqmodAdaptationError::HttpClientSendResponseFailed)?;
164+
}
165+
status => {
166+
state.mark_ups_recv_header();
167+
return if let Some(body) = ups_recv_rsp.take_body() {
168+
let (headers, _) = ups_rsp.into_parts();
169+
if ups_body_transfer.finished() {
170+
self.icap_read_finished = true;
171+
}
172+
let ups_rsp = Response::from_parts(headers, body);
173+
Ok(ReqmodAdaptationEndState::AdaptedTransferred(http_req, ups_rsp))
174+
} else {
175+
Err(H2ReqmodAdaptationError::UnsupportedInformationalResponse(
176+
status,
177+
))
178+
};
179+
}
147180
}
148-
Ok(ReqmodAdaptationEndState::AdaptedTransferred(http_req, ups_rsp))
149181
}
150-
Err(e) => Err(H2ReqmodAdaptationError::HttpUpstreamRecvResponseFailed(e)),
151-
};
182+
Err(e) => return Err(H2ReqmodAdaptationError::HttpUpstreamRecvResponseFailed(e)),
183+
}
152184
}
153185
r = &mut clt_body_transfer => {
154186
return match r {
@@ -157,7 +189,7 @@ impl<I: IdleCheck> BidirectionalRecvHttpRequest<'_, I> {
157189
Ok(_) => {
158190
state.mark_ups_send_all();
159191
self.icap_read_finished = true;
160-
let ups_rsp = recv_ups_response_head_after_transfer(ups_recv_rsp, self.http_rsp_head_recv_timeout).await?;
192+
let ups_rsp = recv_ups_response_head_after_transfer(&mut ups_recv_rsp, clt_send_rsp, allow_continue, self.http_rsp_head_recv_timeout).await?;
161193
Ok(ReqmodAdaptationEndState::AdaptedTransferred(http_req, ups_rsp))
162194
}
163195
Err(H2StreamFromChunkedTransferError::ReadError(e)) => Err(H2ReqmodAdaptationError::IcapServerReadFailed(e)),
@@ -176,7 +208,7 @@ impl<I: IdleCheck> BidirectionalRecvHttpRequest<'_, I> {
176208
Ok(_) => {
177209
state.mark_ups_send_all();
178210
self.icap_read_finished = true;
179-
let ups_rsp = recv_ups_response_head_after_transfer(ups_recv_rsp, self.http_rsp_head_recv_timeout).await?;
211+
let ups_rsp = recv_ups_response_head_after_transfer(&mut ups_recv_rsp, clt_send_rsp, allow_continue, self.http_rsp_head_recv_timeout).await?;
180212
Ok(ReqmodAdaptationEndState::AdaptedTransferred(http_req, ups_rsp))
181213
}
182214
Err(H2StreamFromChunkedTransferError::ReadError(e)) => Err(H2ReqmodAdaptationError::IcapServerReadFailed(e)),

lib/vey-icap-client/src/reqmod/h2/error.rs

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@
55

66
use std::io;
77

8+
use http::StatusCode;
89
use thiserror::Error;
910

1011
use vey_h2::H2PreviewError;
@@ -47,6 +48,12 @@ pub enum H2ReqmodAdaptationError {
4748
HttpUpstreamRecvResponseFailed(h2::Error),
4849
#[error("recv response from http upstream timeout")]
4950
HttpUpstreamRecvResponseTimeout,
51+
#[error("invalid http upstream 100-continue response")]
52+
InvalidUpstreamContinueResponse,
53+
#[error("unsupported http upstream informational response {0}")]
54+
UnsupportedInformationalResponse(StatusCode),
55+
#[error("send response to client failed: {0}")]
56+
HttpClientSendResponseFailed(h2::Error),
5057
#[error("internal server error: {0}")]
5158
InternalServerError(&'static str),
5259
#[error("force quit from idle checker: {0:?}")]

lib/vey-icap-client/src/reqmod/h2/forward_body.rs

Lines changed: 21 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@ use std::io::{IoSlice, Write};
88
use bytes::{BufMut, Bytes};
99
use h2::RecvStream;
1010
use h2::client::SendRequest;
11+
use h2::server::SendResponse;
1112
use http::Request;
1213

1314
use vey_h2::{
@@ -42,7 +43,8 @@ impl<I: IdleCheck> H2RequestAdapter<I> {
4243
http_request: Request<()>,
4344
preview_data: H2PreviewData,
4445
clt_body: RecvStream,
45-
ups_send_request: SendRequest<Bytes>,
46+
ups_send_req: SendRequest<Bytes>,
47+
clt_send_rsp: &mut SendResponse<Bytes>,
4648
) -> Result<ReqmodAdaptationEndState, H2ReqmodAdaptationError> {
4749
let http_header = http_request.serialize_for_adapter();
4850
let icap_header = self.build_forward_all_request(http_header.len());
@@ -99,7 +101,8 @@ impl<I: IdleCheck> H2RequestAdapter<I> {
99101
rsp,
100102
header_size,
101103
http_request,
102-
ups_send_request,
104+
ups_send_req,
105+
clt_send_rsp,
103106
)
104107
.await
105108
}
@@ -109,7 +112,8 @@ impl<I: IdleCheck> H2RequestAdapter<I> {
109112
rsp,
110113
header_size,
111114
http_request,
112-
ups_send_request,
115+
ups_send_req,
116+
clt_send_rsp,
113117
)
114118
.await
115119
}
@@ -179,7 +183,8 @@ impl<I: IdleCheck> H2RequestAdapter<I> {
179183
state: &mut ReqmodAdaptationRunState,
180184
http_request: Request<()>,
181185
mut clt_body: RecvStream,
182-
ups_send_request: SendRequest<Bytes>,
186+
ups_send_req: SendRequest<Bytes>,
187+
clt_send_rsp: &mut SendResponse<Bytes>,
183188
) -> Result<ReqmodAdaptationEndState, H2ReqmodAdaptationError> {
184189
let http_header = http_request.serialize_for_adapter();
185190
let icap_header = self.build_forward_all_request(http_header.len());
@@ -242,7 +247,8 @@ impl<I: IdleCheck> H2RequestAdapter<I> {
242247
rsp,
243248
header_size,
244249
http_request,
245-
ups_send_request,
250+
ups_send_req,
251+
clt_send_rsp,
246252
)
247253
.await
248254
}
@@ -254,7 +260,8 @@ impl<I: IdleCheck> H2RequestAdapter<I> {
254260
rsp,
255261
header_size,
256262
http_request,
257-
ups_send_request,
263+
ups_send_req,
264+
clt_send_rsp,
258265
)
259266
.await
260267
} else {
@@ -270,7 +277,14 @@ impl<I: IdleCheck> H2RequestAdapter<I> {
270277
icap_read_finished: false,
271278
};
272279
let r = bidirectional_transfer
273-
.transfer(state, &mut body_transfer, http_request, ups_send_request)
280+
.transfer(
281+
state,
282+
&mut body_transfer,
283+
http_request,
284+
ups_send_req,
285+
clt_send_rsp,
286+
self.allow_continue,
287+
)
274288
.await?;
275289

276290
let icap_read_finished = bidirectional_transfer.icap_read_finished;

0 commit comments

Comments
 (0)