mod context;
mod continuation;
mod sanitize;
mod streaming;
#[cfg(test)]
#[expect(clippy::allow_attributes, reason = "blanket test suppressions")]
#[allow(
clippy::unwrap_used,
clippy::expect_used,
clippy::panic,
clippy::too_many_lines,
reason = "tests"
)]
mod tests;
mod transport;
mod types;
use std::{
sync::{
Arc,
atomic::{AtomicBool, Ordering},
},
time::{Duration, Instant},
};
use bytes::Bytes;
use http::HeaderMap;
use praxis_core::{
connectivity::Upstream,
subrequest::{FrameworkHeaders, StreamLimits, SubRequestError},
};
use tracing::{Instrument as _, warn};
pub use self::types::{CalloutOutcome, CalloutResponse, StagedUpstream, StagedUpstreamFallback, SubrequestRuntime};
use self::{
context::{SubrequestRuntimeResources, build_sub_filter_context},
sanitize::{
apply_pre_read_header_mutations, apply_request_header_mutations, body_exceeds_limit, ensure_destination_host,
response_body_overflow_limit, sanitize_subrequest_headers, sanitize_subresponse_headers, set_authority_host,
streaming_transport_limit, strip_reserved_headers, subresponse_from_rejection,
},
transport::{build_peer, classify_transport_failure, stream_termination_cause},
};
pub(crate) use self::{
continuation::FilteredSubrequestContinuation,
streaming::{CalloutStreamingBody, FilteredStreamingBody},
types::{
FilteredSubrequestError, FilteredSubrequestInput, OpenedResponse, OpenedSubrequest, RawResponse,
ResponseOrigin, ResponseTooLargeInfo, SubrequestOutcome, TransportFailure,
},
};
#[cfg(feature = "iterative-request-router")]
pub(crate) use self::{continuation::SubrequestCompletion, sanitize::normalize_response_status};
#[cfg(feature = "chain-binding")]
use crate::credentials::{PendingCredentials, ResolvedDestination};
#[cfg(feature = "bound-upstream-request-body")]
use crate::extensions::BoundRequestBodyRewrite;
#[cfg(feature = "upstream-binding")]
use crate::extensions::{BoundUpstream, BoundUpstreamFrozen};
use crate::{
FilterAction, FilterError, FilterPipeline, StreamTermination, StreamTerminationCause, SubRequest,
SubRequestResponseMode, SubResponse,
actions::Rejection,
context::PendingStreamChunks,
extensions::{RequestExtensions, SelectedClusterApplication},
results::RetainedFilterResults,
};
#[cfg(feature = "upstream-binding")]
struct ParentUpstreamState {
binding: Option<BoundUpstream>,
frozen: bool,
#[cfg(feature = "bound-upstream-request-body")]
body_rewrite: Option<BoundRequestBodyRewrite>,
}
#[cfg(feature = "upstream-binding")]
#[derive(Default)]
struct ParentUpstreamStates(Vec<ParentUpstreamState>);
#[cfg(feature = "upstream-binding")]
fn enter_nested_upstream_scope(extensions: &mut RequestExtensions, inherits_binding: bool) {
let (binding, frozen) = if inherits_binding {
(
extensions.get::<BoundUpstream>().cloned(),
extensions.get::<BoundUpstreamFrozen>().is_some(),
)
} else {
(
extensions.remove::<BoundUpstream>(),
extensions.remove::<BoundUpstreamFrozen>().is_some(),
)
};
let state = ParentUpstreamState {
binding,
frozen,
#[cfg(feature = "bound-upstream-request-body")]
body_rewrite: extensions.remove::<BoundRequestBodyRewrite>(),
};
let mut states = extensions.remove::<ParentUpstreamStates>().unwrap_or_default();
states.0.push(state);
extensions.insert(states);
extensions.remove::<SelectedClusterApplication>();
}
#[cfg(not(feature = "upstream-binding"))]
fn enter_nested_upstream_scope(extensions: &mut RequestExtensions, _inherits_binding: bool) {
extensions.remove::<SelectedClusterApplication>();
}
#[cfg(feature = "upstream-binding")]
fn restore_parent_upstream_scope(extensions: &mut RequestExtensions) {
extensions.remove::<SelectedClusterApplication>();
let Some(mut states) = extensions.remove::<ParentUpstreamStates>() else {
return;
};
let Some(state) = states.0.pop() else {
return;
};
extensions.remove::<BoundUpstream>();
extensions.remove::<BoundUpstreamFrozen>();
if let Some(binding) = state.binding {
extensions.insert(binding);
}
if state.frozen {
extensions.insert(BoundUpstreamFrozen);
}
#[cfg(feature = "bound-upstream-request-body")]
{
extensions.remove::<BoundRequestBodyRewrite>();
if let Some(body_rewrite) = state.body_rewrite {
extensions.insert(body_rewrite);
}
}
if !states.0.is_empty() {
extensions.insert(states);
}
}
#[cfg(not(feature = "upstream-binding"))]
fn restore_parent_upstream_scope(extensions: &mut RequestExtensions) {
extensions.remove::<SelectedClusterApplication>();
}
pub(crate) const STREAMING_IDLE_TIMEOUT: Duration = Duration::from_secs(30);
pub(crate) trait RetainedStateAccounting {
fn exceeds_limit(&self, extensions: &RequestExtensions) -> bool;
}
impl FilteredSubrequestError {
fn new(error: FilterError, extensions: RequestExtensions) -> Self {
Self {
error,
extensions,
too_large: None,
}
}
fn capture(error: FilterError, ctx: &mut crate::HttpFilterContext<'_>) -> Self {
ctx.extensions.remove::<PendingStreamChunks>();
ctx.extensions.remove::<RetainedFilterResults>();
ctx.extensions.remove::<StreamTermination>();
Self::new(error, std::mem::take(&mut ctx.extensions))
}
#[must_use]
fn with_too_large(mut self, info: ResponseTooLargeInfo) -> Self {
self.too_large = Some(info);
self
}
pub(crate) fn too_large(&self) -> Option<ResponseTooLargeInfo> {
self.too_large
}
pub(crate) fn into_parts(mut self) -> (FilterError, RequestExtensions) {
restore_parent_upstream_scope(&mut self.extensions);
(self.error, self.extensions)
}
}
struct NoRetainedState;
impl RetainedStateAccounting for NoRetainedState {
fn exceeds_limit(&self, _extensions: &RequestExtensions) -> bool {
false
}
}
fn repin_staged_upstream(ctx: &mut crate::HttpFilterContext<'_>, pinned_upstream: Option<&Upstream>) {
if let Some(upstream) = pinned_upstream {
ctx.upstream = Some(upstream.clone());
ctx.extensions.remove::<SelectedClusterApplication>();
}
}
pub struct FilteredSubrequestExecutor {
accounting: Box<dyn RetainedStateAccounting + Send + Sync>,
client: praxis_core::subrequest::SubRequestClient,
depth: u8,
downstream: SubrequestRuntime,
max_response_bytes: usize,
max_state_bytes: usize,
step_timeout: Duration,
}
impl FilteredSubrequestExecutor {
#[expect(
clippy::too_many_arguments,
reason = "executor owns explicit subrequest limits and resources"
)]
pub(crate) fn new(
accounting: Box<dyn RetainedStateAccounting + Send + Sync>,
client: praxis_core::subrequest::SubRequestClient,
depth: u8,
downstream: SubrequestRuntime,
max_response_bytes: usize,
max_state_bytes: usize,
step_timeout: Duration,
) -> Self {
Self {
accounting,
client,
depth,
downstream,
max_response_bytes,
max_state_bytes,
step_timeout,
}
}
#[must_use]
pub fn for_callout(
client: praxis_core::subrequest::SubRequestClient,
downstream: SubrequestRuntime,
depth: u8,
max_response_bytes: usize,
step_timeout: Duration,
) -> Self {
Self::new(
Box::new(NoRetainedState),
client,
depth,
downstream,
max_response_bytes,
max_response_bytes,
step_timeout,
)
}
#[expect(clippy::large_futures, reason = "delegates to the full step future")]
#[expect(
clippy::large_stack_frames,
reason = "delegates to execute, which reconstructs a full filter context"
)]
pub async fn run(
&self,
pipeline: &Arc<FilterPipeline>,
request: &SubRequest,
extensions: RequestExtensions,
deadline: Instant,
) -> Result<CalloutResponse, FilterError> {
let input = FilteredSubrequestInput::callout(pipeline, request, deadline, extensions);
let opened = self.execute(input).await.map_err(|error| error.into_parts().0)?;
Ok(self.callout_response_from(opened))
}
#[expect(clippy::large_futures, reason = "delegates to the full step future")]
#[expect(
clippy::large_stack_frames,
reason = "delegates to execute, which reconstructs a full filter context"
)]
pub async fn run_classified(
&self,
pipeline: &Arc<FilterPipeline>,
request: &SubRequest,
extensions: RequestExtensions,
deadline: Instant,
) -> Result<CalloutOutcome, FilterError> {
let input = FilteredSubrequestInput::callout(pipeline, request, deadline, extensions);
match self.execute(input).await {
Ok(opened) => {
if let OpenedResponse::Complete(outcome) = &opened.kind
&& let Some(TransportFailure::ResponseTooLarge { actual, limit }) = outcome.transport_error
{
return Ok(CalloutOutcome::ResponseTooLarge {
actual: Some(actual),
limit,
});
}
Ok(CalloutOutcome::Response(self.callout_response_from(opened)))
},
Err(error) => match error.too_large() {
Some(info) => Ok(CalloutOutcome::ResponseTooLarge {
actual: info.actual,
limit: info.limit,
}),
None => Err(error.into_parts().0),
},
}
}
fn callout_response_from(&self, opened: OpenedSubrequest) -> CalloutResponse {
let OpenedSubrequest { continuation, kind } = opened;
match kind {
OpenedResponse::Complete(outcome) => CalloutResponse::Buffered(outcome.response),
OpenedResponse::Streaming { body, outcome } => CalloutResponse::Streaming {
response: outcome.response,
body: Box::new(CalloutStreamingBody::new(
FilteredStreamingBody::new(body, continuation),
self.max_response_bytes,
)),
},
}
}
#[expect(
clippy::too_many_lines,
reason = "one sub-request owns the complete filter and transport lifecycle"
)]
#[expect(clippy::large_futures, reason = "step execution spans filter and transport futures")]
#[expect(
clippy::large_stack_frames,
reason = "step execution reconstructs a full filter context"
)]
pub(crate) async fn execute(
&self,
input: FilteredSubrequestInput<'_>,
) -> Result<OpenedSubrequest, FilteredSubrequestError> {
let FilteredSubrequestInput {
pipeline,
request: current_request,
label,
iteration,
deadline,
mut extensions,
inherits_binding,
} = input;
enter_nested_upstream_scope(&mut extensions, inherits_binding);
let remaining = deadline
.checked_duration_since(Instant::now())
.unwrap_or(Duration::ZERO);
if remaining.is_zero() {
return Err(FilteredSubrequestError::new(
"filtered_subrequest: overall deadline exceeded".to_owned().into(),
extensions,
));
}
let mut sub_headers = current_request.headers.clone();
strip_reserved_headers(&mut sub_headers);
let sub_req = crate::Request {
method: current_request.method.clone(),
uri: current_request.uri.clone(),
headers: sub_headers.clone(),
};
let mut routed_req = sub_req.clone();
let mut response_header = crate::Response {
headers: HeaderMap::new(),
status: http::StatusCode::OK,
};
let resources = SubrequestRuntimeResources {
client_addr: self.downstream.client_addr,
downstream_tls: self.downstream.downstream_tls,
health_registry: pipeline.health_registry(),
id_generator: pipeline.id_generator(),
kv_stores: pipeline.kv_stores(),
session_stores: pipeline.session_stores(),
peer_identity: self.downstream.peer_identity.as_ref(),
request_start: self.downstream.request_start,
subrequest_client: Some(&self.client),
time_source: pipeline.time_source(),
};
let mut filter_ctx = build_sub_filter_context(pipeline, &sub_req, resources);
filter_ctx.extensions = std::mem::take(&mut extensions);
filter_ctx.extensions.insert(RetainedFilterResults::default());
filter_ctx.enable_stream_chunk_emission(self.max_state_bytes);
let pinned_upstream = filter_ctx.extensions.remove::<StagedUpstream>().map(|staged| staged.0);
if let Some(upstream) = &pinned_upstream {
filter_ctx.upstream = Some(upstream.clone());
}
let fallback_addresses = filter_ctx
.extensions
.remove::<StagedUpstreamFallback>()
.map(|fallback| fallback.0)
.unwrap_or_default();
let step_budget = remaining.min(self.step_timeout);
let step_started = Instant::now();
let step_deadline = step_started.checked_add(step_budget).unwrap_or(deadline);
let in_transport = Arc::new(AtomicBool::new(false));
let in_transport_inner = Arc::clone(&in_transport);
let step_span = tracing::info_span!("filtered_subrequest", step = label, iteration = iteration);
let timed: Result<Result<RawResponse, FilterError>, tokio::time::error::Elapsed> =
tokio::time::timeout(step_budget, async {
let mut request_body = Some(current_request.body.clone());
if body_exceeds_limit(
pipeline.body_capabilities().request_body_mode,
request_body.as_ref().map_or(0, Bytes::len),
) {
return Ok(RawResponse::Rejected(Rejection::status(413)));
}
let pre_read_body = matches!(
pipeline.body_capabilities().request_body_mode,
crate::BodyMode::StreamBuffer { .. }
);
if pre_read_body {
let action = pipeline
.execute_http_request_body(&mut filter_ctx, &mut request_body, true)
.await?;
if let FilterAction::Reject(rejection) = action {
return Ok(RawResponse::Rejected(rejection));
}
if self.accounting.exceeds_limit(&filter_ctx.extensions) {
return Ok(RawResponse::Rejected(Rejection::status(413)));
}
apply_pre_read_header_mutations(&mut routed_req.headers, &filter_ctx);
filter_ctx.extra_request_headers.clear();
filter_ctx.request_headers_to_remove.clear();
filter_ctx.request_headers_to_set.clear();
filter_ctx.pre_read_mutations.clear();
sub_headers.clone_from(&routed_req.headers);
filter_ctx.request = &routed_req;
filter_ctx.buffered_request_body.clone_from(&request_body);
}
let action = pipeline.execute_http_request(&mut filter_ctx).await?;
if let FilterAction::Reject(rejection) = action {
return Ok(RawResponse::Rejected(rejection));
}
if self.accounting.exceeds_limit(&filter_ctx.extensions) {
return Ok(RawResponse::Rejected(Rejection::status(413)));
}
if let Some(rewritten) = filter_ctx.take_bound_request_body_rewrite() {
request_body = Some(rewritten);
}
if !pre_read_body {
let action = pipeline
.execute_http_request_body(&mut filter_ctx, &mut request_body, true)
.await?;
if let FilterAction::Reject(rejection) = action {
return Ok(RawResponse::Rejected(rejection));
}
if self.accounting.exceeds_limit(&filter_ctx.extensions) {
return Ok(RawResponse::Rejected(Rejection::status(413)));
}
}
repin_staged_upstream(&mut filter_ctx, pinned_upstream.as_ref());
if pipeline.body_capabilities().needs_selected_upstream_request_body
&& filter_ctx.upstream.is_some()
{
let action = pipeline
.execute_http_selected_upstream_request_body(&mut filter_ctx, &mut request_body)
.await?;
if let FilterAction::Reject(rejection) = action {
return Ok(RawResponse::Rejected(rejection));
}
if pipeline.body_capabilities().any_selected_upstream_request_body_writer
&& request_body.as_ref().map_or(0, Bytes::len)
> pipeline.selected_upstream_request_body_limit()
{
return Ok(RawResponse::Rejected(Rejection::status(413)));
}
if self.accounting.exceeds_limit(&filter_ctx.extensions) {
return Ok(RawResponse::Rejected(Rejection::status(413)));
}
repin_staged_upstream(&mut filter_ctx, pinned_upstream.as_ref());
}
let upstream = filter_ctx.upstream.as_ref().ok_or_else(|| -> FilterError {
format!("filtered_subrequest: step '{label}' did not resolve an upstream").into()
})?;
let destination_authority = Arc::clone(&upstream.address);
let authority_override: Option<Arc<str>> = upstream
.authority
.as_ref()
.and_then(|value| value.to_str().ok())
.map(Arc::from);
#[cfg(feature = "chain-binding")]
let logical_authority: Arc<str> = authority_override
.clone()
.unwrap_or_else(|| Arc::clone(&destination_authority));
in_transport_inner.store(true, Ordering::Release);
let peers = if fallback_addresses.len() > 1 {
let mut built = Vec::with_capacity(fallback_addresses.len());
for address in &fallback_addresses {
let mut per_address = upstream.clone();
per_address.address = Arc::from(address.to_string().as_str());
built.push(build_peer(&per_address, pipeline.allow_private_upstreams()).await);
}
built
} else {
vec![build_peer(upstream, pipeline.allow_private_upstreams()).await]
};
apply_request_header_mutations(&mut sub_headers, &filter_ctx);
match authority_override.as_deref() {
Some(authority) => set_authority_host(&mut sub_headers, authority)?,
None => ensure_destination_host(&mut sub_headers, &destination_authority)?,
}
sanitize_subrequest_headers(&mut sub_headers);
#[cfg(feature = "chain-binding")]
if let Some(pending) = filter_ctx.extensions.remove::<PendingCredentials>() {
let destination = ResolvedDestination {
authority: &logical_authority,
transport: &destination_authority,
};
let injected = pending.inject_authorized(&destination, &mut sub_headers);
if injected > 0 {
set_authority_host(&mut sub_headers, &logical_authority)?;
tracing::debug!(
authority = %logical_authority,
transport = %destination_authority,
injected,
"materialized destination-bound credentials"
);
}
}
let request = SubRequest {
method: current_request.method.clone(),
uri: filter_ctx.rewritten_path.as_ref().map_or_else(
|| current_request.uri.clone(),
|path| http::Uri::try_from(path.as_str()).unwrap_or_else(|_| current_request.uri.clone()),
),
headers: sub_headers,
body: request_body.unwrap_or_default(),
};
let mut framework_headers = FrameworkHeaders::new();
framework_headers.set_depth(self.depth + 1);
filter_ctx.apply_trace_propagation(&mut framework_headers);
let transport_budget = step_budget
.checked_sub(step_started.elapsed())
.unwrap_or(Duration::ZERO);
if transport_budget.is_zero() {
return Ok(RawResponse::Rejected(Rejection::status(504)));
}
match filter_ctx.subrequest_response_mode {
SubRequestResponseMode::Streaming => {
if matches!(
pipeline.body_capabilities().response_body_mode,
crate::BodyMode::StreamBuffer { .. }
) {
return Err(format!(
"filtered_subrequest: step '{label}' selected streaming despite a StreamBuffer response mode"
)
.into());
}
let limits = StreamLimits {
idle_timeout: STREAMING_IDLE_TIMEOUT,
max_stream_duration: None,
max_total_bytes: streaming_transport_limit(
pipeline.body_capabilities().response_body_mode,
),
};
let mut response =
Err(SubRequestError::Connect("filtered_subrequest: no addresses to dial".to_owned()));
for peer in &peers {
let attempt_budget = step_deadline
.checked_duration_since(Instant::now())
.unwrap_or(Duration::ZERO);
if attempt_budget.is_zero() {
break;
}
response = match peer {
Ok(peer) => self
.client
.send_streaming(peer, &request, attempt_budget, limits.clone(), Some(&framework_headers))
.await,
Err(error) => Err(SubRequestError::Connect(error.to_string())),
};
if matches!(response, Err(SubRequestError::Connect(_))) {
continue;
}
break;
}
in_transport_inner.store(false, Ordering::Release);
match response {
Ok(response) => {
let status = response.status;
let mut headers = response.headers;
sanitize_subresponse_headers(&mut headers);
response_header.status = http::StatusCode::from_u16(status)
.map_err(|error| -> FilterError { format!("invalid upstream status: {error}").into() })?;
response_header.headers.clone_from(&headers);
filter_ctx.response_header = Some(&mut response_header);
let response_action = pipeline.execute_http_response(&mut filter_ctx).await?;
if let FilterAction::Reject(rejection) = response_action {
response.body.cancel().await;
return Ok(RawResponse::Rejected(rejection));
}
if self.accounting.exceeds_limit(&filter_ctx.extensions) {
response.body.cancel().await;
return Ok(RawResponse::Rejected(Rejection::status(413)));
}
let metadata = filter_ctx.response_header.as_deref().ok_or_else(|| -> FilterError {
"filtered_subrequest: response metadata missing after header filters"
.to_owned()
.into()
})?;
let status = metadata.status;
let mut headers = metadata.headers.clone();
sanitize_subresponse_headers(&mut headers);
Ok(RawResponse::Streaming {
body: Box::new(response.body),
outcome: SubrequestOutcome {
response: SubResponse { status: status.as_u16(), headers, body: Bytes::new() },
origin: ResponseOrigin::Upstream,
transport_error: None,
},
})
},
Err(error) => {
let (status, kind) = classify_transport_failure(&error);
warn!(step = label, %error, status, "filtered sub-request streaming transport failure");
let response = SubResponse { status, headers: HeaderMap::new(), body: Bytes::new() };
response_header.status = http::StatusCode::from_u16(status)
.map_err(|source| -> FilterError { source.into() })?;
filter_ctx.response_header = Some(&mut response_header);
let response_action = pipeline.execute_http_response(&mut filter_ctx).await?;
if let FilterAction::Reject(rejection) = response_action {
return Ok(RawResponse::Rejected(rejection));
}
if self.accounting.exceeds_limit(&filter_ctx.extensions) {
return Ok(RawResponse::Rejected(Rejection::status(413)));
}
let metadata = filter_ctx.response_header.as_deref().ok_or_else(|| -> FilterError {
"filtered_subrequest: response metadata missing after header filters"
.to_owned()
.into()
})?;
let mut headers = metadata.headers.clone();
sanitize_subresponse_headers(&mut headers);
Ok(RawResponse::Complete(SubrequestOutcome {
response: SubResponse {
status: metadata.status.as_u16(),
headers,
body: response.body,
},
origin: ResponseOrigin::Transport,
transport_error: Some(kind),
}))
},
}
},
SubRequestResponseMode::Buffered => {
let mut attempt = None;
for peer in &peers {
let attempt_budget = step_deadline
.checked_duration_since(Instant::now())
.unwrap_or(Duration::ZERO);
if attempt_budget.is_zero() {
break;
}
let outcome = match peer {
Ok(peer) => match self
.client
.execute(peer, &request, self.max_response_bytes, attempt_budget, Some(&framework_headers))
.await
{
Ok(response) => (response, ResponseOrigin::Upstream, None),
Err(error) => {
let (status, kind) = classify_transport_failure(&error);
warn!(step = label, %error, status, "filtered sub-request buffered transport failure");
let response =
SubResponse { status, headers: HeaderMap::new(), body: Bytes::new() };
if matches!(error, SubRequestError::Connect(_)) {
attempt = Some((response, ResponseOrigin::Transport, Some(kind)));
continue;
}
(response, ResponseOrigin::Transport, Some(kind))
},
},
Err(error) => {
warn!(step = label, %error, status = 502_u16, "filtered sub-request buffered transport failure");
attempt = Some((
SubResponse { status: 502, headers: HeaderMap::new(), body: Bytes::new() },
ResponseOrigin::Transport,
Some(TransportFailure::Connect),
));
continue;
},
};
attempt = Some(outcome);
break;
}
let (mut response, origin, transport_error) = attempt.unwrap_or_else(|| {
(
SubResponse { status: 502, headers: HeaderMap::new(), body: Bytes::new() },
ResponseOrigin::Transport,
Some(TransportFailure::Connect),
)
});
in_transport_inner.store(false, Ordering::Release);
sanitize_subresponse_headers(&mut response.headers);
response_header.status = http::StatusCode::from_u16(response.status)
.map_err(|error| -> FilterError { error.into() })?;
response_header.headers.clone_from(&response.headers);
filter_ctx.response_header = Some(&mut response_header);
let response_action = pipeline.execute_http_response(&mut filter_ctx).await?;
if let FilterAction::Reject(rejection) = response_action {
return Ok(RawResponse::Rejected(rejection));
}
if self.accounting.exceeds_limit(&filter_ctx.extensions) {
return Ok(RawResponse::Rejected(Rejection::status(413)));
}
let mut body = Some(std::mem::take(&mut response.body));
if let Some(limit) = response_body_overflow_limit(
pipeline.body_capabilities().response_body_mode,
self.max_response_bytes,
body.as_ref().map_or(0, Bytes::len),
) {
return Ok(RawResponse::ResponseTooLarge {
actual: body.as_ref().map_or(0, Bytes::len),
limit,
message: "filtered_subrequest: step response exceeds configured body limit",
});
}
let body_action = pipeline.execute_http_response_body(&mut filter_ctx, &mut body, true)?;
if let FilterAction::Reject(rejection) = body_action {
return Ok(RawResponse::Rejected(rejection));
}
if self.accounting.exceeds_limit(&filter_ctx.extensions) {
return Ok(RawResponse::Rejected(Rejection::status(413)));
}
if let Some(limit) = response_body_overflow_limit(
pipeline.body_capabilities().response_body_mode,
self.max_response_bytes,
body.as_ref().map_or(0, Bytes::len),
) {
return Ok(RawResponse::ResponseTooLarge {
actual: body.as_ref().map_or(0, Bytes::len),
limit,
message: "filtered_subrequest: transformed step response exceeds configured body limit",
});
}
response.body = body.unwrap_or_default();
if let Some(metadata) = filter_ctx.response_header.as_deref() {
response.status = metadata.status.as_u16();
response.headers.clone_from(&metadata.headers);
}
sanitize_subresponse_headers(&mut response.headers);
Ok(RawResponse::Complete(SubrequestOutcome { response, origin, transport_error }))
},
}
}
.instrument(step_span))
.await;
let mut raw = match timed {
Ok(Ok(raw)) => raw,
Ok(Err(error)) => return Err(FilteredSubrequestError::capture(error, &mut filter_ctx)),
Err(_) if in_transport.load(Ordering::Acquire) => RawResponse::Complete(SubrequestOutcome {
response: SubResponse {
status: 504,
headers: HeaderMap::new(),
body: Bytes::new(),
},
origin: ResponseOrigin::Transport,
transport_error: Some(TransportFailure::DeadlineExceeded),
}),
Err(_) => RawResponse::Rejected(Rejection::status(504)),
};
filter_ctx.response_header = None;
if filter_ctx.subrequest_response_mode == SubRequestResponseMode::Streaming
&& let RawResponse::Complete(outcome) = &mut raw
&& outcome.origin == ResponseOrigin::Transport
{
let cause = outcome
.transport_error
.map_or(StreamTerminationCause::Io, stream_termination_cause);
filter_ctx.extensions.insert(StreamTermination::new(cause));
let response_snapshot = crate::Response {
status: http::StatusCode::from_u16(outcome.response.status).unwrap_or(http::StatusCode::BAD_GATEWAY),
headers: outcome.response.headers.clone(),
};
let mut completion_body = None;
let completion_action = pipeline
.execute_http_response_body_with_response_header(
&mut filter_ctx,
&mut completion_body,
true,
Some(&response_snapshot),
)
.map_err(|error| FilteredSubrequestError::capture(error, &mut filter_ctx))?;
if let FilterAction::Reject(_) = completion_action {
let error = "filtered_subrequest: step completion filter rejected an abnormal stream"
.to_owned()
.into();
return Err(FilteredSubrequestError::capture(error, &mut filter_ctx));
}
if self.accounting.exceeds_limit(&filter_ctx.extensions) {
let error = "filtered_subrequest: retained state limit exceeded during stream completion"
.to_owned()
.into();
return Err(FilteredSubrequestError::capture(error, &mut filter_ctx));
}
if let Some(limit) = response_body_overflow_limit(
pipeline.body_capabilities().response_body_mode,
self.max_response_bytes,
completion_body.as_ref().map_or(0, Bytes::len),
) {
let error = FilteredSubrequestError::capture(
"filtered_subrequest: abnormal completion exceeds response body limit"
.to_owned()
.into(),
&mut filter_ctx,
)
.with_too_large(ResponseTooLargeInfo {
actual: Some(completion_body.as_ref().map_or(0, Bytes::len)),
limit,
});
return Err(error);
}
outcome.response.body = completion_body.unwrap_or_default();
}
let (kind, response_snapshot, completed) = match raw {
RawResponse::Complete(outcome) => {
let snapshot = crate::Response {
status: http::StatusCode::from_u16(outcome.response.status)
.unwrap_or(http::StatusCode::BAD_GATEWAY),
headers: outcome.response.headers.clone(),
};
(OpenedResponse::Complete(outcome), snapshot, true)
},
RawResponse::Rejected(rejection) => {
let response = subresponse_from_rejection(rejection);
let snapshot = crate::Response {
status: http::StatusCode::from_u16(response.status).unwrap_or(http::StatusCode::BAD_GATEWAY),
headers: response.headers.clone(),
};
(
OpenedResponse::Complete(SubrequestOutcome {
response,
origin: ResponseOrigin::Local,
transport_error: None,
}),
snapshot,
true,
)
},
RawResponse::Streaming { body, outcome } => {
let snapshot = crate::Response {
status: http::StatusCode::from_u16(outcome.response.status)
.unwrap_or(http::StatusCode::BAD_GATEWAY),
headers: outcome.response.headers.clone(),
};
(OpenedResponse::Streaming { body, outcome }, snapshot, false)
},
RawResponse::ResponseTooLarge { actual, limit, message } => {
let error = FilteredSubrequestError::capture(message.to_owned().into(), &mut filter_ctx)
.with_too_large(ResponseTooLargeInfo {
actual: Some(actual),
limit,
});
return Err(error);
},
};
let request_snapshot = crate::Request {
method: filter_ctx.request.method.clone(),
uri: filter_ctx.request.uri.clone(),
headers: filter_ctx.request.headers.clone(),
};
let continuation = FilteredSubrequestContinuation::capture(
Arc::clone(pipeline),
request_snapshot,
response_snapshot,
&mut filter_ctx,
completed,
step_deadline,
);
Ok(OpenedSubrequest { continuation, kind })
}
}