Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 4 additions & 1 deletion lib/vey-icap-client/src/reqmod/h1/bidirectional.rs
Original file line number Diff line number Diff line change
Expand Up @@ -66,7 +66,10 @@ impl<I: IdleCheck> BidirectionalRecvIcapResponse<'_, I> {
}
r = self.icap_reader.fill_wait_data() => {
return match r {
Ok(true) => self.recv_icap_response().await,
Ok(true) => {
state.clt_req_body_size = Some(body_transfer.body_size());
self.recv_icap_response().await
}
Ok(false) => {
state.clt_req_body_size = Some(body_transfer.body_size());
Err(H1ReqmodAdaptationError::IcapServerConnectionClosed)
Expand Down
5 changes: 4 additions & 1 deletion lib/vey-icap-client/src/reqmod/h2/bidirectional.rs
Original file line number Diff line number Diff line change
Expand Up @@ -67,7 +67,10 @@ impl<I: IdleCheck> BidirectionalRecvIcapResponse<'_, I> {
}
r = self.icap_reader.fill_wait_data() => {
return match r {
Ok(true) => self.recv_icap_response().await,
Ok(true) => {
state.clt_req_body_size = Some(body_transfer.received_size());
self.recv_icap_response().await
}
Ok(false) => {
state.clt_req_body_size = Some(body_transfer.received_size());
Err(H2ReqmodAdaptationError::IcapServerConnectionClosed)
Expand Down
5 changes: 4 additions & 1 deletion lib/vey-icap-client/src/reqmod/h2_to_h1/bidirectional.rs
Original file line number Diff line number Diff line change
Expand Up @@ -52,7 +52,10 @@ impl<I: IdleCheck> BidirectionalRecvIcapResponse<'_, I> {
}
r = self.icap_reader.fill_wait_data() => {
return match r {
Ok(true) => self.recv_icap_response().await,
Ok(true) => {
state.record_clt_body_progress(body_transfer);
self.recv_icap_response().await
}
Ok(false) => {
state.record_clt_body_progress(body_transfer);
Err(H2ToH1ReqmodAdaptationError::IcapServerConnectionClosed)
Expand Down
5 changes: 4 additions & 1 deletion lib/vey-icap-client/src/respmod/h1/bidirectional.rs
Original file line number Diff line number Diff line change
Expand Up @@ -66,7 +66,10 @@ impl<I: IdleCheck> BidirectionalRecvIcapResponse<'_, I> {
}
r = self.icap_reader.fill_wait_data() => {
return match r {
Ok(true) => self.recv_icap_response().await,
Ok(true) => {
state.ups_rsp_body_size = Some(body_transfer.body_size());
self.recv_icap_response().await
}
Ok(false) => {
state.ups_rsp_body_size = Some(body_transfer.body_size());
Err(H1RespmodAdaptationError::IcapServerConnectionClosed)
Expand Down
5 changes: 4 additions & 1 deletion lib/vey-icap-client/src/respmod/h1_to_h2/bidirectional.rs
Original file line number Diff line number Diff line change
Expand Up @@ -65,7 +65,10 @@ impl<I: IdleCheck> BidirectionalRecvIcapResponse<'_, I> {
}
r = self.icap_reader.fill_wait_data() => {
return match r {
Ok(true) => self.recv_icap_response().await,
Ok(true) => {
state.ups_rsp_body_size = Some(body_transfer.body_size());
self.recv_icap_response().await
}
Ok(false) => {
state.ups_rsp_body_size = Some(body_transfer.body_size());
Err(H1ToH2RespmodAdaptationError::IcapServerConnectionClosed)
Expand Down
5 changes: 4 additions & 1 deletion lib/vey-icap-client/src/respmod/h2/bidirectional.rs
Original file line number Diff line number Diff line change
Expand Up @@ -68,7 +68,10 @@ impl<I: IdleCheck> BidirectionalRecvIcapResponse<'_, I> {
}
r = self.icap_reader.fill_wait_data() => {
return match r {
Ok(true) => self.recv_icap_response().await,
Ok(true) => {
state.ups_rsp_body_size = Some(body_transfer.received_size());
self.recv_icap_response().await
}
Ok(false) => {
state.ups_rsp_body_size = Some(body_transfer.received_size());
Err(H2RespmodAdaptationError::IcapServerConnectionClosed)
Expand Down
8 changes: 0 additions & 8 deletions vey-proxy/src/auth/user.rs
Original file line number Diff line number Diff line change
Expand Up @@ -801,14 +801,6 @@ impl TenantContext {
pub(crate) fn acquire_request_semaphore(&self) -> Result<GaugeSemaphorePermit, ()> {
self.user.acquire_request_semaphore(&self.forbid_stats)
}

pub(crate) fn check_http_user_agent(
&self,
user_agents: impl IntoIterator<Item = impl AsRef<str>>,
) -> Option<AclAction> {
self.user
.check_http_user_agent(user_agents, &self.forbid_stats)
}
}

#[derive(Clone)]
Expand Down
16 changes: 10 additions & 6 deletions vey-proxy/src/inspect/http/v1/forward/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -443,12 +443,15 @@ impl<'a, SC: ServerConfig> H1ForwardTask<'a, SC> {
clt_w,
&self.ctx.server_config.limited_copy_config(),
);
(&mut copy_to_clt).await.map_err(|e| match e {
StreamCopyError::ReadFailed(e) => ServerTaskError::InternalAdapterError(anyhow!(
"read http error response from adapter failed: {e:?}"
)),
StreamCopyError::WriteFailed(e) => ServerTaskError::ClientTcpWriteFailed(e),
})?;
if let Err(e) = (&mut copy_to_clt).await {
self.http_notes.clt_rsp_body_size = Some(copy_to_clt.reader().body_size());
return Err(match e {
StreamCopyError::ReadFailed(e) => ServerTaskError::InternalAdapterError(
anyhow!("read http error response from adapter failed: {e:?}"),
),
StreamCopyError::WriteFailed(e) => ServerTaskError::ClientTcpWriteFailed(e),
});
}
self.http_notes.clt_rsp_body_size = Some(copy_to_clt.reader().body_size());
recv_body.save_connection().await;
} else {
Expand Down Expand Up @@ -621,6 +624,7 @@ impl<'a, SC: ServerConfig> H1ForwardTask<'a, SC> {
let copy_done = clt_to_ups.finished();
let rsp_head = match rsp_head {
Some(header) => {
record_progress!();
if !clt_body_reader.finished() {
// not all client data read in, drop the client connection
self.should_close = true;
Expand Down
1 change: 1 addition & 0 deletions vey-proxy/src/inspect/http/v2/forward/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -504,6 +504,7 @@ where
match r {
Ok(rsp) => {
if let Some(final_rsp) = self.check_out_final_response(rsp, clt_send_rsp, &mut ups_recv_rsp)? {
record_progress!();
ups_rsp = Some(final_rsp);
break;
}
Expand Down
72 changes: 49 additions & 23 deletions vey-proxy/src/serve/http_expose/task/forward/task.rs
Original file line number Diff line number Diff line change
Expand Up @@ -443,26 +443,19 @@ impl<'a> HttpExposeForwardTask<'a> {
CDR: AsyncRead + Unpin,
CDW: AsyncWrite + Unpin,
{
let tcp_client_misc_opts;

if self.task_notes.check_layered_rate_limit().is_err() {
self.reply_too_many_requests(clt_w).await;
return Err(ServerTaskError::ForbiddenByRule(
ServerTaskForbiddenError::RateLimited,
));
}

if self.task_notes.acquire_site_request_semaphores().is_err() {
self.reply_too_many_requests(clt_w).await;
return Err(ServerTaskError::ForbiddenByRule(
ServerTaskForbiddenError::FullyLoaded,
));
}

if let Some(user_ctx) = self.task_notes.user_ctx() {
if user_ctx.check_rate_limit().is_err() {
self.reply_too_many_requests(clt_w).await;
return Err(ServerTaskError::ForbiddenByRule(
ServerTaskForbiddenError::RateLimited,
));
}
let user_ctx = user_ctx.clone();

if self.task_notes.acquire_user_request_semaphore().is_err() {
if self
.task_notes
.acquire_user_request_semaphore(&user_ctx)
.is_err()
{
self.reply_too_many_requests(clt_w).await;
return Err(ServerTaskError::ForbiddenByRule(
ServerTaskForbiddenError::FullyLoaded,
Expand All @@ -478,13 +471,40 @@ impl<'a> HttpExposeForwardTask<'a> {
) {
self.handle_user_ua_acl_action(action, clt_w).await?;
}
}
if let Some(site_ctx) = self.task_notes.site_ctx() {
if site_ctx.check_rate_limit().is_err() {
self.reply_too_many_requests(clt_w).await;
return Err(ServerTaskError::ForbiddenByRule(
ServerTaskForbiddenError::RateLimited,
));
}
let site_ctx = site_ctx.clone();
if self
.task_notes
.acquire_site_request_semaphores(&site_ctx)
.is_err()
{
self.reply_too_many_requests(clt_w).await;
return Err(ServerTaskError::ForbiddenByRule(
ServerTaskForbiddenError::FullyLoaded,
));
}
}

tcp_client_misc_opts = user_ctx
let tcp_client_misc_opts = if let Some(user_ctx) = self.task_notes.user_ctx() {
user_ctx
.user_config()
.tcp_client_misc_opts(&self.ctx.server_config.tcp_misc_opts);
.tcp_client_misc_opts(&self.ctx.server_config.tcp_misc_opts)
} else if let Some(site_ctx) = self.task_notes.site_ctx()
&& let Some(tenant) = site_ctx.tenant_ctx()
{
tenant
.user_config()
.tcp_client_misc_opts(&self.ctx.server_config.tcp_misc_opts)
} else {
tcp_client_misc_opts = Cow::Borrowed(&self.ctx.server_config.tcp_misc_opts);
}
Cow::Borrowed(&self.ctx.server_config.tcp_misc_opts)
};

// set client side socket options
self.ctx
Expand Down Expand Up @@ -910,7 +930,12 @@ impl<'a> HttpExposeForwardTask<'a> {
.map_err(ServerTaskError::UpstreamWriteFailed)?;
self.http_notes.mark_req_send_hdr();
self.http_notes.mark_req_send_all();
self.http_notes.ups_req_body_size = Some(body.len() as u64);
// Chunked bodies are buffered on the wire, while clt_req_body_size is the
// decoded payload. A fully read body was already counted that way.
self.http_notes.ups_req_body_size = self
.http_notes
.clt_req_body_size
.or(Some(body.len() as u64));

match tokio::time::timeout(
self.rsp_hdr_recv_timeout(),
Expand Down Expand Up @@ -1140,6 +1165,7 @@ impl<'a> HttpExposeForwardTask<'a> {
let copy_done = clt_to_ups.finished();
let mut rsp_header = match rsp_header {
Some(header) => {
record_progress!();
if !clt_body_reader.finished() {
// not all client data read in, drop the client connection
self.should_close = true;
Expand Down
22 changes: 16 additions & 6 deletions vey-proxy/src/serve/http_expose/task/untrusted/task.rs
Original file line number Diff line number Diff line change
Expand Up @@ -68,12 +68,22 @@ impl<'a> HttpExposeUntrustedTask<'a> {
CDR: AsyncRead + Unpin,
CDW: AsyncWrite + Unpin,
{
if self.task_notes.check_layered_rate_limit().is_err()
|| self.task_notes.acquire_site_request_semaphores().is_err()
{
self.should_close = true;
self.reply_too_many_requests(clt_w).await;
return;
if let Some(site_ctx) = self.task_notes.site_ctx() {
if site_ctx.check_rate_limit().is_err() {
self.should_close = true;
self.reply_too_many_requests(clt_w).await;
return;
}
let site_ctx = site_ctx.clone();
if self
.task_notes
.acquire_site_request_semaphores(&site_ctx)
.is_err()
{
self.should_close = true;
self.reply_too_many_requests(clt_w).await;
return;
}
}

let site_io = self.site_ctx.fetch_traffic_stats(
Expand Down
Loading
Loading