pub(crate) mod body;
pub(crate) mod branch;
mod build;
pub(crate) mod build_branch;
#[cfg(feature = "upstream-binding")]
pub(crate) mod catalog;
mod checks;
mod clusters;
pub(crate) mod evaluate;
mod extension;
pub(crate) mod filter;
mod http;
mod http_utils;
pub(crate) mod introspection;
pub(crate) mod subrequest;
mod tcp;
#[cfg(test)]
mod test_filters;
#[cfg(test)]
#[expect(clippy::allow_attributes, reason = "blanket test suppressions")]
#[allow(
clippy::unwrap_used,
clippy::expect_used,
clippy::indexing_slicing,
clippy::panic,
clippy::field_reassign_with_default,
clippy::type_complexity,
clippy::too_many_lines,
clippy::redundant_closure_for_method_calls,
clippy::significant_drop_tightening,
clippy::doc_markdown,
reason = "tests"
)]
mod tests;
use std::sync::Arc;
pub use extension::PipelineExtension;
use praxis_core::{
config::{ABSOLUTE_MAX_BODY_BYTES, FailureMode, InsecureOptions},
health::HealthRegistry,
id::IdGenerator,
kv::KvStoreRegistry,
time::TimeSource,
};
use tracing::{error, warn};
use self::filter::PipelineFilter;
use crate::{
FilterError,
any_filter::AnyFilter,
body::{BodyCapabilities, BodyMode},
builtins::http::payload_processing::compression_config::CompressionConfig,
extensions::RequestExtensions,
};
#[expect(
clippy::struct_excessive_bools,
reason = "independent server-injected runtime toggles, not a state machine"
)]
pub struct FilterPipeline {
body_capabilities: BodyCapabilities,
compression: Option<CompressionConfig>,
pub(crate) filters: Vec<PipelineFilter>,
record_filter_duration_metrics: bool,
route_templates: Arc<praxis_core::config::RouteTemplates>,
health_registry: Option<HealthRegistry>,
id_generator: Arc<IdGenerator>,
kv_stores: Option<KvStoreRegistry>,
session_stores: Option<Arc<crate::SessionStoreRegistry>>,
subrequest_client: Option<praxis_core::subrequest::SubRequestClient>,
may_select_streaming_subrequest_response: bool,
trace_context_filter_indices: Vec<usize>,
pipeline_extensions: Vec<Box<dyn PipelineExtension>>,
time_source: Arc<dyn TimeSource>,
request_body_ceiling: Option<usize>,
response_body_ceiling: Option<usize>,
request_body_filter_indices: Vec<usize>,
response_body_filter_indices: Vec<usize>,
selected_upstream_request_body_filter_indices: Vec<usize>,
#[cfg(feature = "bound-upstream-request-body")]
bound_upstream_request_body_filter_indices: Vec<usize>,
allow_private_upstreams: bool,
response_trailer_filter_indices: Vec<usize>,
}
#[expect(
clippy::multiple_inherent_impl,
reason = "pipeline concerns are split across modules"
)]
impl FilterPipeline {
pub fn apply_body_limits(
&mut self,
max_request: Option<usize>,
max_response: Option<usize>,
allow_unbounded: bool,
) -> Result<(), FilterError> {
self.apply_nested_body_limits(max_request, max_response, allow_unbounded)?;
if let Some(ceiling) = max_request {
self.body_capabilities.request_body_mode = clamp_body_mode(
self.body_capabilities.request_body_mode,
ceiling,
self.body_capabilities.needs_request_body,
);
self.body_capabilities.needs_request_body = true;
}
if let Some(ceiling) = max_response {
self.body_capabilities.response_body_mode = clamp_body_mode(
self.body_capabilities.response_body_mode,
ceiling,
self.body_capabilities.needs_response_body,
);
self.body_capabilities.needs_response_body = true;
}
check_unbounded_stream_buffer(
"request",
&mut self.body_capabilities.request_body_mode,
allow_unbounded,
)?;
check_unbounded_stream_buffer(
"response",
&mut self.body_capabilities.response_body_mode,
allow_unbounded,
)?;
self.request_body_ceiling = max_request;
self.response_body_ceiling = max_response;
Ok(())
}
fn apply_nested_body_limits(
&mut self,
max_request: Option<usize>,
max_response: Option<usize>,
allow_unbounded: bool,
) -> Result<(), FilterError> {
let mut nested_error = None;
self.visit_nested_pipelines(&mut |pipeline| {
if nested_error.is_none()
&& let Err(error) = pipeline.apply_body_limits(max_request, max_response, allow_unbounded)
{
nested_error = Some(error);
}
});
if let Some(error) = nested_error {
return Err(error);
}
Ok(())
}
#[must_use]
pub fn request_body_ceiling(&self) -> Option<usize> {
self.request_body_ceiling
}
#[must_use]
pub fn selected_upstream_request_body_limit(&self) -> usize {
let mode_limit = match self.body_capabilities().request_body_mode {
BodyMode::StreamBuffer { max_bytes } => max_bytes.unwrap_or(ABSOLUTE_MAX_BODY_BYTES),
BodyMode::SizeLimit { max_bytes } => max_bytes,
_ => ABSOLUTE_MAX_BODY_BYTES,
};
self.request_body_ceiling()
.map_or(mode_limit, |ceiling| mode_limit.min(ceiling))
}
#[must_use]
pub fn response_body_ceiling(&self) -> Option<usize> {
self.response_body_ceiling
}
pub fn body_capabilities(&self) -> &BodyCapabilities {
&self.body_capabilities
}
pub fn clear_request_body_done(&self, body_done_indices: &mut [bool]) {
for &idx in &self.request_body_filter_indices {
if let Some(done) = body_done_indices.get_mut(idx) {
*done = false;
}
}
}
pub fn needs_body_filters(&self) -> bool {
self.body_capabilities.needs_request_body || self.body_capabilities.needs_response_body
}
pub fn len(&self) -> usize {
self.filters.len()
}
pub fn is_empty(&self) -> bool {
self.filters.is_empty()
}
pub fn contains_filter(&self, type_name: &str) -> bool {
self.filters.iter().any(|pf| pf.filter.name() == type_name)
}
pub(crate) fn enables_trace_propagation(&self, request: &crate::Request) -> bool {
self.trace_context_filter_indices.iter().any(|&idx| {
self.filters
.get(idx)
.is_some_and(|pf| crate::condition::should_execute(&pf.conditions, request))
})
}
pub(crate) fn enables_trace_propagation_from<S: crate::condition::HeaderSource>(
&self,
request: &crate::Request,
headers: &S,
) -> Result<bool, S::Error> {
for &idx in &self.trace_context_filter_indices {
let Some(pf) = self.filters.get(idx) else {
continue;
};
if crate::condition::should_execute_from(
&pf.conditions,
request,
headers,
crate::condition::BoundUpstreamView::default(),
crate::condition::SelectedUpstream::none(),
)? {
return Ok(true);
}
}
Ok(false)
}
pub fn filters_unsupported_by(&self, listener_protocol: praxis_core::config::ProtocolKind) -> Vec<&'static str> {
self.filters
.iter()
.filter(|pf| !listener_protocol.supports(pf.filter.protocol_level()))
.map(|pf| pf.filter.name())
.collect()
}
pub fn non_http_filters(&self) -> Vec<&'static str> {
let mut names = Vec::new();
for_each_pipeline_filter(&self.filters, &mut |pf| {
if pf.filter.protocol_level() != praxis_core::config::ProtocolKind::Http {
names.push(pf.filter.name());
}
});
names
}
pub fn terminal_filters(&self) -> Vec<&'static str> {
let mut names = Vec::new();
for_each_pipeline_filter(&self.filters, &mut |pf| {
if let AnyFilter::Http(filter) = &pf.filter
&& filter.produces_terminal_response()
{
names.push(filter.name());
}
});
names
}
pub fn filter_request_conditions_match(&self, type_name: &str, ctx: &crate::HttpFilterContext<'_>) -> bool {
let bound = ctx.bound_upstream_view();
let selected = http_utils::ctx_selected_upstream(ctx);
self.filters
.iter()
.filter(|pf| pf.filter.name() == type_name)
.any(|pf| crate::condition::should_execute_bound_selected(&pf.conditions, ctx.request, bound, selected))
}
pub fn emit_deferred_records(&self, ctx: &crate::HttpFilterContext<'_>, status: u16) -> bool {
self.filters
.iter()
.any(|pf| pf.filter.emit_deferred_record(ctx, status))
}
pub fn needs_response_trailers(&self) -> bool {
self.body_capabilities.needs_response_trailers
}
pub fn compression_config(&self) -> Option<&CompressionConfig> {
self.compression.as_ref()
}
pub fn set_health_registry(&mut self, registry: HealthRegistry) {
self.visit_nested_pipelines(&mut |pipeline| pipeline.set_health_registry(Arc::clone(®istry)));
self.health_registry = Some(registry);
}
pub fn set_record_filter_duration_metrics(&mut self, enabled: bool) {
self.visit_nested_pipelines(&mut |pipeline| pipeline.set_record_filter_duration_metrics(enabled));
self.record_filter_duration_metrics = enabled;
}
pub fn records_filter_duration_metrics(&self) -> bool {
self.record_filter_duration_metrics
}
pub fn set_route_templates(&mut self, templates: Arc<praxis_core::config::RouteTemplates>) {
self.visit_nested_pipelines(&mut |pipeline| pipeline.set_route_templates(Arc::clone(&templates)));
self.route_templates = templates;
}
pub fn route_templates(&self) -> &praxis_core::config::RouteTemplates {
&self.route_templates
}
pub fn health_registry(&self) -> Option<&HealthRegistry> {
self.health_registry.as_ref()
}
pub fn id_generator(&self) -> &IdGenerator {
&self.id_generator
}
pub fn set_id_generator(&mut self, generator: Arc<IdGenerator>) {
self.visit_nested_pipelines(&mut |pipeline| pipeline.set_id_generator(Arc::clone(&generator)));
self.id_generator = generator;
}
pub fn kv_stores(&self) -> Option<&KvStoreRegistry> {
self.kv_stores.as_ref()
}
pub fn set_kv_stores(&mut self, stores: KvStoreRegistry) {
self.visit_nested_pipelines(&mut |pipeline| pipeline.set_kv_stores(stores.clone()));
self.kv_stores = Some(stores);
}
pub fn session_stores(&self) -> Option<&Arc<crate::SessionStoreRegistry>> {
self.session_stores.as_ref()
}
pub fn set_session_stores(&mut self, stores: Arc<crate::SessionStoreRegistry>) {
self.visit_nested_pipelines(&mut |pipeline| pipeline.set_session_stores(Arc::clone(&stores)));
self.session_stores = Some(stores);
}
pub fn subrequest_client(&self) -> Option<&praxis_core::subrequest::SubRequestClient> {
self.subrequest_client.as_ref()
}
pub fn may_select_streaming_subrequest_response(&self) -> bool {
self.may_select_streaming_subrequest_response
}
pub fn set_subrequest_client(&mut self, client: praxis_core::subrequest::SubRequestClient) {
self.visit_nested_pipelines(&mut |pipeline| pipeline.set_subrequest_client(client.clone()));
self.subrequest_client = Some(client);
}
pub fn add_pipeline_extension(&mut self, ext: Box<dyn PipelineExtension>) {
self.pipeline_extensions.push(ext);
}
pub fn prepare_extensions(&self, extensions: &mut RequestExtensions) {
for ext in &self.pipeline_extensions {
ext.prepare(extensions);
}
}
pub fn time_source(&self) -> &dyn TimeSource {
&*self.time_source
}
pub fn set_time_source(&mut self, source: Arc<dyn TimeSource>) {
self.visit_nested_pipelines(&mut |pipeline| pipeline.set_time_source(Arc::clone(&source)));
self.time_source = source;
}
pub fn referenced_files(&self) -> Vec<std::path::PathBuf> {
let mut files = Vec::new();
for_each_pipeline_filter(&self.filters, &mut |pf| {
if let AnyFilter::Http(f) = &pf.filter {
files.extend(f.referenced_files());
}
});
files
}
pub fn allow_private_upstreams(&self) -> bool {
self.allow_private_upstreams
}
pub fn set_allow_private_upstreams(&mut self, allow: bool) {
self.visit_nested_pipelines(&mut |pipeline| pipeline.set_allow_private_upstreams(allow));
self.allow_private_upstreams = allow;
}
pub fn apply_insecure_options(&self, options: &InsecureOptions) {
for_each_pipeline_filter(&self.filters, &mut |pf| {
if let AnyFilter::Http(f) = &pf.filter {
f.apply_insecure_options(options);
}
});
}
fn visit_nested_pipelines(&mut self, visitor: &mut dyn FnMut(&mut FilterPipeline)) {
visit_branch_nested_pipelines(&mut self.filters, visitor);
}
#[cfg(feature = "iterative-request-router")]
pub(crate) fn consumes_bound_upstream(&self) -> bool {
let mut found = false;
for_each_pipeline_filter(&self.filters, &mut |pf| {
if let AnyFilter::Http(f) = &pf.filter {
found = found || f.consumes_bound_upstream();
}
});
found
}
#[cfg(feature = "iterative-request-router")]
pub(crate) fn uses_bound_upstream(&self) -> bool {
checks::uses_bound_upstream(&self.filters)
}
#[cfg(feature = "iterative-request-router")]
pub(crate) fn bound_upstream_candidates(&self) -> std::collections::HashSet<String> {
checks::bound_consumer_clusters(&self.filters)
}
#[cfg(feature = "iterative-request-router")]
pub(crate) fn bound_when_matchers(&self) -> Vec<praxis_core::config::ApplicationMatch> {
checks::bound_when_matchers(&self.filters)
}
#[cfg(feature = "iterative-request-router")]
pub(crate) fn serves_bound_cluster(
&self,
cluster: &str,
metadata: Option<&catalog::ClusterApplicationMetadata>,
) -> bool {
checks::serves_bound_cluster(&self.filters, cluster, metadata)
}
#[cfg(feature = "iterative-request-router")]
pub(crate) fn cluster_metadata_declarations(&self) -> Vec<catalog::ClusterMetadataDeclaration> {
collect_cluster_declarations(&self.filters)
}
}
fn for_each_pipeline_filter(filters: &[PipelineFilter], visit: &mut dyn FnMut(&PipelineFilter)) {
for pf in filters {
visit(pf);
for branch in &pf.branches {
for_each_pipeline_filter(&branch.filters, visit);
}
}
}
#[cfg(feature = "upstream-binding")]
pub(super) fn collect_cluster_declarations(filters: &[PipelineFilter]) -> Vec<catalog::ClusterMetadataDeclaration> {
let mut declarations = Vec::new();
for_each_pipeline_filter(filters, &mut |pf| {
if let AnyFilter::Http(f) = &pf.filter {
declarations.extend(f.declared_cluster_metadata());
}
});
declarations
}
fn visit_branch_nested_pipelines(filters: &mut [PipelineFilter], visitor: &mut dyn FnMut(&mut FilterPipeline)) {
for pf in filters {
if let AnyFilter::Http(filter) = &mut pf.filter {
filter.visit_nested_pipelines(visitor);
}
for branch in &mut pf.branches {
visit_branch_nested_pipelines(&mut branch.filters, visitor);
}
}
}
fn clamp_body_mode(mode: BodyMode, ceiling: usize, filter_declared: bool) -> BodyMode {
match mode {
BodyMode::StreamBuffer { max_bytes } => BodyMode::StreamBuffer {
max_bytes: Some(max_bytes.map_or(ceiling, |m| m.min(ceiling))),
},
BodyMode::SizeLimit { max_bytes } => BodyMode::SizeLimit {
max_bytes: max_bytes.min(ceiling),
},
BodyMode::Stream if filter_declared => BodyMode::Stream,
BodyMode::Stream => BodyMode::SizeLimit { max_bytes: ceiling },
}
}
fn check_unbounded_stream_buffer(
direction: &str,
mode: &mut BodyMode,
allow_unbounded: bool,
) -> Result<(), FilterError> {
if let BodyMode::StreamBuffer { max_bytes: max @ None } = mode {
if allow_unbounded {
warn!(
direction = direction,
ceiling = ABSOLUTE_MAX_BODY_BYTES,
"StreamBuffer body mode has no per-filter size limit; \
clamped to absolute ceiling ({} MiB)",
ABSOLUTE_MAX_BODY_BYTES / 1_048_576
);
*max = Some(ABSOLUTE_MAX_BODY_BYTES);
} else {
return Err(format!(
"StreamBuffer {direction} body mode has no size limit; \
set max_{direction}_body_bytes or set \
insecure_options.allow_unbounded_body: true to allow"
)
.into());
}
}
Ok(())
}
pub(crate) fn check_failure_mode(
filter_name: &str,
error: FilterError,
phase: &str,
failure_mode: FailureMode,
) -> Result<(), FilterError> {
match failure_mode {
FailureMode::Open => {
warn!(
filter = filter_name,
error = %error,
phase,
failure_mode = "open",
"filter error, continuing"
);
Ok(())
},
FailureMode::Closed => {
error!(
filter = filter_name,
error = %error,
phase,
failure_mode = "closed",
"filter error, aborting request"
);
Err(error)
},
}
}