Skip to main content

praxis_protocol/http/pingora/
context.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2024 Praxis Contributors
3
4//! Per-request context that carries filter pipeline results through Pingora's request/response lifecycle hooks.
5
6use 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// -----------------------------------------------------------------------------
15// PingoraRequestCtx
16// -----------------------------------------------------------------------------
17
18/// Per-request context carrying filter pipeline results through Pingora hooks.
19///
20/// ```
21/// use std::sync::Arc;
22///
23/// use praxis_protocol::http::pingora::context::PingoraRequestCtx;
24///
25/// let mut ctx = PingoraRequestCtx::default();
26/// ctx.cluster = Some(Arc::from("api-cluster"));
27/// assert_eq!(ctx.cluster.as_deref(), Some("api-cluster"));
28/// ```
29#[expect(clippy::struct_excessive_bools, reason = "lifecycle flags")]
30pub struct PingoraRequestCtx {
31    /// Connection permit from the per-listener semaphore.
32    ///
33    /// Held for the lifetime of the request. RAII drop
34    /// releases the permit when the context is dropped,
35    /// including error and timeout paths.
36    pub _connection_permit: Option<OwnedSemaphorePermit>,
37
38    /// Permit from the process-wide connection semaphore.
39    ///
40    /// Present only when `runtime.max_connections` is configured.
41    pub _global_connection_permit: Option<OwnedSemaphorePermit>,
42
43    /// Downstream client IP address.
44    pub client_addr: Option<IpAddr>,
45
46    /// HTTP version of the downstream client request.
47    ///
48    /// Captured during `request_filter` so the response-phase Via
49    /// header can reflect the protocol the client used.
50    pub client_http_version: Option<http::Version>,
51
52    /// Name of the cluster selected by a cluster-selecting filter.
53    pub cluster: Option<Arc<str>>,
54
55    /// Cached per-filter body-done indices. Swapped into each
56    /// [`HttpFilterContext`] and written back after execution so
57    /// that the heap allocation is reused across pipeline phases.
58    ///
59    /// [`HttpFilterContext`]: praxis_filter::HttpFilterContext
60    pub cached_body_done_indices: Vec<bool>,
61
62    /// Cached per-filter execution indices. Same lifecycle as
63    /// [`cached_body_done_indices`].
64    ///
65    /// [`cached_body_done_indices`]: Self::cached_body_done_indices
66    pub cached_executed_filter_indices: Vec<bool>,
67
68    /// Whether the downstream connection uses TLS.
69    ///
70    /// Derived from the Pingora session's SSL digest during
71    /// `request_filter`. Used by the forwarded headers filter
72    /// to set `X-Forwarded-Proto` correctly for HTTP/1.1
73    /// connections where the URI lacks a scheme.
74    pub downstream_tls: bool,
75
76    /// Verified downstream TLS peer identity.
77    ///
78    /// Set once from the SSL digest in `request_filter` before
79    /// the first filter runs.  Shared via [`Arc`] with each
80    /// `HttpFilterContext` (an `Arc` clone per phase, not a deep
81    /// copy per body chunk) so it is available in both pre-read
82    /// body phases and the main filter pipeline.  `None` for
83    /// non-mTLS or no-client-cert connections.
84    pub peer_identity: Option<Arc<praxis_tls::TlsPeerIdentity>>,
85
86    /// Whether the connection was upgraded via 101 Switching Protocols.
87    ///
88    /// Set during `response_filter` when the upstream returns 101.
89    /// Body filter hooks skip processing when true, since post-upgrade
90    /// bytes are raw protocol frames (e.g. `WebSocket`), not HTTP bodies.
91    pub connection_upgraded: bool,
92
93    /// Type-safe request-scoped extension container. Swapped into each
94    /// [`HttpFilterContext`] and written back after filter execution,
95    /// following the same lifecycle as [`filter_metadata`].
96    ///
97    /// [`HttpFilterContext`]: praxis_filter::HttpFilterContext
98    /// [`filter_metadata`]: PingoraRequestCtx::filter_metadata
99    pub extensions: praxis_filter::RequestExtensions,
100
101    /// Durable per-request metadata that persists across all lifecycle
102    /// phases. Swapped into each [`HttpFilterContext`] and written back
103    /// after filter execution.
104    ///
105    /// [`HttpFilterContext`]: praxis_filter::HttpFilterContext
106    pub filter_metadata: std::collections::HashMap<String, String>,
107
108    /// Ordered log of trusted header mutations from pre-read body
109    /// processing. Swapped into each [`HttpFilterContext`] and written
110    /// back after filter execution.
111    ///
112    /// [`HttpFilterContext`]: praxis_filter::HttpFilterContext
113    pub pre_read_mutations: Vec<TrustedHeaderMutation>,
114
115    /// Structured per-request metadata keyed by namespace.
116    /// Swapped into each [`HttpFilterContext`] and written back
117    /// after filter execution.
118    ///
119    /// [`HttpFilterContext`]: praxis_filter::HttpFilterContext
120    pub structured_metadata: std::collections::HashMap<String, serde_json::Value>,
121
122    /// Post-mutation request body length produced during `StreamBuffer`
123    /// pre-read.
124    ///
125    /// Stored so `upstream_request_filter` can repair request framing
126    /// before Pingora sends headers to the backend.
127    pub mutated_request_body_len: Option<usize>,
128
129    /// Pipeline pinned for this request's entire lifecycle.
130    ///
131    /// Set once during `request_filter` by cloning the [`Arc`] from the
132    /// listener's [`ArcSwap`]. All subsequent hooks (request body, response,
133    /// response body, logging) use this reference instead of re-loading
134    /// from the [`ArcSwap`], ensuring that a hot configuration reload
135    /// cannot change the pipeline mid-request.
136    ///
137    /// [`ArcSwap`]: arc_swap::ArcSwap
138    pub pinned_pipeline: Option<Arc<FilterPipeline>>,
139
140    /// Filter results from body pre-read. Carried into the next
141    /// [`HttpFilterContext`] so that branch chains attached to the
142    /// first `on_request` filter can evaluate body-derived results.
143    ///
144    /// [`HttpFilterContext`]: praxis_filter::HttpFilterContext
145    pub filter_results: std::collections::HashMap<&'static str, praxis_filter::FilterResultSet>,
146
147    /// Typed per-filter state that persists across all lifecycle
148    /// phases. Keyed by stable filter invocation ID, unique within
149    /// the request's pinned [`FilterPipeline`]. Swapped into each
150    /// [`HttpFilterContext`] and written back after filter execution,
151    /// following the same pattern as [`filter_metadata`].
152    ///
153    /// [`FilterPipeline`]: praxis_filter::FilterPipeline
154    /// [`HttpFilterContext`]: praxis_filter::HttpFilterContext
155    /// [`filter_metadata`]: Self::filter_metadata
156    pub filter_state: std::collections::HashMap<usize, Box<dyn std::any::Any + Send + Sync>>,
157
158    /// Cluster name snapshot retained for metrics emission in the
159    /// `logging()` hook, after `cluster` has been consumed by filter
160    /// context construction.
161    pub metrics_cluster: Option<Arc<str>>,
162
163    /// Pre-built [`SharedString`] for the metrics cluster label.
164    ///
165    /// Cached when `metrics_cluster` is set so that
166    /// `emit_request_metrics` avoids an `Arc` clone per request.
167    ///
168    /// [`SharedString`]: ::metrics::SharedString
169    pub metrics_cluster_shared: Option<::metrics::SharedString>,
170
171    /// Matched route path-match pattern for the `route` metric label.
172    pub metrics_route: Option<::metrics::SharedString>,
173
174    /// First proxy error cause seen for this request, for the
175    /// `praxis_errors_total` `type` label.
176    ///
177    /// Stamped first-cause-wins: one failing request can traverse several
178    /// classification sites across a retry fan-out, so the counter is
179    /// incremented once from the logging hook rather than at each site.
180    pub(crate) error_type: Option<&'static str>,
181
182    /// RAII guard that decrements `praxis_http_active_requests` on drop.
183    pub(crate) _active_request: Option<crate::http::pingora::metrics::ActiveRequestGuard>,
184
185    /// When the current upstream connect attempt started.
186    pub upstream_connect_start: Option<Instant>,
187
188    /// Pre-read body chunks (`StreamBuffer` mode). When `StreamBuffer` is
189    /// active, the body is read during `request_filter` (before upstream
190    /// selection) so that body-based routing can influence `upstream_peer`.
191    /// The `request_body_filter` hook then forwards these stored chunks
192    /// instead of reading from the session.
193    ///
194    /// Uses `VecDeque` so that draining from the front is O(1).
195    pub pre_read_body: Option<VecDeque<Bytes>>,
196
197    /// Retained copy of the mutated pre-read body for retry replay.
198    ///
199    /// The first attempt drains `pre_read_body` as it forwards the mutated
200    /// body. A retry replays from Pingora's fixed retry buffer, which holds the
201    /// ORIGINAL (pre-mutation) bytes, while `apply_mutated_content_length`
202    /// re-stamps the mutated length. This retained copy re-seeds `pre_read_body`
203    /// on each retry so the replayed body matches the stamped Content-Length,
204    /// closing a request-smuggling mismatch. Set only when a body writer ran.
205    pub retained_pre_read_body: Option<VecDeque<Bytes>>,
206
207    /// Buffer for request body accumulation in [`StreamBuffer`] mode.
208    ///
209    /// [`StreamBuffer`]: praxis_filter::BodyMode::StreamBuffer
210    pub request_body_buffer: Option<BodyBuffer>,
211
212    /// Accumulated request body bytes seen so far.
213    pub request_body_bytes: u64,
214
215    /// Per-request body delivery mode for the request direction.
216    /// Seeded from static pipeline capabilities, then potentially
217    /// upgraded by filters during `on_request`.
218    pub request_body_mode: BodyMode,
219
220    /// Whether the request body has been released (`StreamBuffer` mode).
221    /// Once true, remaining chunks bypass buffering and stream through.
222    pub request_body_released: bool,
223
224    /// Whether the request method is idempotent (GET, HEAD, OPTIONS).
225    pub request_is_idempotent: bool,
226
227    /// Snapshot of the original request for body/response body phases.
228    pub request_snapshot: Option<Request>,
229
230    /// Root tracing span for this request's lifecycle.
231    ///
232    /// Created during `request_filter` with OpenTelemetry HTTP semantic
233    /// convention attributes. Response-phase attributes
234    /// (`http.response.status_code`, `http.route`, `error.type`,
235    /// `upstream.address`, `upstream.cluster`) are recorded in the
236    /// `logging` hook before the span is dropped.
237    pub request_span: Span,
238
239    /// Child span covering upstream request/response exchange.
240    ///
241    /// Created in `connected_to_upstream` after the connection is
242    /// established (or reused). Response-phase attributes
243    /// (`http.response.status_code`, `http.response.body.size`) are
244    /// recorded in the `logging` hook before the span is dropped.
245    pub upstream_exchange_span: Span,
246
247    /// When this request was received.
248    pub request_start: Instant,
249
250    /// Buffer for response body accumulation in [`StreamBuffer`] mode.
251    ///
252    /// [`StreamBuffer`]: praxis_filter::BodyMode::StreamBuffer
253    pub response_body_buffer: Option<BodyBuffer>,
254
255    /// Accumulated response body bytes seen so far.
256    pub response_body_bytes: u64,
257
258    /// Per-request body delivery mode for the response direction.
259    /// Seeded from static pipeline capabilities, then potentially
260    /// upgraded by filters during `on_response`.
261    pub response_body_mode: BodyMode,
262
263    /// Whether the response body has been released (`StreamBuffer` mode).
264    pub response_body_released: bool,
265
266    /// Snapshot of response headers after response-phase filters.
267    ///
268    /// Used to evaluate `response_conditions` consistently during
269    /// response-body hooks, where Pingora no longer exposes mutable
270    /// response headers.
271    pub response_header_snapshot: Option<Response>,
272
273    /// Upstream response status code, captured during `response_filter`
274    /// for passive health recording in the `logging` hook.
275    pub upstream_response_status: Option<u16>,
276
277    /// Whether the response phase has been executed. Used to ensure
278    /// cleanup (e.g. least-connections counter release) in the
279    /// `logging()` hook when errors bypass `response_filter`.
280    pub response_phase_done: bool,
281
282    /// Rejection raised during the response phase, carried across the
283    /// Pingora error boundary so `fail_to_proxy` can deliver its full
284    /// configured headers and body instead of a bare status envelope.
285    pub pending_rejection: Option<praxis_filter::Rejection>,
286
287    /// Whether the response was delivered to completion: a bodyless
288    /// response finished its response phase, or the response body
289    /// reached end-of-stream. When still `false` in the `logging()`
290    /// hook, the request ended early (rejection, upstream failure, or
291    /// aborted stream) and a fallback access record is emitted.
292    pub response_delivery_complete: bool,
293
294    /// Number of upstream connection retries attempted.
295    pub retries: u32,
296
297    /// Index of the selected endpoint in the cluster's
298    /// endpoint list. Set during load balancing; used
299    /// for passive health recording in the logging hook.
300    pub selected_endpoint_index: Option<usize>,
301
302    /// Endpoints already attempted (for alternate-host retry).
303    pub attempted_endpoints: Vec<Arc<str>>,
304
305    /// Snapshot of the resolved retry policy for this request.
306    pub retry_policy: Option<Arc<praxis_core::config::RetryPolicy>>,
307
308    /// Optional route-level retry policy override (merged by the load balancer).
309    pub route_retry_policy: Option<Arc<praxis_core::config::RetryPolicy>>,
310
311    /// Shared cluster retry state (budget + active requests).
312    pub cluster_retry_state: Option<Arc<praxis_core::retry::ClusterRetryState>>,
313
314    /// Whether `cluster_retry_state.leave()` has already been called.
315    pub cluster_retry_state_released: bool,
316
317    /// Reselector for alternate-host selection on retry.
318    pub endpoint_reselector: Option<Arc<praxis_filter::EndpointReselector>>,
319
320    /// Pending backoff delay to apply before the next `upstream_peer` call.
321    pub pending_backoff: Option<std::time::Duration>,
322
323    /// Whether the next upstream attempt should re-select (alternate host).
324    pub reselect_on_retry: bool,
325
326    /// Rewritten URI path for the upstream request.
327    ///
328    /// Set by the `path_rewrite` filter via [`HttpFilterContext`] and
329    /// applied in `upstream_request_filter`.
330    ///
331    /// [`HttpFilterContext`]: praxis_filter::HttpFilterContext
332    pub rewritten_path: Option<String>,
333
334    /// Upstream endpoint selected by the load balancer filter.
335    pub upstream: Option<Upstream>,
336
337    /// Saved upstream for retry (cloned before first use).
338    pub upstream_for_retry: Option<Upstream>,
339
340    /// Whether the upstream was contacted (a peer was resolved for at least
341    /// one attempt) during this request. Unlike `upstream_for_retry`, which a
342    /// retry decision clears to force reselection, this stays set once the
343    /// upstream has been reached, so response-phase health accounting (passive
344    /// health, circuit breaker) can tell a genuine connect/read failure from a
345    /// request that never reached the cluster. Reset per request.
346    pub upstream_contacted: bool,
347}
348
349/// Build an [`HttpFilterContext`] from a `PingoraRequestCtx`.
350///
351/// Macro (not a function) so Rust's disjoint field borrowing
352/// works: `filter_context_for` borrows `self.request_snapshot`
353/// immutably while `cluster`, `upstream`, and `rewritten_path`
354/// are taken mutably. A function call with `&mut self` would
355/// collapse these into a single mutable borrow.
356///
357/// [`HttpFilterContext`]: praxis_filter::HttpFilterContext
358macro_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            // Seeded per pre-read pass by `pre_read_body`; empty otherwise.
380            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    /// Build an [`HttpFilterContext`] using an external request reference.
417    ///
418    /// Takes `cluster` and `upstream` from `self` (leaving `None`
419    /// behind) so that filters can reassign them. The caller must
420    /// write those fields back after filter execution.
421    ///
422    /// ```
423    /// use praxis_filter::{FilterPipeline, FilterRegistry, Request};
424    /// use praxis_protocol::http::pingora::context::PingoraRequestCtx;
425    ///
426    /// let registry = FilterRegistry::with_builtins();
427    /// let pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
428    /// let request = Request {
429    ///     method: http::Method::GET,
430    ///     uri: http::Uri::from_static("/"),
431    ///     headers: http::HeaderMap::new(),
432    /// };
433    /// let mut ctx = PingoraRequestCtx::default();
434    /// let filter_ctx = ctx.build_filter_context(&pipeline, &request, None);
435    /// assert!(filter_ctx.cluster.is_none());
436    /// ```
437    ///
438    /// [`HttpFilterContext`]: praxis_filter::HttpFilterContext
439    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    /// Build an [`HttpFilterContext`] from the stored [`request_snapshot`].
449    ///
450    /// Uses disjoint field borrowing so that `request_snapshot` is
451    /// borrowed immutably while `cluster` and `upstream` are taken
452    /// mutably.
453    ///
454    /// Returns `None` when `request_snapshot` is not set.
455    ///
456    /// ```
457    /// use praxis_filter::{FilterPipeline, FilterRegistry, Request};
458    /// use praxis_protocol::http::pingora::context::PingoraRequestCtx;
459    ///
460    /// let registry = FilterRegistry::with_builtins();
461    /// let pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
462    /// let mut ctx = PingoraRequestCtx::default();
463    /// ctx.request_snapshot = Some(Request {
464    ///     method: http::Method::GET,
465    ///     uri: http::Uri::from_static("/"),
466    ///     headers: http::HeaderMap::new(),
467    /// });
468    /// let filter_ctx = ctx.filter_context_for(&pipeline, None);
469    /// assert!(filter_ctx.is_some());
470    /// ```
471    ///
472    /// [`HttpFilterContext`]: praxis_filter::HttpFilterContext
473    /// [`request_snapshot`]: PingoraRequestCtx::request_snapshot
474    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    /// Build an [`HttpFilterContext`] plus the saved response header for body conditions.
484    ///
485    /// Response body hooks do not receive mutable headers, but their
486    /// `response_conditions` still need the response status and headers
487    /// captured during the response phase.
488    ///
489    /// [`HttpFilterContext`]: praxis_filter::HttpFilterContext
490    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    /// Pin the current pipeline for this request's entire lifecycle.
500    ///
501    /// Clones the [`Arc`] from the [`ArcSwap`] and stores it in
502    /// [`pinned_pipeline`]. All subsequent hooks should call
503    /// [`pipeline`] instead of re-loading from the [`ArcSwap`].
504    ///
505    /// Called once by `request_filter` in both body-capable and
506    /// no-body handlers.
507    ///
508    /// [`ArcSwap`]: arc_swap::ArcSwap
509    /// [`pinned_pipeline`]: Self::pinned_pipeline
510    /// [`pipeline`]: Self::pipeline
511    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        // Extension preparation is once per request, when the pipeline is
517        // pinned — not per filter context (body chunks build one each).
518        pipeline.prepare_extensions(&mut self.extensions);
519        self.pinned_pipeline = Some(Arc::clone(&pipeline));
520        pipeline
521    }
522
523    /// Record the first proxy error cause seen for this request.
524    ///
525    /// Later causes are ignored: a single failing request can pass through
526    /// several classification sites (connect failure, retry, final proxy
527    /// error), and the counter must reflect the root cause once.
528    pub(crate) fn stamp_error_type(&mut self, error_type: &'static str) {
529        self.error_type.get_or_insert(error_type);
530    }
531
532    /// Return the pinned pipeline, falling back to a fresh
533    /// [`ArcSwap`] load when no pipeline was pinned.
534    ///
535    /// The fallback covers early-failure paths where a lifecycle
536    /// hook runs before `request_filter` (e.g. after
537    /// `early_request_filter` rejection triggers `logging`).
538    ///
539    /// Per-body-chunk hooks (`request_body_filter`,
540    /// `response_body_filter`) call this on every chunk,
541    /// incurring one [`Arc::clone`] per invocation.
542    ///
543    /// [`ArcSwap`]: arc_swap::ArcSwap
544    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// -----------------------------------------------------------------------------
621// Tests
622// -----------------------------------------------------------------------------
623
624#[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    // -------------------------------------------------------------------------
897    // Metadata Roundtrip Tests
898    // -------------------------------------------------------------------------
899
900    #[test]
901    fn metadata_roundtrip_through_filter_context() {
902        let registry = FilterRegistry::with_builtins();
903        let pipeline = FilterPipeline::build(&mut [], &registry).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 [], &registry).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    // -------------------------------------------------------------------------
956    // Hot-Reload Pipeline Pinning (via production helpers)
957    // -------------------------------------------------------------------------
958
959    #[test]
960    fn pin_pipeline_captures_current_arc() {
961        let registry = FilterRegistry::with_builtins();
962        let pipeline_a = Arc::new(FilterPipeline::build(&mut [], &registry).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        // A single request can pass through more than one classification
980        // site (for example a request-filter reject followed by the
981        // terminal fail_to_proxy hook). The first stamp must win so the
982        // recorded cause is stable; the once-per-request increment in the
983        // logging hook is what keeps one error from being multiplied across
984        // a retry fan-out.
985        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 [], &registry).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 [], &registry).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 [], &registry).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 [], &registry).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 [], &registry).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 [], &registry).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 [], &registry).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 [], &registry).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 [], &registry).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 [], &registry).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 [], &registry).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 [], &registry).unwrap());
1122        let pipeline_b = Arc::new(FilterPipeline::build(&mut [], &registry).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    // -------------------------------------------------------------------------
1157    // Test Utilities
1158    // -------------------------------------------------------------------------
1159
1160    /// Create a default request context for tests.
1161    fn default_ctx() -> PingoraRequestCtx {
1162        PingoraRequestCtx::default()
1163    }
1164}