use std::{collections::HashMap, sync::Arc, time::Duration};
use arc_swap::ArcSwap;
use bytes::Bytes;
use pingora_core::{
Result, apps::HttpServerOptions, protocols::http::v2::server::H2Options, server::Server,
services::listening::Service,
};
use pingora_proxy::{Session, http_proxy};
use praxis_core::{config::ABSOLUTE_MAX_BODY_BYTES, connectivity::Upstream};
use praxis_filter::{BodyBuffer, BodyMode, FilterPipeline, HttpFilterContext, RequestExtensions};
use tokio::sync::Semaphore;
use tracing::{debug, warn};
use super::{context::PingoraRequestCtx, metrics};
mod compression;
mod connected_to_upstream;
mod fail_to_proxy;
mod hop_by_hop;
mod normalize;
mod request_body_filter;
mod request_filter;
mod reserved_headers;
mod response_body_filter;
mod response_filter;
mod retry;
mod upstream_peer;
mod upstream_request;
mod upstream_response;
mod via;
mod with_body;
pub use upstream_peer::{UpstreamRetryGateRelease, arm_upstream_retry_gate, lock_upstream_retry_gate_tests};
pub use with_body::PingoraHttpHandler;
pub fn load_http_handler(
server: &mut Server,
listener: &praxis_core::config::Listener,
pipeline: Arc<ArcSwap<FilterPipeline>>,
cert_watcher_shutdowns: &mut Vec<tokio::sync::watch::Sender<bool>>,
) -> Result<(), praxis_core::ProxyError> {
let downstream_read_timeout = listener.downstream_read_timeout_ms.map(Duration::from_millis);
let connection_semaphore = listener
.max_connections
.map(|max| Arc::new(Semaphore::new(max as usize)));
debug!(listener = %listener.name, "loading HTTP handler with body filters");
let handler = PingoraHttpHandler::new(
pipeline,
downstream_read_timeout,
connection_semaphore,
::metrics::SharedString::from_shared(Arc::from(listener.name.as_str())),
);
wire_service(server, listener, handler, cert_watcher_shutdowns)?;
Ok(())
}
fn wire_service<H>(
server: &mut Server,
listener: &praxis_core::config::Listener,
handler: H,
cert_watcher_shutdowns: &mut Vec<tokio::sync::watch::Sender<bool>>,
) -> Result<(), praxis_core::ProxyError>
where
H: pingora_proxy::ProxyHttp + Send + Sync + 'static,
H::CTX: Send + Sync,
{
let service_name = format!("http-proxy:{name}", name = listener.name);
let mut proxy = http_proxy(&server.configuration, handler);
proxy.server_options = Some(h2c_server_options());
proxy.h2_options = Some(h2_server_options());
let mut service = Service::new(service_name, proxy);
if let Some(tx) = super::listener::add_listener(&mut service, listener)? {
cert_watcher_shutdowns.push(tx);
}
server.add_service(service);
Ok(())
}
fn clamp_body_mode_to_ceiling(mode: BodyMode, baseline: BodyMode) -> BodyMode {
let ceiling = match baseline {
BodyMode::StreamBuffer { max_bytes: Some(v) } | BodyMode::SizeLimit { max_bytes: v } => Some(v),
_ => None,
};
match (mode, ceiling) {
(BodyMode::StreamBuffer { max_bytes }, Some(limit)) => BodyMode::StreamBuffer {
max_bytes: Some(max_bytes.map_or(limit, |v| v.min(limit))),
},
(BodyMode::SizeLimit { max_bytes }, Some(limit)) => BodyMode::SizeLimit {
max_bytes: max_bytes.min(limit),
},
(m, None | Some(_)) => m,
}
}
fn legacy_default_policy() -> Arc<praxis_core::config::RetryPolicy> {
static LEGACY_DEFAULT: std::sync::LazyLock<Arc<praxis_core::config::RetryPolicy>> =
std::sync::LazyLock::new(|| Arc::new(praxis_core::config::RetryPolicy::legacy_default()));
Arc::clone(&LEGACY_DEFAULT)
}
#[expect(clippy::too_many_lines, reason = "sequential guard checks")]
fn handle_connect_failure(ctx: &mut PingoraRequestCtx, e: Box<pingora_core::Error>) -> Box<pingora_core::Error> {
let cluster = ctx.metrics_cluster_shared.clone().unwrap_or_else(metrics::cluster_none);
if let Some(start) = ctx.upstream_connect_start.take() {
metrics::record_upstream_connect_duration(cluster.clone(), start.elapsed().as_secs_f64());
}
metrics::record_upstream_connect_failure(cluster.clone());
let policy = ctx.retry_policy.clone().unwrap_or_else(legacy_default_policy);
let outcome = retry::classify_error(&e);
let decision = retry::should_retry(ctx, &policy, outcome, ctx.cluster_retry_state.as_deref());
match decision {
retry::RetryDecision::Retry { backoff } => {
ctx.retries += 1;
ctx.pending_backoff = Some(backoff);
ctx.reselect_on_retry = policy.configured;
if let Some(upstream) = ctx.upstream_for_retry.as_ref() {
let addr = Arc::clone(&upstream.address);
if !ctx.attempted_endpoints.iter().any(|e| e.as_ref() == addr.as_ref()) {
ctx.attempted_endpoints.push(addr);
}
}
if policy.configured {
if let Some(upstream) = ctx.upstream_for_retry.as_ref()
&& let Some(reselector) = ctx.endpoint_reselector.as_ref()
{
reselector.release(&upstream.address);
}
ctx.upstream_for_retry = None;
}
let upstream_address = ctx
.upstream_for_retry
.as_ref()
.map_or("unknown", |u| u.address.as_ref());
debug!(
retries = ctx.retries,
max = policy.effective_max_retries(),
?backoff,
upstream_address,
"retrying after connect failure"
);
let mut e = e;
e.set_retry(true);
e
},
retry::RetryDecision::DoNotRetry => {
if ctx.retries > 0 {
warn!(
retries = ctx.retries,
max = policy.effective_max_retries(),
upstream_address = ctx
.upstream_for_retry
.as_ref()
.map_or("unknown", |u| u.address.as_ref()),
"retry limit exhausted"
);
}
record_retry_exhausted_if_attempted(ctx, cluster);
let mut e = e;
e.set_retry(false);
e
},
}
}
#[expect(clippy::too_many_lines, reason = "sequential guard checks")]
fn maybe_retry_response(ctx: &mut PingoraRequestCtx, status: u16) -> Option<Box<pingora_core::Error>> {
let policy = ctx.retry_policy.clone().unwrap_or_else(legacy_default_policy);
let outcome = retry::RetryOutcome::StatusCode(status);
let decision = retry::should_retry(ctx, &policy, outcome, ctx.cluster_retry_state.as_deref());
match decision {
retry::RetryDecision::Retry { backoff } => {
ctx.retries += 1;
ctx.pending_backoff = Some(backoff);
ctx.reselect_on_retry = policy.configured;
if let Some(upstream) = ctx.upstream_for_retry.as_ref() {
let addr = Arc::clone(&upstream.address);
if !ctx.attempted_endpoints.iter().any(|e| e.as_ref() == addr.as_ref()) {
ctx.attempted_endpoints.push(addr);
}
}
if policy.configured {
if let Some(upstream) = ctx.upstream_for_retry.as_ref()
&& let Some(reselector) = ctx.endpoint_reselector.as_ref()
{
reselector.release(&upstream.address);
}
ctx.upstream_for_retry = None;
}
debug!(
status,
retries = ctx.retries,
max = policy.effective_max_retries(),
?backoff,
"retrying after retriable response status"
);
let mut e =
pingora_core::Error::explain(pingora_core::ErrorType::HTTPStatus(status), "retriable upstream status");
e.set_retry(true);
Some(e)
},
retry::RetryDecision::DoNotRetry => None,
}
}
fn release_retry_state(ctx: &mut PingoraRequestCtx) {
if !ctx.cluster_retry_state_released
&& let Some(state) = ctx.cluster_retry_state.take()
{
state.leave();
ctx.cluster_retry_state_released = true;
}
}
fn record_retry_exhausted_if_attempted(ctx: &PingoraRequestCtx, cluster: ::metrics::SharedString) {
if ctx.retries > 0 {
metrics::record_upstream_retry(cluster, metrics::RETRY_RESULT_EXHAUSTED);
}
}
fn maybe_emit_fallback_access_log(pipeline: &FilterPipeline, status: u16, ctx: &mut PingoraRequestCtx) {
if ctx.response_delivery_complete || ctx.connection_upgraded || !pipeline.contains_filter("access_log") {
return;
}
if let Some(filter_ctx) = ctx.filter_context_for(pipeline, None) {
if praxis_filter::access_record_already_emitted(&filter_ctx) {
return;
}
if !pipeline.filter_request_conditions_match("access_log", filter_ctx.request) {
return;
}
praxis_filter::emit_access_record(&filter_ctx, status);
}
}
async fn logging_cleanup(pipeline: &FilterPipeline, ctx: &mut PingoraRequestCtx) {
if !ctx.response_phase_done
&& let Some(mut filter_ctx) = ctx.filter_context_for(pipeline, None)
{
let _result = pipeline.execute_http_response(&mut filter_ctx).await;
let extensions = filter_ctx.extensions;
let metadata = filter_ctx.filter_metadata;
let state = filter_ctx.filter_state;
let exec_idx = filter_ctx.executed_filter_indices;
let body_idx = filter_ctx.body_done_indices;
let cluster = filter_ctx.cluster;
let upstream = filter_ctx.upstream;
ctx.extensions = extensions;
ctx.filter_metadata = metadata;
ctx.filter_state = state;
ctx.cached_executed_filter_indices = exec_idx;
ctx.cached_body_done_indices = body_idx;
ctx.cluster = cluster;
ctx.upstream = upstream;
}
}
fn emit_request_metrics(session: &Session, ctx: &PingoraRequestCtx) {
if !metrics::is_recorder_installed() {
return;
}
let status_code = session.response_written().map_or(0, |resp| resp.status.as_u16());
let status_class = metrics::status_class(status_code);
let method = request_method_label(session, ctx);
let cluster = ctx.metrics_cluster_shared.clone().unwrap_or_else(metrics::cluster_none);
let route = ctx.metrics_route.clone().unwrap_or_else(metrics::route_unknown);
let labels = metrics::RequestMetricLabels {
cluster: cluster.clone(),
method,
route,
status_class,
};
emit_upstream_request_metric(ctx, &cluster);
if let Some(error_type) = ctx.error_type {
metrics::record_error(error_type);
}
let duration_secs = ctx.request_start.elapsed().as_secs_f64();
metrics::record_request_metrics(labels, duration_secs);
metrics::record_body_size_metrics(
method,
status_class,
cluster,
ctx.request_body_bytes,
ctx.response_body_bytes,
);
}
fn request_method_label(session: &Session, ctx: &PingoraRequestCtx) -> &'static str {
let request_method = session.req_header().method.as_str();
let raw_method = if request_method.is_empty() {
ctx.request_snapshot.as_ref().map_or("UNKNOWN", |r| r.method.as_str())
} else {
request_method
};
metrics::method_label(raw_method)
}
fn emit_upstream_request_metric(ctx: &PingoraRequestCtx, cluster: &::metrics::SharedString) {
if let Some(upstream_status) = ctx.upstream_response_status
&& let Some(upstream) = ctx.upstream_for_retry.as_ref()
{
metrics::record_upstream_request(
cluster.clone(),
::metrics::SharedString::from(Arc::clone(&upstream.address)),
metrics::status_class(upstream_status),
);
}
}
fn record_passive_health(pipeline: &FilterPipeline, error: Option<&pingora_core::Error>, ctx: &PingoraRequestCtx) {
let cluster_name = ctx.cluster.as_ref().or(ctx.metrics_cluster.as_ref());
let Some(cluster_name) = cluster_name else {
return;
};
let Some(idx) = ctx.selected_endpoint_index else {
return;
};
let Some(registry) = pipeline.health_registry() else {
return;
};
let Some(health) = registry.get(cluster_name) else {
return;
};
if !ctx.upstream_contacted {
return;
}
let is_downstream_error = error.is_some_and(|e| matches!(e.esource(), pingora_core::ErrorSource::Downstream));
if is_downstream_error && ctx.upstream_response_status.is_none() {
return;
}
let is_failure =
ctx.upstream_response_status.is_some_and(|s| s >= 500) || (error.is_some() && !is_downstream_error);
apply_passive_threshold(health, idx, cluster_name, is_failure);
}
fn apply_passive_threshold(
health: &praxis_core::health::ClusterHealthEntry,
idx: usize,
cluster_name: &Arc<str>,
is_failure: bool,
) {
if is_failure {
if let Some(threshold) = health.passive_unhealthy_threshold()
&& health
.endpoints()
.get(idx)
.is_some_and(|ep| ep.record_failure(threshold))
{
tracing::warn!(
cluster = %cluster_name,
endpoint_index = idx,
threshold,
"passive health: endpoint marked unhealthy"
);
emit_passive_health_transition(health, cluster_name, metrics::HEALTH_RESULT_UNHEALTHY);
}
} else if let Some(threshold) = health.passive_healthy_threshold()
&& health
.endpoints()
.get(idx)
.is_some_and(|ep| ep.record_success(threshold))
{
tracing::info!(
cluster = %cluster_name,
endpoint_index = idx,
threshold,
"passive health: endpoint recovered"
);
emit_passive_health_transition(health, cluster_name, metrics::HEALTH_RESULT_HEALTHY);
}
}
fn emit_passive_health_transition(
health: &praxis_core::health::ClusterHealthEntry,
cluster_name: &Arc<str>,
result: &'static str,
) {
let (healthy, total) = metrics::count_healthy_endpoints(health);
metrics::record_health_transition(
::metrics::SharedString::from(Arc::clone(cluster_name)),
result,
healthy,
total,
);
}
pub(super) fn http_version_label(version: http::Version) -> &'static str {
match version {
http::Version::HTTP_09 => "0.9",
http::Version::HTTP_10 => "1.0",
http::Version::HTTP_11 => "1.1",
http::Version::HTTP_2 => "2",
http::Version::HTTP_3 => "3",
_ => "unknown",
}
}
fn record_response_span_attributes(session: &Session, ctx: &PingoraRequestCtx) {
if ctx.request_span.is_disabled() {
return;
}
let response = session.response_written();
let status = response.map(|resp| resp.status);
let method = session.req_header().method.as_str();
record_response_span_fields(status, method, response, ctx);
}
fn record_response_span_fields(
status: Option<http::StatusCode>,
method: &str,
response: Option<&pingora_http::ResponseHeader>,
ctx: &PingoraRequestCtx,
) {
if let Some(status) = status {
let code = status.as_u16();
if code > 0 {
ctx.request_span.record("http.response.status_code", code);
}
if status.is_server_error() {
ctx.request_span.record("otel.status_code", "ERROR");
ctx.request_span.record("error.type", code.to_string().as_str());
}
}
if let Some(route) = &ctx.metrics_route {
ctx.request_span.record("http.route", route.as_ref());
ctx.request_span
.record("otel.name", format!("{method} {route}").as_str());
}
if let Some(upstream) = &ctx.upstream_for_retry {
ctx.request_span.record("upstream.address", upstream.address.as_ref());
}
if let Some(cluster) = &ctx.metrics_cluster {
ctx.request_span.record("upstream.cluster", cluster.as_ref());
}
record_upstream_exchange_span(ctx, response);
}
fn record_upstream_exchange_span(ctx: &PingoraRequestCtx, response: Option<&pingora_http::ResponseHeader>) {
if ctx.upstream_exchange_span.is_disabled() {
return;
}
if let Some(status) = ctx
.upstream_response_status
.or_else(|| response.map(|resp| resp.status.as_u16()))
{
ctx.upstream_exchange_span.record("http.response.status_code", status);
}
ctx.upstream_exchange_span
.record("http.response.body.size", ctx.response_body_bytes);
}
fn h2c_server_options() -> HttpServerOptions {
let mut opts = HttpServerOptions::default();
opts.h2c = true;
opts
}
fn h2_server_options() -> H2Options {
let mut opts = H2Options::new();
opts.max_header_list_size(65_536); opts.max_concurrent_streams(128);
opts
}
fn check_body_size_limit(body: Option<&Bytes>, accumulated_bytes: &mut u64, max_bytes: usize) -> bool {
if let Some(chunk) = body {
let chunk_len = chunk.len() as u64;
*accumulated_bytes += chunk_len;
let limit = max_bytes as u64;
return *accumulated_bytes > limit;
}
false
}
fn accumulate_stream_buffer(
body: &mut Option<Bytes>,
body_buffer: &mut Option<BodyBuffer>,
end_of_stream: bool,
max_bytes: Option<usize>,
) -> bool {
if let Some(chunk) = &*body {
let limit = max_bytes.unwrap_or(ABSOLUTE_MAX_BODY_BYTES);
let buf = body_buffer.get_or_insert_with(|| BodyBuffer::new(limit));
if buf.push(chunk.clone()).is_err() {
return true;
}
}
if end_of_stream {
tracing::trace!("stream buffer: freezing accumulated body before pipeline at EOS");
*body = body_buffer.take().map(BodyBuffer::freeze);
} else {
tracing::trace!("stream buffer: filters see the original chunk");
}
false
}
#[expect(
clippy::fn_params_excessive_bools,
reason = "mirrors the caller's existing condition flags"
)]
fn suppress_stream_buffer_chunk(body: &mut Option<Bytes>, is_stream_buffer: bool, released: bool, end_of_stream: bool) {
if is_stream_buffer && !released && !end_of_stream {
*body = None;
}
}
fn release_stream_buffer(
body: &mut Option<Bytes>,
is_stream_buffer: bool,
released: &mut bool,
body_buffer: &mut Option<BodyBuffer>,
end_of_stream: bool,
) {
if is_stream_buffer && !*released {
*released = true;
if !end_of_stream {
*body = body_buffer.take().map(BodyBuffer::freeze);
}
}
}
struct BodyFilterOutput {
cluster: Option<Arc<str>>,
upstream: Option<Upstream>,
extensions: RequestExtensions,
filter_metadata: HashMap<String, String>,
filter_state: HashMap<usize, Box<dyn std::any::Any + Send + Sync>>,
executed_filter_indices: Vec<bool>,
body_done_indices: Vec<bool>,
attempted_endpoints: Vec<Arc<str>>,
}
impl BodyFilterOutput {
fn take_from(fctx: &mut HttpFilterContext<'_>) -> Self {
Self {
cluster: fctx.cluster.take(),
upstream: fctx.upstream.take(),
extensions: std::mem::take(&mut fctx.extensions),
filter_metadata: std::mem::take(&mut fctx.filter_metadata),
filter_state: std::mem::take(&mut fctx.filter_state),
executed_filter_indices: std::mem::take(&mut fctx.executed_filter_indices),
body_done_indices: std::mem::take(&mut fctx.body_done_indices),
attempted_endpoints: std::mem::take(&mut fctx.attempted_endpoints),
}
}
fn write_back(self, ctx: &mut PingoraRequestCtx) {
ctx.cluster = self.cluster;
ctx.upstream = self.upstream;
ctx.extensions = self.extensions;
ctx.filter_metadata = self.filter_metadata;
ctx.filter_state = self.filter_state;
ctx.cached_executed_filter_indices = self.executed_filter_indices;
ctx.cached_body_done_indices = self.body_done_indices;
ctx.attempted_endpoints = self.attempted_endpoints;
}
}
#[cfg(test)]
#[expect(clippy::allow_attributes, reason = "blanket test suppressions")]
#[allow(
clippy::unwrap_used,
clippy::expect_used,
clippy::indexing_slicing,
clippy::field_reassign_with_default,
clippy::too_many_lines,
clippy::cast_possible_truncation,
clippy::significant_drop_tightening,
reason = "tests"
)]
mod tests {
use praxis_core::connectivity::ConnectionOptions;
use super::*;
const MAX_RETRIES: usize = praxis_core::config::DEFAULT_MAX_RETRIES as usize;
const RETRY_BODY_LIMIT: u64 = praxis_core::config::DEFAULT_RETRY_BODY_LIMIT_BYTES;
#[test]
fn first_failure_idempotent_sets_retry() {
let mut ctx = PingoraRequestCtx::default();
ctx.request_is_idempotent = true;
let e = handle_connect_failure(&mut ctx, make_error());
assert!(e.retry(), "first failure should set retry flag");
assert_eq!(ctx.retries, 1);
}
#[test]
fn large_body_skips_retry() {
let mut ctx = PingoraRequestCtx::default();
ctx.request_is_idempotent = true;
ctx.request_body_bytes = RETRY_BODY_LIMIT + 1;
let e = handle_connect_failure(&mut ctx, make_error());
assert!(!e.retry(), "should not retry when body exceeds retry buffer limit");
assert_eq!(ctx.retries, 0, "retry counter should not increment");
}
#[test]
fn mutated_body_exceeding_limit_skips_retry() {
let mut ctx = PingoraRequestCtx::default();
ctx.request_is_idempotent = true;
ctx.request_body_bytes = 1024;
ctx.mutated_request_body_len = Some((RETRY_BODY_LIMIT + 1) as usize);
let e = handle_connect_failure(&mut ctx, make_error());
assert!(
!e.retry(),
"should not retry when mutated body exceeds retry buffer limit"
);
assert_eq!(ctx.retries, 0);
}
#[test]
fn body_at_limit_allows_retry() {
let mut ctx = PingoraRequestCtx::default();
ctx.request_is_idempotent = true;
ctx.request_body_bytes = RETRY_BODY_LIMIT;
let e = handle_connect_failure(&mut ctx, make_error());
assert!(e.retry(), "body exactly at limit should allow retry");
assert_eq!(ctx.retries, 1);
}
#[test]
fn zero_body_allows_retry() {
let mut ctx = PingoraRequestCtx::default();
ctx.request_is_idempotent = true;
ctx.request_body_bytes = 0;
let e = handle_connect_failure(&mut ctx, make_error());
assert!(e.retry(), "zero-length body should allow retry");
assert_eq!(ctx.retries, 1);
}
#[test]
fn max_retries_exhausted_does_not_retry() {
let mut ctx = PingoraRequestCtx::default();
ctx.request_is_idempotent = true;
ctx.retries = MAX_RETRIES as u32;
let e = handle_connect_failure(&mut ctx, make_error());
assert!(!e.retry(), "should not retry after MAX_RETRIES");
assert_eq!(ctx.retries as usize, MAX_RETRIES);
}
#[test]
fn counter_increments_across_calls() {
let mut ctx = PingoraRequestCtx::default();
ctx.request_is_idempotent = true;
for expected in 1..=MAX_RETRIES {
let _result = handle_connect_failure(&mut ctx, make_error());
assert_eq!(ctx.retries as usize, expected);
}
let e = handle_connect_failure(&mut ctx, make_error());
assert!(!e.retry(), "should not retry after reaching MAX_RETRIES");
assert_eq!(ctx.retries as usize, MAX_RETRIES);
}
#[test]
fn non_idempotent_request_never_retries() {
let mut ctx = PingoraRequestCtx::default();
ctx.request_is_idempotent = false;
let e = handle_connect_failure(&mut ctx, make_error());
assert!(!e.retry(), "non-idempotent request should never retry");
assert_eq!(ctx.retries, 0);
}
#[test]
fn connect_failure_clears_upstream_connect_start() {
let mut ctx = PingoraRequestCtx::default();
ctx.upstream_connect_start = Some(std::time::Instant::now());
let _e = handle_connect_failure(&mut ctx, make_error());
assert!(
ctx.upstream_connect_start.is_none(),
"failed connect should consume upstream_connect_start for duration recording"
);
}
#[test]
fn response_503_retries_when_status5xx_enabled() {
let mut ctx = PingoraRequestCtx::default();
ctx.request_is_idempotent = true;
ctx.retry_policy = Some(Arc::new(praxis_core::config::RetryPolicy {
configured: true,
retriable_conditions: vec![praxis_core::config::RetriableCondition::Status5xx],
..praxis_core::config::RetryPolicy::legacy_default()
}));
let e = maybe_retry_response(&mut ctx, 503).expect("503 should be retriable");
assert!(e.retry(), "503 under Status5xx should set retry");
assert_eq!(ctx.retries, 1);
assert!(ctx.reselect_on_retry);
assert!(ctx.pending_backoff.is_some());
}
#[test]
fn response_404_does_not_retry() {
let mut ctx = PingoraRequestCtx::default();
ctx.request_is_idempotent = true;
ctx.retry_policy = Some(Arc::new(praxis_core::config::RetryPolicy {
retriable_conditions: vec![praxis_core::config::RetriableCondition::Status5xx],
..praxis_core::config::RetryPolicy::legacy_default()
}));
assert!(
maybe_retry_response(&mut ctx, 404).is_none(),
"404 must never trigger status-based retry"
);
assert_eq!(ctx.retries, 0);
}
#[test]
fn response_502_does_not_retry_under_legacy_default() {
let mut ctx = PingoraRequestCtx::default();
ctx.request_is_idempotent = true;
assert!(
maybe_retry_response(&mut ctx, 502).is_none(),
"legacy default must forward 5xx without retry"
);
}
#[test]
fn max_retries_zero_disables_connect_retry() {
let mut ctx = PingoraRequestCtx::default();
ctx.request_is_idempotent = true;
ctx.retry_policy = Some(Arc::new(praxis_core::config::RetryPolicy {
max_retries: Some(0),
..praxis_core::config::RetryPolicy::legacy_default()
}));
let e = handle_connect_failure(&mut ctx, make_error());
assert!(!e.retry(), "max_retries: 0 must disable retries");
assert_eq!(ctx.retries, 0);
}
#[test]
fn non_idempotent_clears_pingora_default_retry_flag() {
let mut ctx = PingoraRequestCtx::default();
ctx.request_is_idempotent = false;
let mut e = make_error();
e.set_retry(true);
let e = handle_connect_failure(&mut ctx, e);
assert!(!e.retry(), "policy denial must clear Pingora's default retry flag");
assert_eq!(ctx.retries, 0);
}
#[tokio::test]
async fn logging_cleanup_noop_when_response_phase_done() {
let registry = praxis_filter::FilterRegistry::with_builtins();
let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
let mut ctx = PingoraRequestCtx::default();
ctx.response_phase_done = true;
ctx.request_snapshot = Some(praxis_filter::Request {
method: http::Method::GET,
uri: "/".parse().unwrap(),
headers: http::HeaderMap::new(),
});
logging_cleanup(&pipeline, &mut ctx).await;
}
#[tokio::test]
async fn logging_cleanup_noop_when_no_snapshot() {
let registry = praxis_filter::FilterRegistry::with_builtins();
let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
let mut ctx = PingoraRequestCtx::default();
ctx.response_phase_done = false;
ctx.request_snapshot = None;
logging_cleanup(&pipeline, &mut ctx).await;
}
#[tokio::test]
async fn logging_cleanup_runs_response_pipeline_when_needed() {
let registry = praxis_filter::FilterRegistry::with_builtins();
let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
let mut ctx = PingoraRequestCtx::default();
ctx.response_phase_done = false;
ctx.cluster = Some(Arc::from("test-cluster"));
ctx.request_snapshot = Some(praxis_filter::Request {
method: http::Method::GET,
uri: "/test".parse().unwrap(),
headers: http::HeaderMap::new(),
});
logging_cleanup(&pipeline, &mut ctx).await;
assert_eq!(
ctx.cluster.as_deref(),
Some("test-cluster"),
"cluster must be restored so the fallback access record can attribute the failure"
);
}
#[tokio::test]
async fn logging_cleanup_preserves_filter_metadata() {
let registry = praxis_filter::FilterRegistry::with_builtins();
let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
let mut ctx = PingoraRequestCtx::default();
ctx.response_phase_done = false;
ctx.filter_metadata
.insert("json_rpc.method".to_owned(), "service/invoke".to_owned());
ctx.request_snapshot = Some(praxis_filter::Request {
method: http::Method::POST,
uri: "/api".parse().unwrap(),
headers: http::HeaderMap::new(),
});
logging_cleanup(&pipeline, &mut ctx).await;
assert_eq!(
ctx.filter_metadata.get("json_rpc.method").map(String::as_str),
Some("service/invoke"),
"filter_metadata should survive logging_cleanup"
);
}
#[tokio::test]
async fn logging_cleanup_preserves_extensions() {
let registry = praxis_filter::FilterRegistry::with_builtins();
let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
let mut ctx = PingoraRequestCtx::default();
ctx.response_phase_done = false;
ctx.extensions.insert(42_u32);
ctx.request_snapshot = Some(praxis_filter::Request {
method: http::Method::POST,
uri: "/test".parse().unwrap(),
headers: http::HeaderMap::new(),
});
logging_cleanup(&pipeline, &mut ctx).await;
assert_eq!(
ctx.extensions.get::<u32>(),
Some(&42),
"extensions should survive logging_cleanup"
);
}
#[test]
fn passive_health_error_is_failure() {
let (pipeline, ctx) = make_passive_scenario(Some(3), Some(2));
let error = make_error();
record_passive_health(&pipeline, Some(&error), &ctx);
let registry = pipeline.health_registry().unwrap();
let entry = registry.get("test-cluster").unwrap();
assert!(
entry.endpoints()[0].is_healthy(),
"single failure should not yet mark unhealthy (threshold=3)"
);
}
#[test]
fn passive_health_downstream_error_is_not_failure() {
let (pipeline, ctx) = make_passive_scenario(Some(1), Some(1));
let error = make_error().into_down();
record_passive_health(&pipeline, Some(&error), &ctx);
let registry = pipeline.health_registry().unwrap();
let entry = registry.get("test-cluster").unwrap();
assert!(
entry.endpoints()[0].is_healthy(),
"a downstream/client error must not mark the endpoint unhealthy"
);
}
#[test]
fn passive_health_downstream_error_with_5xx_still_counts_as_failure() {
let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
ctx.upstream_response_status = Some(503);
let error = make_error().into_down();
record_passive_health(&pipeline, Some(&error), &ctx);
let registry = pipeline.health_registry().unwrap();
let entry = registry.get("test-cluster").unwrap();
assert!(
!entry.endpoints()[0].is_healthy(),
"a downstream error with an upstream 503 must still mark the endpoint unhealthy"
);
}
#[test]
fn passive_health_downstream_error_does_not_reset_failure_streak() {
let (pipeline, ctx) = make_passive_scenario(Some(2), Some(1));
let mut upstream_err = make_error();
upstream_err.as_up();
let downstream_err = make_error().into_down();
record_passive_health(&pipeline, Some(&upstream_err), &ctx);
record_passive_health(&pipeline, Some(&downstream_err), &ctx);
record_passive_health(&pipeline, Some(&upstream_err), &ctx);
let registry = pipeline.health_registry().unwrap();
let entry = registry.get("test-cluster").unwrap();
assert!(
!entry.endpoints()[0].is_healthy(),
"two upstream failures must eject the endpoint even with an interleaved client error"
);
}
#[test]
fn passive_health_upstream_error_is_failure() {
let (pipeline, ctx) = make_passive_scenario(Some(1), Some(1));
let mut error = make_error();
error.as_up();
record_passive_health(&pipeline, Some(&error), &ctx);
let registry = pipeline.health_registry().unwrap();
let entry = registry.get("test-cluster").unwrap();
assert!(
!entry.endpoints()[0].is_healthy(),
"an upstream error at unhealthy-threshold 1 must mark the endpoint unhealthy"
);
}
#[test]
fn passive_health_skips_observations_without_upstream_contact() {
let (pipeline, mut ctx) = make_passive_scenario(Some(2), Some(1));
let mut upstream_err = make_error();
upstream_err.as_up();
record_passive_health(&pipeline, Some(&upstream_err), &ctx);
ctx.upstream_contacted = false;
record_passive_health(&pipeline, None, &ctx);
ctx.upstream_response_status = Some(200);
record_passive_health(&pipeline, None, &ctx);
ctx.upstream_response_status = None;
ctx.upstream_contacted = true;
record_passive_health(&pipeline, Some(&upstream_err), &ctx);
let registry = pipeline.health_registry().unwrap();
let entry = registry.get("test-cluster").unwrap();
assert!(
!entry.endpoints()[0].is_healthy(),
"observations without upstream contact must not reset the failure streak"
);
}
#[test]
fn passive_health_records_connect_failure_after_reselect_clears_upstream() {
let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
ctx.upstream_for_retry = None;
ctx.upstream_contacted = true;
let mut error = make_error();
error.as_up();
record_passive_health(&pipeline, Some(&error), &ctx);
let registry = pipeline.health_registry().unwrap();
let entry = registry.get("test-cluster").unwrap();
assert!(
!entry.endpoints()[0].is_healthy(),
"a connect failure after reselection cleared upstream_for_retry must still count"
);
}
#[test]
fn passive_health_status_500_is_failure() {
let (pipeline, mut ctx) = make_passive_scenario(Some(3), Some(2));
ctx.upstream_response_status = Some(500);
record_passive_health(&pipeline, None, &ctx);
let registry = pipeline.health_registry().unwrap();
let entry = registry.get("test-cluster").unwrap();
assert!(
entry.endpoints()[0].is_healthy(),
"single 500 should not yet mark unhealthy (threshold=3)"
);
}
#[test]
fn passive_health_status_below_500_is_success() {
let (pipeline, mut ctx) = make_passive_scenario(Some(2), Some(1));
ctx.upstream_response_status = Some(499);
record_passive_health(&pipeline, None, &ctx);
let registry = pipeline.health_registry().unwrap();
let entry = registry.get("test-cluster").unwrap();
assert!(entry.endpoints()[0].is_healthy(), "status 499 should count as success");
}
#[test]
fn passive_unhealthy_threshold_transition() {
let (pipeline, ctx) = make_passive_scenario(Some(2), Some(1));
let error = make_error();
record_passive_health(&pipeline, Some(&error), &ctx);
record_passive_health(&pipeline, Some(&error), &ctx);
let registry = pipeline.health_registry().unwrap();
let entry = registry.get("test-cluster").unwrap();
assert!(
!entry.endpoints()[0].is_healthy(),
"2 consecutive failures should mark unhealthy (threshold=2)"
);
}
#[test]
fn passive_healthy_threshold_recovery() {
let (pipeline, ctx) = make_passive_scenario(Some(1), Some(2));
let error = make_error();
record_passive_health(&pipeline, Some(&error), &ctx);
let registry = pipeline.health_registry().unwrap();
let entry = registry.get("test-cluster").unwrap();
assert!(
!entry.endpoints()[0].is_healthy(),
"should be unhealthy after 1 failure"
);
let ctx_ok = make_passive_ctx("test-cluster", 0, Some(200));
record_passive_health(&pipeline, None, &ctx_ok);
assert!(
!entry.endpoints()[0].is_healthy(),
"one success should not recover (threshold=2)"
);
record_passive_health(&pipeline, None, &ctx_ok);
assert!(
entry.endpoints()[0].is_healthy(),
"2 consecutive successes should recover (threshold=2)"
);
}
#[test]
fn passive_health_no_thresholds_is_noop() {
let (pipeline, ctx) = make_passive_scenario(None, None);
let error = make_error();
record_passive_health(&pipeline, Some(&error), &ctx);
let registry = pipeline.health_registry().unwrap();
let entry = registry.get("test-cluster").unwrap();
assert!(
entry.endpoints()[0].is_healthy(),
"no passive thresholds means failures are no-op"
);
}
#[test]
fn passive_health_endpoint_index_out_of_bounds() {
let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
ctx.selected_endpoint_index = Some(999);
let error = make_error();
record_passive_health(&pipeline, Some(&error), &ctx);
let registry = pipeline.health_registry().unwrap();
let entry = registry.get("test-cluster").unwrap();
assert!(entry.endpoints()[0].is_healthy(), "out-of-bounds index should be no-op");
}
#[test]
fn passive_health_missing_cluster_is_noop() {
let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
ctx.cluster = None;
ctx.metrics_cluster = None;
let error = make_error();
record_passive_health(&pipeline, Some(&error), &ctx);
}
#[test]
fn passive_health_falls_back_to_metrics_cluster() {
let (pipeline, mut ctx) = make_passive_scenario(Some(2), Some(1));
ctx.cluster = None;
ctx.metrics_cluster = Some(Arc::from("test-cluster"));
let error = make_error();
record_passive_health(&pipeline, Some(&error), &ctx);
record_passive_health(&pipeline, Some(&error), &ctx);
let registry = pipeline.health_registry().unwrap();
let entry = registry.get("test-cluster").unwrap();
assert!(
!entry.endpoints()[0].is_healthy(),
"fallback to metrics_cluster should still record passive health"
);
}
#[test]
fn passive_health_missing_endpoint_index_is_noop() {
let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
ctx.selected_endpoint_index = None;
let error = make_error();
record_passive_health(&pipeline, Some(&error), &ctx);
}
#[test]
fn passive_health_missing_registry_is_noop() {
let registry = praxis_filter::FilterRegistry::with_builtins();
let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
let mut ctx = PingoraRequestCtx::default();
ctx.cluster = Some(Arc::from("test-cluster"));
ctx.selected_endpoint_index = Some(0);
let error = make_error();
record_passive_health(&pipeline, Some(&error), &ctx);
}
#[test]
fn passive_health_unknown_cluster_is_noop() {
let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
ctx.cluster = Some(Arc::from("nonexistent"));
let error = make_error();
record_passive_health(&pipeline, Some(&error), &ctx);
}
#[test]
fn size_limit_none_body_returns_false() {
let mut bytes = 0_u64;
assert!(!check_body_size_limit(None, &mut bytes, 100));
assert_eq!(bytes, 0, "accumulated bytes unchanged for None body");
}
#[test]
fn size_limit_within_limit() {
let mut bytes = 0_u64;
let body = Some(Bytes::from_static(b"hello"));
assert!(!check_body_size_limit(body.as_ref(), &mut bytes, 10));
assert_eq!(bytes, 5);
}
#[test]
fn size_limit_at_exact_limit() {
let mut bytes = 0_u64;
let body = Some(Bytes::from_static(b"exact"));
assert!(!check_body_size_limit(body.as_ref(), &mut bytes, 5));
assert_eq!(bytes, 5);
}
#[test]
fn size_limit_exceeds_limit() {
let mut bytes = 0_u64;
let body = Some(Bytes::from_static(b"toolong"));
assert!(check_body_size_limit(body.as_ref(), &mut bytes, 3));
}
#[test]
fn size_limit_cumulative_overflow() {
let mut bytes = 0_u64;
let first = Some(Bytes::from_static(b"aaa"));
assert!(!check_body_size_limit(first.as_ref(), &mut bytes, 5));
let second = Some(Bytes::from_static(b"bbb"));
assert!(check_body_size_limit(second.as_ref(), &mut bytes, 5));
assert_eq!(bytes, 6);
}
#[test]
fn stream_buffer_accumulates_chunks() {
let mut body = Some(Bytes::from_static(b"hello "));
let mut buf: Option<BodyBuffer> = None;
assert!(!accumulate_stream_buffer(&mut body, &mut buf, false, Some(100)));
assert!(buf.is_some());
body = Some(Bytes::from_static(b"world"));
assert!(!accumulate_stream_buffer(&mut body, &mut buf, false, Some(100)));
let frozen = buf.take().unwrap().freeze();
assert_eq!(frozen, Bytes::from_static(b"hello world"));
}
#[test]
fn stream_buffer_freezes_at_eos() {
let mut body = Some(Bytes::from_static(b"data"));
let mut buf: Option<BodyBuffer> = None;
assert!(!accumulate_stream_buffer(&mut body, &mut buf, false, Some(100)));
body = Some(Bytes::from_static(b" end"));
assert!(!accumulate_stream_buffer(&mut body, &mut buf, true, Some(100)));
assert!(buf.is_none(), "buffer should be taken at EOS");
assert_eq!(body.unwrap(), Bytes::from_static(b"data end"));
}
#[test]
fn stream_buffer_overflow() {
let mut body = Some(Bytes::from_static(b"too long"));
let mut buf: Option<BodyBuffer> = None;
assert!(accumulate_stream_buffer(&mut body, &mut buf, false, Some(5)));
}
#[test]
fn stream_buffer_none_body() {
let mut body: Option<Bytes> = None;
let mut buf: Option<BodyBuffer> = None;
assert!(!accumulate_stream_buffer(&mut body, &mut buf, false, Some(100)));
assert!(buf.is_none());
}
#[test]
fn stream_buffer_uses_absolute_max_when_none() {
let mut body = Some(Bytes::from_static(b"data"));
let mut buf: Option<BodyBuffer> = None;
assert!(!accumulate_stream_buffer(&mut body, &mut buf, false, None));
assert!(buf.is_some(), "should create buffer with absolute max");
}
#[test]
fn suppress_clears_body_when_buffering() {
let mut body = Some(Bytes::from_static(b"data"));
suppress_stream_buffer_chunk(&mut body, true, false, false);
assert!(body.is_none());
}
#[test]
fn suppress_noop_when_not_stream_buffer() {
let mut body = Some(Bytes::from_static(b"data"));
suppress_stream_buffer_chunk(&mut body, false, false, false);
assert!(body.is_some());
}
#[test]
fn suppress_noop_when_released() {
let mut body = Some(Bytes::from_static(b"data"));
suppress_stream_buffer_chunk(&mut body, true, true, false);
assert!(body.is_some());
}
#[test]
fn suppress_noop_at_eos() {
let mut body = Some(Bytes::from_static(b"data"));
suppress_stream_buffer_chunk(&mut body, true, false, true);
assert!(body.is_some());
}
#[test]
fn release_sets_flag_and_flushes_buffer() {
let mut body: Option<Bytes> = None;
let mut released = false;
let mut buf = Some(BodyBuffer::new(100));
buf.as_mut().unwrap().push(Bytes::from_static(b"buffered")).unwrap();
release_stream_buffer(&mut body, true, &mut released, &mut buf, false);
assert!(released);
assert_eq!(body.unwrap(), Bytes::from_static(b"buffered"));
assert!(buf.is_none());
}
#[test]
fn release_noop_when_already_released() {
let mut body: Option<Bytes> = None;
let mut released = true;
let mut buf: Option<BodyBuffer> = None;
release_stream_buffer(&mut body, true, &mut released, &mut buf, false);
assert!(body.is_none(), "body should be unchanged when already released");
}
#[test]
fn release_noop_when_not_stream_buffer() {
let mut body: Option<Bytes> = None;
let mut released = false;
let mut buf: Option<BodyBuffer> = None;
release_stream_buffer(&mut body, false, &mut released, &mut buf, false);
assert!(!released, "released flag should be unchanged for non-stream-buffer");
}
#[test]
fn release_at_eos_sets_flag_but_no_flush() {
let mut body: Option<Bytes> = None;
let mut released = false;
let mut buf = Some(BodyBuffer::new(100));
buf.as_mut().unwrap().push(Bytes::from_static(b"data")).unwrap();
release_stream_buffer(&mut body, true, &mut released, &mut buf, true);
assert!(released);
assert!(body.is_none(), "body should not be overwritten at EOS");
assert!(buf.is_some(), "buffer should not be taken at EOS");
}
#[test]
fn write_back_transfers_fields() {
let mut ctx = PingoraRequestCtx::default();
let mut extensions = RequestExtensions::new();
extensions.insert(42_u32);
let state_val: Box<dyn std::any::Any + Send + Sync> = Box::new(99_i32);
let filter_state = HashMap::from([(0_usize, state_val)]);
let output = BodyFilterOutput {
cluster: Some(Arc::from("test-cluster")),
upstream: Some(Upstream {
address: Arc::from("10.0.0.1:80"),
authority: None,
connection: Arc::new(ConnectionOptions::default()),
tls: None,
}),
extensions,
attempted_endpoints: Vec::new(),
filter_metadata: HashMap::from([("key".to_owned(), "val".to_owned())]),
filter_state,
executed_filter_indices: vec![true, false],
body_done_indices: vec![false, true],
};
output.write_back(&mut ctx);
assert_eq!(ctx.cluster.as_deref(), Some("test-cluster"));
assert!(ctx.upstream.is_some(), "upstream should transfer");
assert_eq!(ctx.upstream.as_ref().unwrap().address.as_ref(), "10.0.0.1:80");
assert_eq!(ctx.extensions.get::<u32>(), Some(&42));
assert_eq!(ctx.filter_metadata.get("key").map(String::as_str), Some("val"));
assert_eq!(ctx.filter_state.len(), 1, "filter_state should transfer");
assert_eq!(
ctx.filter_state.get(&0).and_then(|v| v.downcast_ref::<i32>()),
Some(&99)
);
assert_eq!(ctx.cached_executed_filter_indices, vec![true, false]);
assert_eq!(ctx.cached_body_done_indices, vec![false, true]);
}
#[test]
fn fallback_access_log_emits_for_incomplete_request() {
let pipeline = access_log_pipeline();
let mut ctx = make_fallback_ctx();
let events = capture_access_events(|| maybe_emit_fallback_access_log(&pipeline, 502, &mut ctx));
assert_eq!(
events.len(),
1,
"incomplete request must produce a fallback access record"
);
}
#[test]
fn fallback_access_log_skips_completed_delivery() {
let pipeline = access_log_pipeline();
let mut ctx = make_fallback_ctx();
ctx.response_delivery_complete = true;
let events = capture_access_events(|| maybe_emit_fallback_access_log(&pipeline, 200, &mut ctx));
assert!(events.is_empty(), "completed delivery already logged via the filter");
}
#[test]
fn fallback_access_log_skips_upgraded_connections() {
let pipeline = access_log_pipeline();
let mut ctx = make_fallback_ctx();
ctx.connection_upgraded = true;
let events = capture_access_events(|| maybe_emit_fallback_access_log(&pipeline, 101, &mut ctx));
assert!(events.is_empty(), "upgraded connections have no body completion");
}
#[test]
fn fallback_access_log_skips_without_access_log_filter() {
let registry = praxis_filter::FilterRegistry::with_builtins();
let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
let mut ctx = make_fallback_ctx();
let events = capture_access_events(|| maybe_emit_fallback_access_log(&pipeline, 502, &mut ctx));
assert!(events.is_empty(), "no access_log filter means no fallback record");
}
#[test]
fn fallback_access_log_honors_entry_conditions() {
let registry = praxis_filter::FilterRegistry::with_builtins();
let mut entries = vec![praxis_filter::FilterEntry {
branch_chains: None,
conditions: vec![serde_yaml::from_str("when:\n path_prefix: /api\n").unwrap()],
failure_mode: praxis_filter::FailureMode::default(),
filter_type: "access_log".to_owned(),
config: serde_yaml::Value::Null,
name: None,
response_conditions: vec![],
}];
let pipeline = FilterPipeline::build(&mut entries, ®istry).unwrap();
let mut excluded = make_fallback_ctx();
let events = capture_access_events(|| maybe_emit_fallback_access_log(&pipeline, 502, &mut excluded));
assert!(
events.is_empty(),
"requests the operator scoped out must not gain fallback records"
);
let mut included = make_fallback_ctx();
if let Some(snapshot) = included.request_snapshot.as_mut() {
snapshot.uri = "/api/users".parse().unwrap();
}
let events = capture_access_events(|| maybe_emit_fallback_access_log(&pipeline, 502, &mut included));
assert_eq!(events.len(), 1, "in-scope incomplete requests still get a record");
}
#[test]
fn aborted_response_body_at_eos_is_not_marked_delivered() {
let pipeline = access_log_pipeline();
let mut ctx = make_fallback_ctx();
ctx.response_body_mode = BodyMode::SizeLimit { max_bytes: 4 };
let mut body = Some(Bytes::from_static(b"exceeds the limit"));
let result = response_body_filter::execute(&pipeline, &mut body, true, &mut ctx);
assert!(result.is_err(), "over-limit body must abort");
assert!(
!ctx.response_delivery_complete,
"a response aborted at end-of-stream was not delivered; the fallback record must fire"
);
}
#[test]
fn response_body_eos_marks_delivery_complete() {
let registry = praxis_filter::FilterRegistry::with_builtins();
let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
let mut ctx = PingoraRequestCtx::default();
let mut body: Option<Bytes> = None;
let _timeout = response_body_filter::execute(&pipeline, &mut body, false, &mut ctx).unwrap();
assert!(
!ctx.response_delivery_complete,
"mid-stream chunks must not mark delivery complete"
);
let _timeout = response_body_filter::execute(&pipeline, &mut body, true, &mut ctx).unwrap();
assert!(ctx.response_delivery_complete, "end-of-stream marks delivery complete");
}
#[test]
fn http_version_label_http_09() {
assert_eq!(
http_version_label(http::Version::HTTP_09),
"0.9",
"HTTP/0.9 should map to '0.9'"
);
}
#[test]
fn http_version_label_http_10() {
assert_eq!(
http_version_label(http::Version::HTTP_10),
"1.0",
"HTTP/1.0 should map to '1.0'"
);
}
#[test]
fn http_version_label_http_11() {
assert_eq!(
http_version_label(http::Version::HTTP_11),
"1.1",
"HTTP/1.1 should map to '1.1'"
);
}
#[test]
fn http_version_label_http_2() {
assert_eq!(
http_version_label(http::Version::HTTP_2),
"2",
"HTTP/2 should map to '2'"
);
}
#[test]
fn http_version_label_http_3() {
assert_eq!(
http_version_label(http::Version::HTTP_3),
"3",
"HTTP/3 should map to '3'"
);
}
#[test]
fn record_response_span_attributes_noop_for_disabled_span() {
let ctx = PingoraRequestCtx::default();
assert!(ctx.request_span.is_disabled(), "default span should be disabled");
}
#[derive(Clone, Default)]
struct RecordCapture(Arc<std::sync::Mutex<Vec<(String, String)>>>);
impl<S> tracing_subscriber::Layer<S> for RecordCapture
where
S: tracing::Subscriber + for<'a> tracing_subscriber::registry::LookupSpan<'a>,
{
fn on_record(
&self,
_id: &tracing::span::Id,
values: &tracing::span::Record<'_>,
_ctx: tracing_subscriber::layer::Context<'_, S>,
) {
struct Visitor<'a>(&'a mut Vec<(String, String)>);
impl tracing::field::Visit for Visitor<'_> {
fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
self.0.push((field.name().to_owned(), format!("{value:?}")));
}
}
let mut captured = self.0.lock().expect("capture lock");
values.record(&mut Visitor(&mut captured));
}
}
#[test]
fn record_response_span_fields_records_status_upstream_and_cluster() {
use tracing_subscriber::layer::SubscriberExt as _;
let capture = RecordCapture::default();
let subscriber = tracing_subscriber::registry().with(capture.clone());
let _guard = tracing::subscriber::set_default(subscriber);
let mut ctx = PingoraRequestCtx::default();
ctx.metrics_cluster = Some(Arc::from("api-cluster"));
ctx.upstream_for_retry = Some(Upstream {
address: Arc::from("10.0.0.1:80"),
authority: None,
connection: Arc::new(ConnectionOptions::default()),
tls: None,
});
ctx.request_span = tracing::info_span!(
"test_span",
"http.response.status_code" = tracing::field::Empty,
"otel.status_code" = tracing::field::Empty,
"upstream.address" = tracing::field::Empty,
"upstream.cluster" = tracing::field::Empty,
);
record_response_span_fields(Some(http::StatusCode::SERVICE_UNAVAILABLE), "GET", None, &ctx);
let captured = capture.0.lock().expect("capture lock");
let get = |name: &str| {
let value = captured.iter().find(|(f, _)| f == name).map(|(_, v)| v.clone());
assert!(value.is_some(), "field {name} not recorded; got {captured:?}");
value.unwrap_or_default()
};
assert_eq!(get("http.response.status_code"), "503");
assert_eq!(get("otel.status_code"), "\"ERROR\"", "5xx must set otel error status");
assert_eq!(get("upstream.address"), "\"10.0.0.1:80\"");
assert_eq!(get("upstream.cluster"), "\"api-cluster\"");
}
#[test]
fn record_response_span_fields_success_has_no_error_status() {
use tracing_subscriber::layer::SubscriberExt as _;
let capture = RecordCapture::default();
let subscriber = tracing_subscriber::registry().with(capture.clone());
let _guard = tracing::subscriber::set_default(subscriber);
let mut ctx = PingoraRequestCtx::default();
ctx.request_span = tracing::info_span!(
"test_span",
"http.response.status_code" = tracing::field::Empty,
"otel.status_code" = tracing::field::Empty,
);
record_response_span_fields(Some(http::StatusCode::OK), "GET", None, &ctx);
let captured = capture.0.lock().expect("capture lock");
assert!(
captured
.iter()
.any(|(f, v)| f == "http.response.status_code" && v == "200"),
"status should be recorded: {captured:?}"
);
assert!(
!captured.iter().any(|(f, _)| f == "otel.status_code"),
"2xx must not set otel error status: {captured:?}"
);
}
#[test]
fn record_response_span_attributes_records_exchange_span_fields() {
let mut ctx = PingoraRequestCtx::default();
ctx.request_span = tracing::info_span!(
"test_request",
"http.response.status_code" = tracing::field::Empty,
"server.address" = tracing::field::Empty,
"upstream.cluster" = tracing::field::Empty,
);
ctx.upstream_exchange_span = tracing::info_span!(
parent: &ctx.request_span,
"upstream_exchange",
"http.response.status_code" = tracing::field::Empty,
"http.response.body.size" = tracing::field::Empty,
);
ctx.response_body_bytes = 4096;
ctx.upstream_exchange_span.record("http.response.status_code", 200_u16);
ctx.upstream_exchange_span.record("http.response.body.size", 4096_u64);
}
#[test]
fn record_response_span_attributes_skips_exchange_when_disabled() {
let mut ctx = PingoraRequestCtx::default();
ctx.request_span = tracing::info_span!(
"test_request",
"http.response.status_code" = tracing::field::Empty,
"server.address" = tracing::field::Empty,
"upstream.cluster" = tracing::field::Empty,
);
assert!(
ctx.upstream_exchange_span.is_disabled(),
"exchange span should be disabled by default"
);
}
#[test]
fn retry_with_upstream_address_sets_retry_flag() {
let mut ctx = PingoraRequestCtx::default();
ctx.request_is_idempotent = true;
ctx.upstream_for_retry = Some(Upstream {
address: Arc::from("10.0.0.1:8080"),
connection: Arc::new(ConnectionOptions::default()),
tls: None,
authority: None,
});
let e = handle_connect_failure(&mut ctx, make_error());
assert!(e.retry(), "should retry with upstream address present");
assert_eq!(ctx.retries, 1, "retry counter should increment to 1");
}
#[test]
fn retry_without_upstream_address_uses_fallback() {
let mut ctx = PingoraRequestCtx::default();
ctx.request_is_idempotent = true;
ctx.upstream_for_retry = None;
let e = handle_connect_failure(&mut ctx, make_error());
assert!(
e.retry(),
"should retry even when upstream_for_retry is None (address defaults to unknown)"
);
assert_eq!(ctx.retries, 1, "retry counter should increment to 1");
}
#[test]
fn retry_exhausted_with_upstream_address_does_not_retry() {
let mut ctx = PingoraRequestCtx::default();
ctx.request_is_idempotent = true;
ctx.retries = MAX_RETRIES as u32;
ctx.upstream_for_retry = Some(Upstream {
address: Arc::from("10.0.0.2:443"),
connection: Arc::new(ConnectionOptions::default()),
tls: None,
authority: None,
});
let e = handle_connect_failure(&mut ctx, make_error());
assert!(
!e.retry(),
"should not retry after MAX_RETRIES even with upstream address"
);
}
#[test]
fn large_body_skip_with_upstream_address() {
let mut ctx = PingoraRequestCtx::default();
ctx.request_is_idempotent = true;
ctx.request_body_bytes = RETRY_BODY_LIMIT + 1;
ctx.upstream_for_retry = Some(Upstream {
address: Arc::from("10.0.0.3:8080"),
connection: Arc::new(ConnectionOptions::default()),
tls: None,
authority: None,
});
let e = handle_connect_failure(&mut ctx, make_error());
assert!(!e.retry(), "should not retry large body even with upstream address");
assert_eq!(ctx.retries, 0, "retry counter should not increment");
}
fn make_error() -> Box<pingora_core::Error> {
pingora_core::Error::explain(pingora_core::ErrorType::ConnectError, "test connect failure")
}
fn access_log_pipeline() -> FilterPipeline {
let registry = praxis_filter::FilterRegistry::with_builtins();
let mut entries = vec![praxis_filter::FilterEntry {
branch_chains: None,
conditions: vec![],
failure_mode: praxis_filter::FailureMode::default(),
filter_type: "access_log".to_owned(),
config: serde_yaml::Value::Null,
name: None,
response_conditions: vec![],
}];
FilterPipeline::build(&mut entries, ®istry).unwrap()
}
fn make_fallback_ctx() -> PingoraRequestCtx {
let mut ctx = PingoraRequestCtx::default();
ctx.request_snapshot = Some(praxis_filter::Request {
method: http::Method::GET,
uri: "/incomplete".parse().unwrap(),
headers: http::HeaderMap::new(),
});
ctx
}
fn capture_access_events<F: FnOnce()>(f: F) -> Vec<String> {
use tracing_subscriber::layer::SubscriberExt as _;
let messages = Arc::new(std::sync::Mutex::new(Vec::<String>::new()));
let capture = AccessCapture(Arc::clone(&messages));
let subscriber = tracing_subscriber::registry().with(capture);
tracing::subscriber::with_default(subscriber, f);
let mut guard = messages.lock().unwrap();
std::mem::take(&mut *guard)
}
struct AccessCapture(Arc<std::sync::Mutex<Vec<String>>>);
impl<S: tracing::Subscriber> tracing_subscriber::Layer<S> for AccessCapture {
fn on_event(&self, event: &tracing::Event<'_>, _ctx: tracing_subscriber::layer::Context<'_, S>) {
let mut visitor = AccessMessageVisitor(String::new());
event.record(&mut visitor);
if visitor.0.contains("access") {
self.0.lock().unwrap().push(visitor.0);
}
}
}
struct AccessMessageVisitor(String);
impl tracing::field::Visit for AccessMessageVisitor {
fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
if field.name() == "message" {
self.0 = format!("{value:?}");
}
}
}
fn make_passive_ctx(cluster: &str, endpoint_idx: usize, status: Option<u16>) -> PingoraRequestCtx {
let mut ctx = PingoraRequestCtx::default();
ctx.cluster = Some(Arc::from(cluster));
ctx.selected_endpoint_index = Some(endpoint_idx);
ctx.upstream_response_status = status;
ctx.upstream_contacted = true;
ctx
}
fn make_passive_scenario(
passive_unhealthy: Option<u32>,
passive_healthy: Option<u32>,
) -> (FilterPipeline, PingoraRequestCtx) {
use praxis_core::health::{ClusterHealthEntry, EndpointHealth};
let entry = ClusterHealthEntry::new(
vec![EndpointHealth::new()],
vec![Arc::from("10.0.0.1:80")],
passive_unhealthy,
passive_healthy,
);
let mut map = HashMap::new();
map.insert(Arc::from("test-cluster"), Arc::new(entry));
let health_registry = Arc::new(map);
let registry = praxis_filter::FilterRegistry::with_builtins();
let mut pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
pipeline.set_health_registry(health_registry);
let ctx = make_passive_ctx("test-cluster", 0, None);
(pipeline, ctx)
}
}