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_filter_indices: Vec<bool>,
65
66 pub downstream_tls: bool,
72
73 pub peer_identity: Option<Arc<praxis_tls::TlsPeerIdentity>>,
80
81 pub connection_upgraded: bool,
87
88 pub extensions: praxis_filter::RequestExtensions,
95
96 pub filter_metadata: std::collections::HashMap<String, String>,
102
103 pub pre_read_mutations: Vec<TrustedHeaderMutation>,
109
110 pub structured_metadata: std::collections::HashMap<String, serde_json::Value>,
116
117 pub mutated_request_body_len: Option<usize>,
123
124 pub pinned_pipeline: Option<Arc<FilterPipeline>>,
134
135 pub filter_results: std::collections::HashMap<&'static str, praxis_filter::FilterResultSet>,
141
142 pub filter_state: std::collections::HashMap<usize, Box<dyn std::any::Any + Send + Sync>>,
152
153 pub metrics_cluster: Option<Arc<str>>,
157
158 pub metrics_cluster_shared: Option<::metrics::SharedString>,
163
164 pub metrics_route: Option<::metrics::SharedString>,
166
167 pub(crate) error_type: Option<&'static str>,
174
175 pub(crate) _active_request: Option<crate::http::pingora::metrics::ActiveRequestGuard>,
177
178 pub upstream_connect_start: Option<Instant>,
180
181 pub pre_read_body: Option<VecDeque<Bytes>>,
187
188 pub retained_pre_read_body: Option<VecDeque<Bytes>>,
197
198 pub adapted_request_body: Option<VecDeque<Bytes>>,
208
209 pub retained_adapted_request_body: Option<VecDeque<Bytes>>,
217
218 pub adapted_request_body_len: Option<usize>,
222
223 pub request_body_buffer: Option<BodyBuffer>,
227
228 pub request_body_bytes: u64,
230
231 pub request_body_mode: BodyMode,
235
236 pub request_body_released: bool,
239
240 pub request_is_idempotent: bool,
242
243 pub request_snapshot: Option<Request>,
245
246 pub request_span: Span,
254
255 pub upstream_exchange_span: Span,
262
263 pub request_start: Instant,
265
266 pub response_body_buffer: Option<BodyBuffer>,
270
271 pub response_body_bytes: u64,
273
274 pub response_body_mode: BodyMode,
278
279 pub response_body_released: bool,
281
282 pub response_header_snapshot: Option<Response>,
288
289 pub upstream_response_status: Option<u16>,
292
293 pub grpc_completion: Option<praxis_core::grpc::GrpcCompletion>,
300
301 pub response_phase_done: bool,
305
306 pub pending_rejection: Option<praxis_filter::Rejection>,
310
311 pub response_delivery_complete: bool,
317
318 pub retries: u32,
320
321 pub selected_endpoint_index: Option<usize>,
325
326 pub attempted_endpoints: Vec<Arc<str>>,
328
329 pub retry_policy: Option<Arc<praxis_core::config::RetryPolicy>>,
331
332 pub route_retry_policy: Option<Arc<praxis_core::config::RetryPolicy>>,
334
335 pub cluster_retry_state: Option<Arc<praxis_core::retry::ClusterRetryState>>,
337
338 pub cluster_retry_state_released: bool,
340
341 pub endpoint_reselector: Option<Arc<praxis_filter::EndpointReselector>>,
343
344 pub pending_backoff: Option<std::time::Duration>,
346
347 pub reselect_on_retry: bool,
349
350 pub rewritten_path: Option<String>,
357
358 pub upstream: Option<Upstream>,
360
361 pub upstream_for_retry: Option<Upstream>,
363
364 pub upstream_contacted: bool,
371}
372
373macro_rules! filter_context {
383 ($ctx:expr, $pipeline:expr, $request:expr, $response_header:expr) => {{
384 praxis_filter::HttpFilterContext {
385 buffered_request_body: $ctx
386 .pre_read_body
387 .as_ref()
388 .map(|chunks| chunks.front().cloned().unwrap_or_default()),
389 body_done_indices: std::mem::take(&mut $ctx.cached_body_done_indices),
390 branch_iterations: std::collections::HashMap::new(),
391 client_addr: $ctx.client_addr,
392 cluster: $ctx.cluster.take(),
393 current_filter_id: None,
394 downstream_tls: $ctx.downstream_tls,
395 metrics_route: $ctx.metrics_route.clone(),
396 peer_identity: $ctx.peer_identity.clone(),
397 extensions: std::mem::take(&mut $ctx.extensions),
398 executed_filter_indices: std::mem::take(&mut $ctx.cached_executed_filter_indices),
399 extra_request_headers: Vec::new(),
400 request_headers_to_remove: Vec::new(),
401 request_headers_to_set: Vec::new(),
402 filter_metadata: std::mem::take(&mut $ctx.filter_metadata),
403 prior_pre_read_mutations: Vec::new(),
405 grpc_completion: $ctx.grpc_completion.clone(),
406 pre_read_mutations: std::mem::take(&mut $ctx.pre_read_mutations),
407 structured_metadata: std::mem::take(&mut $ctx.structured_metadata),
408 filter_results: std::mem::take(&mut $ctx.filter_results),
409 filter_state: std::mem::take(&mut $ctx.filter_state),
410 health_registry: $pipeline.health_registry(),
411 id_generator: $pipeline.id_generator(),
412 kv_stores: $pipeline.kv_stores(),
413 session_stores: $pipeline.session_stores(),
414 subrequest_client: $pipeline.subrequest_client(),
415 subrequest_response_mode: praxis_filter::SubRequestResponseMode::Buffered,
416 request: $request,
417 request_body_bytes: $ctx.request_body_bytes,
418 request_body_mode: $ctx.request_body_mode,
419 request_start: $ctx.request_start,
420 response_body_bytes: $ctx.response_body_bytes,
421 response_body_mode: $ctx.response_body_mode,
422 response_header: $response_header,
423 response_headers_modified: false,
424 upstream_reached: $ctx.upstream_contacted,
425 rewritten_path: $ctx.rewritten_path.take(),
426 selected_endpoint_index: $ctx.selected_endpoint_index,
427 attempted_endpoints: std::mem::take(&mut $ctx.attempted_endpoints),
428 retry_policy: $ctx.retry_policy.clone(),
429 route_retry_policy: $ctx.route_retry_policy.clone(),
430 cluster_retry_state: $ctx.cluster_retry_state.clone(),
431 cluster_retry_state_released: $ctx.cluster_retry_state_released,
432 endpoint_reselector: $ctx.endpoint_reselector.clone(),
433 pinned_endpoint_address: None,
434 time_source: $pipeline.time_source(),
435 upstream: $ctx.upstream.take().or_else(|| $ctx.upstream_for_retry.clone()),
436 }
437 }};
438}
439
440impl PingoraRequestCtx {
441 pub fn build_filter_context<'a>(
465 &mut self,
466 pipeline: &'a FilterPipeline,
467 request: &'a Request,
468 response_header: Option<&'a mut Response>,
469 ) -> praxis_filter::HttpFilterContext<'a> {
470 filter_context!(self, pipeline, request, response_header)
471 }
472
473 pub fn filter_context_for<'a>(
496 &'a mut self,
497 pipeline: &'a FilterPipeline,
498 response_header: Option<&'a mut Response>,
499 ) -> Option<praxis_filter::HttpFilterContext<'a>> {
500 let request = self.request_snapshot.as_ref()?;
501 Some(filter_context!(self, pipeline, request, response_header))
502 }
503
504 pub fn response_body_context_for<'a>(
512 &'a mut self,
513 pipeline: &'a FilterPipeline,
514 ) -> Option<(praxis_filter::HttpFilterContext<'a>, Option<&'a Response>)> {
515 let request = self.request_snapshot.as_ref()?;
516 let response_header = self.response_header_snapshot.as_ref();
517 Some((filter_context!(self, pipeline, request, None), response_header))
518 }
519
520 pub fn pin_pipeline(&mut self, swap: &arc_swap::ArcSwap<FilterPipeline>) -> Arc<FilterPipeline> {
533 if let Some(existing) = &self.pinned_pipeline {
534 return Arc::clone(existing);
535 }
536 let pipeline = swap.load_full();
537 pipeline.prepare_extensions(&mut self.extensions);
540 self.pinned_pipeline = Some(Arc::clone(&pipeline));
541 pipeline
542 }
543
544 pub(crate) fn stamp_error_type(&mut self, error_type: &'static str) {
550 self.error_type.get_or_insert(error_type);
551 }
552
553 pub fn pipeline(&self, swap: &arc_swap::ArcSwap<FilterPipeline>) -> Arc<FilterPipeline> {
562 self.pinned_pipeline
563 .as_ref()
564 .map_or_else(|| swap.load_full(), Arc::clone)
565 }
566}
567
568impl Default for PingoraRequestCtx {
569 #[expect(
570 clippy::too_many_lines,
571 reason = "context default enumerates all lifecycle fields explicitly"
572 )]
573 fn default() -> Self {
574 Self {
575 _connection_permit: None,
576 _global_connection_permit: None,
577 cached_body_done_indices: Vec::new(),
578 cached_executed_filter_indices: Vec::new(),
579 client_addr: None,
580 client_http_version: None,
581 cluster: None,
582 connection_upgraded: false,
583 downstream_tls: false,
584 peer_identity: None,
585 extensions: praxis_filter::RequestExtensions::new(),
586 filter_metadata: std::collections::HashMap::new(),
587 pre_read_mutations: Vec::new(),
588 structured_metadata: std::collections::HashMap::new(),
589 mutated_request_body_len: None,
590 pinned_pipeline: None,
591 filter_results: std::collections::HashMap::new(),
592 filter_state: std::collections::HashMap::new(),
593 metrics_cluster: None,
594 metrics_cluster_shared: None,
595 metrics_route: None,
596 error_type: None,
597 _active_request: None,
598 upstream_connect_start: None,
599 pre_read_body: None,
600 retained_pre_read_body: None,
601 adapted_request_body: None,
602 retained_adapted_request_body: None,
603 adapted_request_body_len: None,
604 request_body_buffer: None,
605 request_body_bytes: 0,
606 request_body_mode: BodyMode::Stream,
607 request_body_released: false,
608 request_is_idempotent: false,
609 request_snapshot: None,
610 request_span: Span::none(),
611 request_start: Instant::now(),
612 upstream_exchange_span: Span::none(),
613 response_body_buffer: None,
614 response_body_bytes: 0,
615 response_body_mode: BodyMode::Stream,
616 response_body_released: false,
617 response_header_snapshot: None,
618 upstream_response_status: None,
619 grpc_completion: None,
620 response_phase_done: false,
621 response_delivery_complete: false,
622 pending_rejection: None,
623 retries: 0,
624 rewritten_path: None,
625 selected_endpoint_index: None,
626 attempted_endpoints: Vec::new(),
627 retry_policy: None,
628 route_retry_policy: None,
629 cluster_retry_state: None,
630 cluster_retry_state_released: false,
631 endpoint_reselector: None,
632 pending_backoff: None,
633 reselect_on_retry: false,
634 upstream: None,
635 upstream_for_retry: None,
636 upstream_contacted: false,
637 }
638 }
639}
640
641#[cfg(test)]
646#[expect(clippy::allow_attributes, reason = "blanket test suppressions")]
647#[allow(
648 clippy::unwrap_used,
649 clippy::expect_used,
650 clippy::indexing_slicing,
651 clippy::significant_drop_tightening,
652 clippy::too_many_lines,
653 reason = "tests"
654)]
655mod tests {
656 use std::net::Ipv4Addr;
657
658 use http::{HeaderMap, Method, Uri};
659 use praxis_filter::FilterRegistry;
660
661 use super::*;
662
663 #[test]
664 fn default_state_has_no_client_addr() {
665 let ctx = default_ctx();
666 assert!(ctx.client_addr.is_none(), "default client_addr should be None");
667 }
668
669 #[test]
670 fn default_state_has_no_cluster() {
671 let ctx = default_ctx();
672 assert!(ctx.cluster.is_none(), "default cluster should be None");
673 }
674
675 #[test]
676 fn default_state_has_zero_retries() {
677 let ctx = default_ctx();
678 assert_eq!(ctx.retries, 0, "default retries should be zero");
679 }
680
681 #[test]
682 fn default_state_flags_are_false() {
683 let ctx = default_ctx();
684 assert!(
685 !ctx.request_body_released,
686 "default request_body_released should be false"
687 );
688 assert!(
689 !ctx.response_body_released,
690 "default response_body_released should be false"
691 );
692 assert!(
693 !ctx.request_is_idempotent,
694 "default request_is_idempotent should be false"
695 );
696 assert!(!ctx.response_phase_done, "default response_phase_done should be false");
697 }
698
699 #[test]
700 fn default_state_buffers_are_none() {
701 let ctx = default_ctx();
702 assert!(
703 ctx.request_body_buffer.is_none(),
704 "default request_body_buffer should be None"
705 );
706 assert!(
707 ctx.response_body_buffer.is_none(),
708 "default response_body_buffer should be None"
709 );
710 assert!(ctx.pre_read_body.is_none(), "default pre_read_body should be None");
711 }
712
713 #[test]
714 fn default_state_request_span_is_disabled() {
715 let ctx = default_ctx();
716 assert!(
717 ctx.request_span.is_disabled(),
718 "default request_span should be a disabled (none) span"
719 );
720 }
721
722 #[test]
723 fn default_state_upstream_exchange_span_is_disabled() {
724 let ctx = default_ctx();
725 assert!(
726 ctx.upstream_exchange_span.is_disabled(),
727 "default upstream_exchange_span should be a disabled (none) span"
728 );
729 }
730
731 #[test]
732 fn default_state_snapshots_are_none() {
733 let ctx = default_ctx();
734 assert!(
735 ctx.request_snapshot.is_none(),
736 "default request_snapshot should be None"
737 );
738 assert!(ctx.upstream.is_none(), "default upstream should be None");
739 assert!(
740 ctx.upstream_for_retry.is_none(),
741 "default upstream_for_retry should be None"
742 );
743 }
744
745 #[test]
746 fn set_client_addr() {
747 let mut ctx = default_ctx();
748 let addr = IpAddr::V4(Ipv4Addr::new(192, 168, 1, 100));
749 ctx.client_addr = Some(addr);
750 assert_eq!(
751 ctx.client_addr.unwrap(),
752 addr,
753 "client_addr should match assigned value"
754 );
755 }
756
757 #[test]
758 fn set_cluster() {
759 let mut ctx = default_ctx();
760 ctx.cluster = Some(Arc::from("api-cluster"));
761 assert_eq!(
762 ctx.cluster.as_deref(),
763 Some("api-cluster"),
764 "cluster should match assigned value"
765 );
766 }
767
768 #[test]
769 fn set_upstream() {
770 let mut ctx = default_ctx();
771 let upstream = Upstream {
772 address: Arc::from("10.0.0.1:80"),
773 authority: None,
774 tls: None,
775 connection: Arc::new(praxis_core::connectivity::ConnectionOptions::default()),
776 };
777 ctx.upstream = Some(upstream.clone());
778 assert_eq!(
779 &*ctx.upstream.as_ref().unwrap().address,
780 "10.0.0.1:80",
781 "upstream address should match assigned value"
782 );
783 }
784
785 #[test]
786 fn increment_retries() {
787 let mut ctx = default_ctx();
788 ctx.retries += 1;
789 ctx.retries += 1;
790 assert_eq!(ctx.retries, 2, "retries should be 2 after two increments");
791 }
792
793 #[test]
794 fn release_request_body_flag() {
795 let mut ctx = default_ctx();
796 assert!(!ctx.request_body_released, "request_body_released should start false");
797 ctx.request_body_released = true;
798 assert!(
799 ctx.request_body_released,
800 "request_body_released should be true after setting"
801 );
802 }
803
804 #[test]
805 fn release_response_body_flag() {
806 let mut ctx = default_ctx();
807 assert!(!ctx.response_body_released, "response_body_released should start false");
808 ctx.response_body_released = true;
809 assert!(
810 ctx.response_body_released,
811 "response_body_released should be true after setting"
812 );
813 }
814
815 #[test]
816 fn response_phase_done_flag() {
817 let mut ctx = default_ctx();
818 assert!(!ctx.response_phase_done, "response_phase_done should start false");
819 ctx.response_phase_done = true;
820 assert!(
821 ctx.response_phase_done,
822 "response_phase_done should be true after setting"
823 );
824 }
825
826 #[test]
827 fn set_pre_read_body() {
828 let mut ctx = default_ctx();
829 let chunks = VecDeque::from([Bytes::from_static(b"chunk1"), Bytes::from_static(b"chunk2")]);
830 ctx.pre_read_body = Some(chunks);
831 let body = ctx.pre_read_body.as_ref().unwrap();
832 assert_eq!(body.len(), 2, "pre_read_body should contain 2 chunks");
833 assert_eq!(body[0], Bytes::from_static(b"chunk1"), "first chunk should be 'chunk1'");
834 assert_eq!(
835 body[1],
836 Bytes::from_static(b"chunk2"),
837 "second chunk should be 'chunk2'"
838 );
839 }
840
841 #[test]
842 fn set_request_snapshot() {
843 let mut ctx = default_ctx();
844 let snapshot = Request {
845 method: Method::POST,
846 uri: "/api/data".parse::<Uri>().unwrap(),
847 headers: HeaderMap::new(),
848 };
849 ctx.request_snapshot = Some(snapshot);
850 let snap = ctx.request_snapshot.as_ref().unwrap();
851 assert_eq!(snap.method, Method::POST, "snapshot method should be POST");
852 assert_eq!(snap.uri.path(), "/api/data", "snapshot URI path should be /api/data");
853 }
854
855 #[test]
856 fn request_body_buffer_lifecycle() {
857 let mut ctx = default_ctx();
858 let mut buf = BodyBuffer::new(100);
859 buf.push(Bytes::from_static(b"data")).unwrap();
860 ctx.request_body_buffer = Some(buf);
861
862 assert!(
863 ctx.request_body_buffer.is_some(),
864 "buffer should be present after assignment"
865 );
866 let taken = ctx.request_body_buffer.take().unwrap();
867 assert_eq!(
868 taken.freeze(),
869 Bytes::from_static(b"data"),
870 "frozen buffer should contain pushed data"
871 );
872 assert!(ctx.request_body_buffer.is_none(), "buffer should be None after take");
873 }
874
875 #[test]
876 fn default_request_body_mode_is_stream() {
877 let ctx = default_ctx();
878 assert_eq!(
879 ctx.request_body_mode,
880 BodyMode::Stream,
881 "default request_body_mode should be Stream"
882 );
883 }
884
885 #[test]
886 fn default_response_body_mode_is_stream() {
887 let ctx = default_ctx();
888 assert_eq!(
889 ctx.response_body_mode,
890 BodyMode::Stream,
891 "default response_body_mode should be Stream"
892 );
893 }
894
895 #[test]
896 fn set_request_body_mode() {
897 let mut ctx = default_ctx();
898 ctx.request_body_mode = BodyMode::StreamBuffer { max_bytes: Some(4096) };
899 assert_eq!(
900 ctx.request_body_mode,
901 BodyMode::StreamBuffer { max_bytes: Some(4096) },
902 "request_body_mode should match assigned value"
903 );
904 }
905
906 #[test]
907 fn set_response_body_mode() {
908 let mut ctx = default_ctx();
909 ctx.response_body_mode = BodyMode::StreamBuffer { max_bytes: Some(8192) };
910 assert_eq!(
911 ctx.response_body_mode,
912 BodyMode::StreamBuffer { max_bytes: Some(8192) },
913 "response_body_mode should match assigned value"
914 );
915 }
916
917 #[test]
922 fn metadata_roundtrip_through_filter_context() {
923 let registry = FilterRegistry::with_builtins();
924 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
925 let request = Request {
926 method: Method::GET,
927 uri: "/".parse::<Uri>().unwrap(),
928 headers: HeaderMap::new(),
929 };
930
931 let mut ctx = default_ctx();
932 ctx.filter_metadata.insert("rpc.method".to_owned(), "echo".to_owned());
933 ctx.filter_metadata.insert("rpc.status".to_owned(), "ok".to_owned());
934
935 let fctx = ctx.build_filter_context(&pipeline, &request, None);
936 assert_eq!(
937 fctx.get_metadata("rpc.method"),
938 Some("echo"),
939 "metadata written before build should survive into filter context"
940 );
941 assert_eq!(
942 fctx.get_metadata("rpc.status"),
943 Some("ok"),
944 "multiple metadata keys should round-trip"
945 );
946 }
947
948 #[test]
949 fn metadata_written_in_filter_context_persists_back() {
950 let registry = FilterRegistry::with_builtins();
951 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
952 let request = Request {
953 method: Method::GET,
954 uri: "/".parse::<Uri>().unwrap(),
955 headers: HeaderMap::new(),
956 };
957
958 let mut ctx = default_ctx();
959 let mut fctx = ctx.build_filter_context(&pipeline, &request, None);
960 fctx.set_metadata("trace.id", "abc-123");
961 fctx.set_metadata("trace.span", "42");
962
963 ctx.filter_metadata = fctx.filter_metadata;
964 assert_eq!(
965 ctx.filter_metadata.get("trace.id").map(String::as_str),
966 Some("abc-123"),
967 "metadata set in filter context should persist back to protocol context"
968 );
969 assert_eq!(
970 ctx.filter_metadata.get("trace.span").map(String::as_str),
971 Some("42"),
972 "multiple metadata keys should persist back"
973 );
974 }
975
976 #[test]
981 fn pin_pipeline_captures_current_arc() {
982 let registry = FilterRegistry::with_builtins();
983 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
984 let swap = arc_swap::ArcSwap::from(Arc::clone(&pipeline_a));
985
986 let mut ctx = default_ctx();
987 let pinned = ctx.pin_pipeline(&swap);
988 assert!(
989 Arc::ptr_eq(&pinned, &pipeline_a),
990 "pin_pipeline should return the current pipeline"
991 );
992 assert!(
993 Arc::ptr_eq(ctx.pinned_pipeline.as_ref().unwrap(), &pipeline_a),
994 "pinned_pipeline should be stored in ctx"
995 );
996 }
997
998 #[test]
999 fn stamp_error_type_is_first_write_wins() {
1000 let mut ctx = PingoraRequestCtx::default();
1001 ctx.stamp_error_type(crate::http::pingora::metrics::ERROR_TYPE_FILTER_REJECT);
1002 ctx.stamp_error_type(crate::http::pingora::metrics::ERROR_TYPE_INTERNAL);
1003 assert_eq!(
1004 ctx.error_type,
1005 Some(crate::http::pingora::metrics::ERROR_TYPE_FILTER_REJECT),
1006 "the first classification wins; a later site must not overwrite it"
1007 );
1008 }
1009
1010 #[test]
1011 fn error_type_is_unset_until_stamped() {
1012 let ctx = PingoraRequestCtx::default();
1013 assert_eq!(ctx.error_type, None, "a fresh context carries no error cause");
1014 }
1015
1016 #[test]
1017 fn pipeline_returns_pinned_after_reload() {
1018 let registry = FilterRegistry::with_builtins();
1019 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1020 let swap = arc_swap::ArcSwap::from(Arc::clone(&pipeline_a));
1021
1022 let mut ctx = default_ctx();
1023 ctx.pin_pipeline(&swap);
1024
1025 let pipeline_b = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1026 swap.store(pipeline_b);
1027
1028 let later = ctx.pipeline(&swap);
1029 assert!(
1030 Arc::ptr_eq(&later, &pipeline_a),
1031 "later hooks should still return pipeline A after reload"
1032 );
1033 }
1034
1035 #[test]
1036 fn new_request_after_reload_pins_new_pipeline() {
1037 let registry = FilterRegistry::with_builtins();
1038 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1039 let swap = arc_swap::ArcSwap::from(Arc::clone(&pipeline_a));
1040
1041 let mut ctx_a = default_ctx();
1042 ctx_a.pin_pipeline(&swap);
1043
1044 let pipeline_b = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1045 swap.store(Arc::clone(&pipeline_b));
1046
1047 let mut ctx_b = default_ctx();
1048 ctx_b.pin_pipeline(&swap);
1049
1050 assert!(
1051 Arc::ptr_eq(&ctx_a.pipeline(&swap), &pipeline_a),
1052 "request A should use old pipeline"
1053 );
1054 assert!(
1055 Arc::ptr_eq(&ctx_b.pipeline(&swap), &pipeline_b),
1056 "request B should use new pipeline"
1057 );
1058 }
1059
1060 #[test]
1061 fn old_pipeline_drops_after_ctx_drops() {
1062 let registry = FilterRegistry::with_builtins();
1063 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1064 let weak_a = Arc::downgrade(&pipeline_a);
1065 let swap = arc_swap::ArcSwap::from(pipeline_a);
1066
1067 let mut ctx = default_ctx();
1068 ctx.pin_pipeline(&swap);
1069
1070 swap.store(Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap()));
1071
1072 assert!(weak_a.upgrade().is_some(), "old pipeline alive while ctx exists");
1073
1074 drop(ctx);
1075
1076 assert!(weak_a.upgrade().is_none(), "old pipeline drops with ctx");
1077 }
1078
1079 #[test]
1080 fn pipeline_utility_returns_pinned_for_every_phase() {
1081 let registry = FilterRegistry::with_builtins();
1082 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1083 let swap = arc_swap::ArcSwap::from(Arc::clone(&pipeline_a));
1084
1085 let mut ctx = default_ctx();
1086 ctx.pin_pipeline(&swap);
1087
1088 swap.store(Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap()));
1089
1090 for phase in ["request_body", "response", "response_body", "logging"] {
1091 let p = ctx.pipeline(&swap);
1092 assert!(
1093 Arc::ptr_eq(&p, &pipeline_a),
1094 "{phase}: should still return pinned pipeline A"
1095 );
1096 }
1097 }
1098
1099 #[test]
1100 fn pipeline_fallback_when_not_pinned() {
1101 let registry = FilterRegistry::with_builtins();
1102 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1103 let swap = arc_swap::ArcSwap::from(Arc::clone(&pipeline_a));
1104
1105 let ctx = default_ctx();
1106
1107 let loaded = ctx.pipeline(&swap);
1108 assert!(
1109 Arc::ptr_eq(&loaded, &pipeline_a),
1110 "unpinned ctx should fall back to current ArcSwap value"
1111 );
1112 }
1113
1114 #[test]
1115 fn pin_pipeline_is_idempotent_after_reload() {
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 mut ctx = default_ctx();
1121 ctx.pin_pipeline(&swap);
1122
1123 let pipeline_b = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1124 swap.store(pipeline_b);
1125
1126 let second_pin = ctx.pin_pipeline(&swap);
1127 assert!(
1128 Arc::ptr_eq(&second_pin, &pipeline_a),
1129 "repeated pin_pipeline after reload should return the original pin"
1130 );
1131 }
1132
1133 #[test]
1134 fn filter_state_isolated_across_pipelines_with_same_ids() {
1135 let registry = FilterRegistry::with_builtins();
1136 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1137 let pipeline_b = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1138
1139 let request = Request {
1140 method: Method::GET,
1141 uri: "/".parse::<Uri>().unwrap(),
1142 headers: HeaderMap::new(),
1143 };
1144
1145 let mut ctx_a = default_ctx();
1146 ctx_a.pinned_pipeline = Some(Arc::clone(&pipeline_a));
1147 ctx_a.request_snapshot = Some(request.clone());
1148 let mut fctx_a = ctx_a.build_filter_context(&pipeline_a, &request, None);
1149 fctx_a.current_filter_id = Some(0);
1150 fctx_a.insert_filter_state(String::from("from_pipeline_a"));
1151 ctx_a.filter_state = fctx_a.filter_state;
1152
1153 let mut ctx_b = default_ctx();
1154 ctx_b.pinned_pipeline = Some(Arc::clone(&pipeline_b));
1155 ctx_b.request_snapshot = Some(request.clone());
1156 let fctx_b = ctx_b.build_filter_context(&pipeline_b, &request, None);
1157
1158 assert!(
1159 fctx_b.filter_state.is_empty(),
1160 "request B should have its own empty state map"
1161 );
1162
1163 let fctx_a2 = ctx_a.filter_context_for(&pipeline_a, None).unwrap();
1164 assert_eq!(
1165 fctx_a2.filter_state.get(&0).and_then(|v| v.downcast_ref::<String>()),
1166 Some(&String::from("from_pipeline_a")),
1167 "request A should still see its own state in a later phase"
1168 );
1169 }
1170
1171 #[test]
1172 fn default_ctx_has_no_adapted_body() {
1173 let ctx = PingoraRequestCtx::default();
1174 assert!(ctx.adapted_request_body.is_none(), "adapted body defaults to None");
1175 assert!(
1176 ctx.retained_adapted_request_body.is_none(),
1177 "retained adapted body defaults to None"
1178 );
1179 assert!(
1180 ctx.adapted_request_body_len.is_none(),
1181 "adapted body length defaults to None"
1182 );
1183 }
1184
1185 fn default_ctx() -> PingoraRequestCtx {
1191 PingoraRequestCtx::default()
1192 }
1193}