1use std::{collections::HashMap, sync::Arc, time::Duration};
17
18use arc_swap::ArcSwap;
19use bytes::Bytes;
20use pingora_core::{
21 Result, apps::HttpServerOptions, protocols::http::v2::server::H2Options, server::Server,
22 services::listening::Service,
23};
24use pingora_proxy::{Session, http_proxy};
25use praxis_core::{config::ABSOLUTE_MAX_BODY_BYTES, connectivity::Upstream};
26use praxis_filter::{BodyBuffer, BodyMode, FilterPipeline, HttpFilterContext, RequestExtensions};
27use tokio::sync::Semaphore;
28use tracing::{debug, warn};
29
30use super::{context::PingoraRequestCtx, metrics};
31
32mod compression;
34mod connected_to_upstream;
36mod fail_to_proxy;
38mod hop_by_hop;
40mod normalize;
42mod request_body_filter;
44mod request_filter;
46mod reserved_headers;
48mod response_body_filter;
50mod response_filter;
52mod retry;
54mod upstream_peer;
56mod upstream_request;
58mod upstream_response;
60mod via;
62mod with_body;
64
65pub use upstream_peer::{UpstreamRetryGateRelease, arm_upstream_retry_gate, lock_upstream_retry_gate_tests};
66pub use with_body::PingoraHttpHandler;
67
68pub fn load_http_handler(
114 server: &mut Server,
115 listener: &praxis_core::config::Listener,
116 pipeline: Arc<ArcSwap<FilterPipeline>>,
117 cert_watcher_shutdowns: &mut Vec<tokio::sync::watch::Sender<bool>>,
118) -> Result<(), praxis_core::ProxyError> {
119 let downstream_read_timeout = listener.downstream_read_timeout_ms.map(Duration::from_millis);
120 let connection_semaphore = listener
121 .max_connections
122 .map(|max| Arc::new(Semaphore::new(max as usize)));
123
124 debug!(listener = %listener.name, "loading HTTP handler with body filters");
127 let handler = PingoraHttpHandler::new(
128 pipeline,
129 downstream_read_timeout,
130 connection_semaphore,
131 ::metrics::SharedString::from_shared(Arc::from(listener.name.as_str())),
135 );
136 wire_service(server, listener, handler, cert_watcher_shutdowns)?;
137 Ok(())
138}
139
140fn wire_service<H>(
142 server: &mut Server,
143 listener: &praxis_core::config::Listener,
144 handler: H,
145 cert_watcher_shutdowns: &mut Vec<tokio::sync::watch::Sender<bool>>,
146) -> Result<(), praxis_core::ProxyError>
147where
148 H: pingora_proxy::ProxyHttp + Send + Sync + 'static,
149 H::CTX: Send + Sync,
150{
151 let service_name = format!("http-proxy:{name}", name = listener.name);
152 let mut proxy = http_proxy(&server.configuration, handler);
153 proxy.server_options = Some(h2c_server_options());
154 proxy.h2_options = Some(h2_server_options());
155 let mut service = Service::new(service_name, proxy);
156 if let Some(tx) = super::listener::add_listener(&mut service, listener)? {
157 cert_watcher_shutdowns.push(tx);
158 }
159 server.add_service(service);
160 Ok(())
161}
162
163fn clamp_body_mode_to_ceiling(mode: BodyMode, baseline: BodyMode) -> BodyMode {
179 let ceiling = match baseline {
180 BodyMode::StreamBuffer { max_bytes: Some(v) } | BodyMode::SizeLimit { max_bytes: v } => Some(v),
181 _ => None,
182 };
183
184 match (mode, ceiling) {
185 (BodyMode::StreamBuffer { max_bytes }, Some(limit)) => BodyMode::StreamBuffer {
186 max_bytes: Some(max_bytes.map_or(limit, |v| v.min(limit))),
187 },
188 (BodyMode::SizeLimit { max_bytes }, Some(limit)) => BodyMode::SizeLimit {
189 max_bytes: max_bytes.min(limit),
190 },
191 (m, None | Some(_)) => m,
194 }
195}
196
197fn legacy_default_policy() -> Arc<praxis_core::config::RetryPolicy> {
199 static LEGACY_DEFAULT: std::sync::LazyLock<Arc<praxis_core::config::RetryPolicy>> =
200 std::sync::LazyLock::new(|| Arc::new(praxis_core::config::RetryPolicy::legacy_default()));
201 Arc::clone(&LEGACY_DEFAULT)
202}
203
204#[expect(clippy::too_many_lines, reason = "sequential guard checks")]
210fn handle_connect_failure(ctx: &mut PingoraRequestCtx, e: Box<pingora_core::Error>) -> Box<pingora_core::Error> {
211 let cluster = ctx.metrics_cluster_shared.clone().unwrap_or_else(metrics::cluster_none);
212 if let Some(start) = ctx.upstream_connect_start.take() {
213 metrics::record_upstream_connect_duration(cluster.clone(), start.elapsed().as_secs_f64());
214 }
215 metrics::record_upstream_connect_failure(cluster.clone());
216
217 let policy = ctx.retry_policy.clone().unwrap_or_else(legacy_default_policy);
218 let outcome = retry::classify_error(&e);
219 let decision = retry::should_retry(ctx, &policy, outcome, ctx.cluster_retry_state.as_deref());
220
221 match decision {
222 retry::RetryDecision::Retry { backoff } => {
223 ctx.retries += 1;
224 ctx.pending_backoff = Some(backoff);
225 ctx.reselect_on_retry = policy.configured;
229 if let Some(upstream) = ctx.upstream_for_retry.as_ref() {
230 let addr = Arc::clone(&upstream.address);
231 if !ctx.attempted_endpoints.iter().any(|e| e.as_ref() == addr.as_ref()) {
232 ctx.attempted_endpoints.push(addr);
233 }
234 }
235 if policy.configured {
240 if let Some(upstream) = ctx.upstream_for_retry.as_ref()
241 && let Some(reselector) = ctx.endpoint_reselector.as_ref()
242 {
243 reselector.release(&upstream.address);
244 }
245 ctx.upstream_for_retry = None;
246 }
247 let upstream_address = ctx
248 .upstream_for_retry
249 .as_ref()
250 .map_or("unknown", |u| u.address.as_ref());
251 debug!(
252 retries = ctx.retries,
253 max = policy.effective_max_retries(),
254 ?backoff,
255 upstream_address,
256 "retrying after connect failure"
257 );
258 let mut e = e;
259 e.set_retry(true);
260 e
261 },
262 retry::RetryDecision::DoNotRetry => {
263 if ctx.retries > 0 {
264 warn!(
265 retries = ctx.retries,
266 max = policy.effective_max_retries(),
267 upstream_address = ctx
268 .upstream_for_retry
269 .as_ref()
270 .map_or("unknown", |u| u.address.as_ref()),
271 "retry limit exhausted"
272 );
273 }
274 record_retry_exhausted_if_attempted(ctx, cluster);
275 let mut e = e;
278 e.set_retry(false);
279 e
280 },
281 }
282}
283
284#[expect(clippy::too_many_lines, reason = "sequential guard checks")]
289fn maybe_retry_response(ctx: &mut PingoraRequestCtx, status: u16) -> Option<Box<pingora_core::Error>> {
290 let policy = ctx.retry_policy.clone().unwrap_or_else(legacy_default_policy);
291 let outcome = retry::RetryOutcome::StatusCode(status);
292 let decision = retry::should_retry(ctx, &policy, outcome, ctx.cluster_retry_state.as_deref());
293 match decision {
294 retry::RetryDecision::Retry { backoff } => {
295 ctx.retries += 1;
296 ctx.pending_backoff = Some(backoff);
297 ctx.reselect_on_retry = policy.configured;
301 if let Some(upstream) = ctx.upstream_for_retry.as_ref() {
302 let addr = Arc::clone(&upstream.address);
303 if !ctx.attempted_endpoints.iter().any(|e| e.as_ref() == addr.as_ref()) {
304 ctx.attempted_endpoints.push(addr);
305 }
306 }
307 if policy.configured {
312 if let Some(upstream) = ctx.upstream_for_retry.as_ref()
313 && let Some(reselector) = ctx.endpoint_reselector.as_ref()
314 {
315 reselector.release(&upstream.address);
316 }
317 ctx.upstream_for_retry = None;
318 }
319 debug!(
320 status,
321 retries = ctx.retries,
322 max = policy.effective_max_retries(),
323 ?backoff,
324 "retrying after retriable response status"
325 );
326 let mut e =
327 pingora_core::Error::explain(pingora_core::ErrorType::HTTPStatus(status), "retriable upstream status");
328 e.set_retry(true);
329 Some(e)
330 },
331 retry::RetryDecision::DoNotRetry => None,
332 }
333}
334
335fn release_retry_state(ctx: &mut PingoraRequestCtx) {
337 if !ctx.cluster_retry_state_released
338 && let Some(state) = ctx.cluster_retry_state.take()
339 {
340 state.leave();
341 ctx.cluster_retry_state_released = true;
342 }
343}
344
345fn record_retry_exhausted_if_attempted(ctx: &PingoraRequestCtx, cluster: ::metrics::SharedString) {
347 if ctx.retries > 0 {
348 metrics::record_upstream_retry(cluster, metrics::RETRY_RESULT_EXHAUSTED);
349 }
350}
351
352fn maybe_emit_fallback_access_log(pipeline: &FilterPipeline, status: u16, ctx: &mut PingoraRequestCtx) {
362 if ctx.response_delivery_complete || ctx.connection_upgraded || !pipeline.contains_filter("access_log") {
363 return;
364 }
365 if let Some(filter_ctx) = ctx.filter_context_for(pipeline, None) {
366 if praxis_filter::access_record_already_emitted(&filter_ctx) {
370 return;
371 }
372 if !pipeline.filter_request_conditions_match("access_log", filter_ctx.request) {
376 return;
377 }
378 praxis_filter::emit_access_record(&filter_ctx, status);
379 }
380}
381
382async fn logging_cleanup(pipeline: &FilterPipeline, ctx: &mut PingoraRequestCtx) {
386 if !ctx.response_phase_done
387 && let Some(mut filter_ctx) = ctx.filter_context_for(pipeline, None)
388 {
389 let _result = pipeline.execute_http_response(&mut filter_ctx).await;
390 let extensions = filter_ctx.extensions;
391 let metadata = filter_ctx.filter_metadata;
392 let state = filter_ctx.filter_state;
393 let exec_idx = filter_ctx.executed_filter_indices;
394 let body_idx = filter_ctx.body_done_indices;
395 let cluster = filter_ctx.cluster;
399 let upstream = filter_ctx.upstream;
400 ctx.extensions = extensions;
401 ctx.filter_metadata = metadata;
402 ctx.filter_state = state;
403 ctx.cached_executed_filter_indices = exec_idx;
404 ctx.cached_body_done_indices = body_idx;
405 ctx.cluster = cluster;
406 ctx.upstream = upstream;
407 }
408}
409
410fn emit_request_metrics(session: &Session, ctx: &PingoraRequestCtx) {
414 if !metrics::is_recorder_installed() {
415 return;
416 }
417
418 let status_code = session.response_written().map_or(0, |resp| resp.status.as_u16());
419 let status_class = metrics::status_class(status_code);
420
421 let method = request_method_label(session, ctx);
422
423 let cluster = ctx.metrics_cluster_shared.clone().unwrap_or_else(metrics::cluster_none);
424
425 let route = ctx.metrics_route.clone().unwrap_or_else(metrics::route_unknown);
426
427 let labels = metrics::RequestMetricLabels {
428 cluster: cluster.clone(),
429 method,
430 route,
431 status_class,
432 };
433
434 emit_upstream_request_metric(ctx, &cluster);
435
436 if let Some(error_type) = ctx.error_type {
437 metrics::record_error(error_type);
438 }
439
440 let duration_secs = ctx.request_start.elapsed().as_secs_f64();
441 metrics::record_request_metrics(labels, duration_secs);
442 metrics::record_body_size_metrics(
443 method,
444 status_class,
445 cluster,
446 ctx.request_body_bytes,
447 ctx.response_body_bytes,
448 );
449}
450
451fn request_method_label(session: &Session, ctx: &PingoraRequestCtx) -> &'static str {
456 let request_method = session.req_header().method.as_str();
457 let raw_method = if request_method.is_empty() {
458 ctx.request_snapshot.as_ref().map_or("UNKNOWN", |r| r.method.as_str())
459 } else {
460 request_method
461 };
462 metrics::method_label(raw_method)
463}
464
465fn emit_upstream_request_metric(ctx: &PingoraRequestCtx, cluster: &::metrics::SharedString) {
471 if let Some(upstream_status) = ctx.upstream_response_status
472 && let Some(upstream) = ctx.upstream_for_retry.as_ref()
473 {
474 metrics::record_upstream_request(
475 cluster.clone(),
476 ::metrics::SharedString::from(Arc::clone(&upstream.address)),
477 metrics::status_class(upstream_status),
478 );
479 }
480}
481
482fn record_passive_health(pipeline: &FilterPipeline, error: Option<&pingora_core::Error>, ctx: &PingoraRequestCtx) {
491 let cluster_name = ctx.cluster.as_ref().or(ctx.metrics_cluster.as_ref());
492 let Some(cluster_name) = cluster_name else {
493 return;
494 };
495 let Some(idx) = ctx.selected_endpoint_index else {
496 return;
497 };
498 let Some(registry) = pipeline.health_registry() else {
499 return;
500 };
501 let Some(health) = registry.get(cluster_name) else {
502 return;
503 };
504
505 if !ctx.upstream_contacted {
514 return;
515 }
516
517 let is_downstream_error = error.is_some_and(|e| matches!(e.esource(), pingora_core::ErrorSource::Downstream));
527 if is_downstream_error && ctx.upstream_response_status.is_none() {
528 return;
529 }
530 let is_failure =
531 ctx.upstream_response_status.is_some_and(|s| s >= 500) || (error.is_some() && !is_downstream_error);
532 apply_passive_threshold(health, idx, cluster_name, is_failure);
533}
534
535fn apply_passive_threshold(
537 health: &praxis_core::health::ClusterHealthEntry,
538 idx: usize,
539 cluster_name: &Arc<str>,
540 is_failure: bool,
541) {
542 if is_failure {
543 if let Some(threshold) = health.passive_unhealthy_threshold()
544 && health
545 .endpoints()
546 .get(idx)
547 .is_some_and(|ep| ep.record_failure(threshold))
548 {
549 tracing::warn!(
550 cluster = %cluster_name,
551 endpoint_index = idx,
552 threshold,
553 "passive health: endpoint marked unhealthy"
554 );
555 emit_passive_health_transition(health, cluster_name, metrics::HEALTH_RESULT_UNHEALTHY);
556 }
557 } else if let Some(threshold) = health.passive_healthy_threshold()
558 && health
559 .endpoints()
560 .get(idx)
561 .is_some_and(|ep| ep.record_success(threshold))
562 {
563 tracing::info!(
564 cluster = %cluster_name,
565 endpoint_index = idx,
566 threshold,
567 "passive health: endpoint recovered"
568 );
569 emit_passive_health_transition(health, cluster_name, metrics::HEALTH_RESULT_HEALTHY);
570 }
571}
572
573fn emit_passive_health_transition(
575 health: &praxis_core::health::ClusterHealthEntry,
576 cluster_name: &Arc<str>,
577 result: &'static str,
578) {
579 let (healthy, total) = metrics::count_healthy_endpoints(health);
580 metrics::record_health_transition(
581 ::metrics::SharedString::from(Arc::clone(cluster_name)),
582 result,
583 healthy,
584 total,
585 );
586}
587
588pub(super) fn http_version_label(version: http::Version) -> &'static str {
592 match version {
593 http::Version::HTTP_09 => "0.9",
594 http::Version::HTTP_10 => "1.0",
595 http::Version::HTTP_11 => "1.1",
596 http::Version::HTTP_2 => "2",
597 http::Version::HTTP_3 => "3",
598 _ => "unknown",
599 }
600}
601
602fn record_response_span_attributes(session: &Session, ctx: &PingoraRequestCtx) {
611 if ctx.request_span.is_disabled() {
612 return;
613 }
614 let response = session.response_written();
615 let status = response.map(|resp| resp.status);
616 let method = session.req_header().method.as_str();
617 record_response_span_fields(status, method, response, ctx);
618}
619
620fn record_response_span_fields(
625 status: Option<http::StatusCode>,
626 method: &str,
627 response: Option<&pingora_http::ResponseHeader>,
628 ctx: &PingoraRequestCtx,
629) {
630 if let Some(status) = status {
631 let code = status.as_u16();
632 if code > 0 {
633 ctx.request_span.record("http.response.status_code", code);
634 }
635 if status.is_server_error() {
636 ctx.request_span.record("otel.status_code", "ERROR");
637 ctx.request_span.record("error.type", code.to_string().as_str());
640 }
641 }
642
643 if let Some(route) = &ctx.metrics_route {
644 ctx.request_span.record("http.route", route.as_ref());
645 ctx.request_span
646 .record("otel.name", format!("{method} {route}").as_str());
647 }
648
649 if let Some(upstream) = &ctx.upstream_for_retry {
650 ctx.request_span.record("upstream.address", upstream.address.as_ref());
651 }
652
653 if let Some(cluster) = &ctx.metrics_cluster {
654 ctx.request_span.record("upstream.cluster", cluster.as_ref());
655 }
656
657 record_upstream_exchange_span(ctx, response);
658}
659
660fn record_upstream_exchange_span(ctx: &PingoraRequestCtx, response: Option<&pingora_http::ResponseHeader>) {
662 if ctx.upstream_exchange_span.is_disabled() {
663 return;
664 }
665 if let Some(status) = ctx
668 .upstream_response_status
669 .or_else(|| response.map(|resp| resp.status.as_u16()))
670 {
671 ctx.upstream_exchange_span.record("http.response.status_code", status);
672 }
673 ctx.upstream_exchange_span
674 .record("http.response.body.size", ctx.response_body_bytes);
675}
676
677fn h2c_server_options() -> HttpServerOptions {
681 let mut opts = HttpServerOptions::default();
682 opts.h2c = true;
683 opts
684}
685
686fn h2_server_options() -> H2Options {
695 let mut opts = H2Options::new();
696 opts.max_header_list_size(65_536); opts.max_concurrent_streams(128);
698 opts
699}
700
701fn check_body_size_limit(body: Option<&Bytes>, accumulated_bytes: &mut u64, max_bytes: usize) -> bool {
704 if let Some(chunk) = body {
705 let chunk_len = chunk.len() as u64;
706 *accumulated_bytes += chunk_len;
707
708 let limit = max_bytes as u64;
709 return *accumulated_bytes > limit;
710 }
711 false
712}
713
714fn accumulate_stream_buffer(
717 body: &mut Option<Bytes>,
718 body_buffer: &mut Option<BodyBuffer>,
719 end_of_stream: bool,
720 max_bytes: Option<usize>,
721) -> bool {
722 if let Some(chunk) = &*body {
723 let limit = max_bytes.unwrap_or(ABSOLUTE_MAX_BODY_BYTES);
724 let buf = body_buffer.get_or_insert_with(|| BodyBuffer::new(limit));
725
726 if buf.push(chunk.clone()).is_err() {
727 return true;
728 }
729 }
730
731 if end_of_stream {
732 tracing::trace!("stream buffer: freezing accumulated body before pipeline at EOS");
733 *body = body_buffer.take().map(BodyBuffer::freeze);
734 } else {
735 tracing::trace!("stream buffer: filters see the original chunk");
736 }
737 false
738}
739
740#[expect(
743 clippy::fn_params_excessive_bools,
744 reason = "mirrors the caller's existing condition flags"
745)]
746fn suppress_stream_buffer_chunk(body: &mut Option<Bytes>, is_stream_buffer: bool, released: bool, end_of_stream: bool) {
747 if is_stream_buffer && !released && !end_of_stream {
748 *body = None;
749 }
750}
751
752fn release_stream_buffer(
754 body: &mut Option<Bytes>,
755 is_stream_buffer: bool,
756 released: &mut bool,
757 body_buffer: &mut Option<BodyBuffer>,
758 end_of_stream: bool,
759) {
760 if is_stream_buffer && !*released {
761 *released = true;
762 if !end_of_stream {
763 *body = body_buffer.take().map(BodyBuffer::freeze);
764 }
765 }
766}
767
768struct BodyFilterOutput {
773 cluster: Option<Arc<str>>,
775 upstream: Option<Upstream>,
777 extensions: RequestExtensions,
779 filter_metadata: HashMap<String, String>,
781 filter_state: HashMap<usize, Box<dyn std::any::Any + Send + Sync>>,
783 executed_filter_indices: Vec<bool>,
785 body_done_indices: Vec<bool>,
787 attempted_endpoints: Vec<Arc<str>>,
789}
790
791impl BodyFilterOutput {
792 fn take_from(fctx: &mut HttpFilterContext<'_>) -> Self {
795 Self {
796 cluster: fctx.cluster.take(),
797 upstream: fctx.upstream.take(),
798 extensions: std::mem::take(&mut fctx.extensions),
799 filter_metadata: std::mem::take(&mut fctx.filter_metadata),
800 filter_state: std::mem::take(&mut fctx.filter_state),
801 executed_filter_indices: std::mem::take(&mut fctx.executed_filter_indices),
802 body_done_indices: std::mem::take(&mut fctx.body_done_indices),
803 attempted_endpoints: std::mem::take(&mut fctx.attempted_endpoints),
804 }
805 }
806
807 fn write_back(self, ctx: &mut PingoraRequestCtx) {
809 ctx.cluster = self.cluster;
810 ctx.upstream = self.upstream;
811 ctx.extensions = self.extensions;
812 ctx.filter_metadata = self.filter_metadata;
813 ctx.filter_state = self.filter_state;
814 ctx.cached_executed_filter_indices = self.executed_filter_indices;
815 ctx.cached_body_done_indices = self.body_done_indices;
816 ctx.attempted_endpoints = self.attempted_endpoints;
817 }
818}
819
820#[cfg(test)]
825#[expect(clippy::allow_attributes, reason = "blanket test suppressions")]
826#[allow(
827 clippy::unwrap_used,
828 clippy::expect_used,
829 clippy::indexing_slicing,
830 clippy::field_reassign_with_default,
831 clippy::too_many_lines,
832 clippy::cast_possible_truncation,
833 clippy::significant_drop_tightening,
834 reason = "tests"
835)]
836mod tests {
837 use praxis_core::connectivity::ConnectionOptions;
838
839 use super::*;
840
841 const MAX_RETRIES: usize = praxis_core::config::DEFAULT_MAX_RETRIES as usize;
843
844 const RETRY_BODY_LIMIT: u64 = praxis_core::config::DEFAULT_RETRY_BODY_LIMIT_BYTES;
846
847 #[test]
848 fn first_failure_idempotent_sets_retry() {
849 let mut ctx = PingoraRequestCtx::default();
850 ctx.request_is_idempotent = true;
851 let e = handle_connect_failure(&mut ctx, make_error());
852 assert!(e.retry(), "first failure should set retry flag");
853 assert_eq!(ctx.retries, 1);
854 }
855
856 #[test]
857 fn large_body_skips_retry() {
858 let mut ctx = PingoraRequestCtx::default();
859 ctx.request_is_idempotent = true;
860 ctx.request_body_bytes = RETRY_BODY_LIMIT + 1;
861 let e = handle_connect_failure(&mut ctx, make_error());
862 assert!(!e.retry(), "should not retry when body exceeds retry buffer limit");
863 assert_eq!(ctx.retries, 0, "retry counter should not increment");
864 }
865
866 #[test]
867 fn mutated_body_exceeding_limit_skips_retry() {
868 let mut ctx = PingoraRequestCtx::default();
869 ctx.request_is_idempotent = true;
870 ctx.request_body_bytes = 1024;
871 ctx.mutated_request_body_len = Some((RETRY_BODY_LIMIT + 1) as usize);
872 let e = handle_connect_failure(&mut ctx, make_error());
873 assert!(
874 !e.retry(),
875 "should not retry when mutated body exceeds retry buffer limit"
876 );
877 assert_eq!(ctx.retries, 0);
878 }
879
880 #[test]
881 fn body_at_limit_allows_retry() {
882 let mut ctx = PingoraRequestCtx::default();
883 ctx.request_is_idempotent = true;
884 ctx.request_body_bytes = RETRY_BODY_LIMIT;
885 let e = handle_connect_failure(&mut ctx, make_error());
886 assert!(e.retry(), "body exactly at limit should allow retry");
887 assert_eq!(ctx.retries, 1);
888 }
889
890 #[test]
891 fn zero_body_allows_retry() {
892 let mut ctx = PingoraRequestCtx::default();
893 ctx.request_is_idempotent = true;
894 ctx.request_body_bytes = 0;
895 let e = handle_connect_failure(&mut ctx, make_error());
896 assert!(e.retry(), "zero-length body should allow retry");
897 assert_eq!(ctx.retries, 1);
898 }
899
900 #[test]
901 fn max_retries_exhausted_does_not_retry() {
902 let mut ctx = PingoraRequestCtx::default();
903 ctx.request_is_idempotent = true;
904 ctx.retries = MAX_RETRIES as u32;
905 let e = handle_connect_failure(&mut ctx, make_error());
906 assert!(!e.retry(), "should not retry after MAX_RETRIES");
907 assert_eq!(ctx.retries as usize, MAX_RETRIES);
908 }
909
910 #[test]
911 fn counter_increments_across_calls() {
912 let mut ctx = PingoraRequestCtx::default();
913 ctx.request_is_idempotent = true;
914 for expected in 1..=MAX_RETRIES {
915 let _result = handle_connect_failure(&mut ctx, make_error());
916 assert_eq!(ctx.retries as usize, expected);
917 }
918 let e = handle_connect_failure(&mut ctx, make_error());
919 assert!(!e.retry(), "should not retry after reaching MAX_RETRIES");
920 assert_eq!(ctx.retries as usize, MAX_RETRIES);
921 }
922
923 #[test]
924 fn non_idempotent_request_never_retries() {
925 let mut ctx = PingoraRequestCtx::default();
926 ctx.request_is_idempotent = false;
927 let e = handle_connect_failure(&mut ctx, make_error());
928 assert!(!e.retry(), "non-idempotent request should never retry");
929 assert_eq!(ctx.retries, 0);
930 }
931
932 #[test]
933 fn connect_failure_clears_upstream_connect_start() {
934 let mut ctx = PingoraRequestCtx::default();
935 ctx.upstream_connect_start = Some(std::time::Instant::now());
936 let _e = handle_connect_failure(&mut ctx, make_error());
937 assert!(
938 ctx.upstream_connect_start.is_none(),
939 "failed connect should consume upstream_connect_start for duration recording"
940 );
941 }
942
943 #[test]
944 fn response_503_retries_when_status5xx_enabled() {
945 let mut ctx = PingoraRequestCtx::default();
946 ctx.request_is_idempotent = true;
947 ctx.retry_policy = Some(Arc::new(praxis_core::config::RetryPolicy {
948 configured: true,
949 retriable_conditions: vec![praxis_core::config::RetriableCondition::Status5xx],
950 ..praxis_core::config::RetryPolicy::legacy_default()
951 }));
952 let e = maybe_retry_response(&mut ctx, 503).expect("503 should be retriable");
953 assert!(e.retry(), "503 under Status5xx should set retry");
954 assert_eq!(ctx.retries, 1);
955 assert!(ctx.reselect_on_retry);
956 assert!(ctx.pending_backoff.is_some());
957 }
958
959 #[test]
960 fn response_404_does_not_retry() {
961 let mut ctx = PingoraRequestCtx::default();
962 ctx.request_is_idempotent = true;
963 ctx.retry_policy = Some(Arc::new(praxis_core::config::RetryPolicy {
964 retriable_conditions: vec![praxis_core::config::RetriableCondition::Status5xx],
965 ..praxis_core::config::RetryPolicy::legacy_default()
966 }));
967 assert!(
968 maybe_retry_response(&mut ctx, 404).is_none(),
969 "404 must never trigger status-based retry"
970 );
971 assert_eq!(ctx.retries, 0);
972 }
973
974 #[test]
975 fn response_502_does_not_retry_under_legacy_default() {
976 let mut ctx = PingoraRequestCtx::default();
977 ctx.request_is_idempotent = true;
978 assert!(
979 maybe_retry_response(&mut ctx, 502).is_none(),
980 "legacy default must forward 5xx without retry"
981 );
982 }
983
984 #[test]
985 fn max_retries_zero_disables_connect_retry() {
986 let mut ctx = PingoraRequestCtx::default();
987 ctx.request_is_idempotent = true;
988 ctx.retry_policy = Some(Arc::new(praxis_core::config::RetryPolicy {
989 max_retries: Some(0),
990 ..praxis_core::config::RetryPolicy::legacy_default()
991 }));
992 let e = handle_connect_failure(&mut ctx, make_error());
993 assert!(!e.retry(), "max_retries: 0 must disable retries");
994 assert_eq!(ctx.retries, 0);
995 }
996
997 #[test]
998 fn non_idempotent_clears_pingora_default_retry_flag() {
999 let mut ctx = PingoraRequestCtx::default();
1000 ctx.request_is_idempotent = false;
1001 let mut e = make_error();
1002 e.set_retry(true);
1003 let e = handle_connect_failure(&mut ctx, e);
1004 assert!(!e.retry(), "policy denial must clear Pingora's default retry flag");
1005 assert_eq!(ctx.retries, 0);
1006 }
1007
1008 #[tokio::test]
1009 async fn logging_cleanup_noop_when_response_phase_done() {
1010 let registry = praxis_filter::FilterRegistry::with_builtins();
1011 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
1012 let mut ctx = PingoraRequestCtx::default();
1013 ctx.response_phase_done = true;
1014 ctx.request_snapshot = Some(praxis_filter::Request {
1015 method: http::Method::GET,
1016 uri: "/".parse().unwrap(),
1017 headers: http::HeaderMap::new(),
1018 });
1019 logging_cleanup(&pipeline, &mut ctx).await;
1020 }
1021
1022 #[tokio::test]
1023 async fn logging_cleanup_noop_when_no_snapshot() {
1024 let registry = praxis_filter::FilterRegistry::with_builtins();
1025 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
1026 let mut ctx = PingoraRequestCtx::default();
1027 ctx.response_phase_done = false;
1028 ctx.request_snapshot = None;
1029 logging_cleanup(&pipeline, &mut ctx).await;
1030 }
1031
1032 #[tokio::test]
1033 async fn logging_cleanup_runs_response_pipeline_when_needed() {
1034 let registry = praxis_filter::FilterRegistry::with_builtins();
1035 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
1036 let mut ctx = PingoraRequestCtx::default();
1037 ctx.response_phase_done = false;
1038 ctx.cluster = Some(Arc::from("test-cluster"));
1039 ctx.request_snapshot = Some(praxis_filter::Request {
1040 method: http::Method::GET,
1041 uri: "/test".parse().unwrap(),
1042 headers: http::HeaderMap::new(),
1043 });
1044 logging_cleanup(&pipeline, &mut ctx).await;
1045 assert_eq!(
1046 ctx.cluster.as_deref(),
1047 Some("test-cluster"),
1048 "cluster must be restored so the fallback access record can attribute the failure"
1049 );
1050 }
1051
1052 #[tokio::test]
1053 async fn logging_cleanup_preserves_filter_metadata() {
1054 let registry = praxis_filter::FilterRegistry::with_builtins();
1055 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
1056 let mut ctx = PingoraRequestCtx::default();
1057 ctx.response_phase_done = false;
1058 ctx.filter_metadata
1059 .insert("json_rpc.method".to_owned(), "service/invoke".to_owned());
1060 ctx.request_snapshot = Some(praxis_filter::Request {
1061 method: http::Method::POST,
1062 uri: "/api".parse().unwrap(),
1063 headers: http::HeaderMap::new(),
1064 });
1065 logging_cleanup(&pipeline, &mut ctx).await;
1066 assert_eq!(
1067 ctx.filter_metadata.get("json_rpc.method").map(String::as_str),
1068 Some("service/invoke"),
1069 "filter_metadata should survive logging_cleanup"
1070 );
1071 }
1072
1073 #[tokio::test]
1074 async fn logging_cleanup_preserves_extensions() {
1075 let registry = praxis_filter::FilterRegistry::with_builtins();
1076 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
1077 let mut ctx = PingoraRequestCtx::default();
1078 ctx.response_phase_done = false;
1079 ctx.extensions.insert(42_u32);
1080 ctx.request_snapshot = Some(praxis_filter::Request {
1081 method: http::Method::POST,
1082 uri: "/test".parse().unwrap(),
1083 headers: http::HeaderMap::new(),
1084 });
1085 logging_cleanup(&pipeline, &mut ctx).await;
1086 assert_eq!(
1087 ctx.extensions.get::<u32>(),
1088 Some(&42),
1089 "extensions should survive logging_cleanup"
1090 );
1091 }
1092
1093 #[test]
1094 fn passive_health_error_is_failure() {
1095 let (pipeline, ctx) = make_passive_scenario(Some(3), Some(2));
1096 let error = make_error();
1097 record_passive_health(&pipeline, Some(&error), &ctx);
1098
1099 let registry = pipeline.health_registry().unwrap();
1100 let entry = registry.get("test-cluster").unwrap();
1101 assert!(
1102 entry.endpoints()[0].is_healthy(),
1103 "single failure should not yet mark unhealthy (threshold=3)"
1104 );
1105 }
1106
1107 #[test]
1108 fn passive_health_downstream_error_is_not_failure() {
1109 let (pipeline, ctx) = make_passive_scenario(Some(1), Some(1));
1110 let error = make_error().into_down();
1111 record_passive_health(&pipeline, Some(&error), &ctx);
1112
1113 let registry = pipeline.health_registry().unwrap();
1114 let entry = registry.get("test-cluster").unwrap();
1115 assert!(
1116 entry.endpoints()[0].is_healthy(),
1117 "a downstream/client error must not mark the endpoint unhealthy"
1118 );
1119 }
1120
1121 #[test]
1122 fn passive_health_downstream_error_with_5xx_still_counts_as_failure() {
1123 let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
1124 ctx.upstream_response_status = Some(503);
1125 let error = make_error().into_down();
1126 record_passive_health(&pipeline, Some(&error), &ctx);
1127
1128 let registry = pipeline.health_registry().unwrap();
1129 let entry = registry.get("test-cluster").unwrap();
1130 assert!(
1131 !entry.endpoints()[0].is_healthy(),
1132 "a downstream error with an upstream 503 must still mark the endpoint unhealthy"
1133 );
1134 }
1135
1136 #[test]
1137 fn passive_health_downstream_error_does_not_reset_failure_streak() {
1138 let (pipeline, ctx) = make_passive_scenario(Some(2), Some(1));
1139 let mut upstream_err = make_error();
1140 upstream_err.as_up();
1141 let downstream_err = make_error().into_down();
1142
1143 record_passive_health(&pipeline, Some(&upstream_err), &ctx);
1144 record_passive_health(&pipeline, Some(&downstream_err), &ctx);
1145 record_passive_health(&pipeline, Some(&upstream_err), &ctx);
1146
1147 let registry = pipeline.health_registry().unwrap();
1148 let entry = registry.get("test-cluster").unwrap();
1149 assert!(
1150 !entry.endpoints()[0].is_healthy(),
1151 "two upstream failures must eject the endpoint even with an interleaved client error"
1152 );
1153 }
1154
1155 #[test]
1156 fn passive_health_upstream_error_is_failure() {
1157 let (pipeline, ctx) = make_passive_scenario(Some(1), Some(1));
1158 let mut error = make_error();
1159 error.as_up();
1160 record_passive_health(&pipeline, Some(&error), &ctx);
1161
1162 let registry = pipeline.health_registry().unwrap();
1163 let entry = registry.get("test-cluster").unwrap();
1164 assert!(
1165 !entry.endpoints()[0].is_healthy(),
1166 "an upstream error at unhealthy-threshold 1 must mark the endpoint unhealthy"
1167 );
1168 }
1169
1170 #[test]
1171 fn passive_health_skips_observations_without_upstream_contact() {
1172 let (pipeline, mut ctx) = make_passive_scenario(Some(2), Some(1));
1173 let mut upstream_err = make_error();
1174 upstream_err.as_up();
1175
1176 record_passive_health(&pipeline, Some(&upstream_err), &ctx);
1177
1178 ctx.upstream_contacted = false;
1179 record_passive_health(&pipeline, None, &ctx);
1180
1181 ctx.upstream_response_status = Some(200);
1182 record_passive_health(&pipeline, None, &ctx);
1183 ctx.upstream_response_status = None;
1184
1185 ctx.upstream_contacted = true;
1186 record_passive_health(&pipeline, Some(&upstream_err), &ctx);
1187
1188 let registry = pipeline.health_registry().unwrap();
1189 let entry = registry.get("test-cluster").unwrap();
1190 assert!(
1191 !entry.endpoints()[0].is_healthy(),
1192 "observations without upstream contact must not reset the failure streak"
1193 );
1194 }
1195
1196 #[test]
1197 fn passive_health_records_connect_failure_after_reselect_clears_upstream() {
1198 let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
1199 ctx.upstream_for_retry = None;
1200 ctx.upstream_contacted = true;
1201 let mut error = make_error();
1202 error.as_up();
1203 record_passive_health(&pipeline, Some(&error), &ctx);
1204
1205 let registry = pipeline.health_registry().unwrap();
1206 let entry = registry.get("test-cluster").unwrap();
1207 assert!(
1208 !entry.endpoints()[0].is_healthy(),
1209 "a connect failure after reselection cleared upstream_for_retry must still count"
1210 );
1211 }
1212
1213 #[test]
1214 fn passive_health_status_500_is_failure() {
1215 let (pipeline, mut ctx) = make_passive_scenario(Some(3), Some(2));
1216 ctx.upstream_response_status = Some(500);
1217 record_passive_health(&pipeline, None, &ctx);
1218
1219 let registry = pipeline.health_registry().unwrap();
1220 let entry = registry.get("test-cluster").unwrap();
1221 assert!(
1222 entry.endpoints()[0].is_healthy(),
1223 "single 500 should not yet mark unhealthy (threshold=3)"
1224 );
1225 }
1226
1227 #[test]
1228 fn passive_health_status_below_500_is_success() {
1229 let (pipeline, mut ctx) = make_passive_scenario(Some(2), Some(1));
1230 ctx.upstream_response_status = Some(499);
1231 record_passive_health(&pipeline, None, &ctx);
1232
1233 let registry = pipeline.health_registry().unwrap();
1234 let entry = registry.get("test-cluster").unwrap();
1235 assert!(entry.endpoints()[0].is_healthy(), "status 499 should count as success");
1236 }
1237
1238 #[test]
1239 fn passive_unhealthy_threshold_transition() {
1240 let (pipeline, ctx) = make_passive_scenario(Some(2), Some(1));
1241 let error = make_error();
1242 record_passive_health(&pipeline, Some(&error), &ctx);
1243 record_passive_health(&pipeline, Some(&error), &ctx);
1244
1245 let registry = pipeline.health_registry().unwrap();
1246 let entry = registry.get("test-cluster").unwrap();
1247 assert!(
1248 !entry.endpoints()[0].is_healthy(),
1249 "2 consecutive failures should mark unhealthy (threshold=2)"
1250 );
1251 }
1252
1253 #[test]
1254 fn passive_healthy_threshold_recovery() {
1255 let (pipeline, ctx) = make_passive_scenario(Some(1), Some(2));
1256 let error = make_error();
1257 record_passive_health(&pipeline, Some(&error), &ctx);
1258
1259 let registry = pipeline.health_registry().unwrap();
1260 let entry = registry.get("test-cluster").unwrap();
1261 assert!(
1262 !entry.endpoints()[0].is_healthy(),
1263 "should be unhealthy after 1 failure"
1264 );
1265
1266 let ctx_ok = make_passive_ctx("test-cluster", 0, Some(200));
1267 record_passive_health(&pipeline, None, &ctx_ok);
1268 assert!(
1269 !entry.endpoints()[0].is_healthy(),
1270 "one success should not recover (threshold=2)"
1271 );
1272
1273 record_passive_health(&pipeline, None, &ctx_ok);
1274 assert!(
1275 entry.endpoints()[0].is_healthy(),
1276 "2 consecutive successes should recover (threshold=2)"
1277 );
1278 }
1279
1280 #[test]
1281 fn passive_health_no_thresholds_is_noop() {
1282 let (pipeline, ctx) = make_passive_scenario(None, None);
1283 let error = make_error();
1284 record_passive_health(&pipeline, Some(&error), &ctx);
1285
1286 let registry = pipeline.health_registry().unwrap();
1287 let entry = registry.get("test-cluster").unwrap();
1288 assert!(
1289 entry.endpoints()[0].is_healthy(),
1290 "no passive thresholds means failures are no-op"
1291 );
1292 }
1293
1294 #[test]
1295 fn passive_health_endpoint_index_out_of_bounds() {
1296 let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
1297 ctx.selected_endpoint_index = Some(999);
1298 let error = make_error();
1299 record_passive_health(&pipeline, Some(&error), &ctx);
1300
1301 let registry = pipeline.health_registry().unwrap();
1302 let entry = registry.get("test-cluster").unwrap();
1303 assert!(entry.endpoints()[0].is_healthy(), "out-of-bounds index should be no-op");
1304 }
1305
1306 #[test]
1307 fn passive_health_missing_cluster_is_noop() {
1308 let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
1309 ctx.cluster = None;
1310 ctx.metrics_cluster = None;
1311 let error = make_error();
1312 record_passive_health(&pipeline, Some(&error), &ctx);
1313 }
1314
1315 #[test]
1316 fn passive_health_falls_back_to_metrics_cluster() {
1317 let (pipeline, mut ctx) = make_passive_scenario(Some(2), Some(1));
1318 ctx.cluster = None;
1319 ctx.metrics_cluster = Some(Arc::from("test-cluster"));
1320 let error = make_error();
1321 record_passive_health(&pipeline, Some(&error), &ctx);
1322 record_passive_health(&pipeline, Some(&error), &ctx);
1323
1324 let registry = pipeline.health_registry().unwrap();
1325 let entry = registry.get("test-cluster").unwrap();
1326 assert!(
1327 !entry.endpoints()[0].is_healthy(),
1328 "fallback to metrics_cluster should still record passive health"
1329 );
1330 }
1331
1332 #[test]
1333 fn passive_health_missing_endpoint_index_is_noop() {
1334 let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
1335 ctx.selected_endpoint_index = None;
1336 let error = make_error();
1337 record_passive_health(&pipeline, Some(&error), &ctx);
1338 }
1339
1340 #[test]
1341 fn passive_health_missing_registry_is_noop() {
1342 let registry = praxis_filter::FilterRegistry::with_builtins();
1343 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
1344 let mut ctx = PingoraRequestCtx::default();
1345 ctx.cluster = Some(Arc::from("test-cluster"));
1346 ctx.selected_endpoint_index = Some(0);
1347 let error = make_error();
1348 record_passive_health(&pipeline, Some(&error), &ctx);
1349 }
1350
1351 #[test]
1352 fn passive_health_unknown_cluster_is_noop() {
1353 let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
1354 ctx.cluster = Some(Arc::from("nonexistent"));
1355 let error = make_error();
1356 record_passive_health(&pipeline, Some(&error), &ctx);
1357 }
1358
1359 #[test]
1360 fn size_limit_none_body_returns_false() {
1361 let mut bytes = 0_u64;
1362 assert!(!check_body_size_limit(None, &mut bytes, 100));
1363 assert_eq!(bytes, 0, "accumulated bytes unchanged for None body");
1364 }
1365
1366 #[test]
1367 fn size_limit_within_limit() {
1368 let mut bytes = 0_u64;
1369 let body = Some(Bytes::from_static(b"hello"));
1370 assert!(!check_body_size_limit(body.as_ref(), &mut bytes, 10));
1371 assert_eq!(bytes, 5);
1372 }
1373
1374 #[test]
1375 fn size_limit_at_exact_limit() {
1376 let mut bytes = 0_u64;
1377 let body = Some(Bytes::from_static(b"exact"));
1378 assert!(!check_body_size_limit(body.as_ref(), &mut bytes, 5));
1379 assert_eq!(bytes, 5);
1380 }
1381
1382 #[test]
1383 fn size_limit_exceeds_limit() {
1384 let mut bytes = 0_u64;
1385 let body = Some(Bytes::from_static(b"toolong"));
1386 assert!(check_body_size_limit(body.as_ref(), &mut bytes, 3));
1387 }
1388
1389 #[test]
1390 fn size_limit_cumulative_overflow() {
1391 let mut bytes = 0_u64;
1392 let first = Some(Bytes::from_static(b"aaa"));
1393 assert!(!check_body_size_limit(first.as_ref(), &mut bytes, 5));
1394
1395 let second = Some(Bytes::from_static(b"bbb"));
1396 assert!(check_body_size_limit(second.as_ref(), &mut bytes, 5));
1397 assert_eq!(bytes, 6);
1398 }
1399
1400 #[test]
1401 fn stream_buffer_accumulates_chunks() {
1402 let mut body = Some(Bytes::from_static(b"hello "));
1403 let mut buf: Option<BodyBuffer> = None;
1404 assert!(!accumulate_stream_buffer(&mut body, &mut buf, false, Some(100)));
1405 assert!(buf.is_some());
1406
1407 body = Some(Bytes::from_static(b"world"));
1408 assert!(!accumulate_stream_buffer(&mut body, &mut buf, false, Some(100)));
1409
1410 let frozen = buf.take().unwrap().freeze();
1411 assert_eq!(frozen, Bytes::from_static(b"hello world"));
1412 }
1413
1414 #[test]
1415 fn stream_buffer_freezes_at_eos() {
1416 let mut body = Some(Bytes::from_static(b"data"));
1417 let mut buf: Option<BodyBuffer> = None;
1418 assert!(!accumulate_stream_buffer(&mut body, &mut buf, false, Some(100)));
1419
1420 body = Some(Bytes::from_static(b" end"));
1421 assert!(!accumulate_stream_buffer(&mut body, &mut buf, true, Some(100)));
1422 assert!(buf.is_none(), "buffer should be taken at EOS");
1423 assert_eq!(body.unwrap(), Bytes::from_static(b"data end"));
1424 }
1425
1426 #[test]
1427 fn stream_buffer_overflow() {
1428 let mut body = Some(Bytes::from_static(b"too long"));
1429 let mut buf: Option<BodyBuffer> = None;
1430 assert!(accumulate_stream_buffer(&mut body, &mut buf, false, Some(5)));
1431 }
1432
1433 #[test]
1434 fn stream_buffer_none_body() {
1435 let mut body: Option<Bytes> = None;
1436 let mut buf: Option<BodyBuffer> = None;
1437 assert!(!accumulate_stream_buffer(&mut body, &mut buf, false, Some(100)));
1438 assert!(buf.is_none());
1439 }
1440
1441 #[test]
1442 fn stream_buffer_uses_absolute_max_when_none() {
1443 let mut body = Some(Bytes::from_static(b"data"));
1444 let mut buf: Option<BodyBuffer> = None;
1445 assert!(!accumulate_stream_buffer(&mut body, &mut buf, false, None));
1446 assert!(buf.is_some(), "should create buffer with absolute max");
1447 }
1448
1449 #[test]
1450 fn suppress_clears_body_when_buffering() {
1451 let mut body = Some(Bytes::from_static(b"data"));
1452 suppress_stream_buffer_chunk(&mut body, true, false, false);
1453 assert!(body.is_none());
1454 }
1455
1456 #[test]
1457 fn suppress_noop_when_not_stream_buffer() {
1458 let mut body = Some(Bytes::from_static(b"data"));
1459 suppress_stream_buffer_chunk(&mut body, false, false, false);
1460 assert!(body.is_some());
1461 }
1462
1463 #[test]
1464 fn suppress_noop_when_released() {
1465 let mut body = Some(Bytes::from_static(b"data"));
1466 suppress_stream_buffer_chunk(&mut body, true, true, false);
1467 assert!(body.is_some());
1468 }
1469
1470 #[test]
1471 fn suppress_noop_at_eos() {
1472 let mut body = Some(Bytes::from_static(b"data"));
1473 suppress_stream_buffer_chunk(&mut body, true, false, true);
1474 assert!(body.is_some());
1475 }
1476
1477 #[test]
1478 fn release_sets_flag_and_flushes_buffer() {
1479 let mut body: Option<Bytes> = None;
1480 let mut released = false;
1481 let mut buf = Some(BodyBuffer::new(100));
1482 buf.as_mut().unwrap().push(Bytes::from_static(b"buffered")).unwrap();
1483
1484 release_stream_buffer(&mut body, true, &mut released, &mut buf, false);
1485 assert!(released);
1486 assert_eq!(body.unwrap(), Bytes::from_static(b"buffered"));
1487 assert!(buf.is_none());
1488 }
1489
1490 #[test]
1491 fn release_noop_when_already_released() {
1492 let mut body: Option<Bytes> = None;
1493 let mut released = true;
1494 let mut buf: Option<BodyBuffer> = None;
1495
1496 release_stream_buffer(&mut body, true, &mut released, &mut buf, false);
1497 assert!(body.is_none(), "body should be unchanged when already released");
1498 }
1499
1500 #[test]
1501 fn release_noop_when_not_stream_buffer() {
1502 let mut body: Option<Bytes> = None;
1503 let mut released = false;
1504 let mut buf: Option<BodyBuffer> = None;
1505
1506 release_stream_buffer(&mut body, false, &mut released, &mut buf, false);
1507 assert!(!released, "released flag should be unchanged for non-stream-buffer");
1508 }
1509
1510 #[test]
1511 fn release_at_eos_sets_flag_but_no_flush() {
1512 let mut body: Option<Bytes> = None;
1513 let mut released = false;
1514 let mut buf = Some(BodyBuffer::new(100));
1515 buf.as_mut().unwrap().push(Bytes::from_static(b"data")).unwrap();
1516
1517 release_stream_buffer(&mut body, true, &mut released, &mut buf, true);
1518 assert!(released);
1519 assert!(body.is_none(), "body should not be overwritten at EOS");
1520 assert!(buf.is_some(), "buffer should not be taken at EOS");
1521 }
1522
1523 #[test]
1524 fn write_back_transfers_fields() {
1525 let mut ctx = PingoraRequestCtx::default();
1526
1527 let mut extensions = RequestExtensions::new();
1528 extensions.insert(42_u32);
1529
1530 let state_val: Box<dyn std::any::Any + Send + Sync> = Box::new(99_i32);
1531 let filter_state = HashMap::from([(0_usize, state_val)]);
1532
1533 let output = BodyFilterOutput {
1534 cluster: Some(Arc::from("test-cluster")),
1535 upstream: Some(Upstream {
1536 address: Arc::from("10.0.0.1:80"),
1537 authority: None,
1538 connection: Arc::new(ConnectionOptions::default()),
1539 tls: None,
1540 }),
1541 extensions,
1542 attempted_endpoints: Vec::new(),
1543 filter_metadata: HashMap::from([("key".to_owned(), "val".to_owned())]),
1544 filter_state,
1545 executed_filter_indices: vec![true, false],
1546 body_done_indices: vec![false, true],
1547 };
1548 output.write_back(&mut ctx);
1549
1550 assert_eq!(ctx.cluster.as_deref(), Some("test-cluster"));
1551 assert!(ctx.upstream.is_some(), "upstream should transfer");
1552 assert_eq!(ctx.upstream.as_ref().unwrap().address.as_ref(), "10.0.0.1:80");
1553 assert_eq!(ctx.extensions.get::<u32>(), Some(&42));
1554 assert_eq!(ctx.filter_metadata.get("key").map(String::as_str), Some("val"));
1555 assert_eq!(ctx.filter_state.len(), 1, "filter_state should transfer");
1556 assert_eq!(
1557 ctx.filter_state.get(&0).and_then(|v| v.downcast_ref::<i32>()),
1558 Some(&99)
1559 );
1560 assert_eq!(ctx.cached_executed_filter_indices, vec![true, false]);
1561 assert_eq!(ctx.cached_body_done_indices, vec![false, true]);
1562 }
1563
1564 #[test]
1569 fn fallback_access_log_emits_for_incomplete_request() {
1570 let pipeline = access_log_pipeline();
1571 let mut ctx = make_fallback_ctx();
1572
1573 let events = capture_access_events(|| maybe_emit_fallback_access_log(&pipeline, 502, &mut ctx));
1574 assert_eq!(
1575 events.len(),
1576 1,
1577 "incomplete request must produce a fallback access record"
1578 );
1579 }
1580
1581 #[test]
1582 fn fallback_access_log_skips_completed_delivery() {
1583 let pipeline = access_log_pipeline();
1584 let mut ctx = make_fallback_ctx();
1585 ctx.response_delivery_complete = true;
1586
1587 let events = capture_access_events(|| maybe_emit_fallback_access_log(&pipeline, 200, &mut ctx));
1588 assert!(events.is_empty(), "completed delivery already logged via the filter");
1589 }
1590
1591 #[test]
1592 fn fallback_access_log_skips_upgraded_connections() {
1593 let pipeline = access_log_pipeline();
1594 let mut ctx = make_fallback_ctx();
1595 ctx.connection_upgraded = true;
1596
1597 let events = capture_access_events(|| maybe_emit_fallback_access_log(&pipeline, 101, &mut ctx));
1598 assert!(events.is_empty(), "upgraded connections have no body completion");
1599 }
1600
1601 #[test]
1602 fn fallback_access_log_skips_without_access_log_filter() {
1603 let registry = praxis_filter::FilterRegistry::with_builtins();
1604 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
1605 let mut ctx = make_fallback_ctx();
1606
1607 let events = capture_access_events(|| maybe_emit_fallback_access_log(&pipeline, 502, &mut ctx));
1608 assert!(events.is_empty(), "no access_log filter means no fallback record");
1609 }
1610
1611 #[test]
1612 fn fallback_access_log_honors_entry_conditions() {
1613 let registry = praxis_filter::FilterRegistry::with_builtins();
1614 let mut entries = vec![praxis_filter::FilterEntry {
1615 branch_chains: None,
1616 conditions: vec![serde_yaml::from_str("when:\n path_prefix: /api\n").unwrap()],
1617 failure_mode: praxis_filter::FailureMode::default(),
1618 filter_type: "access_log".to_owned(),
1619 config: serde_yaml::Value::Null,
1620 name: None,
1621 response_conditions: vec![],
1622 }];
1623 let pipeline = FilterPipeline::build(&mut entries, ®istry).unwrap();
1624
1625 let mut excluded = make_fallback_ctx();
1626 let events = capture_access_events(|| maybe_emit_fallback_access_log(&pipeline, 502, &mut excluded));
1627 assert!(
1628 events.is_empty(),
1629 "requests the operator scoped out must not gain fallback records"
1630 );
1631
1632 let mut included = make_fallback_ctx();
1633 if let Some(snapshot) = included.request_snapshot.as_mut() {
1634 snapshot.uri = "/api/users".parse().unwrap();
1635 }
1636 let events = capture_access_events(|| maybe_emit_fallback_access_log(&pipeline, 502, &mut included));
1637 assert_eq!(events.len(), 1, "in-scope incomplete requests still get a record");
1638 }
1639
1640 #[test]
1641 fn aborted_response_body_at_eos_is_not_marked_delivered() {
1642 let pipeline = access_log_pipeline();
1643 let mut ctx = make_fallback_ctx();
1644 ctx.response_body_mode = BodyMode::SizeLimit { max_bytes: 4 };
1645 let mut body = Some(Bytes::from_static(b"exceeds the limit"));
1646
1647 let result = response_body_filter::execute(&pipeline, &mut body, true, &mut ctx);
1648 assert!(result.is_err(), "over-limit body must abort");
1649 assert!(
1650 !ctx.response_delivery_complete,
1651 "a response aborted at end-of-stream was not delivered; the fallback record must fire"
1652 );
1653 }
1654
1655 #[test]
1656 fn response_body_eos_marks_delivery_complete() {
1657 let registry = praxis_filter::FilterRegistry::with_builtins();
1658 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
1659 let mut ctx = PingoraRequestCtx::default();
1660 let mut body: Option<Bytes> = None;
1661
1662 let _timeout = response_body_filter::execute(&pipeline, &mut body, false, &mut ctx).unwrap();
1663 assert!(
1664 !ctx.response_delivery_complete,
1665 "mid-stream chunks must not mark delivery complete"
1666 );
1667
1668 let _timeout = response_body_filter::execute(&pipeline, &mut body, true, &mut ctx).unwrap();
1669 assert!(ctx.response_delivery_complete, "end-of-stream marks delivery complete");
1670 }
1671
1672 #[test]
1677 fn http_version_label_http_09() {
1678 assert_eq!(
1679 http_version_label(http::Version::HTTP_09),
1680 "0.9",
1681 "HTTP/0.9 should map to '0.9'"
1682 );
1683 }
1684
1685 #[test]
1686 fn http_version_label_http_10() {
1687 assert_eq!(
1688 http_version_label(http::Version::HTTP_10),
1689 "1.0",
1690 "HTTP/1.0 should map to '1.0'"
1691 );
1692 }
1693
1694 #[test]
1695 fn http_version_label_http_11() {
1696 assert_eq!(
1697 http_version_label(http::Version::HTTP_11),
1698 "1.1",
1699 "HTTP/1.1 should map to '1.1'"
1700 );
1701 }
1702
1703 #[test]
1704 fn http_version_label_http_2() {
1705 assert_eq!(
1706 http_version_label(http::Version::HTTP_2),
1707 "2",
1708 "HTTP/2 should map to '2'"
1709 );
1710 }
1711
1712 #[test]
1713 fn http_version_label_http_3() {
1714 assert_eq!(
1715 http_version_label(http::Version::HTTP_3),
1716 "3",
1717 "HTTP/3 should map to '3'"
1718 );
1719 }
1720
1721 #[test]
1722 fn record_response_span_attributes_noop_for_disabled_span() {
1723 let ctx = PingoraRequestCtx::default();
1724 assert!(ctx.request_span.is_disabled(), "default span should be disabled");
1725 }
1726
1727 #[derive(Clone, Default)]
1729 struct RecordCapture(Arc<std::sync::Mutex<Vec<(String, String)>>>);
1730
1731 impl<S> tracing_subscriber::Layer<S> for RecordCapture
1732 where
1733 S: tracing::Subscriber + for<'a> tracing_subscriber::registry::LookupSpan<'a>,
1734 {
1735 fn on_record(
1736 &self,
1737 _id: &tracing::span::Id,
1738 values: &tracing::span::Record<'_>,
1739 _ctx: tracing_subscriber::layer::Context<'_, S>,
1740 ) {
1741 struct Visitor<'a>(&'a mut Vec<(String, String)>);
1742 impl tracing::field::Visit for Visitor<'_> {
1743 fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
1744 self.0.push((field.name().to_owned(), format!("{value:?}")));
1745 }
1746 }
1747 let mut captured = self.0.lock().expect("capture lock");
1748 values.record(&mut Visitor(&mut captured));
1749 }
1750 }
1751
1752 #[test]
1753 fn record_response_span_fields_records_status_upstream_and_cluster() {
1754 use tracing_subscriber::layer::SubscriberExt as _;
1755
1756 let capture = RecordCapture::default();
1757 let subscriber = tracing_subscriber::registry().with(capture.clone());
1758 let _guard = tracing::subscriber::set_default(subscriber);
1759
1760 let mut ctx = PingoraRequestCtx::default();
1761 ctx.metrics_cluster = Some(Arc::from("api-cluster"));
1762 ctx.upstream_for_retry = Some(Upstream {
1763 address: Arc::from("10.0.0.1:80"),
1764 authority: None,
1765 connection: Arc::new(ConnectionOptions::default()),
1766 tls: None,
1767 });
1768 ctx.request_span = tracing::info_span!(
1769 "test_span",
1770 "http.response.status_code" = tracing::field::Empty,
1771 "otel.status_code" = tracing::field::Empty,
1772 "upstream.address" = tracing::field::Empty,
1773 "upstream.cluster" = tracing::field::Empty,
1774 );
1775
1776 record_response_span_fields(Some(http::StatusCode::SERVICE_UNAVAILABLE), "GET", None, &ctx);
1777
1778 let captured = capture.0.lock().expect("capture lock");
1779 let get = |name: &str| {
1780 let value = captured.iter().find(|(f, _)| f == name).map(|(_, v)| v.clone());
1781 assert!(value.is_some(), "field {name} not recorded; got {captured:?}");
1782 value.unwrap_or_default()
1783 };
1784 assert_eq!(get("http.response.status_code"), "503");
1785 assert_eq!(get("otel.status_code"), "\"ERROR\"", "5xx must set otel error status");
1786 assert_eq!(get("upstream.address"), "\"10.0.0.1:80\"");
1787 assert_eq!(get("upstream.cluster"), "\"api-cluster\"");
1788 }
1789
1790 #[test]
1791 fn record_response_span_fields_success_has_no_error_status() {
1792 use tracing_subscriber::layer::SubscriberExt as _;
1793
1794 let capture = RecordCapture::default();
1795 let subscriber = tracing_subscriber::registry().with(capture.clone());
1796 let _guard = tracing::subscriber::set_default(subscriber);
1797
1798 let mut ctx = PingoraRequestCtx::default();
1799 ctx.request_span = tracing::info_span!(
1800 "test_span",
1801 "http.response.status_code" = tracing::field::Empty,
1802 "otel.status_code" = tracing::field::Empty,
1803 );
1804
1805 record_response_span_fields(Some(http::StatusCode::OK), "GET", None, &ctx);
1806
1807 let captured = capture.0.lock().expect("capture lock");
1808 assert!(
1809 captured
1810 .iter()
1811 .any(|(f, v)| f == "http.response.status_code" && v == "200"),
1812 "status should be recorded: {captured:?}"
1813 );
1814 assert!(
1815 !captured.iter().any(|(f, _)| f == "otel.status_code"),
1816 "2xx must not set otel error status: {captured:?}"
1817 );
1818 }
1819
1820 #[test]
1821 fn record_response_span_attributes_records_exchange_span_fields() {
1822 let mut ctx = PingoraRequestCtx::default();
1823 ctx.request_span = tracing::info_span!(
1824 "test_request",
1825 "http.response.status_code" = tracing::field::Empty,
1826 "server.address" = tracing::field::Empty,
1827 "upstream.cluster" = tracing::field::Empty,
1828 );
1829 ctx.upstream_exchange_span = tracing::info_span!(
1830 parent: &ctx.request_span,
1831 "upstream_exchange",
1832 "http.response.status_code" = tracing::field::Empty,
1833 "http.response.body.size" = tracing::field::Empty,
1834 );
1835 ctx.response_body_bytes = 4096;
1836
1837 ctx.upstream_exchange_span.record("http.response.status_code", 200_u16);
1838 ctx.upstream_exchange_span.record("http.response.body.size", 4096_u64);
1839 }
1840
1841 #[test]
1842 fn record_response_span_attributes_skips_exchange_when_disabled() {
1843 let mut ctx = PingoraRequestCtx::default();
1844 ctx.request_span = tracing::info_span!(
1845 "test_request",
1846 "http.response.status_code" = tracing::field::Empty,
1847 "server.address" = tracing::field::Empty,
1848 "upstream.cluster" = tracing::field::Empty,
1849 );
1850 assert!(
1851 ctx.upstream_exchange_span.is_disabled(),
1852 "exchange span should be disabled by default"
1853 );
1854 }
1855
1856 #[test]
1861 fn retry_with_upstream_address_sets_retry_flag() {
1862 let mut ctx = PingoraRequestCtx::default();
1863 ctx.request_is_idempotent = true;
1864 ctx.upstream_for_retry = Some(Upstream {
1865 address: Arc::from("10.0.0.1:8080"),
1866 connection: Arc::new(ConnectionOptions::default()),
1867 tls: None,
1868 authority: None,
1869 });
1870 let e = handle_connect_failure(&mut ctx, make_error());
1871 assert!(e.retry(), "should retry with upstream address present");
1872 assert_eq!(ctx.retries, 1, "retry counter should increment to 1");
1873 }
1874
1875 #[test]
1876 fn retry_without_upstream_address_uses_fallback() {
1877 let mut ctx = PingoraRequestCtx::default();
1878 ctx.request_is_idempotent = true;
1879 ctx.upstream_for_retry = None;
1880 let e = handle_connect_failure(&mut ctx, make_error());
1881 assert!(
1882 e.retry(),
1883 "should retry even when upstream_for_retry is None (address defaults to unknown)"
1884 );
1885 assert_eq!(ctx.retries, 1, "retry counter should increment to 1");
1886 }
1887
1888 #[test]
1889 fn retry_exhausted_with_upstream_address_does_not_retry() {
1890 let mut ctx = PingoraRequestCtx::default();
1891 ctx.request_is_idempotent = true;
1892 ctx.retries = MAX_RETRIES as u32;
1893 ctx.upstream_for_retry = Some(Upstream {
1894 address: Arc::from("10.0.0.2:443"),
1895 connection: Arc::new(ConnectionOptions::default()),
1896 tls: None,
1897 authority: None,
1898 });
1899 let e = handle_connect_failure(&mut ctx, make_error());
1900 assert!(
1901 !e.retry(),
1902 "should not retry after MAX_RETRIES even with upstream address"
1903 );
1904 }
1905
1906 #[test]
1907 fn large_body_skip_with_upstream_address() {
1908 let mut ctx = PingoraRequestCtx::default();
1909 ctx.request_is_idempotent = true;
1910 ctx.request_body_bytes = RETRY_BODY_LIMIT + 1;
1911 ctx.upstream_for_retry = Some(Upstream {
1912 address: Arc::from("10.0.0.3:8080"),
1913 connection: Arc::new(ConnectionOptions::default()),
1914 tls: None,
1915 authority: None,
1916 });
1917 let e = handle_connect_failure(&mut ctx, make_error());
1918 assert!(!e.retry(), "should not retry large body even with upstream address");
1919 assert_eq!(ctx.retries, 0, "retry counter should not increment");
1920 }
1921
1922 fn make_error() -> Box<pingora_core::Error> {
1928 pingora_core::Error::explain(pingora_core::ErrorType::ConnectError, "test connect failure")
1929 }
1930
1931 fn access_log_pipeline() -> FilterPipeline {
1933 let registry = praxis_filter::FilterRegistry::with_builtins();
1934 let mut entries = vec![praxis_filter::FilterEntry {
1935 branch_chains: None,
1936 conditions: vec![],
1937 failure_mode: praxis_filter::FailureMode::default(),
1938 filter_type: "access_log".to_owned(),
1939 config: serde_yaml::Value::Null,
1940 name: None,
1941 response_conditions: vec![],
1942 }];
1943 FilterPipeline::build(&mut entries, ®istry).unwrap()
1944 }
1945
1946 fn make_fallback_ctx() -> PingoraRequestCtx {
1948 let mut ctx = PingoraRequestCtx::default();
1949 ctx.request_snapshot = Some(praxis_filter::Request {
1950 method: http::Method::GET,
1951 uri: "/incomplete".parse().unwrap(),
1952 headers: http::HeaderMap::new(),
1953 });
1954 ctx
1955 }
1956
1957 fn capture_access_events<F: FnOnce()>(f: F) -> Vec<String> {
1959 use tracing_subscriber::layer::SubscriberExt as _;
1960
1961 let messages = Arc::new(std::sync::Mutex::new(Vec::<String>::new()));
1962 let capture = AccessCapture(Arc::clone(&messages));
1963 let subscriber = tracing_subscriber::registry().with(capture);
1964 tracing::subscriber::with_default(subscriber, f);
1965 let mut guard = messages.lock().unwrap();
1966 std::mem::take(&mut *guard)
1967 }
1968
1969 struct AccessCapture(Arc<std::sync::Mutex<Vec<String>>>);
1971
1972 impl<S: tracing::Subscriber> tracing_subscriber::Layer<S> for AccessCapture {
1973 fn on_event(&self, event: &tracing::Event<'_>, _ctx: tracing_subscriber::layer::Context<'_, S>) {
1974 let mut visitor = AccessMessageVisitor(String::new());
1975 event.record(&mut visitor);
1976 if visitor.0.contains("access") {
1977 self.0.lock().unwrap().push(visitor.0);
1978 }
1979 }
1980 }
1981
1982 struct AccessMessageVisitor(String);
1984
1985 impl tracing::field::Visit for AccessMessageVisitor {
1986 fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
1987 if field.name() == "message" {
1988 self.0 = format!("{value:?}");
1989 }
1990 }
1991 }
1992
1993 fn make_passive_ctx(cluster: &str, endpoint_idx: usize, status: Option<u16>) -> PingoraRequestCtx {
1995 let mut ctx = PingoraRequestCtx::default();
1996 ctx.cluster = Some(Arc::from(cluster));
1997 ctx.selected_endpoint_index = Some(endpoint_idx);
1998 ctx.upstream_response_status = status;
1999 ctx.upstream_contacted = true;
2001 ctx
2002 }
2003
2004 fn make_passive_scenario(
2007 passive_unhealthy: Option<u32>,
2008 passive_healthy: Option<u32>,
2009 ) -> (FilterPipeline, PingoraRequestCtx) {
2010 use praxis_core::health::{ClusterHealthEntry, EndpointHealth};
2011
2012 let entry = ClusterHealthEntry::new(
2013 vec![EndpointHealth::new()],
2014 vec![Arc::from("10.0.0.1:80")],
2015 passive_unhealthy,
2016 passive_healthy,
2017 );
2018 let mut map = HashMap::new();
2019 map.insert(Arc::from("test-cluster"), Arc::new(entry));
2020 let health_registry = Arc::new(map);
2021
2022 let registry = praxis_filter::FilterRegistry::with_builtins();
2023 let mut pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
2024 pipeline.set_health_registry(health_registry);
2025
2026 let ctx = make_passive_ctx("test-cluster", 0, None);
2027
2028 (pipeline, ctx)
2029 }
2030}