Skip to content

Commit 9ca8fe3

Browse files
committed
pd-edge: perf: decouple proxy forward wait from vm path
1 parent f0a20b6 commit 9ca8fe3

6 files changed

Lines changed: 151 additions & 57 deletions

File tree

pd-edge/examples/http_proxy_perf_framework.rs

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -38,7 +38,6 @@ use tokio::{
3838
use tokio_rustls::rustls::pki_types::{CertificateDer, PrivateKeyDer, PrivatePkcs8KeyDer};
3939
#[cfg(all(feature = "http2", feature = "tls"))]
4040
use tokio_rustls::{TlsAcceptor, rustls::ServerConfig};
41-
#[cfg(feature = "http3")]
4241
use url::Url;
4342
use vm::{compile_source, encode_program, validate_program};
4443

pd-edge/src/abi_impl/http/state.rs

Lines changed: 117 additions & 44 deletions
Original file line numberDiff line numberDiff line change
@@ -1867,6 +1867,7 @@ pub(crate) struct DownstreamState {
18671867
pub(crate) post_response_plan: Option<DownstreamPostResponsePlan>,
18681868
pub(crate) native_default_upstream_http_forward: bool,
18691869
native_default_upstream_forward_response: Option<NativeDefaultUpstreamForwardResponse>,
1870+
native_default_upstream_forward_task: Option<NativeDefaultUpstreamForwardTask>,
18701871
inline_http_response_sender: Option<InlineDownstreamHttpResponseSender>,
18711872
}
18721873

@@ -1880,6 +1881,7 @@ impl DownstreamState {
18801881
post_response_plan: None,
18811882
native_default_upstream_http_forward: false,
18821883
native_default_upstream_forward_response: None,
1884+
native_default_upstream_forward_task: None,
18831885
inline_http_response_sender: None,
18841886
}
18851887
}
@@ -1918,6 +1920,7 @@ impl DownstreamState {
19181920
post_response_plan: None,
19191921
native_default_upstream_http_forward: false,
19201922
native_default_upstream_forward_response: None,
1923+
native_default_upstream_forward_task: None,
19211924
inline_http_response_sender: None,
19221925
}
19231926
}
@@ -2483,16 +2486,19 @@ impl ProxyVmContext {
24832486
}
24842487

24852488
pub(crate) fn clear_native_default_upstream_http_forward(&self) {
2486-
self.lock_downstream().native_default_upstream_http_forward = false;
2489+
let mut downstream = self.lock_downstream();
2490+
downstream.native_default_upstream_http_forward = false;
2491+
downstream.native_default_upstream_forward_response = None;
2492+
if let Some(task) = downstream.native_default_upstream_forward_task.take() {
2493+
task.abort();
2494+
}
24872495
}
24882496

