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