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