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 fail_to_proxy;
34mod hop_by_hop;
36mod normalize;
38mod request_body_filter;
40mod request_filter;
42mod reserved_headers;
44mod response_body_filter;
46mod response_filter;
48mod upstream_peer;
50mod upstream_request;
52mod upstream_response;
54mod via;
56mod with_body;
58
59pub use upstream_peer::{UpstreamRetryGateRelease, arm_upstream_retry_gate, lock_upstream_retry_gate_tests};
60pub use with_body::PingoraHttpHandler;
61
62const MAX_RETRIES: usize = 3;
68
69const RETRY_BODY_LIMIT: u64 = 65_536; pub fn load_http_handler(
122 server: &mut Server,
123 listener: &praxis_core::config::Listener,
124 pipeline: Arc<ArcSwap<FilterPipeline>>,
125 cert_watcher_shutdowns: &mut Vec<tokio::sync::watch::Sender<bool>>,
126) -> Result<(), praxis_core::ProxyError> {
127 let downstream_read_timeout = listener.downstream_read_timeout_ms.map(Duration::from_millis);
128 let connection_semaphore = listener
129 .max_connections
130 .map(|max| Arc::new(Semaphore::new(max as usize)));
131
132 debug!(listener = %listener.name, "loading HTTP handler with body filters");
135 let handler = PingoraHttpHandler::new(
136 pipeline,
137 downstream_read_timeout,
138 connection_semaphore,
139 ::metrics::SharedString::from(listener.name.clone()),
140 );
141 wire_service(server, listener, handler, cert_watcher_shutdowns)?;
142 Ok(())
143}
144
145fn wire_service<H>(
147 server: &mut Server,
148 listener: &praxis_core::config::Listener,
149 handler: H,
150 cert_watcher_shutdowns: &mut Vec<tokio::sync::watch::Sender<bool>>,
151) -> Result<(), praxis_core::ProxyError>
152where
153 H: pingora_proxy::ProxyHttp + Send + Sync + 'static,
154 H::CTX: Send + Sync,
155{
156 let service_name = format!("http-proxy:{name}", name = listener.name);
157 let mut proxy = http_proxy(&server.configuration, handler);
158 proxy.server_options = Some(h2c_server_options());
159 proxy.h2_options = Some(h2_server_options());
160 let mut service = Service::new(service_name, proxy);
161 if let Some(tx) = super::listener::add_listener(&mut service, listener)? {
162 cert_watcher_shutdowns.push(tx);
163 }
164 server.add_service(service);
165 Ok(())
166}
167
168fn clamp_body_mode_to_ceiling(mode: BodyMode, baseline: BodyMode) -> BodyMode {
186 let ceiling = match baseline {
187 BodyMode::StreamBuffer { max_bytes: Some(v) } | BodyMode::SizeLimit { max_bytes: v } => Some(v),
188 _ => None,
189 };
190
191 match (mode, ceiling) {
192 (BodyMode::StreamBuffer { max_bytes }, Some(limit)) => BodyMode::StreamBuffer {
193 max_bytes: Some(max_bytes.map_or(limit, |v| v.min(limit))),
194 },
195 (BodyMode::SizeLimit { max_bytes }, Some(limit)) => BodyMode::SizeLimit {
196 max_bytes: max_bytes.min(limit),
197 },
198 (m, None | Some(_)) => m,
201 }
202}
203
204fn adjust_compression(
206 session: &mut Session,
207 upstream_response: &pingora_http::ResponseHeader,
208 compression: Option<&CompressionConfig>,
209) {
210 use pingora_core::{modules::http::compression::ResponseCompression, protocols::http::compression::Algorithm};
211
212 let Some(cfg) = compression else {
213 return;
214 };
215
216 let Some(module) = session.downstream_modules_ctx.get_mut::<ResponseCompression>() else {
217 return;
218 };
219
220 let headers = &upstream_response.headers;
221
222 if !cfg.should_compress(headers) {
223 debug!("disabling compression: response does not qualify");
224 module.adjust_level(0);
225 return;
226 }
227
228 for (enabled, level, algo) in [
229 (cfg.gzip_enabled, cfg.gzip_level, Algorithm::Gzip),
230 (cfg.brotli_enabled, cfg.brotli_level, Algorithm::Brotli),
231 (cfg.zstd_enabled, cfg.zstd_level, Algorithm::Zstd),
232 ] {
233 if !enabled {
234 module.adjust_algorithm_level(algo, 0);
235 } else if let Some(lvl) = level {
236 module.adjust_algorithm_level(algo, lvl);
237 }
238 }
239}
240
241fn handle_connect_failure(ctx: &mut PingoraRequestCtx, e: Box<pingora_core::Error>) -> Box<pingora_core::Error> {
249 let cluster = ctx.metrics_cluster_shared.clone().unwrap_or_else(metrics::cluster_none);
250 if let Some(start) = ctx.upstream_connect_start.take() {
251 metrics::record_upstream_connect_duration(cluster.clone(), start.elapsed().as_secs_f64());
252 }
253 metrics::record_upstream_connect_failure(cluster.clone());
254 if ctx.request_is_idempotent {
255 maybe_retry_idempotent_connect(ctx, cluster, e)
256 } else {
257 e
258 }
259}
260
261fn maybe_retry_idempotent_connect(
263 ctx: &mut PingoraRequestCtx,
264 cluster: ::metrics::SharedString,
265 e: Box<pingora_core::Error>,
266) -> Box<pingora_core::Error> {
267 let mutated_len = ctx.mutated_request_body_len.unwrap_or(0) as u64;
268 if std::cmp::max(ctx.request_body_bytes, mutated_len) > RETRY_BODY_LIMIT {
269 warn!(
270 body_bytes = ctx.request_body_bytes,
271 mutated_len = ?ctx.mutated_request_body_len,
272 limit = RETRY_BODY_LIMIT,
273 "skipping retry: request body exceeds Pingora retry buffer limit"
274 );
275 record_retry_exhausted_if_attempted(ctx, cluster);
276 return e;
277 }
278 if (ctx.retries as usize) < MAX_RETRIES {
279 ctx.retries += 1;
280 debug!(
281 retries = ctx.retries,
282 max = MAX_RETRIES,
283 "retrying idempotent request after connect failure"
284 );
285 let mut e = e;
286 e.set_retry(true);
287 return e;
288 }
289 warn!(
290 retries = ctx.retries,
291 max = MAX_RETRIES,
292 "retry limit reached for idempotent request"
293 );
294 metrics::record_upstream_retry(cluster, metrics::RETRY_RESULT_EXHAUSTED);
295 e
296}
297
298fn record_retry_exhausted_if_attempted(ctx: &PingoraRequestCtx, cluster: ::metrics::SharedString) {
300 if ctx.retries > 0 {
301 metrics::record_upstream_retry(cluster, metrics::RETRY_RESULT_EXHAUSTED);
302 }
303}
304
305async fn logging_cleanup(pipeline: &FilterPipeline, ctx: &mut PingoraRequestCtx) {
309 if !ctx.response_phase_done
310 && let Some(mut filter_ctx) = ctx.filter_context_for(pipeline, None)
311 {
312 let _result = pipeline.execute_http_response(&mut filter_ctx).await;
313 let extensions = filter_ctx.extensions;
314 let metadata = filter_ctx.filter_metadata;
315 let state = filter_ctx.filter_state;
316 let exec_idx = filter_ctx.executed_filter_indices;
317 let body_idx = filter_ctx.body_done_indices;
318 ctx.extensions = extensions;
319 ctx.filter_metadata = metadata;
320 ctx.filter_state = state;
321 ctx.cached_executed_filter_indices = exec_idx;
322 ctx.cached_body_done_indices = body_idx;
323 }
324}
325
326fn emit_request_metrics(session: &Session, ctx: &PingoraRequestCtx) {
330 if !metrics::is_recorder_installed() {
331 return;
332 }
333
334 let status_code = session.response_written().map_or(0, |resp| resp.status.as_u16());
335 let status_class = metrics::status_class(status_code);
336
337 let request_method = session.req_header().method.as_str();
338 let raw_method = if request_method.is_empty() {
339 ctx.request_snapshot.as_ref().map_or("UNKNOWN", |r| r.method.as_str())
340 } else {
341 request_method
342 };
343 let method = metrics::method_label(raw_method);
344
345 let cluster = ctx.metrics_cluster_shared.clone().unwrap_or_else(metrics::cluster_none);
346
347 let route = ctx.metrics_route.clone().unwrap_or_else(metrics::route_unknown);
348
349 let labels = metrics::RequestMetricLabels {
350 cluster: cluster.clone(),
351 method,
352 route,
353 status_class,
354 };
355
356 let duration_secs = ctx.request_start.elapsed().as_secs_f64();
357 metrics::record_request_metrics(labels, duration_secs);
358 metrics::record_body_size_metrics(
359 method,
360 status_class,
361 cluster,
362 ctx.request_body_bytes,
363 ctx.response_body_bytes,
364 );
365}
366
367fn record_passive_health(pipeline: &FilterPipeline, error: Option<&pingora_core::Error>, ctx: &PingoraRequestCtx) {
376 let cluster_name = ctx.cluster.as_ref().or(ctx.metrics_cluster.as_ref());
377 let Some(cluster_name) = cluster_name else {
378 return;
379 };
380 let Some(idx) = ctx.selected_endpoint_index else {
381 return;
382 };
383 let Some(registry) = pipeline.health_registry() else {
384 return;
385 };
386 let Some(health) = registry.get(cluster_name) else {
387 return;
388 };
389
390 let is_failure = error.is_some() || ctx.upstream_response_status.is_some_and(|s| s >= 500);
391 apply_passive_threshold(health, idx, cluster_name, is_failure);
392}
393
394fn apply_passive_threshold(
396 health: &praxis_core::health::ClusterHealthEntry,
397 idx: usize,
398 cluster_name: &Arc<str>,
399 is_failure: bool,
400) {
401 if is_failure {
402 if let Some(threshold) = health.passive_unhealthy_threshold()
403 && health
404 .endpoints()
405 .get(idx)
406 .is_some_and(|ep| ep.record_failure(threshold))
407 {
408 tracing::warn!(
409 cluster = %cluster_name,
410 endpoint_index = idx,
411 threshold,
412 "passive health: endpoint marked unhealthy"
413 );
414 emit_passive_health_transition(health, cluster_name, metrics::HEALTH_RESULT_UNHEALTHY);
415 }
416 } else if let Some(threshold) = health.passive_healthy_threshold()
417 && health
418 .endpoints()
419 .get(idx)
420 .is_some_and(|ep| ep.record_success(threshold))
421 {
422 tracing::info!(
423 cluster = %cluster_name,
424 endpoint_index = idx,
425 threshold,
426 "passive health: endpoint recovered"
427 );
428 emit_passive_health_transition(health, cluster_name, metrics::HEALTH_RESULT_HEALTHY);
429 }
430}
431
432fn emit_passive_health_transition(
434 health: &praxis_core::health::ClusterHealthEntry,
435 cluster_name: &Arc<str>,
436 result: &'static str,
437) {
438 let (healthy, total) = metrics::count_healthy_endpoints(health);
439 metrics::record_health_transition(
440 ::metrics::SharedString::from(Arc::clone(cluster_name)),
441 result,
442 healthy,
443 total,
444 );
445}
446
447pub(super) fn http_version_label(version: http::Version) -> &'static str {
451 match version {
452 http::Version::HTTP_09 => "0.9",
453 http::Version::HTTP_10 => "1.0",
454 http::Version::HTTP_11 => "1.1",
455 http::Version::HTTP_2 => "2",
456 http::Version::HTTP_3 => "3",
457 _ => "unknown",
458 }
459}
460
461fn record_response_span_attributes(session: &Session, ctx: &PingoraRequestCtx) {
467 if ctx.request_span.is_disabled() {
468 return;
469 }
470
471 if let Some(resp) = session.response_written() {
472 let status = resp.status.as_u16();
473 if status > 0 {
474 ctx.request_span.record("http.response.status_code", status);
475 }
476 if resp.status.is_server_error() {
477 ctx.request_span.record("otel.status_code", "ERROR");
478 }
479 }
480
481 if let Some(upstream) = &ctx.upstream_for_retry {
482 ctx.request_span.record("upstream.address", upstream.address.as_ref());
483 }
484
485 if let Some(cluster) = &ctx.metrics_cluster {
486 ctx.request_span.record("upstream.cluster", cluster.as_ref());
487 }
488}
489
490fn h2c_server_options() -> HttpServerOptions {
494 let mut opts = HttpServerOptions::default();
495 opts.h2c = true;
496 opts
497}
498
499fn h2_server_options() -> H2Options {
508 let mut opts = H2Options::new();
509 opts.max_header_list_size(65_536); opts.max_concurrent_streams(128);
511 opts
512}
513
514fn check_body_size_limit(body: Option<&Bytes>, accumulated_bytes: &mut u64, max_bytes: usize) -> bool {
517 if let Some(chunk) = body {
518 #[expect(clippy::allow_attributes, reason = "cast lint is platform-dependent")]
519 #[allow(clippy::cast_possible_truncation, reason = "chunk length fits u64")]
520 let chunk_len = chunk.len() as u64;
521 *accumulated_bytes += chunk_len;
522
523 #[expect(clippy::allow_attributes, reason = "cast lint is platform-dependent")]
524 #[allow(clippy::cast_possible_truncation, reason = "max_bytes fits u64")]
525 let limit = max_bytes as u64;
526 return *accumulated_bytes > limit;
527 }
528 false
529}
530
531fn accumulate_stream_buffer(
534 body: &mut Option<Bytes>,
535 body_buffer: &mut Option<BodyBuffer>,
536 end_of_stream: bool,
537 max_bytes: Option<usize>,
538) -> bool {
539 if let Some(chunk) = &*body {
540 let limit = max_bytes.unwrap_or(ABSOLUTE_MAX_BODY_BYTES);
541 let buf = body_buffer.get_or_insert_with(|| BodyBuffer::new(limit));
542
543 if buf.push(chunk.clone()).is_err() {
544 return true;
545 }
546 }
547
548 if end_of_stream {
549 tracing::trace!("stream buffer: freezing accumulated body before pipeline at EOS");
550 *body = body_buffer.take().map(BodyBuffer::freeze);
551 } else {
552 tracing::trace!("stream buffer: filters see the original chunk");
553 }
554 false
555}
556
557#[expect(
560 clippy::fn_params_excessive_bools,
561 reason = "mirrors the caller's existing condition flags"
562)]
563fn suppress_stream_buffer_chunk(body: &mut Option<Bytes>, is_stream_buffer: bool, released: bool, end_of_stream: bool) {
564 if is_stream_buffer && !released && !end_of_stream {
565 *body = None;
566 }
567}
568
569fn release_stream_buffer(
571 body: &mut Option<Bytes>,
572 is_stream_buffer: bool,
573 released: &mut bool,
574 body_buffer: &mut Option<BodyBuffer>,
575 end_of_stream: bool,
576) {
577 if is_stream_buffer && !*released {
578 *released = true;
579 if !end_of_stream {
580 *body = body_buffer.take().map(BodyBuffer::freeze);
581 }
582 }
583}
584
585struct BodyFilterOutput {
590 cluster: Option<Arc<str>>,
592 upstream: Option<Upstream>,
594 extensions: RequestExtensions,
596 filter_metadata: HashMap<String, String>,
598 filter_state: HashMap<usize, Box<dyn std::any::Any + Send + Sync>>,
600 executed_filter_indices: Vec<bool>,
602 body_done_indices: Vec<bool>,
604}
605
606impl BodyFilterOutput {
607 fn take_from(fctx: &mut HttpFilterContext<'_>) -> Self {
610 Self {
611 cluster: fctx.cluster.take(),
612 upstream: fctx.upstream.take(),
613 extensions: std::mem::take(&mut fctx.extensions),
614 filter_metadata: std::mem::take(&mut fctx.filter_metadata),
615 filter_state: std::mem::take(&mut fctx.filter_state),
616 executed_filter_indices: std::mem::take(&mut fctx.executed_filter_indices),
617 body_done_indices: std::mem::take(&mut fctx.body_done_indices),
618 }
619 }
620
621 fn write_back(self, ctx: &mut PingoraRequestCtx) {
623 ctx.cluster = self.cluster;
624 ctx.upstream = self.upstream;
625 ctx.extensions = self.extensions;
626 ctx.filter_metadata = self.filter_metadata;
627 ctx.filter_state = self.filter_state;
628 ctx.cached_executed_filter_indices = self.executed_filter_indices;
629 ctx.cached_body_done_indices = self.body_done_indices;
630 }
631}
632
633#[cfg(test)]
638#[expect(clippy::allow_attributes, reason = "blanket test suppressions")]
639#[allow(
640 clippy::unwrap_used,
641 clippy::expect_used,
642 clippy::indexing_slicing,
643 clippy::field_reassign_with_default,
644 clippy::too_many_lines,
645 clippy::cast_possible_truncation,
646 clippy::significant_drop_tightening,
647 reason = "tests"
648)]
649mod tests {
650 use praxis_core::connectivity::ConnectionOptions;
651
652 use super::*;
653
654 #[test]
655 fn first_failure_idempotent_sets_retry() {
656 let mut ctx = PingoraRequestCtx::default();
657 ctx.request_is_idempotent = true;
658 let e = handle_connect_failure(&mut ctx, make_error());
659 assert!(e.retry(), "first failure should set retry flag");
660 assert_eq!(ctx.retries, 1);
661 }
662
663 #[test]
664 fn large_body_skips_retry() {
665 let mut ctx = PingoraRequestCtx::default();
666 ctx.request_is_idempotent = true;
667 ctx.request_body_bytes = RETRY_BODY_LIMIT + 1;
668 let e = handle_connect_failure(&mut ctx, make_error());
669 assert!(!e.retry(), "should not retry when body exceeds retry buffer limit");
670 assert_eq!(ctx.retries, 0, "retry counter should not increment");
671 }
672
673 #[test]
674 fn mutated_body_exceeding_limit_skips_retry() {
675 let mut ctx = PingoraRequestCtx::default();
676 ctx.request_is_idempotent = true;
677 ctx.request_body_bytes = 1024;
678 ctx.mutated_request_body_len = Some((RETRY_BODY_LIMIT + 1) as usize);
679 let e = handle_connect_failure(&mut ctx, make_error());
680 assert!(
681 !e.retry(),
682 "should not retry when mutated body exceeds retry buffer limit"
683 );
684 assert_eq!(ctx.retries, 0);
685 }
686
687 #[test]
688 fn body_at_limit_allows_retry() {
689 let mut ctx = PingoraRequestCtx::default();
690 ctx.request_is_idempotent = true;
691 ctx.request_body_bytes = RETRY_BODY_LIMIT;
692 let e = handle_connect_failure(&mut ctx, make_error());
693 assert!(e.retry(), "body exactly at limit should allow retry");
694 assert_eq!(ctx.retries, 1);
695 }
696
697 #[test]
698 fn zero_body_allows_retry() {
699 let mut ctx = PingoraRequestCtx::default();
700 ctx.request_is_idempotent = true;
701 ctx.request_body_bytes = 0;
702 let e = handle_connect_failure(&mut ctx, make_error());
703 assert!(e.retry(), "zero-length body should allow retry");
704 assert_eq!(ctx.retries, 1);
705 }
706
707 #[test]
708 fn max_retries_exhausted_does_not_retry() {
709 let mut ctx = PingoraRequestCtx::default();
710 ctx.request_is_idempotent = true;
711 ctx.retries = MAX_RETRIES as u32;
712 let e = handle_connect_failure(&mut ctx, make_error());
713 assert!(!e.retry(), "should not retry after MAX_RETRIES");
714 assert_eq!(ctx.retries as usize, MAX_RETRIES);
715 }
716
717 #[test]
718 fn counter_increments_across_calls() {
719 let mut ctx = PingoraRequestCtx::default();
720 ctx.request_is_idempotent = true;
721 for expected in 1..=MAX_RETRIES {
722 let _result = handle_connect_failure(&mut ctx, make_error());
723 assert_eq!(ctx.retries as usize, expected);
724 }
725 let e = handle_connect_failure(&mut ctx, make_error());
726 assert!(!e.retry(), "should not retry after reaching MAX_RETRIES");
727 assert_eq!(ctx.retries as usize, MAX_RETRIES);
728 }
729
730 #[test]
731 fn non_idempotent_request_never_retries() {
732 let mut ctx = PingoraRequestCtx::default();
733 ctx.request_is_idempotent = false;
734 let e = handle_connect_failure(&mut ctx, make_error());
735 assert!(!e.retry(), "non-idempotent request should never retry");
736 assert_eq!(ctx.retries, 0);
737 }
738
739 #[test]
740 fn connect_failure_clears_upstream_connect_start() {
741 let mut ctx = PingoraRequestCtx::default();
742 ctx.upstream_connect_start = Some(std::time::Instant::now());
743 let _e = handle_connect_failure(&mut ctx, make_error());
744 assert!(
745 ctx.upstream_connect_start.is_none(),
746 "failed connect should consume upstream_connect_start for duration recording"
747 );
748 }
749
750 #[tokio::test]
751 async fn logging_cleanup_noop_when_response_phase_done() {
752 let registry = praxis_filter::FilterRegistry::with_builtins();
753 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
754 let mut ctx = PingoraRequestCtx::default();
755 ctx.response_phase_done = true;
756 ctx.request_snapshot = Some(praxis_filter::Request {
757 method: http::Method::GET,
758 uri: "/".parse().unwrap(),
759 headers: http::HeaderMap::new(),
760 });
761 logging_cleanup(&pipeline, &mut ctx).await;
762 }
763
764 #[tokio::test]
765 async fn logging_cleanup_noop_when_no_snapshot() {
766 let registry = praxis_filter::FilterRegistry::with_builtins();
767 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
768 let mut ctx = PingoraRequestCtx::default();
769 ctx.response_phase_done = false;
770 ctx.request_snapshot = None;
771 logging_cleanup(&pipeline, &mut ctx).await;
772 }
773
774 #[tokio::test]
775 async fn logging_cleanup_runs_response_pipeline_when_needed() {
776 let registry = praxis_filter::FilterRegistry::with_builtins();
777 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
778 let mut ctx = PingoraRequestCtx::default();
779 ctx.response_phase_done = false;
780 ctx.cluster = Some(Arc::from("test-cluster"));
781 ctx.request_snapshot = Some(praxis_filter::Request {
782 method: http::Method::GET,
783 uri: "/test".parse().unwrap(),
784 headers: http::HeaderMap::new(),
785 });
786 logging_cleanup(&pipeline, &mut ctx).await;
787 assert!(ctx.cluster.is_none(), "cluster should be taken by logging_cleanup");
788 assert!(ctx.upstream.is_none(), "upstream should be taken by logging_cleanup");
789 }
790
791 #[tokio::test]
792 async fn logging_cleanup_preserves_filter_metadata() {
793 let registry = praxis_filter::FilterRegistry::with_builtins();
794 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
795 let mut ctx = PingoraRequestCtx::default();
796 ctx.response_phase_done = false;
797 ctx.filter_metadata
798 .insert("json_rpc.method".to_owned(), "service/invoke".to_owned());
799 ctx.request_snapshot = Some(praxis_filter::Request {
800 method: http::Method::POST,
801 uri: "/api".parse().unwrap(),
802 headers: http::HeaderMap::new(),
803 });
804 logging_cleanup(&pipeline, &mut ctx).await;
805 assert_eq!(
806 ctx.filter_metadata.get("json_rpc.method").map(String::as_str),
807 Some("service/invoke"),
808 "filter_metadata should survive logging_cleanup"
809 );
810 }
811
812 #[tokio::test]
813 async fn logging_cleanup_preserves_extensions() {
814 let registry = praxis_filter::FilterRegistry::with_builtins();
815 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
816 let mut ctx = PingoraRequestCtx::default();
817 ctx.response_phase_done = false;
818 ctx.extensions.insert(42_u32);
819 ctx.request_snapshot = Some(praxis_filter::Request {
820 method: http::Method::POST,
821 uri: "/test".parse().unwrap(),
822 headers: http::HeaderMap::new(),
823 });
824 logging_cleanup(&pipeline, &mut ctx).await;
825 assert_eq!(
826 ctx.extensions.get::<u32>(),
827 Some(&42),
828 "extensions should survive logging_cleanup"
829 );
830 }
831
832 #[test]
833 fn passive_health_error_is_failure() {
834 let (pipeline, ctx) = make_passive_scenario(Some(3), Some(2));
835 let error = make_error();
836 record_passive_health(&pipeline, Some(&error), &ctx);
837
838 let registry = pipeline.health_registry().unwrap();
839 let entry = registry.get("test-cluster").unwrap();
840 assert!(
841 entry.endpoints()[0].is_healthy(),
842 "single failure should not yet mark unhealthy (threshold=3)"
843 );
844 }
845
846 #[test]
847 fn passive_health_status_500_is_failure() {
848 let (pipeline, mut ctx) = make_passive_scenario(Some(3), Some(2));
849 ctx.upstream_response_status = Some(500);
850 record_passive_health(&pipeline, None, &ctx);
851
852 let registry = pipeline.health_registry().unwrap();
853 let entry = registry.get("test-cluster").unwrap();
854 assert!(
855 entry.endpoints()[0].is_healthy(),
856 "single 500 should not yet mark unhealthy (threshold=3)"
857 );
858 }
859
860 #[test]
861 fn passive_health_status_below_500_is_success() {
862 let (pipeline, mut ctx) = make_passive_scenario(Some(2), Some(1));
863 ctx.upstream_response_status = Some(499);
864 record_passive_health(&pipeline, None, &ctx);
865
866 let registry = pipeline.health_registry().unwrap();
867 let entry = registry.get("test-cluster").unwrap();
868 assert!(entry.endpoints()[0].is_healthy(), "status 499 should count as success");
869 }
870
871 #[test]
872 fn passive_unhealthy_threshold_transition() {
873 let (pipeline, ctx) = make_passive_scenario(Some(2), Some(1));
874 let error = make_error();
875 record_passive_health(&pipeline, Some(&error), &ctx);
876 record_passive_health(&pipeline, Some(&error), &ctx);
877
878 let registry = pipeline.health_registry().unwrap();
879 let entry = registry.get("test-cluster").unwrap();
880 assert!(
881 !entry.endpoints()[0].is_healthy(),
882 "2 consecutive failures should mark unhealthy (threshold=2)"
883 );
884 }
885
886 #[test]
887 fn passive_healthy_threshold_recovery() {
888 let (pipeline, ctx) = make_passive_scenario(Some(1), Some(2));
889 let error = make_error();
890 record_passive_health(&pipeline, Some(&error), &ctx);
891
892 let registry = pipeline.health_registry().unwrap();
893 let entry = registry.get("test-cluster").unwrap();
894 assert!(
895 !entry.endpoints()[0].is_healthy(),
896 "should be unhealthy after 1 failure"
897 );
898
899 let ctx_ok = make_passive_ctx("test-cluster", 0, Some(200));
900 record_passive_health(&pipeline, None, &ctx_ok);
901 assert!(
902 !entry.endpoints()[0].is_healthy(),
903 "one success should not recover (threshold=2)"
904 );
905
906 record_passive_health(&pipeline, None, &ctx_ok);
907 assert!(
908 entry.endpoints()[0].is_healthy(),
909 "2 consecutive successes should recover (threshold=2)"
910 );
911 }
912
913 #[test]
914 fn passive_health_no_thresholds_is_noop() {
915 let (pipeline, ctx) = make_passive_scenario(None, None);
916 let error = make_error();
917 record_passive_health(&pipeline, Some(&error), &ctx);
918
919 let registry = pipeline.health_registry().unwrap();
920 let entry = registry.get("test-cluster").unwrap();
921 assert!(
922 entry.endpoints()[0].is_healthy(),
923 "no passive thresholds means failures are no-op"
924 );
925 }
926
927 #[test]
928 fn passive_health_endpoint_index_out_of_bounds() {
929 let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
930 ctx.selected_endpoint_index = Some(999);
931 let error = make_error();
932 record_passive_health(&pipeline, Some(&error), &ctx);
933
934 let registry = pipeline.health_registry().unwrap();
935 let entry = registry.get("test-cluster").unwrap();
936 assert!(entry.endpoints()[0].is_healthy(), "out-of-bounds index should be no-op");
937 }
938
939 #[test]
940 fn passive_health_missing_cluster_is_noop() {
941 let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
942 ctx.cluster = None;
943 ctx.metrics_cluster = None;
944 let error = make_error();
945 record_passive_health(&pipeline, Some(&error), &ctx);
946 }
947
948 #[test]
949 fn passive_health_falls_back_to_metrics_cluster() {
950 let (pipeline, mut ctx) = make_passive_scenario(Some(2), Some(1));
951 ctx.cluster = None;
952 ctx.metrics_cluster = Some(Arc::from("test-cluster"));
953 let error = make_error();
954 record_passive_health(&pipeline, Some(&error), &ctx);
955 record_passive_health(&pipeline, Some(&error), &ctx);
956
957 let registry = pipeline.health_registry().unwrap();
958 let entry = registry.get("test-cluster").unwrap();
959 assert!(
960 !entry.endpoints()[0].is_healthy(),
961 "fallback to metrics_cluster should still record passive health"
962 );
963 }
964
965 #[test]
966 fn passive_health_missing_endpoint_index_is_noop() {
967 let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
968 ctx.selected_endpoint_index = None;
969 let error = make_error();
970 record_passive_health(&pipeline, Some(&error), &ctx);
971 }
972
973 #[test]
974 fn passive_health_missing_registry_is_noop() {
975 let registry = praxis_filter::FilterRegistry::with_builtins();
976 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
977 let mut ctx = PingoraRequestCtx::default();
978 ctx.cluster = Some(Arc::from("test-cluster"));
979 ctx.selected_endpoint_index = Some(0);
980 let error = make_error();
981 record_passive_health(&pipeline, Some(&error), &ctx);
982 }
983
984 #[test]
985 fn passive_health_unknown_cluster_is_noop() {
986 let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
987 ctx.cluster = Some(Arc::from("nonexistent"));
988 let error = make_error();
989 record_passive_health(&pipeline, Some(&error), &ctx);
990 }
991
992 #[test]
993 fn size_limit_none_body_returns_false() {
994 let mut bytes = 0_u64;
995 assert!(!check_body_size_limit(None, &mut bytes, 100));
996 assert_eq!(bytes, 0, "accumulated bytes unchanged for None body");
997 }
998
999 #[test]
1000 fn size_limit_within_limit() {
1001 let mut bytes = 0_u64;
1002 let body = Some(Bytes::from_static(b"hello"));
1003 assert!(!check_body_size_limit(body.as_ref(), &mut bytes, 10));
1004 assert_eq!(bytes, 5);
1005 }
1006
1007 #[test]
1008 fn size_limit_at_exact_limit() {
1009 let mut bytes = 0_u64;
1010 let body = Some(Bytes::from_static(b"exact"));
1011 assert!(!check_body_size_limit(body.as_ref(), &mut bytes, 5));
1012 assert_eq!(bytes, 5);
1013 }
1014
1015 #[test]
1016 fn size_limit_exceeds_limit() {
1017 let mut bytes = 0_u64;
1018 let body = Some(Bytes::from_static(b"toolong"));
1019 assert!(check_body_size_limit(body.as_ref(), &mut bytes, 3));
1020 }
1021
1022 #[test]
1023 fn size_limit_cumulative_overflow() {
1024 let mut bytes = 0_u64;
1025 let first = Some(Bytes::from_static(b"aaa"));
1026 assert!(!check_body_size_limit(first.as_ref(), &mut bytes, 5));
1027
1028 let second = Some(Bytes::from_static(b"bbb"));
1029 assert!(check_body_size_limit(second.as_ref(), &mut bytes, 5));
1030 assert_eq!(bytes, 6);
1031 }
1032
1033 #[test]
1034 fn stream_buffer_accumulates_chunks() {
1035 let mut body = Some(Bytes::from_static(b"hello "));
1036 let mut buf: Option<BodyBuffer> = None;
1037 assert!(!accumulate_stream_buffer(&mut body, &mut buf, false, Some(100)));
1038 assert!(buf.is_some());
1039
1040 body = Some(Bytes::from_static(b"world"));
1041 assert!(!accumulate_stream_buffer(&mut body, &mut buf, false, Some(100)));
1042
1043 let frozen = buf.take().unwrap().freeze();
1044 assert_eq!(frozen, Bytes::from_static(b"hello world"));
1045 }
1046
1047 #[test]
1048 fn stream_buffer_freezes_at_eos() {
1049 let mut body = Some(Bytes::from_static(b"data"));
1050 let mut buf: Option<BodyBuffer> = None;
1051 assert!(!accumulate_stream_buffer(&mut body, &mut buf, false, Some(100)));
1052
1053 body = Some(Bytes::from_static(b" end"));
1054 assert!(!accumulate_stream_buffer(&mut body, &mut buf, true, Some(100)));
1055 assert!(buf.is_none(), "buffer should be taken at EOS");
1056 assert_eq!(body.unwrap(), Bytes::from_static(b"data end"));
1057 }
1058
1059 #[test]
1060 fn stream_buffer_overflow() {
1061 let mut body = Some(Bytes::from_static(b"too long"));
1062 let mut buf: Option<BodyBuffer> = None;
1063 assert!(accumulate_stream_buffer(&mut body, &mut buf, false, Some(5)));
1064 }
1065
1066 #[test]
1067 fn stream_buffer_none_body() {
1068 let mut body: Option<Bytes> = None;
1069 let mut buf: Option<BodyBuffer> = None;
1070 assert!(!accumulate_stream_buffer(&mut body, &mut buf, false, Some(100)));
1071 assert!(buf.is_none());
1072 }
1073
1074 #[test]
1075 fn stream_buffer_uses_absolute_max_when_none() {
1076 let mut body = Some(Bytes::from_static(b"data"));
1077 let mut buf: Option<BodyBuffer> = None;
1078 assert!(!accumulate_stream_buffer(&mut body, &mut buf, false, None));
1079 assert!(buf.is_some(), "should create buffer with absolute max");
1080 }
1081
1082 #[test]
1083 fn suppress_clears_body_when_buffering() {
1084 let mut body = Some(Bytes::from_static(b"data"));
1085 suppress_stream_buffer_chunk(&mut body, true, false, false);
1086 assert!(body.is_none());
1087 }
1088
1089 #[test]
1090 fn suppress_noop_when_not_stream_buffer() {
1091 let mut body = Some(Bytes::from_static(b"data"));
1092 suppress_stream_buffer_chunk(&mut body, false, false, false);
1093 assert!(body.is_some());
1094 }
1095
1096 #[test]
1097 fn suppress_noop_when_released() {
1098 let mut body = Some(Bytes::from_static(b"data"));
1099 suppress_stream_buffer_chunk(&mut body, true, true, false);
1100 assert!(body.is_some());
1101 }
1102
1103 #[test]
1104 fn suppress_noop_at_eos() {
1105 let mut body = Some(Bytes::from_static(b"data"));
1106 suppress_stream_buffer_chunk(&mut body, true, false, true);
1107 assert!(body.is_some());
1108 }
1109
1110 #[test]
1111 fn release_sets_flag_and_flushes_buffer() {
1112 let mut body: Option<Bytes> = None;
1113 let mut released = false;
1114 let mut buf = Some(BodyBuffer::new(100));
1115 buf.as_mut().unwrap().push(Bytes::from_static(b"buffered")).unwrap();
1116
1117 release_stream_buffer(&mut body, true, &mut released, &mut buf, false);
1118 assert!(released);
1119 assert_eq!(body.unwrap(), Bytes::from_static(b"buffered"));
1120 assert!(buf.is_none());
1121 }
1122
1123 #[test]
1124 fn release_noop_when_already_released() {
1125 let mut body: Option<Bytes> = None;
1126 let mut released = true;
1127 let mut buf: Option<BodyBuffer> = None;
1128
1129 release_stream_buffer(&mut body, true, &mut released, &mut buf, false);
1130 assert!(body.is_none(), "body should be unchanged when already released");
1131 }
1132
1133 #[test]
1134 fn release_noop_when_not_stream_buffer() {
1135 let mut body: Option<Bytes> = None;
1136 let mut released = false;
1137 let mut buf: Option<BodyBuffer> = None;
1138
1139 release_stream_buffer(&mut body, false, &mut released, &mut buf, false);
1140 assert!(!released, "released flag should be unchanged for non-stream-buffer");
1141 }
1142
1143 #[test]
1144 fn release_at_eos_sets_flag_but_no_flush() {
1145 let mut body: Option<Bytes> = None;
1146 let mut released = false;
1147 let mut buf = Some(BodyBuffer::new(100));
1148 buf.as_mut().unwrap().push(Bytes::from_static(b"data")).unwrap();
1149
1150 release_stream_buffer(&mut body, true, &mut released, &mut buf, true);
1151 assert!(released);
1152 assert!(body.is_none(), "body should not be overwritten at EOS");
1153 assert!(buf.is_some(), "buffer should not be taken at EOS");
1154 }
1155
1156 #[test]
1157 fn write_back_transfers_fields() {
1158 let mut ctx = PingoraRequestCtx::default();
1159
1160 let mut extensions = RequestExtensions::new();
1161 extensions.insert(42_u32);
1162
1163 let state_val: Box<dyn std::any::Any + Send + Sync> = Box::new(99_i32);
1164 let filter_state = HashMap::from([(0_usize, state_val)]);
1165
1166 let output = BodyFilterOutput {
1167 cluster: Some(Arc::from("test-cluster")),
1168 upstream: Some(Upstream {
1169 address: Arc::from("10.0.0.1:80"),
1170 connection: Arc::new(ConnectionOptions::default()),
1171 tls: None,
1172 }),
1173 extensions,
1174 filter_metadata: HashMap::from([("key".to_owned(), "val".to_owned())]),
1175 filter_state,
1176 executed_filter_indices: vec![true, false],
1177 body_done_indices: vec![false, true],
1178 };
1179 output.write_back(&mut ctx);
1180
1181 assert_eq!(ctx.cluster.as_deref(), Some("test-cluster"));
1182 assert!(ctx.upstream.is_some(), "upstream should transfer");
1183 assert_eq!(ctx.upstream.as_ref().unwrap().address.as_ref(), "10.0.0.1:80");
1184 assert_eq!(ctx.extensions.get::<u32>(), Some(&42));
1185 assert_eq!(ctx.filter_metadata.get("key").map(String::as_str), Some("val"));
1186 assert_eq!(ctx.filter_state.len(), 1, "filter_state should transfer");
1187 assert_eq!(
1188 ctx.filter_state.get(&0).and_then(|v| v.downcast_ref::<i32>()),
1189 Some(&99)
1190 );
1191 assert_eq!(ctx.cached_executed_filter_indices, vec![true, false]);
1192 assert_eq!(ctx.cached_body_done_indices, vec![false, true]);
1193 }
1194
1195 #[test]
1200 fn http_version_label_http_09() {
1201 assert_eq!(
1202 http_version_label(http::Version::HTTP_09),
1203 "0.9",
1204 "HTTP/0.9 should map to '0.9'"
1205 );
1206 }
1207
1208 #[test]
1209 fn http_version_label_http_10() {
1210 assert_eq!(
1211 http_version_label(http::Version::HTTP_10),
1212 "1.0",
1213 "HTTP/1.0 should map to '1.0'"
1214 );
1215 }
1216
1217 #[test]
1218 fn http_version_label_http_11() {
1219 assert_eq!(
1220 http_version_label(http::Version::HTTP_11),
1221 "1.1",
1222 "HTTP/1.1 should map to '1.1'"
1223 );
1224 }
1225
1226 #[test]
1227 fn http_version_label_http_2() {
1228 assert_eq!(
1229 http_version_label(http::Version::HTTP_2),
1230 "2",
1231 "HTTP/2 should map to '2'"
1232 );
1233 }
1234
1235 #[test]
1236 fn http_version_label_http_3() {
1237 assert_eq!(
1238 http_version_label(http::Version::HTTP_3),
1239 "3",
1240 "HTTP/3 should map to '3'"
1241 );
1242 }
1243
1244 #[test]
1245 fn record_response_span_attributes_noop_for_disabled_span() {
1246 let ctx = PingoraRequestCtx::default();
1247 assert!(ctx.request_span.is_disabled(), "default span should be disabled");
1248 }
1249
1250 #[test]
1251 fn record_response_span_attributes_records_upstream_cluster() {
1252 let mut ctx = PingoraRequestCtx::default();
1253 ctx.metrics_cluster = Some(Arc::from("api-cluster"));
1254 ctx.upstream_for_retry = Some(Upstream {
1255 address: Arc::from("10.0.0.1:80"),
1256 connection: Arc::new(ConnectionOptions::default()),
1257 tls: None,
1258 });
1259 ctx.request_span = tracing::info_span!(
1260 "test_span",
1261 "http.response.status_code" = tracing::field::Empty,
1262 "upstream.address" = tracing::field::Empty,
1263 "upstream.cluster" = tracing::field::Empty,
1264 );
1265 ctx.request_span.record("upstream.cluster", "api-cluster");
1266 ctx.request_span.record("upstream.address", "10.0.0.1:80");
1267 }
1268
1269 fn make_error() -> Box<pingora_core::Error> {
1275 pingora_core::Error::explain(pingora_core::ErrorType::ConnectError, "test connect failure")
1276 }
1277
1278 fn make_passive_ctx(cluster: &str, endpoint_idx: usize, status: Option<u16>) -> PingoraRequestCtx {
1280 let mut ctx = PingoraRequestCtx::default();
1281 ctx.cluster = Some(Arc::from(cluster));
1282 ctx.selected_endpoint_index = Some(endpoint_idx);
1283 ctx.upstream_response_status = status;
1284 ctx
1285 }
1286
1287 fn make_passive_scenario(
1290 passive_unhealthy: Option<u32>,
1291 passive_healthy: Option<u32>,
1292 ) -> (FilterPipeline, PingoraRequestCtx) {
1293 use praxis_core::health::{ClusterHealthEntry, EndpointHealth};
1294
1295 let entry = ClusterHealthEntry::new(
1296 vec![EndpointHealth::new()],
1297 vec![Arc::from("10.0.0.1:80")],
1298 passive_unhealthy,
1299 passive_healthy,
1300 );
1301 let mut map = HashMap::new();
1302 map.insert(Arc::from("test-cluster"), Arc::new(entry));
1303 let health_registry = Arc::new(map);
1304
1305 let registry = praxis_filter::FilterRegistry::with_builtins();
1306 let mut pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
1307 pipeline.set_health_registry(health_registry);
1308
1309 let ctx = make_passive_ctx("test-cluster", 0, None);
1310
1311 (pipeline, ctx)
1312 }
1313}