1use std::{collections::VecDeque, net::IpAddr, sync::Arc, time::Instant};
7
8use bytes::Bytes;
9use praxis_core::connectivity::Upstream;
10use praxis_filter::{BodyBuffer, BodyMode, FilterPipeline, Request, Response, TrustedHeaderMutation};
11use tokio::sync::OwnedSemaphorePermit;
12use tracing::Span;
13
14#[expect(clippy::struct_excessive_bools, reason = "lifecycle flags")]
30pub struct PingoraRequestCtx {
31 pub _connection_permit: Option<OwnedSemaphorePermit>,
36
37 pub _global_connection_permit: Option<OwnedSemaphorePermit>,
41
42 pub client_addr: Option<IpAddr>,
44
45 pub client_http_version: Option<http::Version>,
50
51 pub cluster: Option<Arc<str>>,
53
54 pub cached_body_done_indices: Vec<bool>,
59
60 pub cached_executed_branch_filters: Vec<bool>,
65
66 pub cached_executed_filter_indices: Vec<bool>,
71
72 pub downstream_tls: bool,
78
79 pub peer_identity: Option<Arc<praxis_tls::TlsPeerIdentity>>,
86
87 pub connection_upgraded: bool,
93
94 pub extensions: praxis_filter::RequestExtensions,
101
102 pub filter_metadata: std::collections::HashMap<String, String>,
108
109 pub pre_read_mutations: Vec<TrustedHeaderMutation>,
115
116 pub structured_metadata: std::collections::HashMap<String, serde_json::Value>,
122
123 pub mutated_request_body_len: Option<usize>,
129
130 pub pinned_pipeline: Option<Arc<FilterPipeline>>,
140
141 pub filter_results: std::collections::HashMap<&'static str, praxis_filter::FilterResultSet>,
147
148 pub filter_state: std::collections::HashMap<usize, Box<dyn std::any::Any + Send + Sync>>,
158
159 pub metrics_cluster: Option<Arc<str>>,
163
164 pub metrics_cluster_shared: Option<::metrics::SharedString>,
169
170 pub metrics_route: Option<::metrics::SharedString>,
172
173 pub(crate) error_type: Option<&'static str>,
180
181 pub(crate) _active_request: Option<crate::http::pingora::metrics::ActiveRequestGuard>,
183
184 pub upstream_connect_start: Option<Instant>,
186
187 pub pre_read_body: Option<VecDeque<Bytes>>,
193
194 pub retained_pre_read_body: Option<VecDeque<Bytes>>,
203
204 pub adapted_request_body: Option<VecDeque<Bytes>>,
214
215 pub retained_adapted_request_body: Option<VecDeque<Bytes>>,
223
224 pub adapted_request_body_len: Option<usize>,
228
229 pub request_body_buffer: Option<BodyBuffer>,
233
234 pub request_body_bytes: u64,
236
237 pub request_body_mode: BodyMode,
241
242 pub request_body_released: bool,
245
246 pub request_is_idempotent: bool,
248
249 pub request_snapshot: Option<Request>,
251
252 pub request_span: Span,
260
261 pub upstream_exchange_span: Span,
268
269 pub request_start: Instant,
271
272 pub response_body_buffer: Option<BodyBuffer>,
276
277 pub response_body_bytes: u64,
279
280 pub response_body_mode: BodyMode,
284
285 pub response_body_released: bool,
287
288 pub response_header_snapshot: Option<Response>,
294
295 pub upstream_response_status: Option<u16>,
298
299 pub grpc_completion: Option<praxis_core::grpc::GrpcCompletion>,
306
307 pub response_phase_done: bool,
311
312 pub pending_rejection: Option<praxis_filter::Rejection>,
316
317 pub response_delivery_complete: bool,
323
324 pub retries: u32,
326
327 pub selected_endpoint_index: Option<usize>,
331
332 pub attempted_endpoints: Vec<Arc<str>>,
334
335 pub retry_policy: Option<Arc<praxis_core::config::RetryPolicy>>,
337
338 pub route_retry_policy: Option<Arc<praxis_core::config::RetryPolicy>>,
340
341 pub cluster_retry_state: Option<Arc<praxis_core::retry::ClusterRetryState>>,
343
344 pub cluster_retry_state_released: bool,
346
347 pub endpoint_reselector: Option<Arc<praxis_filter::EndpointReselector>>,
349
350 pub pending_backoff: Option<std::time::Duration>,
352
353 pub reselect_on_retry: bool,
355
356 pub rewritten_path: Option<String>,
363
364 pub upstream: Option<Upstream>,
366
367 pub upstream_for_retry: Option<Upstream>,
369
370 pub upstream_contacted: bool,
377
378 pub ended_by_client: bool,
383}
384
385macro_rules! filter_context {
395 ($ctx:expr, $pipeline:expr, $request:expr, $response_header:expr) => {{
396 praxis_filter::HttpFilterContext {
397 buffered_request_body: $ctx
398 .pre_read_body
399 .as_ref()
400 .map(|chunks| chunks.front().cloned().unwrap_or_default()),
401 body_done_indices: std::mem::take(&mut $ctx.cached_body_done_indices),
402 branch_iterations: std::collections::HashMap::new(),
403 client_addr: $ctx.client_addr,
404 cluster: $ctx.cluster.take(),
405 current_filter_id: None,
406 downstream_tls: $ctx.downstream_tls,
407 metrics_route: $ctx.metrics_route.clone(),
408 peer_identity: $ctx.peer_identity.clone(),
409 extensions: std::mem::take(&mut $ctx.extensions),
410 executed_branch_filters: std::mem::take(&mut $ctx.cached_executed_branch_filters),
411 executed_filter_indices: std::mem::take(&mut $ctx.cached_executed_filter_indices),
412 extra_request_headers: Vec::new(),
413 request_headers_to_remove: Vec::new(),
414 request_headers_to_set: Vec::new(),
415 filter_metadata: std::mem::take(&mut $ctx.filter_metadata),
416 prior_pre_read_mutations: Vec::new(),
418 grpc_completion: $ctx.grpc_completion.clone(),
419 pre_read_mutations: std::mem::take(&mut $ctx.pre_read_mutations),
420 structured_metadata: std::mem::take(&mut $ctx.structured_metadata),
421 filter_results: std::mem::take(&mut $ctx.filter_results),
422 filter_state: std::mem::take(&mut $ctx.filter_state),
423 health_registry: $pipeline.health_registry(),
424 id_generator: $pipeline.id_generator(),
425 kv_stores: $pipeline.kv_stores(),
426 session_stores: $pipeline.session_stores(),
427 subrequest_client: $pipeline.subrequest_client(),
428 subrequest_response_mode: praxis_filter::SubRequestResponseMode::Buffered,
429 request: $request,
430 request_body_bytes: $ctx.request_body_bytes,
431 request_body_mode: $ctx.request_body_mode,
432 request_start: $ctx.request_start,
433 response_body_bytes: $ctx.response_body_bytes,
434 response_body_mode: $ctx.response_body_mode,
435 response_header: $response_header,
436 response_headers_modified: false,
437 upstream_reached: $ctx.upstream_contacted && !$ctx.ended_by_client,
438 rewritten_path: $ctx.rewritten_path.take(),
439 selected_endpoint_index: $ctx.selected_endpoint_index,
440 attempted_endpoints: std::mem::take(&mut $ctx.attempted_endpoints),
441 retry_policy: $ctx.retry_policy.clone(),
442 route_retry_policy: $ctx.route_retry_policy.clone(),
443 cluster_retry_state: $ctx.cluster_retry_state.clone(),
444 cluster_retry_state_released: $ctx.cluster_retry_state_released,
445 endpoint_reselector: $ctx.endpoint_reselector.clone(),
446 pinned_endpoint_address: None,
447 time_source: $pipeline.time_source(),
448 upstream: $ctx.upstream.take().or_else(|| $ctx.upstream_for_retry.clone()),
449 }
450 }};
451}
452
453impl PingoraRequestCtx {
454 pub fn build_filter_context<'a>(
478 &mut self,
479 pipeline: &'a FilterPipeline,
480 request: &'a Request,
481 response_header: Option<&'a mut Response>,
482 ) -> praxis_filter::HttpFilterContext<'a> {
483 filter_context!(self, pipeline, request, response_header)
484 }
485
486 pub fn filter_context_for<'a>(
509 &'a mut self,
510 pipeline: &'a FilterPipeline,
511 response_header: Option<&'a mut Response>,
512 ) -> Option<praxis_filter::HttpFilterContext<'a>> {
513 let request = self.request_snapshot.as_ref()?;
514 Some(filter_context!(self, pipeline, request, response_header))
515 }
516
517 pub fn response_body_context_for<'a>(
525 &'a mut self,
526 pipeline: &'a FilterPipeline,
527 ) -> Option<(praxis_filter::HttpFilterContext<'a>, Option<&'a Response>)> {
528 let request = self.request_snapshot.as_ref()?;
529 let response_header = self.response_header_snapshot.as_ref();
530 Some((filter_context!(self, pipeline, request, None), response_header))
531 }
532
533 pub fn pin_pipeline(&mut self, swap: &arc_swap::ArcSwap<FilterPipeline>) -> Arc<FilterPipeline> {
546 if let Some(existing) = &self.pinned_pipeline {
547 return Arc::clone(existing);
548 }
549 let pipeline = swap.load_full();
550 pipeline.prepare_extensions(&mut self.extensions);
553 self.pinned_pipeline = Some(Arc::clone(&pipeline));
554 pipeline
555 }
556
557 pub(crate) fn stamp_error_type(&mut self, error_type: &'static str) {
563 self.error_type.get_or_insert(error_type);
564 }
565
566 pub fn pipeline(&self, swap: &arc_swap::ArcSwap<FilterPipeline>) -> Arc<FilterPipeline> {
575 self.pinned_pipeline
576 .as_ref()
577 .map_or_else(|| swap.load_full(), Arc::clone)
578 }
579}
580
581impl Default for PingoraRequestCtx {
582 #[expect(
583 clippy::too_many_lines,
584 reason = "context default enumerates all lifecycle fields explicitly"
585 )]
586 fn default() -> Self {
587 Self {
588 _connection_permit: None,
589 _global_connection_permit: None,
590 cached_body_done_indices: Vec::new(),
591 cached_executed_branch_filters: Vec::new(),
592 cached_executed_filter_indices: Vec::new(),
593 client_addr: None,
594 client_http_version: None,
595 cluster: None,
596 connection_upgraded: false,
597 downstream_tls: false,
598 peer_identity: None,
599 extensions: praxis_filter::RequestExtensions::new(),
600 filter_metadata: std::collections::HashMap::new(),
601 pre_read_mutations: Vec::new(),
602 structured_metadata: std::collections::HashMap::new(),
603 mutated_request_body_len: None,
604 pinned_pipeline: None,
605 filter_results: std::collections::HashMap::new(),
606 filter_state: std::collections::HashMap::new(),
607 metrics_cluster: None,
608 metrics_cluster_shared: None,
609 metrics_route: None,
610 error_type: None,
611 _active_request: None,
612 upstream_connect_start: None,
613 pre_read_body: None,
614 retained_pre_read_body: None,
615 adapted_request_body: None,
616 retained_adapted_request_body: None,
617 adapted_request_body_len: None,
618 request_body_buffer: None,
619 request_body_bytes: 0,
620 request_body_mode: BodyMode::Stream,
621 request_body_released: false,
622 request_is_idempotent: false,
623 request_snapshot: None,
624 request_span: Span::none(),
625 request_start: Instant::now(),
626 upstream_exchange_span: Span::none(),
627 response_body_buffer: None,
628 response_body_bytes: 0,
629 response_body_mode: BodyMode::Stream,
630 response_body_released: false,
631 response_header_snapshot: None,
632 upstream_response_status: None,
633 grpc_completion: None,
634 response_phase_done: false,
635 response_delivery_complete: false,
636 pending_rejection: None,
637 retries: 0,
638 rewritten_path: None,
639 selected_endpoint_index: None,
640 attempted_endpoints: Vec::new(),
641 retry_policy: None,
642 route_retry_policy: None,
643 cluster_retry_state: None,
644 cluster_retry_state_released: false,
645 endpoint_reselector: None,
646 pending_backoff: None,
647 reselect_on_retry: false,
648 upstream: None,
649 upstream_for_retry: None,
650 upstream_contacted: false,
651 ended_by_client: false,
652 }
653 }
654}
655
656#[cfg(test)]
661#[expect(clippy::allow_attributes, reason = "blanket test suppressions")]
662#[allow(
663 clippy::unwrap_used,
664 clippy::expect_used,
665 clippy::indexing_slicing,
666 clippy::significant_drop_tightening,
667 clippy::too_many_lines,
668 reason = "tests"
669)]
670mod tests {
671 use std::net::Ipv4Addr;
672
673 use http::{HeaderMap, Method, Uri};
674 use praxis_filter::FilterRegistry;
675
676 use super::*;
677
678 #[test]
679 fn default_state_has_no_client_addr() {
680 let ctx = default_ctx();
681 assert!(ctx.client_addr.is_none(), "default client_addr should be None");
682 }
683
684 #[test]
685 fn default_state_has_no_cluster() {
686 let ctx = default_ctx();
687 assert!(ctx.cluster.is_none(), "default cluster should be None");
688 }
689
690 #[test]
691 fn default_state_has_zero_retries() {
692 let ctx = default_ctx();
693 assert_eq!(ctx.retries, 0, "default retries should be zero");
694 }
695
696 #[test]
697 fn default_state_flags_are_false() {
698 let ctx = default_ctx();
699 assert!(
700 !ctx.request_body_released,
701 "default request_body_released should be false"
702 );
703 assert!(
704 !ctx.response_body_released,
705 "default response_body_released should be false"
706 );
707 assert!(
708 !ctx.request_is_idempotent,
709 "default request_is_idempotent should be false"
710 );
711 assert!(!ctx.response_phase_done, "default response_phase_done should be false");
712 }
713
714 #[test]
715 fn default_state_buffers_are_none() {
716 let ctx = default_ctx();
717 assert!(
718 ctx.request_body_buffer.is_none(),
719 "default request_body_buffer should be None"
720 );
721 assert!(
722 ctx.response_body_buffer.is_none(),
723 "default response_body_buffer should be None"
724 );
725 assert!(ctx.pre_read_body.is_none(), "default pre_read_body should be None");
726 }
727
728 #[test]
729 fn default_state_request_span_is_disabled() {
730 let ctx = default_ctx();
731 assert!(
732 ctx.request_span.is_disabled(),
733 "default request_span should be a disabled (none) span"
734 );
735 }
736
737 #[test]
738 fn default_state_upstream_exchange_span_is_disabled() {
739 let ctx = default_ctx();
740 assert!(
741 ctx.upstream_exchange_span.is_disabled(),
742 "default upstream_exchange_span should be a disabled (none) span"
743 );
744 }
745
746 #[test]
747 fn default_state_snapshots_are_none() {
748 let ctx = default_ctx();
749 assert!(
750 ctx.request_snapshot.is_none(),
751 "default request_snapshot should be None"
752 );
753 assert!(ctx.upstream.is_none(), "default upstream should be None");
754 assert!(
755 ctx.upstream_for_retry.is_none(),
756 "default upstream_for_retry should be None"
757 );
758 }
759
760 #[test]
761 fn set_client_addr() {
762 let mut ctx = default_ctx();
763 let addr = IpAddr::V4(Ipv4Addr::new(192, 168, 1, 100));
764 ctx.client_addr = Some(addr);
765 assert_eq!(
766 ctx.client_addr.unwrap(),
767 addr,
768 "client_addr should match assigned value"
769 );
770 }
771
772 #[test]
773 fn set_cluster() {
774 let mut ctx = default_ctx();
775 ctx.cluster = Some(Arc::from("api-cluster"));
776 assert_eq!(
777 ctx.cluster.as_deref(),
778 Some("api-cluster"),
779 "cluster should match assigned value"
780 );
781 }
782
783 #[test]
784 fn set_upstream() {
785 let mut ctx = default_ctx();
786 let upstream = Upstream {
787 address: Arc::from("10.0.0.1:80"),
788 authority: None,
789 tls: None,
790 connection: Arc::new(praxis_core::connectivity::ConnectionOptions::default()),
791 };
792 ctx.upstream = Some(upstream.clone());
793 assert_eq!(
794 &*ctx.upstream.as_ref().unwrap().address,
795 "10.0.0.1:80",
796 "upstream address should match assigned value"
797 );
798 }
799
800 #[test]
801 fn increment_retries() {
802 let mut ctx = default_ctx();
803 ctx.retries += 1;
804 ctx.retries += 1;
805 assert_eq!(ctx.retries, 2, "retries should be 2 after two increments");
806 }
807
808 #[test]
809 fn release_request_body_flag() {
810 let mut ctx = default_ctx();
811 assert!(!ctx.request_body_released, "request_body_released should start false");
812 ctx.request_body_released = true;
813 assert!(
814 ctx.request_body_released,
815 "request_body_released should be true after setting"
816 );
817 }
818
819 #[test]
820 fn release_response_body_flag() {
821 let mut ctx = default_ctx();
822 assert!(!ctx.response_body_released, "response_body_released should start false");
823 ctx.response_body_released = true;
824 assert!(
825 ctx.response_body_released,
826 "response_body_released should be true after setting"
827 );
828 }
829
830 #[test]
831 fn response_phase_done_flag() {
832 let mut ctx = default_ctx();
833 assert!(!ctx.response_phase_done, "response_phase_done should start false");
834 ctx.response_phase_done = true;
835 assert!(
836 ctx.response_phase_done,
837 "response_phase_done should be true after setting"
838 );
839 }
840
841 #[test]
842 fn set_pre_read_body() {
843 let mut ctx = default_ctx();
844 let chunks = VecDeque::from([Bytes::from_static(b"chunk1"), Bytes::from_static(b"chunk2")]);
845 ctx.pre_read_body = Some(chunks);
846 let body = ctx.pre_read_body.as_ref().unwrap();
847 assert_eq!(body.len(), 2, "pre_read_body should contain 2 chunks");
848 assert_eq!(body[0], Bytes::from_static(b"chunk1"), "first chunk should be 'chunk1'");
849 assert_eq!(
850 body[1],
851 Bytes::from_static(b"chunk2"),
852 "second chunk should be 'chunk2'"
853 );
854 }
855
856 #[test]
857 fn set_request_snapshot() {
858 let mut ctx = default_ctx();
859 let snapshot = Request {
860 method: Method::POST,
861 uri: "/api/data".parse::<Uri>().unwrap(),
862 headers: HeaderMap::new(),
863 };
864 ctx.request_snapshot = Some(snapshot);
865 let snap = ctx.request_snapshot.as_ref().unwrap();
866 assert_eq!(snap.method, Method::POST, "snapshot method should be POST");
867 assert_eq!(snap.uri.path(), "/api/data", "snapshot URI path should be /api/data");
868 }
869
870 #[test]
871 fn request_body_buffer_lifecycle() {
872 let mut ctx = default_ctx();
873 let mut buf = BodyBuffer::new(100);
874 buf.push(Bytes::from_static(b"data")).unwrap();
875 ctx.request_body_buffer = Some(buf);
876
877 assert!(
878 ctx.request_body_buffer.is_some(),
879 "buffer should be present after assignment"
880 );
881 let taken = ctx.request_body_buffer.take().unwrap();
882 assert_eq!(
883 taken.freeze(),
884 Bytes::from_static(b"data"),
885 "frozen buffer should contain pushed data"
886 );
887 assert!(ctx.request_body_buffer.is_none(), "buffer should be None after take");
888 }
889
890 #[test]
891 fn default_request_body_mode_is_stream() {
892 let ctx = default_ctx();
893 assert_eq!(
894 ctx.request_body_mode,
895 BodyMode::Stream,
896 "default request_body_mode should be Stream"
897 );
898 }
899
900 #[test]
901 fn default_response_body_mode_is_stream() {
902 let ctx = default_ctx();
903 assert_eq!(
904 ctx.response_body_mode,
905 BodyMode::Stream,
906 "default response_body_mode should be Stream"
907 );
908 }
909
910 #[test]
911 fn set_request_body_mode() {
912 let mut ctx = default_ctx();
913 ctx.request_body_mode = BodyMode::StreamBuffer { max_bytes: Some(4096) };
914 assert_eq!(
915 ctx.request_body_mode,
916 BodyMode::StreamBuffer { max_bytes: Some(4096) },
917 "request_body_mode should match assigned value"
918 );
919 }
920
921 #[test]
922 fn set_response_body_mode() {
923 let mut ctx = default_ctx();
924 ctx.response_body_mode = BodyMode::StreamBuffer { max_bytes: Some(8192) };
925 assert_eq!(
926 ctx.response_body_mode,
927 BodyMode::StreamBuffer { max_bytes: Some(8192) },
928 "response_body_mode should match assigned value"
929 );
930 }
931
932 #[test]
937 fn metadata_roundtrip_through_filter_context() {
938 let registry = FilterRegistry::with_builtins();
939 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
940 let request = Request {
941 method: Method::GET,
942 uri: "/".parse::<Uri>().unwrap(),
943 headers: HeaderMap::new(),
944 };
945
946 let mut ctx = default_ctx();
947 ctx.filter_metadata.insert("rpc.method".to_owned(), "echo".to_owned());
948 ctx.filter_metadata.insert("rpc.status".to_owned(), "ok".to_owned());
949
950 let fctx = ctx.build_filter_context(&pipeline, &request, None);
951 assert_eq!(
952 fctx.get_metadata("rpc.method"),
953 Some("echo"),
954 "metadata written before build should survive into filter context"
955 );
956 assert_eq!(
957 fctx.get_metadata("rpc.status"),
958 Some("ok"),
959 "multiple metadata keys should round-trip"
960 );
961 }
962
963 #[test]
964 fn metadata_written_in_filter_context_persists_back() {
965 let registry = FilterRegistry::with_builtins();
966 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
967 let request = Request {
968 method: Method::GET,
969 uri: "/".parse::<Uri>().unwrap(),
970 headers: HeaderMap::new(),
971 };
972
973 let mut ctx = default_ctx();
974 let mut fctx = ctx.build_filter_context(&pipeline, &request, None);
975 fctx.set_metadata("trace.id", "abc-123");
976 fctx.set_metadata("trace.span", "42");
977
978 ctx.filter_metadata = fctx.filter_metadata;
979 assert_eq!(
980 ctx.filter_metadata.get("trace.id").map(String::as_str),
981 Some("abc-123"),
982 "metadata set in filter context should persist back to protocol context"
983 );
984 assert_eq!(
985 ctx.filter_metadata.get("trace.span").map(String::as_str),
986 Some("42"),
987 "multiple metadata keys should persist back"
988 );
989 }
990
991 #[test]
996 fn pin_pipeline_captures_current_arc() {
997 let registry = FilterRegistry::with_builtins();
998 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
999 let swap = arc_swap::ArcSwap::from(Arc::clone(&pipeline_a));
1000
1001 let mut ctx = default_ctx();
1002 let pinned = ctx.pin_pipeline(&swap);
1003 assert!(
1004 Arc::ptr_eq(&pinned, &pipeline_a),
1005 "pin_pipeline should return the current pipeline"
1006 );
1007 assert!(
1008 Arc::ptr_eq(ctx.pinned_pipeline.as_ref().unwrap(), &pipeline_a),
1009 "pinned_pipeline should be stored in ctx"
1010 );
1011 }
1012
1013 #[test]
1014 fn stamp_error_type_is_first_write_wins() {
1015 let mut ctx = PingoraRequestCtx::default();
1016 ctx.stamp_error_type(crate::http::pingora::metrics::ERROR_TYPE_FILTER_REJECT);
1017 ctx.stamp_error_type(crate::http::pingora::metrics::ERROR_TYPE_INTERNAL);
1018 assert_eq!(
1019 ctx.error_type,
1020 Some(crate::http::pingora::metrics::ERROR_TYPE_FILTER_REJECT),
1021 "the first classification wins; a later site must not overwrite it"
1022 );
1023 }
1024
1025 #[test]
1026 fn error_type_is_unset_until_stamped() {
1027 let ctx = PingoraRequestCtx::default();
1028 assert_eq!(ctx.error_type, None, "a fresh context carries no error cause");
1029 }
1030
1031 #[test]
1032 fn pipeline_returns_pinned_after_reload() {
1033 let registry = FilterRegistry::with_builtins();
1034 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1035 let swap = arc_swap::ArcSwap::from(Arc::clone(&pipeline_a));
1036
1037 let mut ctx = default_ctx();
1038 ctx.pin_pipeline(&swap);
1039
1040 let pipeline_b = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1041 swap.store(pipeline_b);
1042
1043 let later = ctx.pipeline(&swap);
1044 assert!(
1045 Arc::ptr_eq(&later, &pipeline_a),
1046 "later hooks should still return pipeline A after reload"
1047 );
1048 }
1049
1050 #[test]
1051 fn new_request_after_reload_pins_new_pipeline() {
1052 let registry = FilterRegistry::with_builtins();
1053 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1054 let swap = arc_swap::ArcSwap::from(Arc::clone(&pipeline_a));
1055
1056 let mut ctx_a = default_ctx();
1057 ctx_a.pin_pipeline(&swap);
1058
1059 let pipeline_b = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1060 swap.store(Arc::clone(&pipeline_b));
1061
1062 let mut ctx_b = default_ctx();
1063 ctx_b.pin_pipeline(&swap);
1064
1065 assert!(
1066 Arc::ptr_eq(&ctx_a.pipeline(&swap), &pipeline_a),
1067 "request A should use old pipeline"
1068 );
1069 assert!(
1070 Arc::ptr_eq(&ctx_b.pipeline(&swap), &pipeline_b),
1071 "request B should use new pipeline"
1072 );
1073 }
1074
1075 #[test]
1076 fn old_pipeline_drops_after_ctx_drops() {
1077 let registry = FilterRegistry::with_builtins();
1078 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1079 let weak_a = Arc::downgrade(&pipeline_a);
1080 let swap = arc_swap::ArcSwap::from(pipeline_a);
1081
1082 let mut ctx = default_ctx();
1083 ctx.pin_pipeline(&swap);
1084
1085 swap.store(Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap()));
1086
1087 assert!(weak_a.upgrade().is_some(), "old pipeline alive while ctx exists");
1088
1089 drop(ctx);
1090
1091 assert!(weak_a.upgrade().is_none(), "old pipeline drops with ctx");
1092 }
1093
1094 #[test]
1095 fn pipeline_utility_returns_pinned_for_every_phase() {
1096 let registry = FilterRegistry::with_builtins();
1097 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1098 let swap = arc_swap::ArcSwap::from(Arc::clone(&pipeline_a));
1099
1100 let mut ctx = default_ctx();
1101 ctx.pin_pipeline(&swap);
1102
1103 swap.store(Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap()));
1104
1105 for phase in ["request_body", "response", "response_body", "logging"] {
1106 let p = ctx.pipeline(&swap);
1107 assert!(
1108 Arc::ptr_eq(&p, &pipeline_a),
1109 "{phase}: should still return pinned pipeline A"
1110 );
1111 }
1112 }
1113
1114 #[test]
1115 fn pipeline_fallback_when_not_pinned() {
1116 let registry = FilterRegistry::with_builtins();
1117 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1118 let swap = arc_swap::ArcSwap::from(Arc::clone(&pipeline_a));
1119
1120 let ctx = default_ctx();
1121
1122 let loaded = ctx.pipeline(&swap);
1123 assert!(
1124 Arc::ptr_eq(&loaded, &pipeline_a),
1125 "unpinned ctx should fall back to current ArcSwap value"
1126 );
1127 }
1128
1129 #[test]
1130 fn pin_pipeline_is_idempotent_after_reload() {
1131 let registry = FilterRegistry::with_builtins();
1132 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1133 let swap = arc_swap::ArcSwap::from(Arc::clone(&pipeline_a));
1134
1135 let mut ctx = default_ctx();
1136 ctx.pin_pipeline(&swap);
1137
1138 let pipeline_b = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1139 swap.store(pipeline_b);
1140
1141 let second_pin = ctx.pin_pipeline(&swap);
1142 assert!(
1143 Arc::ptr_eq(&second_pin, &pipeline_a),
1144 "repeated pin_pipeline after reload should return the original pin"
1145 );
1146 }
1147
1148 #[test]
1149 fn filter_state_isolated_across_pipelines_with_same_ids() {
1150 let registry = FilterRegistry::with_builtins();
1151 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1152 let pipeline_b = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1153
1154 let request = Request {
1155 method: Method::GET,
1156 uri: "/".parse::<Uri>().unwrap(),
1157 headers: HeaderMap::new(),
1158 };
1159
1160 let mut ctx_a = default_ctx();
1161 ctx_a.pinned_pipeline = Some(Arc::clone(&pipeline_a));
1162 ctx_a.request_snapshot = Some(request.clone());
1163 let mut fctx_a = ctx_a.build_filter_context(&pipeline_a, &request, None);
1164 fctx_a.current_filter_id = Some(0);
1165 fctx_a.insert_filter_state(String::from("from_pipeline_a"));
1166 ctx_a.filter_state = fctx_a.filter_state;
1167
1168 let mut ctx_b = default_ctx();
1169 ctx_b.pinned_pipeline = Some(Arc::clone(&pipeline_b));
1170 ctx_b.request_snapshot = Some(request.clone());
1171 let fctx_b = ctx_b.build_filter_context(&pipeline_b, &request, None);
1172
1173 assert!(
1174 fctx_b.filter_state.is_empty(),
1175 "request B should have its own empty state map"
1176 );
1177
1178 let fctx_a2 = ctx_a.filter_context_for(&pipeline_a, None).unwrap();
1179 assert_eq!(
1180 fctx_a2.filter_state.get(&0).and_then(|v| v.downcast_ref::<String>()),
1181 Some(&String::from("from_pipeline_a")),
1182 "request A should still see its own state in a later phase"
1183 );
1184 }
1185
1186 #[test]
1187 fn default_ctx_has_no_adapted_body() {
1188 let ctx = PingoraRequestCtx::default();
1189 assert!(ctx.adapted_request_body.is_none(), "adapted body defaults to None");
1190 assert!(
1191 ctx.retained_adapted_request_body.is_none(),
1192 "retained adapted body defaults to None"
1193 );
1194 assert!(
1195 ctx.adapted_request_body_len.is_none(),
1196 "adapted body length defaults to None"
1197 );
1198 }
1199
1200 fn default_ctx() -> PingoraRequestCtx {
1206 PingoraRequestCtx::default()
1207 }
1208}