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;
12
13#[expect(clippy::struct_excessive_bools, reason = "lifecycle flags")]
29pub struct PingoraRequestCtx {
30 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>,
60
61 pub cached_executed_filter_indices: Vec<bool>,
66
67 pub downstream_tls: bool,
74
75 pub peer_identity: Option<praxis_tls::TlsPeerIdentity>,
83
84 pub connection_upgraded: bool,
90
91 pub extensions: praxis_filter::RequestExtensions,
98
99 pub filter_metadata: std::collections::HashMap<String, String>,
105
106 pub pre_read_mutations: Vec<TrustedHeaderMutation>,
112
113 pub structured_metadata: std::collections::HashMap<String, serde_json::Value>,
119
120 pub mutated_request_body_len: Option<usize>,
126
127 pub pinned_pipeline: Option<Arc<FilterPipeline>>,
137
138 pub filter_results: std::collections::HashMap<&'static str, praxis_filter::FilterResultSet>,
144
145 pub filter_state: std::collections::HashMap<usize, Box<dyn std::any::Any + Send + Sync>>,
155
156 pub metrics_cluster: Option<Arc<str>>,
160
161 pub metrics_cluster_shared: Option<::metrics::SharedString>,
168
169 pub pre_read_body: Option<VecDeque<Bytes>>,
177
178 pub request_body_buffer: Option<BodyBuffer>,
182
183 pub request_body_bytes: u64,
185
186 pub request_body_mode: BodyMode,
190
191 pub request_body_released: bool,
194
195 pub request_is_idempotent: bool,
197
198 pub request_snapshot: Option<Request>,
200
201 pub request_start: Instant,
203
204 pub response_body_buffer: Option<BodyBuffer>,
208
209 pub response_body_bytes: u64,
211
212 pub response_body_mode: BodyMode,
216
217 pub response_body_released: bool,
219
220 pub response_header_snapshot: Option<Response>,
226
227 pub upstream_response_status: Option<u16>,
230
231 pub response_phase_done: bool,
235
236 pub retries: u32,
238
239 pub selected_endpoint_index: Option<usize>,
243
244 pub rewritten_path: Option<String>,
251
252 pub upstream: Option<Upstream>,
254
255 pub upstream_for_retry: Option<Upstream>,
257}
258
259macro_rules! filter_context {
269 ($ctx:expr, $pipeline:expr, $request:expr, $response_header:expr) => {{
270 $pipeline.prepare_extensions(&mut $ctx.extensions);
271 praxis_filter::HttpFilterContext {
272 buffered_request_body: $ctx
273 .pre_read_body
274 .as_ref()
275 .map(|chunks| chunks.front().cloned().unwrap_or_default()),
276 body_done_indices: std::mem::take(&mut $ctx.cached_body_done_indices),
277 branch_iterations: std::collections::HashMap::new(),
278 client_addr: $ctx.client_addr,
279 cluster: $ctx.cluster.take(),
280 current_filter_id: None,
281 downstream_tls: $ctx.downstream_tls,
282 extensions: std::mem::take(&mut $ctx.extensions),
283 executed_filter_indices: std::mem::take(&mut $ctx.cached_executed_filter_indices),
284 extra_request_headers: Vec::new(),
285 request_headers_to_remove: Vec::new(),
286 request_headers_to_set: Vec::new(),
287 filter_metadata: std::mem::take(&mut $ctx.filter_metadata),
288 pre_read_mutations: std::mem::take(&mut $ctx.pre_read_mutations),
289 structured_metadata: std::mem::take(&mut $ctx.structured_metadata),
290 filter_results: std::mem::take(&mut $ctx.filter_results),
291 filter_state: std::mem::take(&mut $ctx.filter_state),
292 health_registry: $pipeline.health_registry(),
293 peer_identity: $ctx.peer_identity.clone(),
294 id_generator: $pipeline.id_generator(),
295 kv_stores: $pipeline.kv_stores(),
296 subrequest_connector: $pipeline.subrequest_connector(),
297 request: $request,
298 request_body_bytes: $ctx.request_body_bytes,
299 request_body_mode: $ctx.request_body_mode,
300 request_start: $ctx.request_start,
301 response_body_bytes: $ctx.response_body_bytes,
302 response_body_mode: $ctx.response_body_mode,
303 response_header: $response_header,
304 response_headers_modified: false,
305 rewritten_path: $ctx.rewritten_path.take(),
306 selected_endpoint_index: $ctx.selected_endpoint_index,
307 time_source: $pipeline.time_source(),
308 upstream: $ctx.upstream.take(),
309 }
310 }};
311}
312
313impl PingoraRequestCtx {
314 pub fn build_filter_context<'a>(
338 &mut self,
339 pipeline: &'a FilterPipeline,
340 request: &'a Request,
341 response_header: Option<&'a mut Response>,
342 ) -> praxis_filter::HttpFilterContext<'a> {
343 filter_context!(self, pipeline, request, response_header)
344 }
345
346 pub fn filter_context_for<'a>(
373 &'a mut self,
374 pipeline: &'a FilterPipeline,
375 response_header: Option<&'a mut Response>,
376 ) -> Option<praxis_filter::HttpFilterContext<'a>> {
377 let request = self.request_snapshot.as_ref()?;
378 Some(filter_context!(self, pipeline, request, response_header))
379 }
380
381 pub fn response_body_context_for<'a>(
389 &'a mut self,
390 pipeline: &'a FilterPipeline,
391 ) -> Option<(praxis_filter::HttpFilterContext<'a>, Option<&'a Response>)> {
392 let request = self.request_snapshot.as_ref()?;
393 let response_header = self.response_header_snapshot.as_ref();
394 Some((filter_context!(self, pipeline, request, None), response_header))
395 }
396
397 pub fn pin_pipeline(&mut self, swap: &arc_swap::ArcSwap<FilterPipeline>) -> Arc<FilterPipeline> {
410 if let Some(existing) = &self.pinned_pipeline {
411 return Arc::clone(existing);
412 }
413 let pipeline = swap.load_full();
414 self.pinned_pipeline = Some(Arc::clone(&pipeline));
415 pipeline
416 }
417
418 pub fn pipeline(&self, swap: &arc_swap::ArcSwap<FilterPipeline>) -> Arc<FilterPipeline> {
431 self.pinned_pipeline
432 .as_ref()
433 .map_or_else(|| swap.load_full(), Arc::clone)
434 }
435}
436
437impl Default for PingoraRequestCtx {
438 #[expect(
439 clippy::too_many_lines,
440 reason = "context default enumerates all lifecycle fields explicitly"
441 )]
442 fn default() -> Self {
443 Self {
444 _connection_permit: None,
445 _global_connection_permit: None,
446 cached_body_done_indices: Vec::new(),
447 cached_executed_filter_indices: Vec::new(),
448 client_addr: None,
449 client_http_version: None,
450 cluster: None,
451 connection_upgraded: false,
452 downstream_tls: false,
453 peer_identity: None,
454 extensions: praxis_filter::RequestExtensions::new(),
455 filter_metadata: std::collections::HashMap::new(),
456 pre_read_mutations: Vec::new(),
457 structured_metadata: std::collections::HashMap::new(),
458 mutated_request_body_len: None,
459 pinned_pipeline: None,
460 filter_results: std::collections::HashMap::new(),
461 filter_state: std::collections::HashMap::new(),
462 metrics_cluster: None,
463 metrics_cluster_shared: None,
464 pre_read_body: None,
465 request_body_buffer: None,
466 request_body_bytes: 0,
467 request_body_mode: BodyMode::Stream,
468 request_body_released: false,
469 request_is_idempotent: false,
470 request_snapshot: None,
471 request_start: Instant::now(),
472 response_body_buffer: None,
473 response_body_bytes: 0,
474 response_body_mode: BodyMode::Stream,
475 response_body_released: false,
476 response_header_snapshot: None,
477 upstream_response_status: None,
478 response_phase_done: false,
479 retries: 0,
480 rewritten_path: None,
481 selected_endpoint_index: None,
482 upstream: None,
483 upstream_for_retry: None,
484 }
485 }
486}
487
488#[cfg(test)]
493#[expect(clippy::allow_attributes, reason = "blanket test suppressions")]
494#[allow(
495 clippy::unwrap_used,
496 clippy::expect_used,
497 clippy::indexing_slicing,
498 clippy::significant_drop_tightening,
499 clippy::too_many_lines,
500 reason = "tests"
501)]
502mod tests {
503 use std::{
504 collections::VecDeque,
505 net::{IpAddr, Ipv4Addr},
506 sync::Arc,
507 };
508
509 use bytes::Bytes;
510 use http::{HeaderMap, Method, Uri};
511 use praxis_core::connectivity::Upstream;
512 use praxis_filter::{BodyBuffer, BodyMode, FilterPipeline, FilterRegistry};
513
514 use super::*;
515
516 #[test]
517 fn default_state_has_no_client_addr() {
518 let ctx = default_ctx();
519 assert!(ctx.client_addr.is_none(), "default client_addr should be None");
520 }
521
522 #[test]
523 fn default_state_has_no_cluster() {
524 let ctx = default_ctx();
525 assert!(ctx.cluster.is_none(), "default cluster should be None");
526 }
527
528 #[test]
529 fn default_state_has_zero_retries() {
530 let ctx = default_ctx();
531 assert_eq!(ctx.retries, 0, "default retries should be zero");
532 }
533
534 #[test]
535 fn default_state_flags_are_false() {
536 let ctx = default_ctx();
537 assert!(
538 !ctx.request_body_released,
539 "default request_body_released should be false"
540 );
541 assert!(
542 !ctx.response_body_released,
543 "default response_body_released should be false"
544 );
545 assert!(
546 !ctx.request_is_idempotent,
547 "default request_is_idempotent should be false"
548 );
549 assert!(!ctx.response_phase_done, "default response_phase_done should be false");
550 }
551
552 #[test]
553 fn default_state_buffers_are_none() {
554 let ctx = default_ctx();
555 assert!(
556 ctx.request_body_buffer.is_none(),
557 "default request_body_buffer should be None"
558 );
559 assert!(
560 ctx.response_body_buffer.is_none(),
561 "default response_body_buffer should be None"
562 );
563 assert!(ctx.pre_read_body.is_none(), "default pre_read_body should be None");
564 }
565
566 #[test]
567 fn default_state_snapshots_are_none() {
568 let ctx = default_ctx();
569 assert!(
570 ctx.request_snapshot.is_none(),
571 "default request_snapshot should be None"
572 );
573 assert!(ctx.upstream.is_none(), "default upstream should be None");
574 assert!(
575 ctx.upstream_for_retry.is_none(),
576 "default upstream_for_retry should be None"
577 );
578 }
579
580 #[test]
581 fn set_client_addr() {
582 let mut ctx = default_ctx();
583 let addr = IpAddr::V4(Ipv4Addr::new(192, 168, 1, 100));
584 ctx.client_addr = Some(addr);
585 assert_eq!(
586 ctx.client_addr.unwrap(),
587 addr,
588 "client_addr should match assigned value"
589 );
590 }
591
592 #[test]
593 fn set_cluster() {
594 let mut ctx = default_ctx();
595 ctx.cluster = Some(Arc::from("api-cluster"));
596 assert_eq!(
597 ctx.cluster.as_deref(),
598 Some("api-cluster"),
599 "cluster should match assigned value"
600 );
601 }
602
603 #[test]
604 fn set_upstream() {
605 let mut ctx = default_ctx();
606 let upstream = Upstream {
607 address: Arc::from("10.0.0.1:80"),
608 tls: None,
609 connection: Arc::new(praxis_core::connectivity::ConnectionOptions::default()),
610 };
611 ctx.upstream = Some(upstream.clone());
612 assert_eq!(
613 &*ctx.upstream.as_ref().unwrap().address,
614 "10.0.0.1:80",
615 "upstream address should match assigned value"
616 );
617 }
618
619 #[test]
620 fn increment_retries() {
621 let mut ctx = default_ctx();
622 ctx.retries += 1;
623 ctx.retries += 1;
624 assert_eq!(ctx.retries, 2, "retries should be 2 after two increments");
625 }
626
627 #[test]
628 fn release_request_body_flag() {
629 let mut ctx = default_ctx();
630 assert!(!ctx.request_body_released, "request_body_released should start false");
631 ctx.request_body_released = true;
632 assert!(
633 ctx.request_body_released,
634 "request_body_released should be true after setting"
635 );
636 }
637
638 #[test]
639 fn release_response_body_flag() {
640 let mut ctx = default_ctx();
641 assert!(!ctx.response_body_released, "response_body_released should start false");
642 ctx.response_body_released = true;
643 assert!(
644 ctx.response_body_released,
645 "response_body_released should be true after setting"
646 );
647 }
648
649 #[test]
650 fn response_phase_done_flag() {
651 let mut ctx = default_ctx();
652 assert!(!ctx.response_phase_done, "response_phase_done should start false");
653 ctx.response_phase_done = true;
654 assert!(
655 ctx.response_phase_done,
656 "response_phase_done should be true after setting"
657 );
658 }
659
660 #[test]
661 fn set_pre_read_body() {
662 let mut ctx = default_ctx();
663 let chunks = VecDeque::from([Bytes::from_static(b"chunk1"), Bytes::from_static(b"chunk2")]);
664 ctx.pre_read_body = Some(chunks);
665 let body = ctx.pre_read_body.as_ref().unwrap();
666 assert_eq!(body.len(), 2, "pre_read_body should contain 2 chunks");
667 assert_eq!(body[0], Bytes::from_static(b"chunk1"), "first chunk should be 'chunk1'");
668 assert_eq!(
669 body[1],
670 Bytes::from_static(b"chunk2"),
671 "second chunk should be 'chunk2'"
672 );
673 }
674
675 #[test]
676 fn set_request_snapshot() {
677 let mut ctx = default_ctx();
678 let snapshot = Request {
679 method: Method::POST,
680 uri: "/api/data".parse::<Uri>().unwrap(),
681 headers: HeaderMap::new(),
682 };
683 ctx.request_snapshot = Some(snapshot);
684 let snap = ctx.request_snapshot.as_ref().unwrap();
685 assert_eq!(snap.method, Method::POST, "snapshot method should be POST");
686 assert_eq!(snap.uri.path(), "/api/data", "snapshot URI path should be /api/data");
687 }
688
689 #[test]
690 fn request_body_buffer_lifecycle() {
691 let mut ctx = default_ctx();
692 let mut buf = BodyBuffer::new(100);
693 buf.push(Bytes::from_static(b"data")).unwrap();
694 ctx.request_body_buffer = Some(buf);
695
696 assert!(
697 ctx.request_body_buffer.is_some(),
698 "buffer should be present after assignment"
699 );
700 let taken = ctx.request_body_buffer.take().unwrap();
701 assert_eq!(
702 taken.freeze(),
703 Bytes::from_static(b"data"),
704 "frozen buffer should contain pushed data"
705 );
706 assert!(ctx.request_body_buffer.is_none(), "buffer should be None after take");
707 }
708
709 #[test]
710 fn default_request_body_mode_is_stream() {
711 let ctx = default_ctx();
712 assert_eq!(
713 ctx.request_body_mode,
714 BodyMode::Stream,
715 "default request_body_mode should be Stream"
716 );
717 }
718
719 #[test]
720 fn default_response_body_mode_is_stream() {
721 let ctx = default_ctx();
722 assert_eq!(
723 ctx.response_body_mode,
724 BodyMode::Stream,
725 "default response_body_mode should be Stream"
726 );
727 }
728
729 #[test]
730 fn set_request_body_mode() {
731 let mut ctx = default_ctx();
732 ctx.request_body_mode = BodyMode::StreamBuffer { max_bytes: Some(4096) };
733 assert_eq!(
734 ctx.request_body_mode,
735 BodyMode::StreamBuffer { max_bytes: Some(4096) },
736 "request_body_mode should match assigned value"
737 );
738 }
739
740 #[test]
741 fn set_response_body_mode() {
742 let mut ctx = default_ctx();
743 ctx.response_body_mode = BodyMode::StreamBuffer { max_bytes: Some(8192) };
744 assert_eq!(
745 ctx.response_body_mode,
746 BodyMode::StreamBuffer { max_bytes: Some(8192) },
747 "response_body_mode should match assigned value"
748 );
749 }
750
751 #[test]
756 fn metadata_roundtrip_through_filter_context() {
757 let registry = FilterRegistry::with_builtins();
758 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
759 let request = Request {
760 method: Method::GET,
761 uri: "/".parse::<Uri>().unwrap(),
762 headers: HeaderMap::new(),
763 };
764
765 let mut ctx = default_ctx();
766 ctx.filter_metadata.insert("rpc.method".to_owned(), "echo".to_owned());
767 ctx.filter_metadata.insert("rpc.status".to_owned(), "ok".to_owned());
768
769 let fctx = ctx.build_filter_context(&pipeline, &request, None);
770 assert_eq!(
771 fctx.get_metadata("rpc.method"),
772 Some("echo"),
773 "metadata written before build should survive into filter context"
774 );
775 assert_eq!(
776 fctx.get_metadata("rpc.status"),
777 Some("ok"),
778 "multiple metadata keys should round-trip"
779 );
780 }
781
782 #[test]
783 fn metadata_written_in_filter_context_persists_back() {
784 let registry = FilterRegistry::with_builtins();
785 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
786 let request = Request {
787 method: Method::GET,
788 uri: "/".parse::<Uri>().unwrap(),
789 headers: HeaderMap::new(),
790 };
791
792 let mut ctx = default_ctx();
793 let mut fctx = ctx.build_filter_context(&pipeline, &request, None);
794 fctx.set_metadata("trace.id", "abc-123");
795 fctx.set_metadata("trace.span", "42");
796
797 ctx.filter_metadata = fctx.filter_metadata;
798 assert_eq!(
799 ctx.filter_metadata.get("trace.id").map(String::as_str),
800 Some("abc-123"),
801 "metadata set in filter context should persist back to protocol context"
802 );
803 assert_eq!(
804 ctx.filter_metadata.get("trace.span").map(String::as_str),
805 Some("42"),
806 "multiple metadata keys should persist back"
807 );
808 }
809
810 #[test]
815 fn pin_pipeline_captures_current_arc() {
816 let registry = FilterRegistry::with_builtins();
817 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
818 let swap = arc_swap::ArcSwap::from(Arc::clone(&pipeline_a));
819
820 let mut ctx = default_ctx();
821 let pinned = ctx.pin_pipeline(&swap);
822 assert!(
823 Arc::ptr_eq(&pinned, &pipeline_a),
824 "pin_pipeline should return the current pipeline"
825 );
826 assert!(
827 Arc::ptr_eq(ctx.pinned_pipeline.as_ref().unwrap(), &pipeline_a),
828 "pinned_pipeline should be stored in ctx"
829 );
830 }
831
832 #[test]
833 fn pipeline_returns_pinned_after_reload() {
834 let registry = FilterRegistry::with_builtins();
835 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
836 let swap = arc_swap::ArcSwap::from(Arc::clone(&pipeline_a));
837
838 let mut ctx = default_ctx();
839 ctx.pin_pipeline(&swap);
840
841 let pipeline_b = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
842 swap.store(pipeline_b);
843
844 let later = ctx.pipeline(&swap);
845 assert!(
846 Arc::ptr_eq(&later, &pipeline_a),
847 "later hooks should still return pipeline A after reload"
848 );
849 }
850
851 #[test]
852 fn new_request_after_reload_pins_new_pipeline() {
853 let registry = FilterRegistry::with_builtins();
854 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
855 let swap = arc_swap::ArcSwap::from(Arc::clone(&pipeline_a));
856
857 let mut ctx_a = default_ctx();
858 ctx_a.pin_pipeline(&swap);
859
860 let pipeline_b = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
861 swap.store(Arc::clone(&pipeline_b));
862
863 let mut ctx_b = default_ctx();
864 ctx_b.pin_pipeline(&swap);
865
866 assert!(
867 Arc::ptr_eq(&ctx_a.pipeline(&swap), &pipeline_a),
868 "request A should use old pipeline"
869 );
870 assert!(
871 Arc::ptr_eq(&ctx_b.pipeline(&swap), &pipeline_b),
872 "request B should use new pipeline"
873 );
874 }
875
876 #[test]
877 fn old_pipeline_drops_after_ctx_drops() {
878 let registry = FilterRegistry::with_builtins();
879 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
880 let weak_a = Arc::downgrade(&pipeline_a);
881 let swap = arc_swap::ArcSwap::from(pipeline_a);
882
883 let mut ctx = default_ctx();
884 ctx.pin_pipeline(&swap);
885
886 swap.store(Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap()));
887
888 assert!(weak_a.upgrade().is_some(), "old pipeline alive while ctx exists");
889
890 drop(ctx);
891
892 assert!(weak_a.upgrade().is_none(), "old pipeline drops with ctx");
893 }
894
895 #[test]
896 fn pipeline_helper_returns_pinned_for_every_phase() {
897 let registry = FilterRegistry::with_builtins();
898 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
899 let swap = arc_swap::ArcSwap::from(Arc::clone(&pipeline_a));
900
901 let mut ctx = default_ctx();
902 ctx.pin_pipeline(&swap);
903
904 swap.store(Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap()));
905
906 for phase in ["request_body", "response", "response_body", "logging"] {
907 let p = ctx.pipeline(&swap);
908 assert!(
909 Arc::ptr_eq(&p, &pipeline_a),
910 "{phase}: should still return pinned pipeline A"
911 );
912 }
913 }
914
915 #[test]
916 fn pipeline_fallback_when_not_pinned() {
917 let registry = FilterRegistry::with_builtins();
918 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
919 let swap = arc_swap::ArcSwap::from(Arc::clone(&pipeline_a));
920
921 let ctx = default_ctx();
922
923 let loaded = ctx.pipeline(&swap);
924 assert!(
925 Arc::ptr_eq(&loaded, &pipeline_a),
926 "unpinned ctx should fall back to current ArcSwap value"
927 );
928 }
929
930 #[test]
931 fn pin_pipeline_is_idempotent_after_reload() {
932 let registry = FilterRegistry::with_builtins();
933 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
934 let swap = arc_swap::ArcSwap::from(Arc::clone(&pipeline_a));
935
936 let mut ctx = default_ctx();
937 ctx.pin_pipeline(&swap);
938
939 let pipeline_b = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
940 swap.store(pipeline_b);
941
942 let second_pin = ctx.pin_pipeline(&swap);
943 assert!(
944 Arc::ptr_eq(&second_pin, &pipeline_a),
945 "repeated pin_pipeline after reload should return the original pin"
946 );
947 }
948
949 #[test]
950 fn filter_state_isolated_across_pipelines_with_same_ids() {
951 let registry = FilterRegistry::with_builtins();
952 let pipeline_a = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
953 let pipeline_b = Arc::new(FilterPipeline::build(&mut [], ®istry).unwrap());
954
955 let request = Request {
956 method: Method::GET,
957 uri: "/".parse::<Uri>().unwrap(),
958 headers: HeaderMap::new(),
959 };
960
961 let mut ctx_a = default_ctx();
962 ctx_a.pinned_pipeline = Some(Arc::clone(&pipeline_a));
963 ctx_a.request_snapshot = Some(request.clone());
964 let mut fctx_a = ctx_a.build_filter_context(&pipeline_a, &request, None);
965 fctx_a.current_filter_id = Some(0);
966 fctx_a.insert_filter_state(String::from("from_pipeline_a"));
967 ctx_a.filter_state = fctx_a.filter_state;
968
969 let mut ctx_b = default_ctx();
970 ctx_b.pinned_pipeline = Some(Arc::clone(&pipeline_b));
971 ctx_b.request_snapshot = Some(request.clone());
972 let fctx_b = ctx_b.build_filter_context(&pipeline_b, &request, None);
973
974 assert!(
975 fctx_b.filter_state.is_empty(),
976 "request B should have its own empty state map"
977 );
978
979 let fctx_a2 = ctx_a.filter_context_for(&pipeline_a, None).unwrap();
980 assert_eq!(
981 fctx_a2.filter_state.get(&0).and_then(|v| v.downcast_ref::<String>()),
982 Some(&String::from("from_pipeline_a")),
983 "request A should still see its own state in a later phase"
984 );
985 }
986
987 fn default_ctx() -> PingoraRequestCtx {
993 PingoraRequestCtx::default()
994 }
995}