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) error_type: Option<&'static str>,
181
182 pub(crate) _active_request: Option<crate::http::pingora::metrics::ActiveRequestGuard>,
184
185 pub upstream_connect_start: Option<Instant>,
187
188 pub pre_read_body: Option<VecDeque<Bytes>>,
196
197 pub retained_pre_read_body: Option<VecDeque<Bytes>>,
206
207 pub request_body_buffer: Option<BodyBuffer>,
211
212 pub request_body_bytes: u64,
214
215 pub request_body_mode: BodyMode,
219
220 pub request_body_released: bool,
223
224 pub request_is_idempotent: bool,
226
227 pub request_snapshot: Option<Request>,
229
230 pub request_span: Span,
238
239 pub upstream_exchange_span: Span,
246
247 pub request_start: Instant,
249
250 pub response_body_buffer: Option<BodyBuffer>,
254
255 pub response_body_bytes: u64,
257
258 pub response_body_mode: BodyMode,
262
263 pub response_body_released: bool,
265
266 pub response_header_snapshot: Option<Response>,
272
273 pub upstream_response_status: Option<u16>,
276
277 pub response_phase_done: bool,
281
282 pub pending_rejection: Option<praxis_filter::Rejection>,
286
287 pub response_delivery_complete: bool,
293
294 pub retries: u32,
296
297 pub selected_endpoint_index: Option<usize>,
301
302 pub attempted_endpoints: Vec<Arc<str>>,
304
305 pub retry_policy: Option<Arc<praxis_core::config::RetryPolicy>>,
307
308 pub route_retry_policy: Option<Arc<praxis_core::config::RetryPolicy>>,
310
311 pub cluster_retry_state: Option<Arc<praxis_core::retry::ClusterRetryState>>,
313
314 pub cluster_retry_state_released: bool,
316
317 pub endpoint_reselector: Option<Arc<praxis_filter::EndpointReselector>>,
319
320 pub pending_backoff: Option<std::time::Duration>,
322
323 pub reselect_on_retry: bool,
325
326 pub rewritten_path: Option<String>,
333
334 pub upstream: Option<Upstream>,
336
337 pub upstream_for_retry: Option<Upstream>,
339
340 pub upstream_contacted: bool,
347}
348
349macro_rules! filter_context {
359 ($ctx:expr, $pipeline:expr, $request:expr, $response_header:expr) => {{
360 praxis_filter::HttpFilterContext {
361 buffered_request_body: $ctx
362 .pre_read_body
363 .as_ref()
364 .map(|chunks| chunks.front().cloned().unwrap_or_default()),
365 body_done_indices: std::mem::take(&mut $ctx.cached_body_done_indices),
366 branch_iterations: std::collections::HashMap::new(),
367 client_addr: $ctx.client_addr,
368 cluster: $ctx.cluster.take(),
369 current_filter_id: None,
370 downstream_tls: $ctx.downstream_tls,
371 metrics_route: $ctx.metrics_route.clone(),
372 peer_identity: $ctx.peer_identity.clone(),
373 extensions: std::mem::take(&mut $ctx.extensions),
374 executed_filter_indices: std::mem::take(&mut $ctx.cached_executed_filter_indices),
375 extra_request_headers: Vec::new(),
376 request_headers_to_remove: Vec::new(),
377 request_headers_to_set: Vec::new(),
378 filter_metadata: std::mem::take(&mut $ctx.filter_metadata),
379 prior_pre_read_mutations: Vec::new(),
381 pre_read_mutations: std::mem::take(&mut $ctx.pre_read_mutations),
382 structured_metadata: std::mem::take(&mut $ctx.structured_metadata),
383 filter_results: std::mem::take(&mut $ctx.filter_results),
384 filter_state: std::mem::take(&mut $ctx.filter_state),
385 health_registry: $pipeline.health_registry(),
386 id_generator: $pipeline.id_generator(),
387 kv_stores: $pipeline.kv_stores(),
388 session_stores: $pipeline.session_stores(),
389 subrequest_client: $pipeline.subrequest_client(),
390 subrequest_response_mode: praxis_filter::SubRequestResponseMode::Buffered,
391 request: $request,
392 request_body_bytes: $ctx.request_body_bytes,
393 request_body_mode: $ctx.request_body_mode,
394 request_start: $ctx.request_start,
395 response_body_bytes: $ctx.response_body_bytes,
396 response_body_mode: $ctx.response_body_mode,
397 response_header: $response_header,
398 response_headers_modified: false,
399 upstream_reached: $ctx.upstream_contacted,
400 rewritten_path: $ctx.rewritten_path.take(),
401 selected_endpoint_index: $ctx.selected_endpoint_index,
402 attempted_endpoints: std::mem::take(&mut $ctx.attempted_endpoints),
403 retry_policy: $ctx.retry_policy.clone(),
404 route_retry_policy: $ctx.route_retry_policy.clone(),
405 cluster_retry_state: $ctx.cluster_retry_state.clone(),
406 cluster_retry_state_released: $ctx.cluster_retry_state_released,
407 endpoint_reselector: $ctx.endpoint_reselector.clone(),
408 pinned_endpoint_address: None,
409 time_source: $pipeline.time_source(),
410 upstream: $ctx.upstream.take().or_else(|| $ctx.upstream_for_retry.clone()),
411 }
412 }};
413}
414
415impl PingoraRequestCtx {
416 pub fn build_filter_context<'a>(
440 &mut self,
441 pipeline: &'a FilterPipeline,
442 request: &'a Request,
443 response_header: Option<&'a mut Response>,
444 ) -> praxis_filter::HttpFilterContext<'a> {
445 filter_context!(self, pipeline, request, response_header)
446 }
447
448 pub fn filter_context_for<'a>(
475 &'a mut self,
476 pipeline: &'a FilterPipeline,
477 response_header: Option<&'a mut Response>,
478 ) -> Option<praxis_filter::HttpFilterContext<'a>> {
479 let request = self.request_snapshot.as_ref()?;
480 Some(filter_context!(self, pipeline, request, response_header))
481 }
482
483 pub fn response_body_context_for<'a>(
491 &'a mut self,
492 pipeline: &'a FilterPipeline,
493 ) -> Option<(praxis_filter::HttpFilterContext<'a>, Option<&'a Response>)> {
494 let request = self.request_snapshot.as_ref()?;
495 let response_header = self.response_header_snapshot.as_ref();
496 Some((filter_context!(self, pipeline, request, None), response_header))
497 }
498
499 pub fn pin_pipeline(&mut self, swap: &arc_swap::ArcSwap<FilterPipeline>) -> Arc<FilterPipeline> {
512 if let Some(existing) = &self.pinned_pipeline {
513 return Arc::clone(existing);
514 }
515 let pipeline = swap.load_full();
516 pipeline.prepare_extensions(&mut self.extensions);
519 self.pinned_pipeline = Some(Arc::clone(&pipeline));
520 pipeline
521 }
522
523 pub(crate) fn stamp_error_type(&mut self, error_type: &'static str) {
529 self.error_type.get_or_insert(error_type);
530 }
531
532 pub fn pipeline(&self, swap: &arc_swap::ArcSwap<FilterPipeline>) -> Arc<FilterPipeline> {
545 self.pinned_pipeline
546 .as_ref()
547 .map_or_else(|| swap.load_full(), Arc::clone)
548 }
549}
550
551impl Default for PingoraRequestCtx {
552 #[expect(
553 clippy::too_many_lines,
554 reason = "context default enumerates all lifecycle fields explicitly"
555 )]
556 fn default() -> Self {
557 Self {
558 _connection_permit: None,
559 _global_connection_permit: None,
560 cached_body_done_indices: Vec::new(),
561 cached_executed_filter_indices: Vec::new(),
562 client_addr: None,
563 client_http_version: None,
564 cluster: None,
565 connection_upgraded: false,
566 downstream_tls: false,
567 peer_identity: None,
568 extensions: praxis_filter::RequestExtensions::new(),
569 filter_metadata: std::collections::HashMap::new(),
570 pre_read_mutations: Vec::new(),
571 structured_metadata: std::collections::HashMap::new(),
572 mutated_request_body_len: None,
573 pinned_pipeline: None,
574 filter_results: std::collections::HashMap::new(),
575 filter_state: std::collections::HashMap::new(),
576 metrics_cluster: None,
577 metrics_cluster_shared: None,
578 metrics_route: None,
579 error_type: None,
580 _active_request: None,
581 upstream_connect_start: None,
582 pre_read_body: None,
583 retained_pre_read_body: None,
584 request_body_buffer: None,
585 request_body_bytes: 0,
586 request_body_mode: BodyMode::Stream,
587 request_body_released: false,
588 request_is_idempotent: false,
589 request_snapshot: None,
590 request_span: Span::none(),
591 request_start: Instant::now(),
592 upstream_exchange_span: Span::none(),
593 response_body_buffer: None,
594 response_body_bytes: 0,
595 response_body_mode: BodyMode::Stream,
596 response_body_released: false,
597 response_header_snapshot: None,
598 upstream_response_status: None,
599 response_phase_done: false,
600 response_delivery_complete: false,
601 pending_rejection: None,
602 retries: 0,
603 rewritten_path: None,
604 selected_endpoint_index: None,
605 attempted_endpoints: Vec::new(),
606 retry_policy: None,
607 route_retry_policy: None,
608 cluster_retry_state: None,
609 cluster_retry_state_released: false,
610 endpoint_reselector: None,
611 pending_backoff: None,
612 reselect_on_retry: false,
613 upstream: None,
614 upstream_for_retry: None,
615 upstream_contacted: false,
616 }
617 }
618}
619
620#[cfg(test)]
625#[expect(clippy::allow_attributes, reason = "blanket test suppressions")]
626#[allow(
627 clippy::unwrap_used,
628 clippy::expect_used,
629 clippy::indexing_slicing,
630 clippy::significant_drop_tightening,
631 clippy::too_many_lines,
632 reason = "tests"
633)]
634mod tests {
635 use std::net::Ipv4Addr;
636
637 use http::{HeaderMap, Method, Uri};
638 use praxis_filter::FilterRegistry;
639
640 use super::*;
641
642 #[test]
643 fn default_state_has_no_client_addr() {
644 let ctx = default_ctx();
645 assert!(ctx.client_addr.is_none(), "default client_addr should be None");
646 }
647
648 #[test]
649 fn default_state_has_no_cluster() {
650 let ctx = default_ctx();
651 assert!(ctx.cluster.is_none(), "default cluster should be None");
652 }
653
654 #[test]
655 fn default_state_has_zero_retries() {
656 let ctx = default_ctx();
657 assert_eq!(ctx.retries, 0, "default retries should be zero");
658 }
659
660 #[test]
661 fn default_state_flags_are_false() {
662 let ctx = default_ctx();
663 assert!(
664 !ctx.request_body_released,
665 "default request_body_released should be false"
666 );
667 assert!(
668 !ctx.response_body_released,
669 "default response_body_released should be false"
670 );
671 assert!(
672 !ctx.request_is_idempotent,
673 "default request_is_idempotent should be false"
674 );
675 assert!(!ctx.response_phase_done, "default response_phase_done should be false");
676 }
677
678 #[test]
679 fn default_state_buffers_are_none() {
680 let ctx = default_ctx();
681 assert!(
682 ctx.request_body_buffer.is_none(),
683 "default request_body_buffer should be None"
684 );
685 assert!(
686 ctx.response_body_buffer.is_none(),
687 "default response_body_buffer should be None"
688 );
689 assert!(ctx.pre_read_body.is_none(), "default pre_read_body should be None");
690 }
691
692 #[test]
693 fn default_state_request_span_is_disabled() {
694 let ctx = default_ctx();
695 assert!(
696 ctx.request_span.is_disabled(),
697 "default request_span should be a disabled (none) span"
698 );
699 }
700
701 #[test]
702 fn default_state_upstream_exchange_span_is_disabled() {
703 let ctx = default_ctx();
704 assert!(
705 ctx.upstream_exchange_span.is_disabled(),
706 "default upstream_exchange_span should be a disabled (none) span"
707 );
708 }
709
710 #[test]
711 fn default_state_snapshots_are_none() {
712 let ctx = default_ctx();
713 assert!(
714 ctx.request_snapshot.is_none(),
715 "default request_snapshot should be None"
716 );
717 assert!(ctx.upstream.is_none(), "default upstream should be None");
718 assert!(
719 ctx.upstream_for_retry.is_none(),
720 "default upstream_for_retry should be None"
721 );
722 }
723
724 #[test]
725 fn set_client_addr() {
726 let mut ctx = default_ctx();
727 let addr = IpAddr::V4(Ipv4Addr::new(192, 168, 1, 100));
728 ctx.client_addr = Some(addr);
729 assert_eq!(
730 ctx.client_addr.unwrap(),
731 addr,
732 "client_addr should match assigned value"
733 );
734 }
735
736 #[test]
737 fn set_cluster() {
738 let mut ctx = default_ctx();
739 ctx.cluster = Some(Arc::from("api-cluster"));
740 assert_eq!(
741 ctx.cluster.as_deref(),
742 Some("api-cluster"),
743 "cluster should match assigned value"
744 );
745 }
746
747 #[test]
748 fn set_upstream() {
749 let mut ctx = default_ctx();
750 let upstream = Upstream {
751 address: Arc::from("10.0.0.1:80"),
752 authority: None,
753 tls: None,
754 connection: Arc::new(praxis_core::connectivity::ConnectionOptions::default()),
755 };
756 ctx.upstream = Some(upstream.clone());
757 assert_eq!(
758 &*ctx.upstream.as_ref().unwrap().address,
759 "10.0.0.1:80",
760 "upstream address should match assigned value"
761 );
762 }
763
764 #[test]
765 fn increment_retries() {
766 let mut ctx = default_ctx();
767 ctx.retries += 1;
768 ctx.retries += 1;
769 assert_eq!(ctx.retries, 2, "retries should be 2 after two increments");
770 }
771
772 #[test]
773 fn release_request_body_flag() {
774 let mut ctx = default_ctx();
775 assert!(!ctx.request_body_released, "request_body_released should start false");
776 ctx.request_body_released = true;
777 assert!(
778 ctx.request_body_released,
779 "request_body_released should be true after setting"
780 );
781 }
782
783 #[test]
784 fn release_response_body_flag() {
785 let mut ctx = default_ctx();
786 assert!(!ctx.response_body_released, "response_body_released should start false");
787 ctx.response_body_released = true;
788 assert!(
789 ctx.response_body_released,
790 "response_body_released should be true after setting"
791 );
792 }
793
794 #[test]
795 fn response_phase_done_flag() {
796 let mut ctx = default_ctx();
797 assert!(!ctx.response_phase_done, "response_phase_done should start false");
798 ctx.response_phase_done = true;
799 assert!(
800 ctx.response_phase_done,
801 "response_phase_done should be true after setting"
802 );
803 }
804
805 #[test]
806 fn set_pre_read_body() {
807 let mut ctx = default_ctx();
808 let chunks = VecDeque::from([Bytes::from_static(b"chunk1"), Bytes::from_static(b"chunk2")]);
809 ctx.pre_read_body = Some(chunks);
810 let body = ctx.pre_read_body.as_ref().unwrap();
811 assert_eq!(body.len(), 2, "pre_read_body should contain 2 chunks");
812 assert_eq!(body[0], Bytes::from_static(b"chunk1"), "first chunk should be 'chunk1'");
813 assert_eq!(
814 body[1],
815 Bytes::from_static(b"chunk2"),
816 "second chunk should be 'chunk2'"
817 );
818 }
819
820 #[test]
821 fn set_request_snapshot() {
822 let mut ctx = default_ctx();
823 let snapshot = Request {
824 method: Method::POST,
825 uri: "/api/data".parse::<Uri>().unwrap(),
826 headers: HeaderMap::new(),
827 };
828 ctx.request_snapshot = Some(snapshot);
829 let snap = ctx.request_snapshot.as_ref().unwrap();
830 assert_eq!(snap.method, Method::POST, "snapshot method should be POST");
831 assert_eq!(snap.uri.path(), "/api/data", "snapshot URI path should be /api/data");
832 }
833
834 #[test]
835 fn request_body_buffer_lifecycle() {
836 let mut ctx = default_ctx();
837 let mut buf = BodyBuffer::new(100);
838 buf.push(Bytes::from_static(b"data")).unwrap();
839 ctx.request_body_buffer = Some(buf);
840
841 assert!(
842 ctx.request_body_buffer.is_some(),
843 "buffer should be present after assignment"
844 );
845 let taken = ctx.request_body_buffer.take().unwrap();
846 assert_eq!(
847 taken.freeze(),
848 Bytes::from_static(b"data"),
849 "frozen buffer should contain pushed data"
850 );
851 assert!(ctx.request_body_buffer.is_none(), "buffer should be None after take");
852 }
853
854 #[test]
855 fn default_request_body_mode_is_stream() {
856 let ctx = default_ctx();
857 assert_eq!(
858 ctx.request_body_mode,
859 BodyMode::Stream,
860 "default request_body_mode should be Stream"
861 );
862 }
863
864 #[test]
865 fn default_response_body_mode_is_stream() {
866 let ctx = default_ctx();
867 assert_eq!(
868 ctx.response_body_mode,
869 BodyMode::Stream,
870 "default response_body_mode should be Stream"
871 );
872 }
873
874 #[test]
875 fn set_request_body_mode() {
876 let mut ctx = default_ctx();
877 ctx.request_body_mode = BodyMode::StreamBuffer { max_bytes: Some(4096) };
878 assert_eq!(
879 ctx.request_body_mode,
880 BodyMode::StreamBuffer { max_bytes: Some(4096) },
881 "request_body_mode should match assigned value"
882 );
883 }
884
885 #[test]
886 fn set_response_body_mode() {
887 let mut ctx = default_ctx();
888 ctx.response_body_mode = BodyMode::StreamBuffer { max_bytes: Some(8192) };
889 assert_eq!(
890 ctx.response_body_mode,
891 BodyMode::StreamBuffer { max_bytes: Some(8192) },
892 "response_body_mode should match assigned value"
893 );
894 }
895
896 #[test]
901 fn metadata_roundtrip_through_filter_context() {
902 let registry = FilterRegistry::with_builtins();
903 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
904 let request = Request {
905 method: Method::GET,
906 uri: "/".parse::<Uri>().unwrap(),
907 headers: HeaderMap::new(),
908 };
909
910 let mut ctx = default_ctx();
911 ctx.filter_metadata.insert("rpc.method".to_owned(), "echo".to_owned());
912 ctx.filter_metadata.insert("rpc.status".to_owned(), "ok".to_owned());
913
914 let fctx = ctx.build_filter_context(&pipeline, &request, None);
915 assert_eq!(
916 fctx.get_metadata("rpc.method"),
917 Some("echo"),
918 "metadata written before build should survive into filter context"
919 );
920 assert_eq!(
921 fctx.get_metadata("rpc.status"),
922 Some("ok"),
923 "multiple metadata keys should round-trip"
924 );
925 }
926
927 #[test]
928 fn metadata_written_in_filter_context_persists_back() {
929 let registry = FilterRegistry::with_builtins();
930 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
931 let request = Request {
932 method: Method::GET,
933 uri: "/".parse::<Uri>().unwrap(),
934 headers: HeaderMap::new(),
935 };
936
937 let mut ctx = default_ctx();
938 let mut fctx = ctx.build_filter_context(&pipeline, &request, None);
939 fctx.set_metadata("trace.id", "abc-123");
940 fctx.set_metadata("trace.span", "42");
941
942 ctx.filter_metadata = fctx.filter_metadata;
943 assert_eq!(
944 ctx.filter_metadata.get("trace.id").map(String::as_str),
945 Some("abc-123"),
946 "metadata set in filter context should persist back to protocol context"
947 );
948 assert_eq!(
949 ctx.filter_metadata.get("trace.span").map(String::as_str),
950 Some("42"),
951 "multiple metadata keys should persist back"
952 );
953 }
954
955 #[test]
960 fn pin_pipeline_captures_current_arc() {
961 let registry = FilterRegistry::with_builtins();
962 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
963 let swap = arc_swap::ArcSwap::from(Arc::clone(&pipeline_a));
964
965 let mut ctx = default_ctx();
966 let pinned = ctx.pin_pipeline(&swap);
967 assert!(
968 Arc::ptr_eq(&pinned, &pipeline_a),
969 "pin_pipeline should return the current pipeline"
970 );
971 assert!(
972 Arc::ptr_eq(ctx.pinned_pipeline.as_ref().unwrap(), &pipeline_a),
973 "pinned_pipeline should be stored in ctx"
974 );
975 }
976
977 #[test]
978 fn stamp_error_type_is_first_write_wins() {
979 let mut ctx = PingoraRequestCtx::default();
986 ctx.stamp_error_type(crate::http::pingora::metrics::ERROR_TYPE_FILTER_REJECT);
987 ctx.stamp_error_type(crate::http::pingora::metrics::ERROR_TYPE_INTERNAL);
988 assert_eq!(
989 ctx.error_type,
990 Some(crate::http::pingora::metrics::ERROR_TYPE_FILTER_REJECT),
991 "the first classification wins; a later site must not overwrite it"
992 );
993 }
994
995 #[test]
996 fn error_type_is_unset_until_stamped() {
997 let ctx = PingoraRequestCtx::default();
998 assert_eq!(ctx.error_type, None, "a fresh context carries no error cause");
999 }
1000
1001 #[test]
1002 fn pipeline_returns_pinned_after_reload() {
1003 let registry = FilterRegistry::with_builtins();
1004 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1005 let swap = arc_swap::ArcSwap::from(Arc::clone(&pipeline_a));
1006
1007 let mut ctx = default_ctx();
1008 ctx.pin_pipeline(&swap);
1009
1010 let pipeline_b = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1011 swap.store(pipeline_b);
1012
1013 let later = ctx.pipeline(&swap);
1014 assert!(
1015 Arc::ptr_eq(&later, &pipeline_a),
1016 "later hooks should still return pipeline A after reload"
1017 );
1018 }
1019
1020 #[test]
1021 fn new_request_after_reload_pins_new_pipeline() {
1022 let registry = FilterRegistry::with_builtins();
1023 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1024 let swap = arc_swap::ArcSwap::from(Arc::clone(&pipeline_a));
1025
1026 let mut ctx_a = default_ctx();
1027 ctx_a.pin_pipeline(&swap);
1028
1029 let pipeline_b = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1030 swap.store(Arc::clone(&pipeline_b));
1031
1032 let mut ctx_b = default_ctx();
1033 ctx_b.pin_pipeline(&swap);
1034
1035 assert!(
1036 Arc::ptr_eq(&ctx_a.pipeline(&swap), &pipeline_a),
1037 "request A should use old pipeline"
1038 );
1039 assert!(
1040 Arc::ptr_eq(&ctx_b.pipeline(&swap), &pipeline_b),
1041 "request B should use new pipeline"
1042 );
1043 }
1044
1045 #[test]
1046 fn old_pipeline_drops_after_ctx_drops() {
1047 let registry = FilterRegistry::with_builtins();
1048 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1049 let weak_a = Arc::downgrade(&pipeline_a);
1050 let swap = arc_swap::ArcSwap::from(pipeline_a);
1051
1052 let mut ctx = default_ctx();
1053 ctx.pin_pipeline(&swap);
1054
1055 swap.store(Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap()));
1056
1057 assert!(weak_a.upgrade().is_some(), "old pipeline alive while ctx exists");
1058
1059 drop(ctx);
1060
1061 assert!(weak_a.upgrade().is_none(), "old pipeline drops with ctx");
1062 }
1063
1064 #[test]
1065 fn pipeline_helper_returns_pinned_for_every_phase() {
1066 let registry = FilterRegistry::with_builtins();
1067 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1068 let swap = arc_swap::ArcSwap::from(Arc::clone(&pipeline_a));
1069
1070 let mut ctx = default_ctx();
1071 ctx.pin_pipeline(&swap);
1072
1073 swap.store(Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap()));
1074
1075 for phase in ["request_body", "response", "response_body", "logging"] {
1076 let p = ctx.pipeline(&swap);
1077 assert!(
1078 Arc::ptr_eq(&p, &pipeline_a),
1079 "{phase}: should still return pinned pipeline A"
1080 );
1081 }
1082 }
1083
1084 #[test]
1085 fn pipeline_fallback_when_not_pinned() {
1086 let registry = FilterRegistry::with_builtins();
1087 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1088 let swap = arc_swap::ArcSwap::from(Arc::clone(&pipeline_a));
1089
1090 let ctx = default_ctx();
1091
1092 let loaded = ctx.pipeline(&swap);
1093 assert!(
1094 Arc::ptr_eq(&loaded, &pipeline_a),
1095 "unpinned ctx should fall back to current ArcSwap value"
1096 );
1097 }
1098
1099 #[test]
1100 fn pin_pipeline_is_idempotent_after_reload() {
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 mut ctx = default_ctx();
1106 ctx.pin_pipeline(&swap);
1107
1108 let pipeline_b = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1109 swap.store(pipeline_b);
1110
1111 let second_pin = ctx.pin_pipeline(&swap);
1112 assert!(
1113 Arc::ptr_eq(&second_pin, &pipeline_a),
1114 "repeated pin_pipeline after reload should return the original pin"
1115 );
1116 }
1117
1118 #[test]
1119 fn filter_state_isolated_across_pipelines_with_same_ids() {
1120 let registry = FilterRegistry::with_builtins();
1121 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1122 let pipeline_b = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
1123
1124 let request = Request {
1125 method: Method::GET,
1126 uri: "/".parse::<Uri>().unwrap(),
1127 headers: HeaderMap::new(),
1128 };
1129
1130 let mut ctx_a = default_ctx();
1131 ctx_a.pinned_pipeline = Some(Arc::clone(&pipeline_a));
1132 ctx_a.request_snapshot = Some(request.clone());
1133 let mut fctx_a = ctx_a.build_filter_context(&pipeline_a, &request, None);
1134 fctx_a.current_filter_id = Some(0);
1135 fctx_a.insert_filter_state(String::from("from_pipeline_a"));
1136 ctx_a.filter_state = fctx_a.filter_state;
1137
1138 let mut ctx_b = default_ctx();
1139 ctx_b.pinned_pipeline = Some(Arc::clone(&pipeline_b));
1140 ctx_b.request_snapshot = Some(request.clone());
1141 let fctx_b = ctx_b.build_filter_context(&pipeline_b, &request, None);
1142
1143 assert!(
1144 fctx_b.filter_state.is_empty(),
1145 "request B should have its own empty state map"
1146 );
1147
1148 let fctx_a2 = ctx_a.filter_context_for(&pipeline_a, None).unwrap();
1149 assert_eq!(
1150 fctx_a2.filter_state.get(&0).and_then(|v| v.downcast_ref::<String>()),
1151 Some(&String::from("from_pipeline_a")),
1152 "request A should still see its own state in a later phase"
1153 );
1154 }
1155
1156 fn default_ctx() -> PingoraRequestCtx {
1162 PingoraRequestCtx::default()
1163 }
1164}