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