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 request_method = session.req_header().method.as_str();
463 let raw_method = if request_method.is_empty() {
464 ctx.request_snapshot.as_ref().map_or("UNKNOWN", |r| r.method.as_str())
465 } else {
466 request_method
467 };
468 let method = metrics::method_label(raw_method);
469
470 let cluster = ctx.metrics_cluster_shared.clone().unwrap_or_else(metrics::cluster_none);
471
472 let route = ctx.metrics_route.clone().unwrap_or_else(metrics::route_unknown);
473
474 let labels = metrics::RequestMetricLabels {
475 cluster: cluster.clone(),
476 method,
477 route,
478 status_class,
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 record_passive_health(pipeline: &FilterPipeline, error: Option<&pingora_core::Error>, ctx: &PingoraRequestCtx) {
501 let cluster_name = ctx.cluster.as_ref().or(ctx.metrics_cluster.as_ref());
502 let Some(cluster_name) = cluster_name else {
503 return;
504 };
505 let Some(idx) = ctx.selected_endpoint_index else {
506 return;
507 };
508 let Some(registry) = pipeline.health_registry() else {
509 return;
510 };
511 let Some(health) = registry.get(cluster_name) else {
512 return;
513 };
514
515 let is_failure = error.is_some() || ctx.upstream_response_status.is_some_and(|s| s >= 500);
516 apply_passive_threshold(health, idx, cluster_name, is_failure);
517}
518
519fn apply_passive_threshold(
521 health: &praxis_core::health::ClusterHealthEntry,
522 idx: usize,
523 cluster_name: &Arc<str>,
524 is_failure: bool,
525) {
526 if is_failure {
527 if let Some(threshold) = health.passive_unhealthy_threshold()
528 && health
529 .endpoints()
530 .get(idx)
531 .is_some_and(|ep| ep.record_failure(threshold))
532 {
533 tracing::warn!(
534 cluster = %cluster_name,
535 endpoint_index = idx,
536 threshold,
537 "passive health: endpoint marked unhealthy"
538 );
539 emit_passive_health_transition(health, cluster_name, metrics::HEALTH_RESULT_UNHEALTHY);
540 }
541 } else if let Some(threshold) = health.passive_healthy_threshold()
542 && health
543 .endpoints()
544 .get(idx)
545 .is_some_and(|ep| ep.record_success(threshold))
546 {
547 tracing::info!(
548 cluster = %cluster_name,
549 endpoint_index = idx,
550 threshold,
551 "passive health: endpoint recovered"
552 );
553 emit_passive_health_transition(health, cluster_name, metrics::HEALTH_RESULT_HEALTHY);
554 }
555}
556
557fn emit_passive_health_transition(
559 health: &praxis_core::health::ClusterHealthEntry,
560 cluster_name: &Arc<str>,
561 result: &'static str,
562) {
563 let (healthy, total) = metrics::count_healthy_endpoints(health);
564 metrics::record_health_transition(
565 ::metrics::SharedString::from(Arc::clone(cluster_name)),
566 result,
567 healthy,
568 total,
569 );
570}
571
572pub(super) fn http_version_label(version: http::Version) -> &'static str {
576 match version {
577 http::Version::HTTP_09 => "0.9",
578 http::Version::HTTP_10 => "1.0",
579 http::Version::HTTP_11 => "1.1",
580 http::Version::HTTP_2 => "2",
581 http::Version::HTTP_3 => "3",
582 _ => "unknown",
583 }
584}
585
586fn record_response_span_attributes(session: &Session, ctx: &PingoraRequestCtx) {
595 if ctx.request_span.is_disabled() {
596 return;
597 }
598 let response = session.response_written();
599 let status = response.map(|resp| resp.status);
600 let method = session.req_header().method.as_str();
601 record_response_span_fields(status, method, response, ctx);
602}
603
604fn record_response_span_fields(
609 status: Option<http::StatusCode>,
610 method: &str,
611 response: Option<&pingora_http::ResponseHeader>,
612 ctx: &PingoraRequestCtx,
613) {
614 if let Some(status) = status {
615 let code = status.as_u16();
616 if code > 0 {
617 ctx.request_span.record("http.response.status_code", code);
618 }
619 if status.is_server_error() {
620 ctx.request_span.record("otel.status_code", "ERROR");
621 ctx.request_span.record("error.type", code.to_string().as_str());
624 }
625 }
626
627 if let Some(route) = &ctx.metrics_route {
628 ctx.request_span.record("http.route", route.as_ref());
629 ctx.request_span
630 .record("otel.name", format!("{method} {route}").as_str());
631 }
632
633 if let Some(upstream) = &ctx.upstream_for_retry {
634 ctx.request_span.record("upstream.address", upstream.address.as_ref());
635 }
636
637 if let Some(cluster) = &ctx.metrics_cluster {
638 ctx.request_span.record("upstream.cluster", cluster.as_ref());
639 }
640
641 record_upstream_exchange_span(ctx, response);
642}
643
644fn record_upstream_exchange_span(ctx: &PingoraRequestCtx, response: Option<&pingora_http::ResponseHeader>) {
646 if ctx.upstream_exchange_span.is_disabled() {
647 return;
648 }
649 if let Some(status) = ctx
652 .upstream_response_status
653 .or_else(|| response.map(|resp| resp.status.as_u16()))
654 {
655 ctx.upstream_exchange_span.record("http.response.status_code", status);
656 }
657 ctx.upstream_exchange_span
658 .record("http.response.body.size", ctx.response_body_bytes);
659}
660
661fn h2c_server_options() -> HttpServerOptions {
665 let mut opts = HttpServerOptions::default();
666 opts.h2c = true;
667 opts
668}
669
670fn h2_server_options() -> H2Options {
679 let mut opts = H2Options::new();
680 opts.max_header_list_size(65_536); opts.max_concurrent_streams(128);
682 opts
683}
684
685fn check_body_size_limit(body: Option<&Bytes>, accumulated_bytes: &mut u64, max_bytes: usize) -> bool {
688 if let Some(chunk) = body {
689 let chunk_len = chunk.len() as u64;
690 *accumulated_bytes += chunk_len;
691
692 let limit = max_bytes as u64;
693 return *accumulated_bytes > limit;
694 }
695 false
696}
697
698fn accumulate_stream_buffer(
701 body: &mut Option<Bytes>,
702 body_buffer: &mut Option<BodyBuffer>,
703 end_of_stream: bool,
704 max_bytes: Option<usize>,
705) -> bool {
706 if let Some(chunk) = &*body {
707 let limit = max_bytes.unwrap_or(ABSOLUTE_MAX_BODY_BYTES);
708 let buf = body_buffer.get_or_insert_with(|| BodyBuffer::new(limit));
709
710 if buf.push(chunk.clone()).is_err() {
711 return true;
712 }
713 }
714
715 if end_of_stream {
716 tracing::trace!("stream buffer: freezing accumulated body before pipeline at EOS");
717 *body = body_buffer.take().map(BodyBuffer::freeze);
718 } else {
719 tracing::trace!("stream buffer: filters see the original chunk");
720 }
721 false
722}
723
724#[expect(
727 clippy::fn_params_excessive_bools,
728 reason = "mirrors the caller's existing condition flags"
729)]
730fn suppress_stream_buffer_chunk(body: &mut Option<Bytes>, is_stream_buffer: bool, released: bool, end_of_stream: bool) {
731 if is_stream_buffer && !released && !end_of_stream {
732 *body = None;
733 }
734}
735
736fn release_stream_buffer(
738 body: &mut Option<Bytes>,
739 is_stream_buffer: bool,
740 released: &mut bool,
741 body_buffer: &mut Option<BodyBuffer>,
742 end_of_stream: bool,
743) {
744 if is_stream_buffer && !*released {
745 *released = true;
746 if !end_of_stream {
747 *body = body_buffer.take().map(BodyBuffer::freeze);
748 }
749 }
750}
751
752struct BodyFilterOutput {
757 cluster: Option<Arc<str>>,
759 upstream: Option<Upstream>,
761 extensions: RequestExtensions,
763 filter_metadata: HashMap<String, String>,
765 filter_state: HashMap<usize, Box<dyn std::any::Any + Send + Sync>>,
767 executed_filter_indices: Vec<bool>,
769 body_done_indices: Vec<bool>,
771 attempted_endpoints: Vec<Arc<str>>,
773}
774
775impl BodyFilterOutput {
776 fn take_from(fctx: &mut HttpFilterContext<'_>) -> Self {
779 Self {
780 cluster: fctx.cluster.take(),
781 upstream: fctx.upstream.take(),
782 extensions: std::mem::take(&mut fctx.extensions),
783 filter_metadata: std::mem::take(&mut fctx.filter_metadata),
784 filter_state: std::mem::take(&mut fctx.filter_state),
785 executed_filter_indices: std::mem::take(&mut fctx.executed_filter_indices),
786 body_done_indices: std::mem::take(&mut fctx.body_done_indices),
787 attempted_endpoints: std::mem::take(&mut fctx.attempted_endpoints),
788 }
789 }
790
791 fn write_back(self, ctx: &mut PingoraRequestCtx) {
793 ctx.cluster = self.cluster;
794 ctx.upstream = self.upstream;
795 ctx.extensions = self.extensions;
796 ctx.filter_metadata = self.filter_metadata;
797 ctx.filter_state = self.filter_state;
798 ctx.cached_executed_filter_indices = self.executed_filter_indices;
799 ctx.cached_body_done_indices = self.body_done_indices;
800 ctx.attempted_endpoints = self.attempted_endpoints;
801 }
802}
803
804#[cfg(test)]
809#[expect(clippy::allow_attributes, reason = "blanket test suppressions")]
810#[allow(
811 clippy::unwrap_used,
812 clippy::expect_used,
813 clippy::indexing_slicing,
814 clippy::field_reassign_with_default,
815 clippy::too_many_lines,
816 clippy::cast_possible_truncation,
817 clippy::significant_drop_tightening,
818 reason = "tests"
819)]
820mod tests {
821 use praxis_core::connectivity::ConnectionOptions;
822
823 use super::*;
824
825 const MAX_RETRIES: usize = praxis_core::config::DEFAULT_MAX_RETRIES as usize;
827
828 const RETRY_BODY_LIMIT: u64 = praxis_core::config::DEFAULT_RETRY_BODY_LIMIT_BYTES;
830
831 #[test]
832 fn first_failure_idempotent_sets_retry() {
833 let mut ctx = PingoraRequestCtx::default();
834 ctx.request_is_idempotent = true;
835 let e = handle_connect_failure(&mut ctx, make_error());
836 assert!(e.retry(), "first failure should set retry flag");
837 assert_eq!(ctx.retries, 1);
838 }
839
840 #[test]
841 fn large_body_skips_retry() {
842 let mut ctx = PingoraRequestCtx::default();
843 ctx.request_is_idempotent = true;
844 ctx.request_body_bytes = RETRY_BODY_LIMIT + 1;
845 let e = handle_connect_failure(&mut ctx, make_error());
846 assert!(!e.retry(), "should not retry when body exceeds retry buffer limit");
847 assert_eq!(ctx.retries, 0, "retry counter should not increment");
848 }
849
850 #[test]
851 fn mutated_body_exceeding_limit_skips_retry() {
852 let mut ctx = PingoraRequestCtx::default();
853 ctx.request_is_idempotent = true;
854 ctx.request_body_bytes = 1024;
855 ctx.mutated_request_body_len = Some((RETRY_BODY_LIMIT + 1) as usize);
856 let e = handle_connect_failure(&mut ctx, make_error());
857 assert!(
858 !e.retry(),
859 "should not retry when mutated body exceeds retry buffer limit"
860 );
861 assert_eq!(ctx.retries, 0);
862 }
863
864 #[test]
865 fn body_at_limit_allows_retry() {
866 let mut ctx = PingoraRequestCtx::default();
867 ctx.request_is_idempotent = true;
868 ctx.request_body_bytes = RETRY_BODY_LIMIT;
869 let e = handle_connect_failure(&mut ctx, make_error());
870 assert!(e.retry(), "body exactly at limit should allow retry");
871 assert_eq!(ctx.retries, 1);
872 }
873
874 #[test]
875 fn zero_body_allows_retry() {
876 let mut ctx = PingoraRequestCtx::default();
877 ctx.request_is_idempotent = true;
878 ctx.request_body_bytes = 0;
879 let e = handle_connect_failure(&mut ctx, make_error());
880 assert!(e.retry(), "zero-length body should allow retry");
881 assert_eq!(ctx.retries, 1);
882 }
883
884 #[test]
885 fn max_retries_exhausted_does_not_retry() {
886 let mut ctx = PingoraRequestCtx::default();
887 ctx.request_is_idempotent = true;
888 ctx.retries = MAX_RETRIES as u32;
889 let e = handle_connect_failure(&mut ctx, make_error());
890 assert!(!e.retry(), "should not retry after MAX_RETRIES");
891 assert_eq!(ctx.retries as usize, MAX_RETRIES);
892 }
893
894 #[test]
895 fn counter_increments_across_calls() {
896 let mut ctx = PingoraRequestCtx::default();
897 ctx.request_is_idempotent = true;
898 for expected in 1..=MAX_RETRIES {
899 let _result = handle_connect_failure(&mut ctx, make_error());
900 assert_eq!(ctx.retries as usize, expected);
901 }
902 let e = handle_connect_failure(&mut ctx, make_error());
903 assert!(!e.retry(), "should not retry after reaching MAX_RETRIES");
904 assert_eq!(ctx.retries as usize, MAX_RETRIES);
905 }
906
907 #[test]
908 fn non_idempotent_request_never_retries() {
909 let mut ctx = PingoraRequestCtx::default();
910 ctx.request_is_idempotent = false;
911 let e = handle_connect_failure(&mut ctx, make_error());
912 assert!(!e.retry(), "non-idempotent request should never retry");
913 assert_eq!(ctx.retries, 0);
914 }
915
916 #[test]
917 fn connect_failure_clears_upstream_connect_start() {
918 let mut ctx = PingoraRequestCtx::default();
919 ctx.upstream_connect_start = Some(std::time::Instant::now());
920 let _e = handle_connect_failure(&mut ctx, make_error());
921 assert!(
922 ctx.upstream_connect_start.is_none(),
923 "failed connect should consume upstream_connect_start for duration recording"
924 );
925 }
926
927 #[test]
928 fn response_503_retries_when_status5xx_enabled() {
929 let mut ctx = PingoraRequestCtx::default();
930 ctx.request_is_idempotent = true;
931 ctx.retry_policy = Some(Arc::new(praxis_core::config::RetryPolicy {
932 configured: true,
933 retriable_conditions: vec![praxis_core::config::RetriableCondition::Status5xx],
934 ..praxis_core::config::RetryPolicy::legacy_default()
935 }));
936 let e = maybe_retry_response(&mut ctx, 503).expect("503 should be retriable");
937 assert!(e.retry(), "503 under Status5xx should set retry");
938 assert_eq!(ctx.retries, 1);
939 assert!(ctx.reselect_on_retry);
940 assert!(ctx.pending_backoff.is_some());
941 }
942
943 #[test]
944 fn response_404_does_not_retry() {
945 let mut ctx = PingoraRequestCtx::default();
946 ctx.request_is_idempotent = true;
947 ctx.retry_policy = Some(Arc::new(praxis_core::config::RetryPolicy {
948 retriable_conditions: vec![praxis_core::config::RetriableCondition::Status5xx],
949 ..praxis_core::config::RetryPolicy::legacy_default()
950 }));
951 assert!(
952 maybe_retry_response(&mut ctx, 404).is_none(),
953 "404 must never trigger status-based retry"
954 );
955 assert_eq!(ctx.retries, 0);
956 }
957
958 #[test]
959 fn response_502_does_not_retry_under_legacy_default() {
960 let mut ctx = PingoraRequestCtx::default();
961 ctx.request_is_idempotent = true;
962 assert!(
964 maybe_retry_response(&mut ctx, 502).is_none(),
965 "legacy default must forward 5xx without retry"
966 );
967 }
968
969 #[test]
970 fn max_retries_zero_disables_connect_retry() {
971 let mut ctx = PingoraRequestCtx::default();
972 ctx.request_is_idempotent = true;
973 ctx.retry_policy = Some(Arc::new(praxis_core::config::RetryPolicy {
974 max_retries: Some(0),
975 ..praxis_core::config::RetryPolicy::legacy_default()
976 }));
977 let e = handle_connect_failure(&mut ctx, make_error());
978 assert!(!e.retry(), "max_retries: 0 must disable retries");
979 assert_eq!(ctx.retries, 0);
980 }
981
982 #[test]
983 fn non_idempotent_clears_pingora_default_retry_flag() {
984 let mut ctx = PingoraRequestCtx::default();
985 ctx.request_is_idempotent = false;
986 let mut e = make_error();
988 e.set_retry(true);
989 let e = handle_connect_failure(&mut ctx, e);
990 assert!(!e.retry(), "policy denial must clear Pingora's default retry flag");
991 assert_eq!(ctx.retries, 0);
992 }
993
994 #[tokio::test]
995 async fn logging_cleanup_noop_when_response_phase_done() {
996 let registry = praxis_filter::FilterRegistry::with_builtins();
997 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
998 let mut ctx = PingoraRequestCtx::default();
999 ctx.response_phase_done = true;
1000 ctx.request_snapshot = Some(praxis_filter::Request {
1001 method: http::Method::GET,
1002 uri: "/".parse().unwrap(),
1003 headers: http::HeaderMap::new(),
1004 });
1005 logging_cleanup(&pipeline, &mut ctx).await;
1006 }
1007
1008 #[tokio::test]
1009 async fn logging_cleanup_noop_when_no_snapshot() {
1010 let registry = praxis_filter::FilterRegistry::with_builtins();
1011 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
1012 let mut ctx = PingoraRequestCtx::default();
1013 ctx.response_phase_done = false;
1014 ctx.request_snapshot = None;
1015 logging_cleanup(&pipeline, &mut ctx).await;
1016 }
1017
1018 #[tokio::test]
1019 async fn logging_cleanup_runs_response_pipeline_when_needed() {
1020 let registry = praxis_filter::FilterRegistry::with_builtins();
1021 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
1022 let mut ctx = PingoraRequestCtx::default();
1023 ctx.response_phase_done = false;
1024 ctx.cluster = Some(Arc::from("test-cluster"));
1025 ctx.request_snapshot = Some(praxis_filter::Request {
1026 method: http::Method::GET,
1027 uri: "/test".parse().unwrap(),
1028 headers: http::HeaderMap::new(),
1029 });
1030 logging_cleanup(&pipeline, &mut ctx).await;
1031 assert_eq!(
1032 ctx.cluster.as_deref(),
1033 Some("test-cluster"),
1034 "cluster must be restored so the fallback access record can attribute the failure"
1035 );
1036 }
1037
1038 #[tokio::test]
1039 async fn logging_cleanup_preserves_filter_metadata() {
1040 let registry = praxis_filter::FilterRegistry::with_builtins();
1041 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
1042 let mut ctx = PingoraRequestCtx::default();
1043 ctx.response_phase_done = false;
1044 ctx.filter_metadata
1045 .insert("json_rpc.method".to_owned(), "service/invoke".to_owned());
1046 ctx.request_snapshot = Some(praxis_filter::Request {
1047 method: http::Method::POST,
1048 uri: "/api".parse().unwrap(),
1049 headers: http::HeaderMap::new(),
1050 });
1051 logging_cleanup(&pipeline, &mut ctx).await;
1052 assert_eq!(
1053 ctx.filter_metadata.get("json_rpc.method").map(String::as_str),
1054 Some("service/invoke"),
1055 "filter_metadata should survive logging_cleanup"
1056 );
1057 }
1058
1059 #[tokio::test]
1060 async fn logging_cleanup_preserves_extensions() {
1061 let registry = praxis_filter::FilterRegistry::with_builtins();
1062 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
1063 let mut ctx = PingoraRequestCtx::default();
1064 ctx.response_phase_done = false;
1065 ctx.extensions.insert(42_u32);
1066 ctx.request_snapshot = Some(praxis_filter::Request {
1067 method: http::Method::POST,
1068 uri: "/test".parse().unwrap(),
1069 headers: http::HeaderMap::new(),
1070 });
1071 logging_cleanup(&pipeline, &mut ctx).await;
1072 assert_eq!(
1073 ctx.extensions.get::<u32>(),
1074 Some(&42),
1075 "extensions should survive logging_cleanup"
1076 );
1077 }
1078
1079 #[test]
1080 fn passive_health_error_is_failure() {
1081 let (pipeline, ctx) = make_passive_scenario(Some(3), Some(2));
1082 let error = make_error();
1083 record_passive_health(&pipeline, Some(&error), &ctx);
1084
1085 let registry = pipeline.health_registry().unwrap();
1086 let entry = registry.get("test-cluster").unwrap();
1087 assert!(
1088 entry.endpoints()[0].is_healthy(),
1089 "single failure should not yet mark unhealthy (threshold=3)"
1090 );
1091 }
1092
1093 #[test]
1094 fn passive_health_status_500_is_failure() {
1095 let (pipeline, mut ctx) = make_passive_scenario(Some(3), Some(2));
1096 ctx.upstream_response_status = Some(500);
1097 record_passive_health(&pipeline, None, &ctx);
1098
1099 let registry = pipeline.health_registry().unwrap();
1100 let entry = registry.get("test-cluster").unwrap();
1101 assert!(
1102 entry.endpoints()[0].is_healthy(),
1103 "single 500 should not yet mark unhealthy (threshold=3)"
1104 );
1105 }
1106
1107 #[test]
1108 fn passive_health_status_below_500_is_success() {
1109 let (pipeline, mut ctx) = make_passive_scenario(Some(2), Some(1));
1110 ctx.upstream_response_status = Some(499);
1111 record_passive_health(&pipeline, None, &ctx);
1112
1113 let registry = pipeline.health_registry().unwrap();
1114 let entry = registry.get("test-cluster").unwrap();
1115 assert!(entry.endpoints()[0].is_healthy(), "status 499 should count as success");
1116 }
1117
1118 #[test]
1119 fn passive_unhealthy_threshold_transition() {
1120 let (pipeline, ctx) = make_passive_scenario(Some(2), Some(1));
1121 let error = make_error();
1122 record_passive_health(&pipeline, Some(&error), &ctx);
1123 record_passive_health(&pipeline, Some(&error), &ctx);
1124
1125 let registry = pipeline.health_registry().unwrap();
1126 let entry = registry.get("test-cluster").unwrap();
1127 assert!(
1128 !entry.endpoints()[0].is_healthy(),
1129 "2 consecutive failures should mark unhealthy (threshold=2)"
1130 );
1131 }
1132
1133 #[test]
1134 fn passive_healthy_threshold_recovery() {
1135 let (pipeline, ctx) = make_passive_scenario(Some(1), Some(2));
1136 let error = make_error();
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 "should be unhealthy after 1 failure"
1144 );
1145
1146 let ctx_ok = make_passive_ctx("test-cluster", 0, Some(200));
1147 record_passive_health(&pipeline, None, &ctx_ok);
1148 assert!(
1149 !entry.endpoints()[0].is_healthy(),
1150 "one success should not recover (threshold=2)"
1151 );
1152
1153 record_passive_health(&pipeline, None, &ctx_ok);
1154 assert!(
1155 entry.endpoints()[0].is_healthy(),
1156 "2 consecutive successes should recover (threshold=2)"
1157 );
1158 }
1159
1160 #[test]
1161 fn passive_health_no_thresholds_is_noop() {
1162 let (pipeline, ctx) = make_passive_scenario(None, None);
1163 let error = make_error();
1164 record_passive_health(&pipeline, Some(&error), &ctx);
1165
1166 let registry = pipeline.health_registry().unwrap();
1167 let entry = registry.get("test-cluster").unwrap();
1168 assert!(
1169 entry.endpoints()[0].is_healthy(),
1170 "no passive thresholds means failures are no-op"
1171 );
1172 }
1173
1174 #[test]
1175 fn passive_health_endpoint_index_out_of_bounds() {
1176 let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
1177 ctx.selected_endpoint_index = Some(999);
1178 let error = make_error();
1179 record_passive_health(&pipeline, Some(&error), &ctx);
1180
1181 let registry = pipeline.health_registry().unwrap();
1182 let entry = registry.get("test-cluster").unwrap();
1183 assert!(entry.endpoints()[0].is_healthy(), "out-of-bounds index should be no-op");
1184 }
1185
1186 #[test]
1187 fn passive_health_missing_cluster_is_noop() {
1188 let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
1189 ctx.cluster = None;
1190 ctx.metrics_cluster = None;
1191 let error = make_error();
1192 record_passive_health(&pipeline, Some(&error), &ctx);
1193 }
1194
1195 #[test]
1196 fn passive_health_falls_back_to_metrics_cluster() {
1197 let (pipeline, mut ctx) = make_passive_scenario(Some(2), Some(1));
1198 ctx.cluster = None;
1199 ctx.metrics_cluster = Some(Arc::from("test-cluster"));
1200 let error = make_error();
1201 record_passive_health(&pipeline, Some(&error), &ctx);
1202 record_passive_health(&pipeline, Some(&error), &ctx);
1203
1204 let registry = pipeline.health_registry().unwrap();
1205 let entry = registry.get("test-cluster").unwrap();
1206 assert!(
1207 !entry.endpoints()[0].is_healthy(),
1208 "fallback to metrics_cluster should still record passive health"
1209 );
1210 }
1211
1212 #[test]
1213 fn passive_health_missing_endpoint_index_is_noop() {
1214 let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
1215 ctx.selected_endpoint_index = None;
1216 let error = make_error();
1217 record_passive_health(&pipeline, Some(&error), &ctx);
1218 }
1219
1220 #[test]
1221 fn passive_health_missing_registry_is_noop() {
1222 let registry = praxis_filter::FilterRegistry::with_builtins();
1223 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
1224 let mut ctx = PingoraRequestCtx::default();
1225 ctx.cluster = Some(Arc::from("test-cluster"));
1226 ctx.selected_endpoint_index = Some(0);
1227 let error = make_error();
1228 record_passive_health(&pipeline, Some(&error), &ctx);
1229 }
1230
1231 #[test]
1232 fn passive_health_unknown_cluster_is_noop() {
1233 let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
1234 ctx.cluster = Some(Arc::from("nonexistent"));
1235 let error = make_error();
1236 record_passive_health(&pipeline, Some(&error), &ctx);
1237 }
1238
1239 #[test]
1240 fn size_limit_none_body_returns_false() {
1241 let mut bytes = 0_u64;
1242 assert!(!check_body_size_limit(None, &mut bytes, 100));
1243 assert_eq!(bytes, 0, "accumulated bytes unchanged for None body");
1244 }
1245
1246 #[test]
1247 fn size_limit_within_limit() {
1248 let mut bytes = 0_u64;
1249 let body = Some(Bytes::from_static(b"hello"));
1250 assert!(!check_body_size_limit(body.as_ref(), &mut bytes, 10));
1251 assert_eq!(bytes, 5);
1252 }
1253
1254 #[test]
1255 fn size_limit_at_exact_limit() {
1256 let mut bytes = 0_u64;
1257 let body = Some(Bytes::from_static(b"exact"));
1258 assert!(!check_body_size_limit(body.as_ref(), &mut bytes, 5));
1259 assert_eq!(bytes, 5);
1260 }
1261
1262 #[test]
1263 fn size_limit_exceeds_limit() {
1264 let mut bytes = 0_u64;
1265 let body = Some(Bytes::from_static(b"toolong"));
1266 assert!(check_body_size_limit(body.as_ref(), &mut bytes, 3));
1267 }
1268
1269 #[test]
1270 fn size_limit_cumulative_overflow() {
1271 let mut bytes = 0_u64;
1272 let first = Some(Bytes::from_static(b"aaa"));
1273 assert!(!check_body_size_limit(first.as_ref(), &mut bytes, 5));
1274
1275 let second = Some(Bytes::from_static(b"bbb"));
1276 assert!(check_body_size_limit(second.as_ref(), &mut bytes, 5));
1277 assert_eq!(bytes, 6);
1278 }
1279
1280 #[test]
1281 fn stream_buffer_accumulates_chunks() {
1282 let mut body = Some(Bytes::from_static(b"hello "));
1283 let mut buf: Option<BodyBuffer> = None;
1284 assert!(!accumulate_stream_buffer(&mut body, &mut buf, false, Some(100)));
1285 assert!(buf.is_some());
1286
1287 body = Some(Bytes::from_static(b"world"));
1288 assert!(!accumulate_stream_buffer(&mut body, &mut buf, false, Some(100)));
1289
1290 let frozen = buf.take().unwrap().freeze();
1291 assert_eq!(frozen, Bytes::from_static(b"hello world"));
1292 }
1293
1294 #[test]
1295 fn stream_buffer_freezes_at_eos() {
1296 let mut body = Some(Bytes::from_static(b"data"));
1297 let mut buf: Option<BodyBuffer> = None;
1298 assert!(!accumulate_stream_buffer(&mut body, &mut buf, false, Some(100)));
1299
1300 body = Some(Bytes::from_static(b" end"));
1301 assert!(!accumulate_stream_buffer(&mut body, &mut buf, true, Some(100)));
1302 assert!(buf.is_none(), "buffer should be taken at EOS");
1303 assert_eq!(body.unwrap(), Bytes::from_static(b"data end"));
1304 }
1305
1306 #[test]
1307 fn stream_buffer_overflow() {
1308 let mut body = Some(Bytes::from_static(b"too long"));
1309 let mut buf: Option<BodyBuffer> = None;
1310 assert!(accumulate_stream_buffer(&mut body, &mut buf, false, Some(5)));
1311 }
1312
1313 #[test]
1314 fn stream_buffer_none_body() {
1315 let mut body: Option<Bytes> = None;
1316 let mut buf: Option<BodyBuffer> = None;
1317 assert!(!accumulate_stream_buffer(&mut body, &mut buf, false, Some(100)));
1318 assert!(buf.is_none());
1319 }
1320
1321 #[test]
1322 fn stream_buffer_uses_absolute_max_when_none() {
1323 let mut body = Some(Bytes::from_static(b"data"));
1324 let mut buf: Option<BodyBuffer> = None;
1325 assert!(!accumulate_stream_buffer(&mut body, &mut buf, false, None));
1326 assert!(buf.is_some(), "should create buffer with absolute max");
1327 }
1328
1329 #[test]
1330 fn suppress_clears_body_when_buffering() {
1331 let mut body = Some(Bytes::from_static(b"data"));
1332 suppress_stream_buffer_chunk(&mut body, true, false, false);
1333 assert!(body.is_none());
1334 }
1335
1336 #[test]
1337 fn suppress_noop_when_not_stream_buffer() {
1338 let mut body = Some(Bytes::from_static(b"data"));
1339 suppress_stream_buffer_chunk(&mut body, false, false, false);
1340 assert!(body.is_some());
1341 }
1342
1343 #[test]
1344 fn suppress_noop_when_released() {
1345 let mut body = Some(Bytes::from_static(b"data"));
1346 suppress_stream_buffer_chunk(&mut body, true, true, false);
1347 assert!(body.is_some());
1348 }
1349
1350 #[test]
1351 fn suppress_noop_at_eos() {
1352 let mut body = Some(Bytes::from_static(b"data"));
1353 suppress_stream_buffer_chunk(&mut body, true, false, true);
1354 assert!(body.is_some());
1355 }
1356
1357 #[test]
1358 fn release_sets_flag_and_flushes_buffer() {
1359 let mut body: Option<Bytes> = None;
1360 let mut released = false;
1361 let mut buf = Some(BodyBuffer::new(100));
1362 buf.as_mut().unwrap().push(Bytes::from_static(b"buffered")).unwrap();
1363
1364 release_stream_buffer(&mut body, true, &mut released, &mut buf, false);
1365 assert!(released);
1366 assert_eq!(body.unwrap(), Bytes::from_static(b"buffered"));
1367 assert!(buf.is_none());
1368 }
1369
1370 #[test]
1371 fn release_noop_when_already_released() {
1372 let mut body: Option<Bytes> = None;
1373 let mut released = true;
1374 let mut buf: Option<BodyBuffer> = None;
1375
1376 release_stream_buffer(&mut body, true, &mut released, &mut buf, false);
1377 assert!(body.is_none(), "body should be unchanged when already released");
1378 }
1379
1380 #[test]
1381 fn release_noop_when_not_stream_buffer() {
1382 let mut body: Option<Bytes> = None;
1383 let mut released = false;
1384 let mut buf: Option<BodyBuffer> = None;
1385
1386 release_stream_buffer(&mut body, false, &mut released, &mut buf, false);
1387 assert!(!released, "released flag should be unchanged for non-stream-buffer");
1388 }
1389
1390 #[test]
1391 fn release_at_eos_sets_flag_but_no_flush() {
1392 let mut body: Option<Bytes> = None;
1393 let mut released = false;
1394 let mut buf = Some(BodyBuffer::new(100));
1395 buf.as_mut().unwrap().push(Bytes::from_static(b"data")).unwrap();
1396
1397 release_stream_buffer(&mut body, true, &mut released, &mut buf, true);
1398 assert!(released);
1399 assert!(body.is_none(), "body should not be overwritten at EOS");
1400 assert!(buf.is_some(), "buffer should not be taken at EOS");
1401 }
1402
1403 #[test]
1404 fn write_back_transfers_fields() {
1405 let mut ctx = PingoraRequestCtx::default();
1406
1407 let mut extensions = RequestExtensions::new();
1408 extensions.insert(42_u32);
1409
1410 let state_val: Box<dyn std::any::Any + Send + Sync> = Box::new(99_i32);
1411 let filter_state = HashMap::from([(0_usize, state_val)]);
1412
1413 let output = BodyFilterOutput {
1414 cluster: Some(Arc::from("test-cluster")),
1415 upstream: Some(Upstream {
1416 address: Arc::from("10.0.0.1:80"),
1417 authority: None,
1418 connection: Arc::new(ConnectionOptions::default()),
1419 tls: None,
1420 }),
1421 extensions,
1422 attempted_endpoints: Vec::new(),
1423 filter_metadata: HashMap::from([("key".to_owned(), "val".to_owned())]),
1424 filter_state,
1425 executed_filter_indices: vec![true, false],
1426 body_done_indices: vec![false, true],
1427 };
1428 output.write_back(&mut ctx);
1429
1430 assert_eq!(ctx.cluster.as_deref(), Some("test-cluster"));
1431 assert!(ctx.upstream.is_some(), "upstream should transfer");
1432 assert_eq!(ctx.upstream.as_ref().unwrap().address.as_ref(), "10.0.0.1:80");
1433 assert_eq!(ctx.extensions.get::<u32>(), Some(&42));
1434 assert_eq!(ctx.filter_metadata.get("key").map(String::as_str), Some("val"));
1435 assert_eq!(ctx.filter_state.len(), 1, "filter_state should transfer");
1436 assert_eq!(
1437 ctx.filter_state.get(&0).and_then(|v| v.downcast_ref::<i32>()),
1438 Some(&99)
1439 );
1440 assert_eq!(ctx.cached_executed_filter_indices, vec![true, false]);
1441 assert_eq!(ctx.cached_body_done_indices, vec![false, true]);
1442 }
1443
1444 #[test]
1449 fn fallback_access_log_emits_for_incomplete_request() {
1450 let pipeline = access_log_pipeline();
1451 let mut ctx = make_fallback_ctx();
1452
1453 let events = capture_access_events(|| maybe_emit_fallback_access_log(&pipeline, 502, &mut ctx));
1454 assert_eq!(
1455 events.len(),
1456 1,
1457 "incomplete request must produce a fallback access record"
1458 );
1459 }
1460
1461 #[test]
1462 fn fallback_access_log_skips_completed_delivery() {
1463 let pipeline = access_log_pipeline();
1464 let mut ctx = make_fallback_ctx();
1465 ctx.response_delivery_complete = true;
1466
1467 let events = capture_access_events(|| maybe_emit_fallback_access_log(&pipeline, 200, &mut ctx));
1468 assert!(events.is_empty(), "completed delivery already logged via the filter");
1469 }
1470
1471 #[test]
1472 fn fallback_access_log_skips_upgraded_connections() {
1473 let pipeline = access_log_pipeline();
1474 let mut ctx = make_fallback_ctx();
1475 ctx.connection_upgraded = true;
1476
1477 let events = capture_access_events(|| maybe_emit_fallback_access_log(&pipeline, 101, &mut ctx));
1478 assert!(events.is_empty(), "upgraded connections have no body completion");
1479 }
1480
1481 #[test]
1482 fn fallback_access_log_skips_without_access_log_filter() {
1483 let registry = praxis_filter::FilterRegistry::with_builtins();
1484 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
1485 let mut ctx = make_fallback_ctx();
1486
1487 let events = capture_access_events(|| maybe_emit_fallback_access_log(&pipeline, 502, &mut ctx));
1488 assert!(events.is_empty(), "no access_log filter means no fallback record");
1489 }
1490
1491 #[test]
1492 fn fallback_access_log_honors_entry_conditions() {
1493 let registry = praxis_filter::FilterRegistry::with_builtins();
1494 let mut entries = vec![praxis_filter::FilterEntry {
1495 branch_chains: None,
1496 conditions: vec![serde_yaml::from_str("when:\n path_prefix: /api\n").unwrap()],
1497 failure_mode: praxis_filter::FailureMode::default(),
1498 filter_type: "access_log".to_owned(),
1499 config: serde_yaml::Value::Null,
1500 name: None,
1501 response_conditions: vec![],
1502 }];
1503 let pipeline = FilterPipeline::build(&mut entries, ®istry).unwrap();
1504
1505 let mut excluded = make_fallback_ctx();
1506 let events = capture_access_events(|| maybe_emit_fallback_access_log(&pipeline, 502, &mut excluded));
1507 assert!(
1508 events.is_empty(),
1509 "requests the operator scoped out must not gain fallback records"
1510 );
1511
1512 let mut included = make_fallback_ctx();
1513 if let Some(snapshot) = included.request_snapshot.as_mut() {
1514 snapshot.uri = "/api/users".parse().unwrap();
1515 }
1516 let events = capture_access_events(|| maybe_emit_fallback_access_log(&pipeline, 502, &mut included));
1517 assert_eq!(events.len(), 1, "in-scope incomplete requests still get a record");
1518 }
1519
1520 #[test]
1521 fn aborted_response_body_at_eos_is_not_marked_delivered() {
1522 let pipeline = access_log_pipeline();
1525 let mut ctx = make_fallback_ctx();
1526 ctx.response_body_mode = BodyMode::SizeLimit { max_bytes: 4 };
1527 let mut body = Some(Bytes::from_static(b"exceeds the limit"));
1528
1529 let result = response_body_filter::execute(&pipeline, &mut body, true, &mut ctx);
1530 assert!(result.is_err(), "over-limit body must abort");
1531 assert!(
1532 !ctx.response_delivery_complete,
1533 "a response aborted at end-of-stream was not delivered; the fallback record must fire"
1534 );
1535 }
1536
1537 #[test]
1538 fn response_body_eos_marks_delivery_complete() {
1539 let registry = praxis_filter::FilterRegistry::with_builtins();
1540 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
1541 let mut ctx = PingoraRequestCtx::default();
1542 let mut body: Option<Bytes> = None;
1543
1544 let _timeout = response_body_filter::execute(&pipeline, &mut body, false, &mut ctx).unwrap();
1545 assert!(
1546 !ctx.response_delivery_complete,
1547 "mid-stream chunks must not mark delivery complete"
1548 );
1549
1550 let _timeout = response_body_filter::execute(&pipeline, &mut body, true, &mut ctx).unwrap();
1551 assert!(ctx.response_delivery_complete, "end-of-stream marks delivery complete");
1552 }
1553
1554 #[test]
1559 fn http_version_label_http_09() {
1560 assert_eq!(
1561 http_version_label(http::Version::HTTP_09),
1562 "0.9",
1563 "HTTP/0.9 should map to '0.9'"
1564 );
1565 }
1566
1567 #[test]
1568 fn http_version_label_http_10() {
1569 assert_eq!(
1570 http_version_label(http::Version::HTTP_10),
1571 "1.0",
1572 "HTTP/1.0 should map to '1.0'"
1573 );
1574 }
1575
1576 #[test]
1577 fn http_version_label_http_11() {
1578 assert_eq!(
1579 http_version_label(http::Version::HTTP_11),
1580 "1.1",
1581 "HTTP/1.1 should map to '1.1'"
1582 );
1583 }
1584
1585 #[test]
1586 fn http_version_label_http_2() {
1587 assert_eq!(
1588 http_version_label(http::Version::HTTP_2),
1589 "2",
1590 "HTTP/2 should map to '2'"
1591 );
1592 }
1593
1594 #[test]
1595 fn http_version_label_http_3() {
1596 assert_eq!(
1597 http_version_label(http::Version::HTTP_3),
1598 "3",
1599 "HTTP/3 should map to '3'"
1600 );
1601 }
1602
1603 #[test]
1604 fn record_response_span_attributes_noop_for_disabled_span() {
1605 let ctx = PingoraRequestCtx::default();
1606 assert!(ctx.request_span.is_disabled(), "default span should be disabled");
1607 }
1608
1609 #[derive(Clone, Default)]
1611 struct RecordCapture(Arc<std::sync::Mutex<Vec<(String, String)>>>);
1612
1613 impl<S> tracing_subscriber::Layer<S> for RecordCapture
1614 where
1615 S: tracing::Subscriber + for<'a> tracing_subscriber::registry::LookupSpan<'a>,
1616 {
1617 fn on_record(
1618 &self,
1619 _id: &tracing::span::Id,
1620 values: &tracing::span::Record<'_>,
1621 _ctx: tracing_subscriber::layer::Context<'_, S>,
1622 ) {
1623 struct Visitor<'a>(&'a mut Vec<(String, String)>);
1624 impl tracing::field::Visit for Visitor<'_> {
1625 fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
1626 self.0.push((field.name().to_owned(), format!("{value:?}")));
1627 }
1628 }
1629 let mut captured = self.0.lock().expect("capture lock");
1630 values.record(&mut Visitor(&mut captured));
1631 }
1632 }
1633
1634 #[test]
1635 fn record_response_span_fields_records_status_upstream_and_cluster() {
1636 use tracing_subscriber::layer::SubscriberExt as _;
1637
1638 let capture = RecordCapture::default();
1639 let subscriber = tracing_subscriber::registry().with(capture.clone());
1640 let _guard = tracing::subscriber::set_default(subscriber);
1641
1642 let mut ctx = PingoraRequestCtx::default();
1643 ctx.metrics_cluster = Some(Arc::from("api-cluster"));
1644 ctx.upstream_for_retry = Some(Upstream {
1645 address: Arc::from("10.0.0.1:80"),
1646 authority: None,
1647 connection: Arc::new(ConnectionOptions::default()),
1648 tls: None,
1649 });
1650 ctx.request_span = tracing::info_span!(
1651 "test_span",
1652 "http.response.status_code" = tracing::field::Empty,
1653 "otel.status_code" = tracing::field::Empty,
1654 "upstream.address" = tracing::field::Empty,
1655 "upstream.cluster" = tracing::field::Empty,
1656 );
1657
1658 record_response_span_fields(Some(http::StatusCode::SERVICE_UNAVAILABLE), "GET", None, &ctx);
1659
1660 let captured = capture.0.lock().expect("capture lock");
1661 let get = |name: &str| {
1662 let value = captured.iter().find(|(f, _)| f == name).map(|(_, v)| v.clone());
1663 assert!(value.is_some(), "field {name} not recorded; got {captured:?}");
1664 value.unwrap_or_default()
1665 };
1666 assert_eq!(get("http.response.status_code"), "503");
1667 assert_eq!(get("otel.status_code"), "\"ERROR\"", "5xx must set otel error status");
1668 assert_eq!(get("upstream.address"), "\"10.0.0.1:80\"");
1669 assert_eq!(get("upstream.cluster"), "\"api-cluster\"");
1670 }
1671
1672 #[test]
1673 fn record_response_span_fields_success_has_no_error_status() {
1674 use tracing_subscriber::layer::SubscriberExt as _;
1675
1676 let capture = RecordCapture::default();
1677 let subscriber = tracing_subscriber::registry().with(capture.clone());
1678 let _guard = tracing::subscriber::set_default(subscriber);
1679
1680 let mut ctx = PingoraRequestCtx::default();
1681 ctx.request_span = tracing::info_span!(
1682 "test_span",
1683 "http.response.status_code" = tracing::field::Empty,
1684 "otel.status_code" = tracing::field::Empty,
1685 );
1686
1687 record_response_span_fields(Some(http::StatusCode::OK), "GET", None, &ctx);
1688
1689 let captured = capture.0.lock().expect("capture lock");
1690 assert!(
1691 captured
1692 .iter()
1693 .any(|(f, v)| f == "http.response.status_code" && v == "200"),
1694 "status should be recorded: {captured:?}"
1695 );
1696 assert!(
1697 !captured.iter().any(|(f, _)| f == "otel.status_code"),
1698 "2xx must not set otel error status: {captured:?}"
1699 );
1700 }
1701
1702 #[test]
1703 fn record_response_span_attributes_records_exchange_span_fields() {
1704 let mut ctx = PingoraRequestCtx::default();
1705 ctx.request_span = tracing::info_span!(
1706 "test_request",
1707 "http.response.status_code" = tracing::field::Empty,
1708 "server.address" = tracing::field::Empty,
1709 "upstream.cluster" = tracing::field::Empty,
1710 );
1711 ctx.upstream_exchange_span = tracing::info_span!(
1712 parent: &ctx.request_span,
1713 "upstream_exchange",
1714 "http.response.status_code" = tracing::field::Empty,
1715 "http.response.body.size" = tracing::field::Empty,
1716 );
1717 ctx.response_body_bytes = 4096;
1718
1719 ctx.upstream_exchange_span.record("http.response.status_code", 200_u16);
1721 ctx.upstream_exchange_span.record("http.response.body.size", 4096_u64);
1722 }
1723
1724 #[test]
1725 fn record_response_span_attributes_skips_exchange_when_disabled() {
1726 let mut ctx = PingoraRequestCtx::default();
1727 ctx.request_span = tracing::info_span!(
1728 "test_request",
1729 "http.response.status_code" = tracing::field::Empty,
1730 "server.address" = tracing::field::Empty,
1731 "upstream.cluster" = tracing::field::Empty,
1732 );
1733 assert!(
1735 ctx.upstream_exchange_span.is_disabled(),
1736 "exchange span should be disabled by default"
1737 );
1738 }
1740
1741 #[test]
1746 fn retry_with_upstream_address_sets_retry_flag() {
1747 let mut ctx = PingoraRequestCtx::default();
1748 ctx.request_is_idempotent = true;
1749 ctx.upstream_for_retry = Some(Upstream {
1750 address: Arc::from("10.0.0.1:8080"),
1751 connection: Arc::new(ConnectionOptions::default()),
1752 tls: None,
1753 authority: None,
1754 });
1755 let e = handle_connect_failure(&mut ctx, make_error());
1756 assert!(e.retry(), "should retry with upstream address present");
1757 assert_eq!(ctx.retries, 1, "retry counter should increment to 1");
1758 }
1759
1760 #[test]
1761 fn retry_without_upstream_address_uses_fallback() {
1762 let mut ctx = PingoraRequestCtx::default();
1763 ctx.request_is_idempotent = true;
1764 ctx.upstream_for_retry = None;
1765 let e = handle_connect_failure(&mut ctx, make_error());
1766 assert!(
1767 e.retry(),
1768 "should retry even when upstream_for_retry is None (address defaults to unknown)"
1769 );
1770 assert_eq!(ctx.retries, 1, "retry counter should increment to 1");
1771 }
1772
1773 #[test]
1774 fn retry_exhausted_with_upstream_address_does_not_retry() {
1775 let mut ctx = PingoraRequestCtx::default();
1776 ctx.request_is_idempotent = true;
1777 ctx.retries = MAX_RETRIES as u32;
1778 ctx.upstream_for_retry = Some(Upstream {
1779 address: Arc::from("10.0.0.2:443"),
1780 connection: Arc::new(ConnectionOptions::default()),
1781 tls: None,
1782 authority: None,
1783 });
1784 let e = handle_connect_failure(&mut ctx, make_error());
1785 assert!(
1786 !e.retry(),
1787 "should not retry after MAX_RETRIES even with upstream address"
1788 );
1789 }
1790
1791 #[test]
1792 fn large_body_skip_with_upstream_address() {
1793 let mut ctx = PingoraRequestCtx::default();
1794 ctx.request_is_idempotent = true;
1795 ctx.request_body_bytes = RETRY_BODY_LIMIT + 1;
1796 ctx.upstream_for_retry = Some(Upstream {
1797 address: Arc::from("10.0.0.3:8080"),
1798 connection: Arc::new(ConnectionOptions::default()),
1799 tls: None,
1800 authority: None,
1801 });
1802 let e = handle_connect_failure(&mut ctx, make_error());
1803 assert!(!e.retry(), "should not retry large body even with upstream address");
1804 assert_eq!(ctx.retries, 0, "retry counter should not increment");
1805 }
1806
1807 fn make_error() -> Box<pingora_core::Error> {
1813 pingora_core::Error::explain(pingora_core::ErrorType::ConnectError, "test connect failure")
1814 }
1815
1816 fn access_log_pipeline() -> FilterPipeline {
1818 let registry = praxis_filter::FilterRegistry::with_builtins();
1819 let mut entries = vec![praxis_filter::FilterEntry {
1820 branch_chains: None,
1821 conditions: vec![],
1822 failure_mode: praxis_filter::FailureMode::default(),
1823 filter_type: "access_log".to_owned(),
1824 config: serde_yaml::Value::Null,
1825 name: None,
1826 response_conditions: vec![],
1827 }];
1828 FilterPipeline::build(&mut entries, ®istry).unwrap()
1829 }
1830
1831 fn make_fallback_ctx() -> PingoraRequestCtx {
1833 let mut ctx = PingoraRequestCtx::default();
1834 ctx.request_snapshot = Some(praxis_filter::Request {
1835 method: http::Method::GET,
1836 uri: "/incomplete".parse().unwrap(),
1837 headers: http::HeaderMap::new(),
1838 });
1839 ctx
1840 }
1841
1842 fn capture_access_events<F: FnOnce()>(f: F) -> Vec<String> {
1844 use tracing_subscriber::layer::SubscriberExt as _;
1845
1846 let messages = Arc::new(std::sync::Mutex::new(Vec::<String>::new()));
1847 let capture = AccessCapture(Arc::clone(&messages));
1848 let subscriber = tracing_subscriber::registry().with(capture);
1849 tracing::subscriber::with_default(subscriber, f);
1850 let mut guard = messages.lock().unwrap();
1851 std::mem::take(&mut *guard)
1852 }
1853
1854 struct AccessCapture(Arc<std::sync::Mutex<Vec<String>>>);
1856
1857 impl<S: tracing::Subscriber> tracing_subscriber::Layer<S> for AccessCapture {
1858 fn on_event(&self, event: &tracing::Event<'_>, _ctx: tracing_subscriber::layer::Context<'_, S>) {
1859 let mut visitor = AccessMessageVisitor(String::new());
1860 event.record(&mut visitor);
1861 if visitor.0.contains("access") {
1862 self.0.lock().unwrap().push(visitor.0);
1863 }
1864 }
1865 }
1866
1867 struct AccessMessageVisitor(String);
1869
1870 impl tracing::field::Visit for AccessMessageVisitor {
1871 fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
1872 if field.name() == "message" {
1873 self.0 = format!("{value:?}");
1874 }
1875 }
1876 }
1877
1878 fn make_passive_ctx(cluster: &str, endpoint_idx: usize, status: Option<u16>) -> PingoraRequestCtx {
1880 let mut ctx = PingoraRequestCtx::default();
1881 ctx.cluster = Some(Arc::from(cluster));
1882 ctx.selected_endpoint_index = Some(endpoint_idx);
1883 ctx.upstream_response_status = status;
1884 ctx
1885 }
1886
1887 fn make_passive_scenario(
1890 passive_unhealthy: Option<u32>,
1891 passive_healthy: Option<u32>,
1892 ) -> (FilterPipeline, PingoraRequestCtx) {
1893 use praxis_core::health::{ClusterHealthEntry, EndpointHealth};
1894
1895 let entry = ClusterHealthEntry::new(
1896 vec![EndpointHealth::new()],
1897 vec![Arc::from("10.0.0.1:80")],
1898 passive_unhealthy,
1899 passive_healthy,
1900 );
1901 let mut map = HashMap::new();
1902 map.insert(Arc::from("test-cluster"), Arc::new(entry));
1903 let health_registry = Arc::new(map);
1904
1905 let registry = praxis_filter::FilterRegistry::with_builtins();
1906 let mut pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
1907 pipeline.set_health_registry(health_registry);
1908
1909 let ctx = make_passive_ctx("test-cluster", 0, None);
1910
1911 (pipeline, ctx)
1912 }
1913}