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>,
37
38 pub _global_connection_permit: Option<OwnedSemaphorePermit>,
42
43 pub client_addr: Option<IpAddr>,
45
46 pub client_http_version: Option<http::Version>,
51
52 pub cluster: Option<Arc<str>>,
54
55 pub cached_body_done_indices: Vec<bool>,
61
62 pub cached_executed_filter_indices: Vec<bool>,
67
68 pub downstream_tls: bool,
75
76 pub peer_identity: Option<Arc<praxis_tls::TlsPeerIdentity>>,
85
86 pub connection_upgraded: bool,
92
93 pub extensions: praxis_filter::RequestExtensions,
100
101 pub filter_metadata: std::collections::HashMap<String, String>,
107
108 pub pre_read_mutations: Vec<TrustedHeaderMutation>,
114
115 pub structured_metadata: std::collections::HashMap<String, serde_json::Value>,
121
122 pub mutated_request_body_len: Option<usize>,
128
129 pub pinned_pipeline: Option<Arc<FilterPipeline>>,
139
140 pub filter_results: std::collections::HashMap<&'static str, praxis_filter::FilterResultSet>,
146
147 pub filter_state: std::collections::HashMap<usize, Box<dyn std::any::Any + Send + Sync>>,
157
158 pub metrics_cluster: Option<Arc<str>>,
162
163 pub metrics_cluster_shared: Option<::metrics::SharedString>,
170
171 pub metrics_route: Option<::metrics::SharedString>,
173
174 pub(crate) _active_connection: Option<crate::http::pingora::metrics::ActiveConnectionGuard>,
176
177 pub upstream_connect_start: Option<Instant>,
179
180 pub pre_read_body: Option<VecDeque<Bytes>>,
188
189 pub request_body_buffer: Option<BodyBuffer>,
193
194 pub request_body_bytes: u64,
196
197 pub request_body_mode: BodyMode,
201
202 pub request_body_released: bool,
205
206 pub request_is_idempotent: bool,
208
209 pub request_snapshot: Option<Request>,
211
212 pub request_span: Span,
220
221 pub upstream_exchange_span: Span,
228
229 pub request_start: Instant,
231
232 pub response_body_buffer: Option<BodyBuffer>,
236
237 pub response_body_bytes: u64,
239
240 pub response_body_mode: BodyMode,
244
245 pub response_body_released: bool,
247
248 pub response_header_snapshot: Option<Response>,
254
255 pub upstream_response_status: Option<u16>,
258
259 pub response_phase_done: bool,
263
264 pub pending_rejection: Option<praxis_filter::Rejection>,
268
269 pub response_delivery_complete: bool,
275
276 pub retries: u32,
278
279 pub selected_endpoint_index: Option<usize>,
283
284 pub attempted_endpoints: Vec<Arc<str>>,
286
287 pub retry_policy: Option<Arc<praxis_core::config::RetryPolicy>>,
289
290 pub route_retry_policy: Option<Arc<praxis_core::config::RetryPolicy>>,
292
293 pub cluster_retry_state: Option<Arc<praxis_core::retry::ClusterRetryState>>,
295
296 pub cluster_retry_state_released: bool,
298
299 pub endpoint_reselector: Option<Arc<praxis_filter::EndpointReselector>>,
301
302 pub pending_backoff: Option<std::time::Duration>,
304
305 pub reselect_on_retry: bool,
307
308 pub rewritten_path: Option<String>,
315
316 pub upstream: Option<Upstream>,
318
319 pub upstream_for_retry: Option<Upstream>,
321}
322
323macro_rules! filter_context {
333 ($ctx:expr, $pipeline:expr, $request:expr, $response_header:expr) => {{
334 praxis_filter::HttpFilterContext {
335 buffered_request_body: $ctx
336 .pre_read_body
337 .as_ref()
338 .map(|chunks| chunks.front().cloned().unwrap_or_default()),
339 body_done_indices: std::mem::take(&mut $ctx.cached_body_done_indices),
340 branch_iterations: std::collections::HashMap::new(),
341 client_addr: $ctx.client_addr,
342 cluster: $ctx.cluster.take(),
343 current_filter_id: None,
344 downstream_tls: $ctx.downstream_tls,
345 metrics_route: $ctx.metrics_route.clone(),
346 peer_identity: $ctx.peer_identity.clone(),
347 extensions: std::mem::take(&mut $ctx.extensions),
348 executed_filter_indices: std::mem::take(&mut $ctx.cached_executed_filter_indices),
349 extra_request_headers: Vec::new(),
350 request_headers_to_remove: Vec::new(),
351 request_headers_to_set: Vec::new(),
352 filter_metadata: std::mem::take(&mut $ctx.filter_metadata),
353 pre_read_mutations: std::mem::take(&mut $ctx.pre_read_mutations),
354 structured_metadata: std::mem::take(&mut $ctx.structured_metadata),
355 filter_results: std::mem::take(&mut $ctx.filter_results),
356 filter_state: std::mem::take(&mut $ctx.filter_state),
357 health_registry: $pipeline.health_registry(),
358 id_generator: $pipeline.id_generator(),
359 kv_stores: $pipeline.kv_stores(),
360 session_stores: $pipeline.session_stores(),
361 subrequest_client: $pipeline.subrequest_client(),
362 subrequest_response_mode: praxis_filter::SubRequestResponseMode::Buffered,
363 request: $request,
364 request_body_bytes: $ctx.request_body_bytes,
365 request_body_mode: $ctx.request_body_mode,
366 request_start: $ctx.request_start,
367 response_body_bytes: $ctx.response_body_bytes,
368 response_body_mode: $ctx.response_body_mode,
369 response_header: $response_header,
370 response_headers_modified: false,
371 rewritten_path: $ctx.rewritten_path.take(),
372 selected_endpoint_index: $ctx.selected_endpoint_index,
373 attempted_endpoints: std::mem::take(&mut $ctx.attempted_endpoints),
374 retry_policy: $ctx.retry_policy.clone(),
375 route_retry_policy: $ctx.route_retry_policy.clone(),
376 cluster_retry_state: $ctx.cluster_retry_state.clone(),
377 cluster_retry_state_released: $ctx.cluster_retry_state_released,
378 endpoint_reselector: $ctx.endpoint_reselector.clone(),
379 pinned_endpoint_address: None,
380 time_source: $pipeline.time_source(),
381 upstream: $ctx.upstream.take().or_else(|| $ctx.upstream_for_retry.clone()),
382 }
383 }};
384}
385
386impl PingoraRequestCtx {
387 pub fn build_filter_context<'a>(
411 &mut self,
412 pipeline: &'a FilterPipeline,
413 request: &'a Request,
414 response_header: Option<&'a mut Response>,
415 ) -> praxis_filter::HttpFilterContext<'a> {
416 filter_context!(self, pipeline, request, response_header)
417 }
418
419 pub fn filter_context_for<'a>(
446 &'a mut self,
447 pipeline: &'a FilterPipeline,
448 response_header: Option<&'a mut Response>,
449 ) -> Option<praxis_filter::HttpFilterContext<'a>> {
450 let request = self.request_snapshot.as_ref()?;
451 Some(filter_context!(self, pipeline, request, response_header))
452 }
453
454 pub fn response_body_context_for<'a>(
462 &'a mut self,
463 pipeline: &'a FilterPipeline,
464 ) -> Option<(praxis_filter::HttpFilterContext<'a>, Option<&'a Response>)> {
465 let request = self.request_snapshot.as_ref()?;
466 let response_header = self.response_header_snapshot.as_ref();
467 Some((filter_context!(self, pipeline, request, None), response_header))
468 }
469
470 pub fn pin_pipeline(&mut self, swap: &arc_swap::ArcSwap<FilterPipeline>) -> Arc<FilterPipeline> {
483 if let Some(existing) = &self.pinned_pipeline {
484 return Arc::clone(existing);
485 }
486 let pipeline = swap.load_full();
487 pipeline.prepare_extensions(&mut self.extensions);
490 self.pinned_pipeline = Some(Arc::clone(&pipeline));
491 pipeline
492 }
493
494 pub fn pipeline(&self, swap: &arc_swap::ArcSwap<FilterPipeline>) -> Arc<FilterPipeline> {
507 self.pinned_pipeline
508 .as_ref()
509 .map_or_else(|| swap.load_full(), Arc::clone)
510 }
511}
512
513impl Default for PingoraRequestCtx {
514 #[expect(
515 clippy::too_many_lines,
516 reason = "context default enumerates all lifecycle fields explicitly"
517 )]
518 fn default() -> Self {
519 Self {
520 _connection_permit: None,
521 _global_connection_permit: None,
522 cached_body_done_indices: Vec::new(),
523 cached_executed_filter_indices: Vec::new(),
524 client_addr: None,
525 client_http_version: None,
526 cluster: None,
527 connection_upgraded: false,
528 downstream_tls: false,
529 peer_identity: None,
530 extensions: praxis_filter::RequestExtensions::new(),
531 filter_metadata: std::collections::HashMap::new(),
532 pre_read_mutations: Vec::new(),
533 structured_metadata: std::collections::HashMap::new(),
534 mutated_request_body_len: None,
535 pinned_pipeline: None,
536 filter_results: std::collections::HashMap::new(),
537 filter_state: std::collections::HashMap::new(),
538 metrics_cluster: None,
539 metrics_cluster_shared: None,
540 metrics_route: None,
541 _active_connection: None,
542 upstream_connect_start: None,
543 pre_read_body: None,
544 request_body_buffer: None,
545 request_body_bytes: 0,
546 request_body_mode: BodyMode::Stream,
547 request_body_released: false,
548 request_is_idempotent: false,
549 request_snapshot: None,
550 request_span: Span::none(),
551 request_start: Instant::now(),
552 upstream_exchange_span: Span::none(),
553 response_body_buffer: None,
554 response_body_bytes: 0,
555 response_body_mode: BodyMode::Stream,
556 response_body_released: false,
557 response_header_snapshot: None,
558 upstream_response_status: None,
559 response_phase_done: false,
560 response_delivery_complete: false,
561 pending_rejection: None,
562 retries: 0,
563 rewritten_path: None,
564 selected_endpoint_index: None,
565 attempted_endpoints: Vec::new(),
566 retry_policy: None,
567 route_retry_policy: None,
568 cluster_retry_state: None,
569 cluster_retry_state_released: false,
570 endpoint_reselector: None,
571 pending_backoff: None,
572 reselect_on_retry: false,
573 upstream: None,
574 upstream_for_retry: None,
575 }
576 }
577}
578
579#[cfg(test)]
584#[expect(clippy::allow_attributes, reason = "blanket test suppressions")]
585#[allow(
586 clippy::unwrap_used,
587 clippy::expect_used,
588 clippy::indexing_slicing,
589 clippy::significant_drop_tightening,
590 clippy::too_many_lines,
591 reason = "tests"
592)]
593mod tests {
594 use std::net::Ipv4Addr;
595
596 use http::{HeaderMap, Method, Uri};
597 use praxis_filter::FilterRegistry;
598
599 use super::*;
600
601 #[test]
602 fn default_state_has_no_client_addr() {
603 let ctx = default_ctx();
604 assert!(ctx.client_addr.is_none(), "default client_addr should be None");
605 }
606
607 #[test]
608 fn default_state_has_no_cluster() {
609 let ctx = default_ctx();
610 assert!(ctx.cluster.is_none(), "default cluster should be None");
611 }
612
613 #[test]
614 fn default_state_has_zero_retries() {
615 let ctx = default_ctx();
616 assert_eq!(ctx.retries, 0, "default retries should be zero");
617 }
618
619 #[test]
620 fn default_state_flags_are_false() {
621 let ctx = default_ctx();
622 assert!(
623 !ctx.request_body_released,
624 "default request_body_released should be false"
625 );
626 assert!(
627 !ctx.response_body_released,
628 "default response_body_released should be false"
629 );
630 assert!(
631 !ctx.request_is_idempotent,
632 "default request_is_idempotent should be false"
633 );
634 assert!(!ctx.response_phase_done, "default response_phase_done should be false");
635 }
636
637 #[test]
638 fn default_state_buffers_are_none() {
639 let ctx = default_ctx();
640 assert!(
641 ctx.request_body_buffer.is_none(),
642 "default request_body_buffer should be None"
643 );
644 assert!(
645 ctx.response_body_buffer.is_none(),
646 "default response_body_buffer should be None"
647 );
648 assert!(ctx.pre_read_body.is_none(), "default pre_read_body should be None");
649 }
650
651 #[test]
652 fn default_state_request_span_is_disabled() {
653 let ctx = default_ctx();
654 assert!(
655 ctx.request_span.is_disabled(),
656 "default request_span should be a disabled (none) span"
657 );
658 }
659
660 #[test]
661 fn default_state_upstream_exchange_span_is_disabled() {
662 let ctx = default_ctx();
663 assert!(
664 ctx.upstream_exchange_span.is_disabled(),
665 "default upstream_exchange_span should be a disabled (none) span"
666 );
667 }
668
669 #[test]
670 fn default_state_snapshots_are_none() {
671 let ctx = default_ctx();
672 assert!(
673 ctx.request_snapshot.is_none(),
674 "default request_snapshot should be None"
675 );
676 assert!(ctx.upstream.is_none(), "default upstream should be None");
677 assert!(
678 ctx.upstream_for_retry.is_none(),
679 "default upstream_for_retry should be None"
680 );
681 }
682
683 #[test]
684 fn set_client_addr() {
685 let mut ctx = default_ctx();
686 let addr = IpAddr::V4(Ipv4Addr::new(192, 168, 1, 100));
687 ctx.client_addr = Some(addr);
688 assert_eq!(
689 ctx.client_addr.unwrap(),
690 addr,
691 "client_addr should match assigned value"
692 );
693 }
694
695 #[test]
696 fn set_cluster() {
697 let mut ctx = default_ctx();
698 ctx.cluster = Some(Arc::from("api-cluster"));
699 assert_eq!(
700 ctx.cluster.as_deref(),
701 Some("api-cluster"),
702 "cluster should match assigned value"
703 );
704 }
705
706 #[test]
707 fn set_upstream() {
708 let mut ctx = default_ctx();
709 let upstream = Upstream {
710 address: Arc::from("10.0.0.1:80"),
711 authority: None,
712 tls: None,
713 connection: Arc::new(praxis_core::connectivity::ConnectionOptions::default()),
714 };
715 ctx.upstream = Some(upstream.clone());
716 assert_eq!(
717 &*ctx.upstream.as_ref().unwrap().address,
718 "10.0.0.1:80",
719 "upstream address should match assigned value"
720 );
721 }
722
723 #[test]
724 fn increment_retries() {
725 let mut ctx = default_ctx();
726 ctx.retries += 1;
727 ctx.retries += 1;
728 assert_eq!(ctx.retries, 2, "retries should be 2 after two increments");
729 }
730
731 #[test]
732 fn release_request_body_flag() {
733 let mut ctx = default_ctx();
734 assert!(!ctx.request_body_released, "request_body_released should start false");
735 ctx.request_body_released = true;
736 assert!(
737 ctx.request_body_released,
738 "request_body_released should be true after setting"
739 );
740 }
741
742 #[test]
743 fn release_response_body_flag() {
744 let mut ctx = default_ctx();
745 assert!(!ctx.response_body_released, "response_body_released should start false");
746 ctx.response_body_released = true;
747 assert!(
748 ctx.response_body_released,
749 "response_body_released should be true after setting"
750 );
751 }
752
753 #[test]
754 fn response_phase_done_flag() {
755 let mut ctx = default_ctx();
756 assert!(!ctx.response_phase_done, "response_phase_done should start false");
757 ctx.response_phase_done = true;
758 assert!(
759 ctx.response_phase_done,
760 "response_phase_done should be true after setting"
761 );
762 }
763
764 #[test]
765 fn set_pre_read_body() {
766 let mut ctx = default_ctx();
767 let chunks = VecDeque::from([Bytes::from_static(b"chunk1"), Bytes::from_static(b"chunk2")]);
768 ctx.pre_read_body = Some(chunks);
769 let body = ctx.pre_read_body.as_ref().unwrap();
770 assert_eq!(body.len(), 2, "pre_read_body should contain 2 chunks");
771 assert_eq!(body[0], Bytes::from_static(b"chunk1"), "first chunk should be 'chunk1'");
772 assert_eq!(
773 body[1],
774 Bytes::from_static(b"chunk2"),
775 "second chunk should be 'chunk2'"
776 );
777 }
778
779 #[test]
780 fn set_request_snapshot() {
781 let mut ctx = default_ctx();
782 let snapshot = Request {
783 method: Method::POST,
784 uri: "/api/data".parse::<Uri>().unwrap(),
785 headers: HeaderMap::new(),
786 };
787 ctx.request_snapshot = Some(snapshot);
788 let snap = ctx.request_snapshot.as_ref().unwrap();
789 assert_eq!(snap.method, Method::POST, "snapshot method should be POST");
790 assert_eq!(snap.uri.path(), "/api/data", "snapshot URI path should be /api/data");
791 }
792
793 #[test]
794 fn request_body_buffer_lifecycle() {
795 let mut ctx = default_ctx();
796 let mut buf = BodyBuffer::new(100);
797 buf.push(Bytes::from_static(b"data")).unwrap();
798 ctx.request_body_buffer = Some(buf);
799
800 assert!(
801 ctx.request_body_buffer.is_some(),
802 "buffer should be present after assignment"
803 );
804 let taken = ctx.request_body_buffer.take().unwrap();
805 assert_eq!(
806 taken.freeze(),
807 Bytes::from_static(b"data"),
808 "frozen buffer should contain pushed data"
809 );
810 assert!(ctx.request_body_buffer.is_none(), "buffer should be None after take");
811 }
812
813 #[test]
814 fn default_request_body_mode_is_stream() {
815 let ctx = default_ctx();
816 assert_eq!(
817 ctx.request_body_mode,
818 BodyMode::Stream,
819 "default request_body_mode should be Stream"
820 );
821 }
822
823 #[test]
824 fn default_response_body_mode_is_stream() {
825 let ctx = default_ctx();
826 assert_eq!(
827 ctx.response_body_mode,
828 BodyMode::Stream,
829 "default response_body_mode should be Stream"
830 );
831 }
832
833 #[test]
834 fn set_request_body_mode() {
835 let mut ctx = default_ctx();
836 ctx.request_body_mode = BodyMode::StreamBuffer { max_bytes: Some(4096) };
837 assert_eq!(
838 ctx.request_body_mode,
839 BodyMode::StreamBuffer { max_bytes: Some(4096) },
840 "request_body_mode should match assigned value"
841 );
842 }
843
844 #[test]
845 fn set_response_body_mode() {
846 let mut ctx = default_ctx();
847 ctx.response_body_mode = BodyMode::StreamBuffer { max_bytes: Some(8192) };
848 assert_eq!(
849 ctx.response_body_mode,
850 BodyMode::StreamBuffer { max_bytes: Some(8192) },
851 "response_body_mode should match assigned value"
852 );
853 }
854
855 #[test]
860 fn metadata_roundtrip_through_filter_context() {
861 let registry = FilterRegistry::with_builtins();
862 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
863 let request = Request {
864 method: Method::GET,
865 uri: "/".parse::<Uri>().unwrap(),
866 headers: HeaderMap::new(),
867 };
868
869 let mut ctx = default_ctx();
870 ctx.filter_metadata.insert("rpc.method".to_owned(), "echo".to_owned());
871 ctx.filter_metadata.insert("rpc.status".to_owned(), "ok".to_owned());
872
873 let fctx = ctx.build_filter_context(&pipeline, &request, None);
874 assert_eq!(
875 fctx.get_metadata("rpc.method"),
876 Some("echo"),
877 "metadata written before build should survive into filter context"
878 );
879 assert_eq!(
880 fctx.get_metadata("rpc.status"),
881 Some("ok"),
882 "multiple metadata keys should round-trip"
883 );
884 }
885
886 #[test]
887 fn metadata_written_in_filter_context_persists_back() {
888 let registry = FilterRegistry::with_builtins();
889 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
890 let request = Request {
891 method: Method::GET,
892 uri: "/".parse::<Uri>().unwrap(),
893 headers: HeaderMap::new(),
894 };
895
896 let mut ctx = default_ctx();
897 let mut fctx = ctx.build_filter_context(&pipeline, &request, None);
898 fctx.set_metadata("trace.id", "abc-123");
899 fctx.set_metadata("trace.span", "42");
900
901 ctx.filter_metadata = fctx.filter_metadata;
902 assert_eq!(
903 ctx.filter_metadata.get("trace.id").map(String::as_str),
904 Some("abc-123"),
905 "metadata set in filter context should persist back to protocol context"
906 );
907 assert_eq!(
908 ctx.filter_metadata.get("trace.span").map(String::as_str),
909 Some("42"),
910 "multiple metadata keys should persist back"
911 );
912 }
913
914 #[test]
919 fn pin_pipeline_captures_current_arc() {
920 let registry = FilterRegistry::with_builtins();
921 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
922 let swap = arc_swap::ArcSwap::from(Arc::clone(&pipeline_a));
923
924 let mut ctx = default_ctx();
925 let pinned = ctx.pin_pipeline(&swap);
926 assert!(
927 Arc::ptr_eq(&pinned, &pipeline_a),
928 "pin_pipeline should return the current pipeline"
929 );
930 assert!(
931 Arc::ptr_eq(ctx.pinned_pipeline.as_ref().unwrap(), &pipeline_a),
932 "pinned_pipeline should be stored in ctx"
933 );
934 }
935
936 #[test]
937 fn pipeline_returns_pinned_after_reload() {
938 let registry = FilterRegistry::with_builtins();
939 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
940 let swap = arc_swap::ArcSwap::from(Arc::clone(&pipeline_a));
941
942 let mut ctx = default_ctx();
943 ctx.pin_pipeline(&swap);
944
945 let pipeline_b = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
946 swap.store(pipeline_b);
947
948 let later = ctx.pipeline(&swap);
949 assert!(
950 Arc::ptr_eq(&later, &pipeline_a),
951 "later hooks should still return pipeline A after reload"
952 );
953 }
954
955 #[test]
956 fn new_request_after_reload_pins_new_pipeline() {
957 let registry = FilterRegistry::with_builtins();
958 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
959 let swap = arc_swap::ArcSwap::from(Arc::clone(&pipeline_a));
960
961 let mut ctx_a = default_ctx();
962 ctx_a.pin_pipeline(&swap);
963
964 let pipeline_b = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
965 swap.store(Arc::clone(&pipeline_b));
966
967 let mut ctx_b = default_ctx();
968 ctx_b.pin_pipeline(&swap);
969
970 assert!(
971 Arc::ptr_eq(&ctx_a.pipeline(&swap), &pipeline_a),
972 "request A should use old pipeline"
973 );
974 assert!(
975 Arc::ptr_eq(&ctx_b.pipeline(&swap), &pipeline_b),
976 "request B should use new pipeline"
977 );
978 }
979
980 #[test]
981 fn old_pipeline_drops_after_ctx_drops() {
982 let registry = FilterRegistry::with_builtins();
983 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
984 let weak_a = Arc::downgrade(&pipeline_a);
985 let swap = arc_swap::ArcSwap::from(pipeline_a);
986
987 let mut ctx = default_ctx();
988 ctx.pin_pipeline(&swap);
989
990 swap.store(Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap()));
991
992 assert!(weak_a.upgrade().is_some(), "old pipeline alive while ctx exists");
993
994 drop(ctx);
995
996 assert!(weak_a.upgrade().is_none(), "old pipeline drops with ctx");
997 }
998
999 #[test]
1000 fn pipeline_helper_returns_pinned_for_every_phase() {
1001 let registry = FilterRegistry::with_builtins();
1002 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1003 let swap = arc_swap::ArcSwap::from(Arc::clone(&pipeline_a));
1004
1005 let mut ctx = default_ctx();
1006 ctx.pin_pipeline(&swap);
1007
1008 swap.store(Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap()));
1009
1010 for phase in ["request_body", "response", "response_body", "logging"] {
1011 let p = ctx.pipeline(&swap);
1012 assert!(
1013 Arc::ptr_eq(&p, &pipeline_a),
1014 "{phase}: should still return pinned pipeline A"
1015 );
1016 }
1017 }
1018
1019 #[test]
1020 fn pipeline_fallback_when_not_pinned() {
1021 let registry = FilterRegistry::with_builtins();
1022 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1023 let swap = arc_swap::ArcSwap::from(Arc::clone(&pipeline_a));
1024
1025 let ctx = default_ctx();
1026
1027 let loaded = ctx.pipeline(&swap);
1028 assert!(
1029 Arc::ptr_eq(&loaded, &pipeline_a),
1030 "unpinned ctx should fall back to current ArcSwap value"
1031 );
1032 }
1033
1034 #[test]
1035 fn pin_pipeline_is_idempotent_after_reload() {
1036 let registry = FilterRegistry::with_builtins();
1037 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1038 let swap = arc_swap::ArcSwap::from(Arc::clone(&pipeline_a));
1039
1040 let mut ctx = default_ctx();
1041 ctx.pin_pipeline(&swap);
1042
1043 let pipeline_b = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1044 swap.store(pipeline_b);
1045
1046 let second_pin = ctx.pin_pipeline(&swap);
1047 assert!(
1048 Arc::ptr_eq(&second_pin, &pipeline_a),
1049 "repeated pin_pipeline after reload should return the original pin"
1050 );
1051 }
1052
1053 #[test]
1054 fn filter_state_isolated_across_pipelines_with_same_ids() {
1055 let registry = FilterRegistry::with_builtins();
1056 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1057 let pipeline_b = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1058
1059 let request = Request {
1060 method: Method::GET,
1061 uri: "/".parse::<Uri>().unwrap(),
1062 headers: HeaderMap::new(),
1063 };
1064
1065 let mut ctx_a = default_ctx();
1066 ctx_a.pinned_pipeline = Some(Arc::clone(&pipeline_a));
1067 ctx_a.request_snapshot = Some(request.clone());
1068 let mut fctx_a = ctx_a.build_filter_context(&pipeline_a, &request, None);
1069 fctx_a.current_filter_id = Some(0);
1070 fctx_a.insert_filter_state(String::from("from_pipeline_a"));
1071 ctx_a.filter_state = fctx_a.filter_state;
1072
1073 let mut ctx_b = default_ctx();
1074 ctx_b.pinned_pipeline = Some(Arc::clone(&pipeline_b));
1075 ctx_b.request_snapshot = Some(request.clone());
1076 let fctx_b = ctx_b.build_filter_context(&pipeline_b, &request, None);
1077
1078 assert!(
1079 fctx_b.filter_state.is_empty(),
1080 "request B should have its own empty state map"
1081 );
1082
1083 let fctx_a2 = ctx_a.filter_context_for(&pipeline_a, None).unwrap();
1084 assert_eq!(
1085 fctx_a2.filter_state.get(&0).and_then(|v| v.downcast_ref::<String>()),
1086 Some(&String::from("from_pipeline_a")),
1087 "request A should still see its own state in a later phase"
1088 );
1089 }
1090
1091 fn default_ctx() -> PingoraRequestCtx {
1097 PingoraRequestCtx::default()
1098 }
1099}