2489-
fn store_native_default_upstream_forward_response(
2490-
&self,
2491-
response: NativeDefaultUpstreamForwardResponse,
2492-
) {
2497+
fn store_native_default_upstream_forward_task(&self, task: NativeDefaultUpstreamForwardTask) {
24932498
let mut downstream = self.lock_downstream();
24942499
downstream.native_default_upstream_http_forward = true;
2495-
downstream.native_default_upstream_forward_response = Some(response);
2500+
downstream.native_default_upstream_forward_response = None;
2501+
downstream.native_default_upstream_forward_task = Some(task);
24962502
}
24972503

24982504
fn take_native_default_upstream_forward_response(
@@ -2509,6 +2515,18 @@ impl ProxyVmContext {
25092515
.is_some()
25102516
}
25112517

2518+
fn take_native_default_upstream_forward_task(&self) -> Option<NativeDefaultUpstreamForwardTask> {
2519+
self.lock_downstream()
2520+
.native_default_upstream_forward_task
2521+
.take()
2522+
}
2523+
2524+
fn native_default_upstream_forward_task_pending(&self) -> bool {
2525+
self.lock_downstream()
2526+
.native_default_upstream_forward_task
2527+
.is_some()
2528+
}
2529+
25122530
fn native_default_upstream_forward_latency_ms(&self) -> Option<u64> {
25132531
self.lock_downstream()
25142532
.native_default_upstream_forward_response
@@ -3303,6 +3321,10 @@ enum NativeDefaultUpstreamForwardBody {
33033321
},
33043322
}
33053323

3324+
type NativeDefaultUpstreamForwardTask = tokio::task::JoinHandle<
3325+
Result<NativeDefaultUpstreamForwardResponse, UpstreamResponseStartError>,
3326+
>;
3327+
33063328
#[derive(Debug)]
33073329
struct NativeDefaultUpstreamForwardResponse {
33083330
status: u16,
@@ -4285,10 +4307,8 @@ async fn start_default_upstream_plain_http1_fast_path(
42854307

42864308
fn materialize_native_default_upstream_forward_response(
42874309
context: &SharedProxyVmContext,
4310+
response: NativeDefaultUpstreamForwardResponse,
42884311
) -> Result<Option<HttpUpstreamResponseSnapshot>, UpstreamResponseStartError> {
4289-
let Some(response) = context.take_native_default_upstream_forward_response() else {
4290-
return Ok(None);
4291-
};
42924312
let NativeDefaultUpstreamForwardResponse {
42934313
status,
42944314
headers,
@@ -4342,6 +4362,53 @@ fn materialize_native_default_upstream_forward_response(
43424362
Ok(Some(snapshot))
43434363
}
43444364

4365+
async fn take_or_await_native_default_upstream_forward_response(
4366+
context: &SharedProxyVmContext,
4367+
) -> Result<Option<NativeDefaultUpstreamForwardResponse>, UpstreamResponseStartError> {
4368+
if let Some(response) = context.take_native_default_upstream_forward_response() {
4369+
return Ok(Some(response));
4370+
}
4371+
4372+
let Some(task) = context.take_native_default_upstream_forward_task() else {
4373+
return Ok(None);
4374+
};
4375+
match task.await {
4376+
Ok(Ok(response)) => Ok(Some(response)),
4377+
Ok(Err(err)) => Err(err),
4378+
Err(err) => Err(UpstreamResponseStartError::Protocol(format!(
4379+
"native default upstream forward task failed: {err}"
4380+
))),
4381+
}
4382+
}
4383+
4384+
async fn try_materialize_ready_or_pending_native_default_upstream_forward_response(
4385+
context: &SharedProxyVmContext,
4386+
) -> Result<Option<HttpUpstreamResponseSnapshot>, UpstreamResponseStartError> {
4387+
let Some(response) = take_or_await_native_default_upstream_forward_response(context).await?
4388+
else {
4389+
return Ok(None);
4390+
};
4391+
materialize_native_default_upstream_forward_response(context, response)
4392+
}
4393+
4394+
async fn try_resolve_ready_or_pending_native_default_upstream_forward_response(
4395+
context: &SharedProxyVmContext,
4396+
response_headers: HeaderMap,
4397+
response_status: Option<u16>,
4398+
) -> Result<Option<ResolvedHttpGraphResponse>, UpstreamResponseStartError> {
4399+
let Some(response) = take_or_await_native_default_upstream_forward_response(context).await?
4400+
else {
4401+
return Ok(None);
4402+
};
4403+
let upstream_latency_ms = response.upstream_latency_ms;
4404+
Ok(Some(ResolvedHttpGraphResponse {
4405+
response: response_from_started_upstream_response(response, response_headers, response_status)
4406+
.await,
4407+
upstream_latency_ms,
4408+
post_response_plan: None,
4409+
}))
4410+
}
4411+
43454412
async fn forward_native_default_upstream_http_via_sender_pool(
43464413
context: &SharedProxyVmContext,
43474414
request: &DefaultUpstreamRequestSnapshot,
@@ -4487,10 +4554,23 @@ async fn forward_native_default_upstream_http_via_sender_pool(
44874554
})
44884555
}
44894556

4557+
fn schedule_native_default_upstream_http_forward_response(
4558+
context: &SharedProxyVmContext,
4559+
request: DefaultUpstreamRequestSnapshot,
4560+
) {
4561+
let task_context = context.clone();
4562+
let task = tokio::spawn(async move {
4563+
forward_native_default_upstream_http_via_sender_pool(&task_context, &request).await
4564+
});
4565+
context.store_native_default_upstream_forward_task(task);
4566+
}
4567+
44904568
pub(crate) async fn start_native_default_upstream_http_forward_response(
44914569
context: &SharedProxyVmContext,
44924570
) -> Result<bool, VmError> {
4493-
if context.native_default_upstream_forward_response_ready() {
4571+
if context.native_default_upstream_forward_response_ready()
4572+
|| context.native_default_upstream_forward_task_pending()
4573+
{
44944574
return Ok(true);
44954575
}
44964576

@@ -4521,13 +4601,8 @@ pub(crate) async fn start_native_default_upstream_http_forward_response(
45214601
return Ok(false);
45224602
}
45234603

4524-
match forward_native_default_upstream_http_via_sender_pool(context, &request).await {
4525-
Ok(response) => {
4526-
context.store_native_default_upstream_forward_response(response);
4527-
Ok(true)
4528-
}
4529-
Err(_) => Ok(false),
4530-
}
4604+
schedule_native_default_upstream_http_forward_response(context, request);
4605+
Ok(true)
45314606
}
45324607

45334608
async fn try_resolve_native_default_upstream_http_forward_response(
@@ -5096,10 +5171,13 @@ async fn start_outbound_exchange_response(
50965171
}
50975172
}
50985173

5099-
if handle == DEFAULT_UPSTREAM_EXCHANGE_HANDLE
5100-
&& let Some(snapshot) = materialize_native_default_upstream_forward_response(context)?
5101-
{
5102-
return Ok(snapshot);
5174+
if handle == DEFAULT_UPSTREAM_EXCHANGE_HANDLE {
5175+
match try_materialize_ready_or_pending_native_default_upstream_forward_response(context)
5176+
.await
5177+
{
5178+
Ok(Some(snapshot)) => return Ok(snapshot),
5179+
Ok(None) | Err(_) => {}
5180+
}
51035181
}
51045182

51055183
if handle == DEFAULT_UPSTREAM_EXCHANGE_HANDLE
@@ -5690,7 +5768,6 @@ pub(crate) async fn resolve_http_graph_response(
56905768
default_upstream_websocket_mode,
56915769
upstream_response,
56925770
native_default_upstream_http_forward,
5693-
native_default_upstream_forward_response_ready,
56945771
) = {
56955772
let downstream = context.lock_downstream();
56965773
let exchanges = context.lock_exchanges();
@@ -5710,9 +5787,6 @@ pub(crate) async fn resolve_http_graph_response(
57105787
HttpUpstreamResponseNode::NotStarted => None,
57115788
},
57125789
downstream.native_default_upstream_http_forward,
5713-
downstream
5714-
.native_default_upstream_forward_response
5715-
.is_some(),
57165790
)
57175791
};
57185792

@@ -5763,24 +5837,20 @@ pub(crate) async fn resolve_http_graph_response(
57635837
};
57645838
}
57655839

5766-
if native_default_upstream_forward_response_ready {
5767-
if let Some(response) = context.take_native_default_upstream_forward_response() {
5768-
let upstream_latency_ms = response.upstream_latency_ms;
5769-
context.clear_native_default_upstream_http_forward();
5770-
return ResolvedHttpGraphResponse {
5771-
response: response_from_started_upstream_response(
5772-
response,
5773-
response_headers,
5774-
response_status,
5775-
)
5776-
.await,
5777-
upstream_latency_ms,
5778-
post_response_plan: None,
5779-
};
5780-
}
5781-
}
5782-
57835840
if native_default_upstream_http_forward && upstream_response.is_none() {
5841+
match try_resolve_ready_or_pending_native_default_upstream_forward_response(
5842+
context,
5843+
response_headers.clone(),
5844+
response_status,
5845+
)
5846+
.await
5847+
{
5848+
Ok(Some(resolved)) => {
5849+
context.clear_native_default_upstream_http_forward();
5850+
return resolved;
5851+
}
5852+
Ok(None) | Err(_) => {}
5853+
}
57845854
match try_resolve_native_default_upstream_http_forward_response(
57855855
context,
57865856
response_headers.clone(),
@@ -6029,10 +6099,13 @@ mod tests {
60296099
resolve_http_graph_response, response_from_upstream_snapshot,
60306100
};
60316101
use crate::abi_impl::RateLimiterStore;
6102+
use crate::abi_impl::http2::{Http2DownstreamStreamAttachment, Http2StreamRef};
6103+
#[cfg(feature = "http2")]
60326104
use crate::abi_impl::http2::{
6033-
Http2DownstreamStreamAttachment, Http2SendRequest, Http2StreamRef, Http2UpstreamMode,
6034-
new_shared_http_upstream_sessions, send_request, total_active_streams,
6105+
Http2SendRequest, Http2UpstreamMode, new_shared_http_upstream_sessions, send_request,
6106+
total_active_streams,
60356107
};
6108+
#[cfg(feature = "http2")]
60366109
use crate::abi_impl::transport::TlsFlowState;
60376110

60386111
fn test_context() -> SharedProxyVmContext {

pd-edge/src/abi_impl/http2/mod.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -12,7 +12,7 @@ pub(crate) use self::model::{
1212
Http2StreamRef, Http2UpstreamMode, configure_reqwest_builder, response_version_label,
1313
select_upstream_mode, supports_response_version,
1414
};
15-
#[cfg(test)]
15+
#[cfg(all(test, feature = "http2"))]
1616
pub(crate) use self::upstream::total_active_streams;
1717
#[cfg(feature = "http2")]
1818
pub(crate) use self::upstream::{Http2RequestError, Http2SendRequest, send_request};

pd-edge/src/runtime/http_plane/proxy_path.rs

Lines changed: 1 addition & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1175,18 +1175,15 @@ fn text_response(status: StatusCode, text: &str) -> Response<Body> {
11751175

11761176
#[cfg(test)]
11771177
mod tests {
1178+
use super::*;
11781179
#[cfg(feature = "http2")]
11791180
use axum::http::Version;
11801181
#[cfg(feature = "http2")]
11811182
use http_body_util::{BodyExt, Full};
1182-
#[cfg(feature = "http2")]
11831183
use vm::{compile_source, encode_program};
11841184

1185-
#[cfg(feature = "http2")]
1186-
use super::*;
11871185
#[cfg(feature = "http2")]
11881186
use crate::abi_impl::Http2SessionFrontier;
1189-
#[cfg(feature = "http2")]
11901187
use crate::runtime::{VmExecutionConfig, VmExecutionMode, apply_program_from_bytes};
11911188

11921189
#[cfg(feature = "http2")]

pd-vm/src/vm/jit/runtime.rs

Lines changed: 6 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -282,13 +282,12 @@ impl Vm {
282282
let Some(trace) = self.jit.trace_clone(trace_id) else {
283283
return Ok(ExecOutcome::Continue);
284284
};
285-
let step_indices_by_ip = trace
286-
.step_ips
287-
.iter()
288-
.copied()
289-
.enumerate()
290-
.map(|(step_index, step_ip)| (step_ip, step_index))
291-
.collect::<HashMap<_, _>>();
285+
let step_indices_by_ip = self.jit.trace_step_index_lookup(trace_id).ok_or_else(|| {
286+
VmError::JitNative(format!(
287+
"trace {} is missing its step-index lookup",
288+
trace_id
289+
))
290+
})?;
292291
let mut step_index = 0usize;
293292
while step_index < trace.steps.len() {
294293
self.ip = trace

0 commit comments

Comments
 (0)