use std::pin::Pin;
use bytes::Bytes;
use tracing::{debug, trace, warn};
use super::{
FilterPipeline,
branch::{BranchOutcome, ResolvedBranch, branch_filter_executed},
filter::PipelineFilter,
http_utils::{
BodyFilterOutcome, HeaderFilterOutcome, accumulate_body_bytes, as_request_body_filter, as_response_body_filter,
released_or_continue, run_request_body_filter, run_request_filter, run_response_body_filter,
run_response_filter, run_selected_upstream_request_body_filter, skip_by_response_conditions,
},
};
use crate::{
FilterError,
actions::{FilterAction, Rejection, SelectedUpstreamBodyOutcome},
any_filter::AnyFilter,
condition::should_execute_bound_selected,
context::{EffectiveHeaders, HttpFilterContext},
trace_context::{TraceContext, ensure_trace_context},
};
#[cfg(feature = "bound-upstream-request-body")]
use crate::{actions::BoundUpstreamBodyOutcome, condition::SelectedUpstream, extensions::BoundRequestBodyRewrite};
#[expect(
clippy::multiple_inherent_impl,
reason = "pipeline concerns are split across modules"
)]
impl FilterPipeline {
#[expect(clippy::indexing_slicing, reason = "while loop bounds idx")]
#[expect(clippy::too_many_lines, reason = "filter identity tracking adds lines per branch")]
pub async fn execute_http_request(&self, ctx: &mut HttpFilterContext<'_>) -> Result<FilterAction, FilterError> {
if self.enables_trace_propagation(ctx.request) {
ensure_trace_context(ctx);
}
ctx.executed_filter_indices.clear();
ctx.executed_filter_indices.resize(self.filters.len(), false);
ctx.executed_branch_filters.clear();
ctx.body_done_indices.clear();
ctx.body_done_indices.resize(self.filters.len(), false);
let mut pre_read_results = std::mem::take(&mut ctx.filter_results);
let mut idx = 0;
while idx < self.filters.len() {
let pf = &self.filters[idx];
let http_filter = match &pf.filter {
AnyFilter::Http(f) => f.as_ref(),
AnyFilter::Tcp(_) => {
idx += 1;
continue;
},
};
if !pf.conditions.is_empty()
&& !should_execute_bound_selected(
&pf.conditions,
ctx.request,
ctx.bound_upstream_view(),
super::http_utils::ctx_selected_upstream(ctx),
)
{
trace!(filter = http_filter.name(), "skipped by conditions");
idx += 1;
continue;
}
if let Some(results) = pre_read_results.remove(http_filter.name()) {
ctx.filter_results.insert(http_filter.name(), results);
}
ctx.current_filter_id = Some(pf.filter_id);
let outcome =
run_request_filter(http_filter, ctx, pf.failure_mode, self.record_filter_duration_metrics).await;
ctx.current_filter_id = None;
match outcome? {
HeaderFilterOutcome::Rejected(r) => {
ctx.executed_filter_indices[idx] = true;
return Ok(FilterAction::Reject(r));
},
HeaderFilterOutcome::TerminalResponse(terminal) => {
ctx.executed_filter_indices[idx] = true;
return Ok(FilterAction::TerminalResponse(terminal));
},
HeaderFilterOutcome::StreamingTerminalResponse(terminal) => {
ctx.executed_filter_indices[idx] = true;
return Ok(FilterAction::StreamingTerminalResponse(terminal));
},
HeaderFilterOutcome::Continue => {},
}
ctx.executed_filter_indices[idx] = true;
#[cfg(feature = "upstream-binding")]
if published_first_binding(http_filter, ctx) {
ctx.freeze_bound_upstream();
#[cfg(feature = "bound-upstream-request-body")]
if let FilterAction::Reject(r) = self.run_bound_upstream_request_body(ctx).await? {
return Ok(FilterAction::Reject(r));
}
}
match super::evaluate::evaluate_branches(&pf.branches, ctx).await? {
BranchOutcome::Continue => idx += 1,
BranchOutcome::Terminal => {
if ctx.cluster.is_some() {
return Ok(FilterAction::Continue);
}
warn!(
filter = http_filter.name(),
"terminal branch produced no response and selected no cluster; \
stopping the pipeline with 500 instead of forwarding upstream"
);
return Ok(FilterAction::Reject(Rejection::status(500)));
},
BranchOutcome::SkipTo(t) => idx = t,
BranchOutcome::ReEnter(t) => {
idx = t;
},
BranchOutcome::Reject(r) => return Ok(FilterAction::Reject(r)),
BranchOutcome::TerminalResponse(t) => return Ok(FilterAction::TerminalResponse(t)),
BranchOutcome::StreamingTerminalResponse(t) => {
return Ok(FilterAction::StreamingTerminalResponse(t));
},
}
}
Ok(FilterAction::Continue)
}
pub async fn execute_http_response(&self, ctx: &mut HttpFilterContext<'_>) -> Result<FilterAction, FilterError> {
ctx.body_done_indices.clear();
ctx.body_done_indices.resize(self.filters.len(), false);
for (idx, pf) in self.filters.iter().enumerate().rev() {
if ctx.executed_filter_indices.get(idx) == Some(&false) {
trace!(
filter = pf.filter.name(),
"skipped on_response (not executed in request phase)"
);
continue;
}
if !pf.branches.is_empty()
&& let Some(rejection) = self.unwind_branches(&pf.branches, ctx).await?
{
return Ok(FilterAction::Reject(rejection));
}
if let Some(rejection) = self.run_response_hook(pf, ctx).await? {
return Ok(FilterAction::Reject(rejection));
}
}
Ok(FilterAction::Continue)
}
async fn unwind_branches(
&self,
branches: &[ResolvedBranch],
ctx: &mut HttpFilterContext<'_>,
) -> Result<Option<Rejection>, FilterError> {
for pf in branches.iter().rev().flat_map(|branch| branch.filters.iter().rev()) {
if !branch_filter_executed(&ctx.executed_branch_filters, pf.filter_id) {
trace!(
filter = pf.filter.name(),
"skipped branch on_response (not executed in request phase)"
);
continue;
}
if !pf.branches.is_empty()
&& let Some(rejection) = self.unwind_branches_boxed(&pf.branches, ctx).await?
{
return Ok(Some(rejection));
}
if let Some(rejection) = self.run_response_hook(pf, ctx).await? {
return Ok(Some(rejection));
}
}
Ok(None)
}
fn unwind_branches_boxed<'a>(
&'a self,
branches: &'a [ResolvedBranch],
ctx: &'a mut HttpFilterContext<'_>,
) -> UnwindFuture<'a> {
Box::pin(self.unwind_branches(branches, ctx))
}
async fn run_response_hook(
&self,
pf: &PipelineFilter,
ctx: &mut HttpFilterContext<'_>,
) -> Result<Option<Rejection>, FilterError> {
let AnyFilter::Http(http_filter) = &pf.filter else {
return Ok(None);
};
if skip_by_response_conditions(http_filter.as_ref(), &pf.response_conditions, ctx) {
return Ok(None);
}
ctx.current_filter_id = Some(pf.filter_id);
let outcome = run_response_filter(
http_filter.as_ref(),
ctx,
pf.failure_mode,
self.record_filter_duration_metrics,
)
.await;
ctx.current_filter_id = None;
match outcome? {
HeaderFilterOutcome::Continue
| HeaderFilterOutcome::TerminalResponse(_)
| HeaderFilterOutcome::StreamingTerminalResponse(_) => Ok(None),
HeaderFilterOutcome::Rejected(rejection) => Ok(Some(rejection)),
}
}
#[expect(clippy::too_many_lines, reason = "body hook loop with metrics dispatch")]
pub async fn execute_http_request_body(
&self,
ctx: &mut HttpFilterContext<'_>,
body: &mut Option<Bytes>,
end_of_stream: bool,
) -> Result<FilterAction, FilterError> {
ensure_body_done_indices(ctx, self.filters.len());
accumulate_body_bytes(&mut ctx.request_body_bytes, body.as_ref());
let request_phase_tracked = request_phase_tracked(ctx, self.filters.len());
let mut released = false;
for &idx in &self.request_body_filter_indices {
if !request_phase_tracked {
self.ensure_matching_trace_context(ctx);
}
let Some(pf) = self.filters.get(idx) else {
continue;
};
if ctx.body_done_indices.get(idx) == Some(&true) {
trace!(filter = pf.filter.name(), "skipped body (body_done)");
continue;
}
if skipped_in_request_phase(ctx, request_phase_tracked, idx) {
trace!(
filter = pf.filter.name(),
"skipped request body (not executed in request phase)"
);
continue;
}
let Some(http_filter) = as_request_body_filter(pf, ctx, request_phase_tracked)? else {
continue;
};
ctx.current_filter_id = Some(pf.filter_id);
let outcome = run_request_body_filter(
http_filter,
ctx,
body,
end_of_stream,
pf.failure_mode,
self.record_filter_duration_metrics,
)
.await;
ctx.current_filter_id = None;
match outcome? {
BodyFilterOutcome::Continue => {},
BodyFilterOutcome::Released => released = true,
BodyFilterOutcome::BodyDone => {
if let Some(done) = ctx.body_done_indices.get_mut(idx) {
*done = true;
}
},
BodyFilterOutcome::Rejected(r) => return Ok(FilterAction::Reject(r)),
}
}
if !request_phase_tracked {
self.ensure_matching_trace_context(ctx);
}
Ok(released_or_continue(released))
}
fn ensure_matching_trace_context(&self, ctx: &mut HttpFilterContext<'_>) {
if ctx.extensions.get::<TraceContext>().is_some() {
return;
}
match self.enables_trace_propagation_from(ctx.request, &EffectiveHeaders(ctx)) {
Ok(true) => ensure_trace_context(ctx),
Ok(false) => {},
Err(error) => debug!(%error, "trace_context: pre-read headers are ambiguous; not starting early"),
}
}
#[expect(clippy::too_many_lines, reason = "body hook loop with per-filter skip checks")]
pub async fn execute_http_selected_upstream_request_body(
&self,
ctx: &mut HttpFilterContext<'_>,
body: &mut Option<Bytes>,
) -> Result<FilterAction, FilterError> {
let request_phase_tracked = request_phase_tracked(ctx, self.filters.len());
for &idx in &self.selected_upstream_request_body_filter_indices {
let Some(pf) = self.filters.get(idx) else {
continue;
};
if skipped_in_request_phase(ctx, request_phase_tracked, idx) {
trace!(
filter = pf.filter.name(),
"skipped selected-upstream request body (not executed in request phase)"
);
continue;
}
let AnyFilter::Http(http_filter) = &pf.filter else {
continue;
};
ctx.current_filter_id = Some(pf.filter_id);
let outcome = run_selected_upstream_request_body_filter(
http_filter.as_ref(),
ctx,
body,
pf.failure_mode,
self.record_filter_duration_metrics,
)
.await;
ctx.current_filter_id = None;
match outcome? {
SelectedUpstreamBodyOutcome::Continue => {},
SelectedUpstreamBodyOutcome::Reject(rejection) => return Ok(FilterAction::Reject(rejection)),
}
}
Ok(FilterAction::Continue)
}
#[cfg(feature = "bound-upstream-request-body")]
async fn run_bound_upstream_request_body(
&self,
ctx: &mut HttpFilterContext<'_>,
) -> Result<FilterAction, FilterError> {
if self.bound_upstream_request_body_filter_indices.is_empty() {
return Ok(FilterAction::Continue);
}
let (action, rewrote) = self.execute_http_bound_upstream_request_body(ctx).await?;
if rewrote && matches!(action, FilterAction::Continue) {
let rewritten = ctx.buffered_request_body.clone().unwrap_or_default();
if rewritten.len() > self.selected_upstream_request_body_limit() {
return Ok(FilterAction::Reject(Rejection::status(413)));
}
ctx.extensions.insert(BoundRequestBodyRewrite(rewritten));
}
Ok(action)
}
#[cfg(feature = "bound-upstream-request-body")]
#[expect(
clippy::too_many_lines,
reason = "body hook loop with take/commit and per-filter skip checks"
)]
async fn execute_http_bound_upstream_request_body(
&self,
ctx: &mut HttpFilterContext<'_>,
) -> Result<(FilterAction, bool), FilterError> {
let had_buffer = ctx.buffered_request_body.is_some();
let mut body = ctx.buffered_request_body.take();
let mut result = Ok(FilterAction::Continue);
let mut rewrote = false;
for &idx in &self.bound_upstream_request_body_filter_indices {
let Some(pf) = self.filters.get(idx) else {
continue;
};
if !pf.conditions.is_empty()
&& !should_execute_bound_selected(
&pf.conditions,
ctx.request,
ctx.bound_upstream_view(),
SelectedUpstream::none(),
)
{
trace!(
filter = pf.filter.name(),
"skipped bound-upstream request body (conditions)"
);
continue;
}
let AnyFilter::Http(http_filter) = &pf.filter else {
continue;
};
ctx.current_filter_id = Some(pf.filter_id);
body = body.filter(|bytes| !bytes.is_empty());
let outcome = super::http_utils::run_bound_upstream_request_body_filter(
http_filter.as_ref(),
ctx,
&mut body,
pf.failure_mode,
self.record_filter_duration_metrics,
)
.await;
ctx.current_filter_id = None;
match outcome {
Ok((BoundUpstreamBodyOutcome::Continue, wrote)) => rewrote |= wrote,
Ok((BoundUpstreamBodyOutcome::Reject(rejection), _)) => {
result = Ok(FilterAction::Reject(rejection));
break;
},
Err(e) => {
result = Err(e);
break;
},
}
}
ctx.buffered_request_body = body.or_else(|| had_buffer.then(Bytes::new));
result.map(|action| (action, rewrote))
}
pub fn execute_http_response_body(
&self,
ctx: &mut HttpFilterContext<'_>,
body: &mut Option<Bytes>,
end_of_stream: bool,
) -> Result<FilterAction, FilterError> {
if !self.body_capabilities.any_response_body_condition {
return self.execute_http_response_body_with_response_header(ctx, body, end_of_stream, None);
}
let response_header = ctx.response_header.take();
let result =
self.execute_http_response_body_with_response_header(ctx, body, end_of_stream, response_header.as_deref());
ctx.response_header = response_header;
result
}
#[expect(clippy::too_many_lines, reason = "body hook loop with per-filter skip checks")]
pub fn execute_http_response_body_with_response_header(
&self,
ctx: &mut HttpFilterContext<'_>,
body: &mut Option<Bytes>,
end_of_stream: bool,
response_header: Option<&crate::context::Response>,
) -> Result<FilterAction, FilterError> {
ensure_body_done_indices(ctx, self.filters.len());
accumulate_body_bytes(&mut ctx.response_body_bytes, body.as_ref());
let request_phase_tracked = request_phase_tracked(ctx, self.filters.len());
let mut released = false;
for &idx in self.response_body_filter_indices.iter().rev() {
let Some(pf) = self.filters.get(idx) else {
continue;
};
if ctx.body_done_indices.get(idx) == Some(&true) {
trace!(filter = pf.filter.name(), "skipped body (body_done)");
continue;
}
if skipped_in_request_phase(ctx, request_phase_tracked, idx) {
trace!(
filter = pf.filter.name(),
"skipped response body (not executed in request phase)"
);
continue;
}
let Some(http_filter) = as_response_body_filter(&pf.filter, &pf.response_conditions, response_header)
else {
continue;
};
ctx.current_filter_id = Some(pf.filter_id);
let outcome = run_response_body_filter(
http_filter,
ctx,
body,
end_of_stream,
pf.failure_mode,
self.record_filter_duration_metrics,
);
ctx.current_filter_id = None;
match outcome? {
BodyFilterOutcome::Continue => {},
BodyFilterOutcome::Released => released = true,
BodyFilterOutcome::BodyDone => {
if let Some(done) = ctx.body_done_indices.get_mut(idx) {
*done = true;
}
},
BodyFilterOutcome::Rejected(r) => return Ok(FilterAction::Reject(r)),
}
}
Ok(released_or_continue(released))
}
pub fn execute_http_response_trailers(
&self,
ctx: &mut HttpFilterContext<'_>,
trailers: &mut http::HeaderMap,
) -> Result<Option<Bytes>, FilterError> {
let request_phase_tracked = request_phase_tracked(ctx, self.filters.len());
for &idx in self.response_trailer_filter_indices.iter().rev() {
let Some(pf) = self.filters.get(idx) else {
continue;
};
if skipped_in_request_phase(ctx, request_phase_tracked, idx) {
trace!(
filter = pf.filter.name(),
"skipped response trailers (not executed in request phase)"
);
continue;
}
let AnyFilter::Http(http_filter) = &pf.filter else {
continue;
};
ctx.current_filter_id = Some(pf.filter_id);
let produced = http_filter.on_response_trailers(ctx, trailers);
ctx.current_filter_id = None;
if let Some(bytes) = produced? {
return Ok(Some(bytes));
}
}
Ok(None)
}
}
type UnwindFuture<'a> = Pin<Box<dyn Future<Output = Result<Option<Rejection>, FilterError>> + Send + 'a>>;
#[cfg(feature = "upstream-binding")]
fn published_first_binding(filter: &dyn crate::filter::HttpFilter, ctx: &HttpFilterContext<'_>) -> bool {
filter.binds_upstream() && !ctx.bound_upstream_frozen() && ctx.bound_cluster().is_some()
}
fn ensure_body_done_indices(ctx: &mut HttpFilterContext<'_>, filter_count: usize) {
if ctx.body_done_indices.len() != filter_count {
ctx.body_done_indices.resize(filter_count, false);
}
}
fn request_phase_tracked(ctx: &HttpFilterContext<'_>, filter_count: usize) -> bool {
ctx.executed_filter_indices.len() == filter_count
}
fn skipped_in_request_phase(ctx: &HttpFilterContext<'_>, tracked: bool, idx: usize) -> bool {
tracked && ctx.executed_filter_indices.get(idx) == Some(&false)
}