pub(crate) mod body;
pub(crate) mod branch;
mod build;
mod build_branch;
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,
body::{BodyCapabilities, BodyMode},
builtins::http::payload_processing::compression_config::CompressionConfig,
extensions::RequestExtensions,
};
pub struct FilterPipeline {
body_capabilities: BodyCapabilities,
compression: Option<CompressionConfig>,
pub(crate) filters: Vec<PipelineFilter>,
record_filter_duration_metrics: bool,
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,
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>,
}
#[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 response_body_ceiling(&self) -> Option<usize> {
self.response_body_ceiling
}
pub fn body_capabilities(&self) -> &BodyCapabilities {
&self.body_capabilities
}
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 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 filter_request_conditions_match(&self, type_name: &str, request: &crate::Request) -> bool {
self.filters
.iter()
.filter(|pf| pf.filter.name() == type_name)
.any(|pf| crate::condition::should_execute(&pf.conditions, request))
}
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 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.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> {
self.filters
.iter()
.filter_map(|pf| match &pf.filter {
crate::any_filter::AnyFilter::Http(f) => Some(f.referenced_files()),
crate::any_filter::AnyFilter::Tcp(_) => None,
})
.flatten()
.collect()
}
pub fn apply_insecure_options(&self, options: &InsecureOptions) {
for pf in &self.filters {
if let crate::any_filter::AnyFilter::Http(f) = &pf.filter {
f.apply_insecure_options(options);
}
}
}
fn visit_nested_pipelines(&mut self, visitor: &mut dyn FnMut(&mut FilterPipeline)) {
for pf in &mut self.filters {
if let crate::any_filter::AnyFilter::Http(filter) = &mut pf.filter {
filter.visit_nested_pipelines(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)
},
}
}