use std::{collections::HashMap, mem, sync::Arc};
use praxis_core::{
config::{FilterEntry, SkipPipelineChecks},
id::IdGenerator,
time::SystemTimeSource,
};
use tracing::{debug, warn};
use super::{
FilterPipeline,
body::{body_filter_indices, compute_body_capabilities},
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)?;
warn_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.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]>,
) -> Result<Self, FilterError> {
let filters = super::build_branch::resolve_chain_filters(entries, registry, chains, 0)?;
Ok(Self::from_filters(filters))
}
fn from_filters(filters: Vec<PipelineFilter>) -> Self {
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 (request_body_filter_indices, response_body_filter_indices) = body_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,
health_registry: None,
id_generator: Arc::clone(&id_generator),
kv_stores: None,
session_stores: None,
pipeline_extensions: Vec::new(),
record_filter_duration_metrics: false,
subrequest_client: None,
may_select_streaming_subrequest_response,
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
}
#[expect(
clippy::too_many_lines,
reason = "streaming capability check adds one validation pass"
)]
pub fn ordering_errors(
&self,
entries: &[FilterEntry],
allow_open_security: bool,
skip: &SkipPipelineChecks,
) -> 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_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_irr_with_router_or_lb(&names, &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(&names, &mut warnings);
super::checks::check_all_routers_conditional(&names, &self.filters, &mut warnings);
warnings
}
}
fn warn_tcp_unsupported_fields(filter: &AnyFilter, entry: &FilterEntry) {
if !matches!(filter, AnyFilter::Tcp(_)) {
return;
}
if !entry.conditions.is_empty() || !entry.response_conditions.is_empty() {
warn!(
filter = filter.name(),
"TCP filter has conditions that will be ignored; \
conditions are only evaluated for HTTP filters"
);
}
if entry.branch_chains.is_some() {
warn!(
filter = filter.name(),
"TCP filter has branch_chains that will be ignored; \
branching is only supported for HTTP filters"
);
}
}
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))
})
}