use std::{collections::HashMap, mem, sync::Arc};
use praxis_core::{
config::{FilterEntry, InsecureOptions, SkipPipelineChecks},
id::IdGenerator,
time::SystemTimeSource,
};
use tracing::debug;
#[cfg(feature = "upstream-binding")]
use super::catalog::ClusterApplicationCatalog;
use super::{
FilterPipeline,
body::{body_filter_indices, compute_body_capabilities, selected_upstream_request_body_indices},
filter::PipelineFilter,
};
use crate::{FilterError, any_filter::AnyFilter, registry::FilterRegistry};
impl FilterPipeline {
pub fn build(entries: &mut [FilterEntry], registry: &FilterRegistry) -> Result<Self, FilterError> {
let mut filters = Vec::with_capacity(entries.len());
for (filter_id, entry) in entries.iter_mut().enumerate() {
let filter = registry.create(&entry.filter_type, &entry.config)?;
reject_tcp_unsupported_fields(&filter, entry)?;
let has_conditions = !entry.conditions.is_empty() || !entry.response_conditions.is_empty();
debug!(
filter = filter.name(),
conditions = has_conditions,
"filter added to pipeline"
);
let mut pf = PipelineFilter::new(
filter_id,
filter,
mem::take(&mut entry.conditions),
mem::take(&mut entry.response_conditions),
);
pf.failure_mode = entry.failure_mode;
pf.is_security = registry.is_security_filter(&entry.filter_type);
pf.name = entry.name.as_ref().map(|n| Arc::from(n.as_str()));
filters.push(pf);
}
Ok(Self::from_filters(filters))
}
pub fn build_with_chains(
entries: &mut [FilterEntry],
registry: &FilterRegistry,
chains: &HashMap<&str, &[FilterEntry]>,
insecure_options: &InsecureOptions,
) -> Result<Self, FilterError> {
let filters = super::build_branch::resolve_chain_filters(entries, registry, chains, 0, insecure_options)?;
Ok(Self::from_filters(filters))
}
#[expect(
clippy::too_many_lines,
reason = "single construction choke point: one precompute per body phase plus the full struct literal"
)]
pub(crate) fn from_filters(
#[cfg_attr(
not(feature = "upstream-binding"),
expect(unused_mut, reason = "only binding enablement mutates the filters")
)]
mut filters: Vec<PipelineFilter>,
) -> Self {
#[cfg(feature = "upstream-binding")]
if super::checks::uses_bound_upstream(&filters) {
let catalog = build_cluster_application_catalog(&filters);
enable_upstream_binding(&mut filters, &catalog);
}
let body_capabilities = compute_body_capabilities(&filters);
let compression = extract_compression_config(&filters);
let may_select_streaming_subrequest_response = filters_may_select_streaming_subrequest_response(&filters);
let trace_context_filter_indices = filters
.iter()
.enumerate()
.filter_map(|(idx, pf)| (pf.filter.name() == "trace_context").then_some(idx))
.collect();
let (request_body_filter_indices, response_body_filter_indices) = body_filter_indices(&filters);
let selected_upstream_request_body_filter_indices = selected_upstream_request_body_indices(&filters);
#[cfg(feature = "bound-upstream-request-body")]
let bound_upstream_request_body_filter_indices = super::body::bound_upstream_request_body_indices(&filters);
let response_trailer_filter_indices = super::body::response_trailer_filter_indices(&filters);
let id_generator = Arc::new(IdGenerator::new());
let time_source: Arc<dyn praxis_core::time::TimeSource> = Arc::new(SystemTimeSource);
let mut pipeline = Self {
body_capabilities,
compression,
filters,
request_body_filter_indices,
response_body_filter_indices,
selected_upstream_request_body_filter_indices,
#[cfg(feature = "bound-upstream-request-body")]
bound_upstream_request_body_filter_indices,
allow_private_upstreams: false,
response_trailer_filter_indices,
health_registry: None,
id_generator: Arc::clone(&id_generator),
kv_stores: None,
session_stores: None,
pipeline_extensions: Vec::new(),
record_filter_duration_metrics: false,
route_templates: Arc::default(),
subrequest_client: None,
may_select_streaming_subrequest_response,
trace_context_filter_indices,
time_source: Arc::clone(&time_source),
request_body_ceiling: None,
response_body_ceiling: None,
};
pipeline.set_id_generator(id_generator);
pipeline.set_time_source(time_source);
pipeline
}
pub fn ordering_errors(
&self,
entries: &[FilterEntry],
allow_open_security: bool,
skip: &SkipPipelineChecks,
) -> Vec<String> {
self.ordering_errors_inner(entries, allow_open_security, skip, false)
}
#[cfg(feature = "upstream-binding")]
fn binding_errors(&self, in_irr_step: bool, names: &[&str], errors: &mut Vec<String>) {
let uses_bound_upstream = super::checks::uses_bound_upstream(&self.filters);
if uses_bound_upstream {
super::checks::check_cluster_metadata_conflicts(&self.filters, errors);
}
super::checks::check_bound_upstream_requires_binding(&self.filters, in_irr_step, errors);
super::checks::check_bound_condition_with_pre_read_body(
&self.filters,
self.body_capabilities.request_body_mode,
errors,
);
#[cfg(feature = "bound-upstream-request-body")]
super::checks::check_bound_upstream_body_participants(&self.filters, in_irr_step, errors);
if uses_bound_upstream {
super::checks::check_no_rebind_after_binding(&self.filters, in_irr_step, errors);
}
super::checks::check_bound_cluster_coverage(&self.filters, errors);
super::checks::check_untagged_bound_cluster_fields(&self.filters, errors);
super::checks::check_irr_coexistence(&self.filters, names, errors);
}
#[cfg(feature = "iterative-request-router")]
pub(crate) fn step_ordering_errors(
&self,
entries: &[FilterEntry],
allow_open_security: bool,
skip: &SkipPipelineChecks,
) -> Vec<String> {
self.ordering_errors_inner(entries, allow_open_security, skip, true)
}
#[expect(
clippy::too_many_lines,
reason = "one sequential invocation per ordering check, including the bound-upstream passes"
)]
fn ordering_errors_inner(
&self,
entries: &[FilterEntry],
allow_open_security: bool,
skip: &SkipPipelineChecks,
in_irr_step: bool,
) -> Vec<String> {
let names: Vec<&str> = self.filters.iter().map(|pf| pf.filter.name()).collect();
let mut errors = Vec::new();
if !skip.lb_without_router {
super::checks::check_lb_without_cluster_selector(&self.filters, &mut errors);
}
if !skip.unreachable_filters {
super::checks::check_unconditional_static_response(&names, &self.filters, &mut errors);
}
if !skip.conditional_security {
super::checks::check_conditional_security(&names, &self.filters, &mut errors);
}
super::checks::check_open_security_filters(&names, &self.filters, allow_open_security, &mut errors);
if !skip.duplicate_routers {
super::checks::check_duplicate_routers(&names, &mut errors);
}
if !skip.duplicate_load_balancers {
super::checks::check_duplicate_load_balancers(&names, &mut errors);
}
if !skip.conflicting_cluster_selectors {
super::checks::check_conflicting_cluster_selectors(&self.filters, &mut errors);
}
if !skip.misaligned_clusters {
super::checks::check_misaligned_clusters(&self.filters, &mut errors);
}
if !skip.duplicate_rewrite_filters {
super::checks::check_duplicate_rewrite_filters(&names, entries, &mut errors);
}
super::checks::check_condition_header_names(&self.filters, &mut errors);
super::checks::check_trace_context_upstream_conditions(&self.filters, &mut errors);
super::checks::check_skip_to_bypasses_security(&self.filters, &mut errors);
super::checks::check_terminal_rejoin_bypasses_security(&self.filters, &mut errors);
super::checks::check_branch_body_filters(&self.filters, &mut errors);
super::checks::check_branch_selected_upstream_body_filters(&self.filters, &mut errors);
super::checks::check_selected_upstream_body_mode(&self.filters, &mut errors);
#[cfg(feature = "upstream-binding")]
self.binding_errors(in_irr_step, &names, &mut errors);
#[cfg(not(feature = "upstream-binding"))]
let _ = in_irr_step;
super::checks::check_selected_upstream_condition_ordering(&self.filters, &mut errors);
super::checks::check_selected_upstream_condition_pre_read(
&self.filters,
self.body_capabilities.request_body_mode,
&mut errors,
);
if self.may_select_streaming_subrequest_response
&& matches!(
self.body_capabilities.response_body_mode,
crate::BodyMode::StreamBuffer { .. }
)
{
errors.push(
"pipeline contains a filter that may select a streaming sub-request response, \
but its response body mode is StreamBuffer"
.to_owned(),
);
}
errors
}
pub fn ordering_warnings(&self) -> Vec<String> {
let names: Vec<&str> = self.filters.iter().map(|pf| pf.filter.name()).collect();
let mut warnings = Vec::new();
super::checks::check_router_without_lb(&self.filters, &names, &mut warnings);
super::checks::check_all_routers_conditional(&names, &self.filters, &mut warnings);
super::checks::check_security_filter_in_conditional_branch(&self.filters, &mut warnings);
warnings
}
}
#[cfg(feature = "upstream-binding")]
fn enable_upstream_binding(filters: &mut [PipelineFilter], catalog: &Arc<ClusterApplicationCatalog>) {
for pf in filters {
if let AnyFilter::Http(filter) = &mut pf.filter {
filter.enable_upstream_binding(Arc::clone(catalog));
}
for branch in &mut pf.branches {
enable_upstream_binding(&mut branch.filters, catalog);
}
}
}
pub(super) fn reject_tcp_unsupported_fields(filter: &AnyFilter, entry: &FilterEntry) -> Result<(), FilterError> {
if !matches!(filter, AnyFilter::Tcp(_)) {
return Ok(());
}
if !entry.conditions.is_empty() || !entry.response_conditions.is_empty() {
return Err(format!(
"filter '{}': conditions are not supported on TCP filters; they apply to HTTP filters only",
filter.name()
)
.into());
}
if entry.branch_chains.is_some() {
return Err(format!(
"filter '{}': branch_chains are not supported on TCP filters; they apply to HTTP filters only",
filter.name()
)
.into());
}
Ok(())
}
#[cfg(feature = "upstream-binding")]
fn build_cluster_application_catalog(filters: &[PipelineFilter]) -> Arc<ClusterApplicationCatalog> {
let (catalog, _conflicts) = super::catalog::build_catalog(super::collect_cluster_declarations(filters));
Arc::new(catalog)
}
fn extract_compression_config(
filters: &[PipelineFilter],
) -> Option<crate::builtins::http::payload_processing::compression_config::CompressionConfig> {
filters.iter().find_map(|pf| match &pf.filter {
AnyFilter::Http(f) => f.compression_config().cloned(),
AnyFilter::Tcp(_) => None,
})
}
fn filters_may_select_streaming_subrequest_response(filters: &[PipelineFilter]) -> bool {
filters.iter().any(|pf| {
let selects_streaming = match &pf.filter {
AnyFilter::Http(filter) => filter.may_select_streaming_subrequest_response(),
AnyFilter::Tcp(_) => false,
};
selects_streaming
|| pf
.branches
.iter()
.any(|branch| filters_may_select_streaming_subrequest_response(&branch.filters))
})
}