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 request's lifetime; released on drop,
34    /// including error and timeout paths.
35    pub _connection_permit: Option<OwnedSemaphorePermit>,
36
37    /// Permit from the process-wide connection semaphore.
38    ///
39    /// Present only when `runtime.max_connections` is configured.
40    pub _global_connection_permit: Option<OwnedSemaphorePermit>,
41
42    /// Downstream client IP address.
43    pub client_addr: Option<IpAddr>,
44
45    /// HTTP version of the downstream client request.
46    ///
47    /// Captured during `request_filter` so the response-phase Via
48    /// header can reflect the protocol the client used.
49    pub client_http_version: Option<http::Version>,
50
51    /// Name of the cluster selected by a cluster-selecting filter.
52    pub cluster: Option<Arc<str>>,
53
54    /// Cached per-filter body-done indices. Swapped into each
55    /// [`HttpFilterContext`] and written back after execution.
56    ///
57    /// [`HttpFilterContext`]: praxis_filter::HttpFilterContext
58    pub cached_body_done_indices: Vec<bool>,
59
60    /// Cached per-filter execution indices. Same lifecycle as
61    /// [`cached_body_done_indices`].
62    ///
63    /// [`cached_body_done_indices`]: Self::cached_body_done_indices
64    pub cached_executed_filter_indices: Vec<bool>,
65
66    /// Whether the downstream connection uses TLS.
67    ///
68    /// Set during `request_filter`. Used by the forwarded headers
69    /// filter to set `X-Forwarded-Proto` correctly for HTTP/1.1
70    /// connections where the URI lacks a scheme.
71    pub downstream_tls: bool,
72
73    /// Verified downstream TLS peer identity.
74    ///
75    /// Set once from the SSL digest in `request_filter` before the first
76    /// filter runs. Shared via [`Arc`] with each `HttpFilterContext` so it
77    /// is available in both pre-read body phases and the main filter
78    /// pipeline. `None` for non-mTLS or no-client-cert connections.
79    pub peer_identity: Option<Arc<praxis_tls::TlsPeerIdentity>>,
80
81    /// Whether the connection was upgraded via 101 Switching Protocols.
82    ///
83    /// Set during `response_filter` when the upstream returns 101.
84    /// Body filter hooks skip processing when true, since post-upgrade
85    /// bytes are raw protocol frames (e.g. `WebSocket`), not HTTP bodies.
86    pub connection_upgraded: bool,
87
88    /// Type-safe request-scoped extension container. Swapped into each
89    /// [`HttpFilterContext`] and written back after filter execution,
90    /// following the same lifecycle as [`filter_metadata`].
91    ///
92    /// [`HttpFilterContext`]: praxis_filter::HttpFilterContext
93    /// [`filter_metadata`]: PingoraRequestCtx::filter_metadata
94    pub extensions: praxis_filter::RequestExtensions,
95
96    /// Durable per-request metadata that persists across all lifecycle
97    /// phases. Swapped into each [`HttpFilterContext`] and written back
98    /// after filter execution.
99    ///
100    /// [`HttpFilterContext`]: praxis_filter::HttpFilterContext
101    pub filter_metadata: std::collections::HashMap<String, String>,
102
103    /// Ordered log of trusted header mutations from pre-read body
104    /// processing. Swapped into each [`HttpFilterContext`] and written
105    /// back after filter execution.
106    ///
107    /// [`HttpFilterContext`]: praxis_filter::HttpFilterContext
108    pub pre_read_mutations: Vec<TrustedHeaderMutation>,
109
110    /// Structured per-request metadata keyed by namespace.
111    /// Swapped into each [`HttpFilterContext`] and written back
112    /// after filter execution.
113    ///
114    /// [`HttpFilterContext`]: praxis_filter::HttpFilterContext
115    pub structured_metadata: std::collections::HashMap<String, serde_json::Value>,
116
117    /// Post-mutation request body length produced during `StreamBuffer`
118    /// pre-read.
119    ///
120    /// Stored so `upstream_request_filter` can repair request framing
121    /// before Pingora sends headers to the backend.
122    pub mutated_request_body_len: Option<usize>,
123
124    /// Pipeline pinned for this request's entire lifecycle.
125    ///
126    /// Set once during `request_filter` by cloning the [`Arc`] from the
127    /// listener's [`ArcSwap`]. All subsequent hooks (request body, response,
128    /// response body, logging) use this reference instead of re-loading
129    /// from the [`ArcSwap`], ensuring that a hot configuration reload
130    /// cannot change the pipeline mid-request.
131    ///
132    /// [`ArcSwap`]: arc_swap::ArcSwap
133    pub pinned_pipeline: Option<Arc<FilterPipeline>>,
134
135    /// Filter results from body pre-read. Carried into the next
136    /// [`HttpFilterContext`] so that branch chains attached to the
137    /// first `on_request` filter can evaluate body-derived results.
138    ///
139    /// [`HttpFilterContext`]: praxis_filter::HttpFilterContext
140    pub filter_results: std::collections::HashMap<&'static str, praxis_filter::FilterResultSet>,
141
142    /// Typed per-filter state that persists across all lifecycle
143    /// phases. Keyed by stable filter invocation ID, unique within
144    /// the request's pinned [`FilterPipeline`]. Swapped into each
145    /// [`HttpFilterContext`] and written back after filter execution,
146    /// following the same pattern as [`filter_metadata`].
147    ///
148    /// [`FilterPipeline`]: praxis_filter::FilterPipeline
149    /// [`HttpFilterContext`]: praxis_filter::HttpFilterContext
150    /// [`filter_metadata`]: Self::filter_metadata
151    pub filter_state: std::collections::HashMap<usize, Box<dyn std::any::Any + Send + Sync>>,
152
153    /// Cluster name snapshot retained for metrics emission in the
154    /// `logging()` hook, after `cluster` has been consumed by filter
155    /// context construction.
156    pub metrics_cluster: Option<Arc<str>>,
157
158    /// Pre-built [`SharedString`] for the metrics cluster label, set
159    /// alongside `metrics_cluster`.
160    ///
161    /// [`SharedString`]: ::metrics::SharedString
162    pub metrics_cluster_shared: Option<::metrics::SharedString>,
163
164    /// Matched route path-match pattern for the `route` metric label.
165    pub metrics_route: Option<::metrics::SharedString>,
166
167    /// First proxy error cause seen for this request, for the
168    /// `praxis_errors_total` `type` label.
169    ///
170    /// Stamped first-cause-wins: one failing request can traverse several
171    /// classification sites across a retry fan-out, so the counter is
172    /// incremented once from the logging hook rather than at each site.
173    pub(crate) error_type: Option<&'static str>,
174
175    /// RAII guard that decrements `praxis_http_active_requests` on drop.
176    pub(crate) _active_request: Option<crate::http::pingora::metrics::ActiveRequestGuard>,
177
178    /// When the current upstream connect attempt started.
179    pub upstream_connect_start: Option<Instant>,
180
181    /// Pre-read body chunks (`StreamBuffer` mode). When `StreamBuffer` is
182    /// active, the body is read during `request_filter` (before upstream
183    /// selection) so that body-based routing can influence `upstream_peer`.
184    /// The `request_body_filter` hook then forwards these stored chunks
185    /// instead of reading from the session.
186    pub pre_read_body: Option<VecDeque<Bytes>>,
187
188    /// Retained copy of the mutated pre-read body for retry replay.
189    ///
190    /// The first attempt drains `pre_read_body` as it forwards the mutated
191    /// body. A retry replays from Pingora's fixed retry buffer, which holds the
192    /// ORIGINAL (pre-mutation) bytes, while `apply_mutated_content_length`
193    /// re-stamps the mutated length. This retained copy re-seeds `pre_read_body`
194    /// on each retry so the replayed body matches the stamped Content-Length,
195    /// closing a request-smuggling mismatch. Set only when a body writer ran.
196    pub retained_pre_read_body: Option<VecDeque<Bytes>>,
197
198    /// Adapted request body captured by the selected-upstream phase (#1139).
199    ///
200    /// When a selected-upstream body writer ran, the phase output is frozen
201    /// here as the exchange-local body Pingora forwards to the selected
202    /// upstream. Drained by `request_body_filter` in preference to
203    /// `pre_read_body`; kept separate from the canonical pre-read body so the
204    /// canonical bytes are never forwarded once adaptation ran (#1138).
205    ///
206    /// Uses `VecDeque` so draining from the front is O(1).
207    pub adapted_request_body: Option<VecDeque<Bytes>>,
208
209    /// Retained copy of the adapted body for retry replay (#1139).
210    ///
211    /// `Some(_)` is the durable marker that adaptation ran for this request; it
212    /// persists for the request's lifetime (unlike `adapted_request_body`,
213    /// which drains to `None`). `reseed_retry_body` restores
214    /// `adapted_request_body` from this copy on each retry so replays match the
215    /// stamped `Content-Length`.
216    pub retained_adapted_request_body: Option<VecDeque<Bytes>>,
217
218    /// Authoritative length of the adapted body sent to the selected upstream
219    /// (#1139). Takes precedence over `mutated_request_body_len` in
220    /// `apply_mutated_content_length` and the retry-body replay guard.
221    pub adapted_request_body_len: Option<usize>,
222
223    /// Buffer for request body accumulation in [`StreamBuffer`] mode.
224    ///
225    /// [`StreamBuffer`]: praxis_filter::BodyMode::StreamBuffer
226    pub request_body_buffer: Option<BodyBuffer>,
227
228    /// Accumulated request body bytes seen so far.
229    pub request_body_bytes: u64,
230
231    /// Per-request body delivery mode for the request direction.
232    /// Seeded from static pipeline capabilities, then potentially
233    /// upgraded by filters during `on_request`.
234    pub request_body_mode: BodyMode,
235
236    /// Whether the request body has been released (`StreamBuffer` mode).
237    /// Once true, remaining chunks bypass buffering and stream through.
238    pub request_body_released: bool,
239
240    /// Whether the request method is idempotent (GET, HEAD, OPTIONS).
241    pub request_is_idempotent: bool,
242
243    /// Snapshot of the original request for body/response body phases.
244    pub request_snapshot: Option<Request>,
245
246    /// Root tracing span for this request's lifecycle.
247    ///
248    /// Created during `request_filter` with OpenTelemetry HTTP semantic
249    /// convention attributes. Response-phase attributes
250    /// (`http.response.status_code`, `http.route`, `error.type`,
251    /// `upstream.address`, `upstream.cluster`) are recorded in the
252    /// `logging` hook before the span is dropped.
253    pub request_span: Span,
254
255    /// Child span covering upstream request/response exchange.
256    ///
257    /// Created in `connected_to_upstream` after the connection is
258    /// established (or reused). Response-phase attributes
259    /// (`http.response.status_code`, `http.response.body.size`) are
260    /// recorded in the `logging` hook before the span is dropped.
261    pub upstream_exchange_span: Span,
262
263    /// When this request was received.
264    pub request_start: Instant,
265
266    /// Buffer for response body accumulation in [`StreamBuffer`] mode.
267    ///
268    /// [`StreamBuffer`]: praxis_filter::BodyMode::StreamBuffer
269    pub response_body_buffer: Option<BodyBuffer>,
270
271    /// Accumulated response body bytes seen so far.
272    pub response_body_bytes: u64,
273
274    /// Per-request body delivery mode for the response direction.
275    /// Seeded from static pipeline capabilities, then potentially
276    /// upgraded by filters during `on_response`.
277    pub response_body_mode: BodyMode,
278
279    /// Whether the response body has been released (`StreamBuffer` mode).
280    pub response_body_released: bool,
281
282    /// Snapshot of response headers after response-phase filters.
283    ///
284    /// Used to evaluate `response_conditions` consistently during
285    /// response-body hooks, where Pingora no longer exposes mutable
286    /// response headers.
287    pub response_header_snapshot: Option<Response>,
288
289    /// Upstream response status code, captured during `response_filter`
290    /// for passive health recording in the `logging` hook.
291    pub upstream_response_status: Option<u16>,
292
293    /// How the upstream gRPC call ended.
294    ///
295    /// Captured from the response trailers, or from the response header
296    /// block of a Trailers-Only response. Held here rather than in the
297    /// filter context because trailers arrive after the response-header
298    /// phase, and the access log runs later still.
299    pub grpc_completion: Option<praxis_core::grpc::GrpcCompletion>,
300
301    /// Whether the response phase has been executed. Used to ensure
302    /// cleanup (e.g. least-connections counter release) in the
303    /// `logging()` hook when errors bypass `response_filter`.
304    pub response_phase_done: bool,
305
306    /// Rejection raised during the response phase, carried across the
307    /// Pingora error boundary so `fail_to_proxy` can deliver its full
308    /// configured headers and body instead of a bare status envelope.
309    pub pending_rejection: Option<praxis_filter::Rejection>,
310
311    /// Whether the response was delivered to completion: a bodyless
312    /// response finished its response phase, or the response body
313    /// reached end-of-stream. When still `false` in the `logging()`
314    /// hook, the request ended early (rejection, upstream failure, or
315    /// aborted stream) and a fallback access record is emitted.
316    pub response_delivery_complete: bool,
317
318    /// Number of upstream connection retries attempted.
319    pub retries: u32,
320
321    /// Index of the selected endpoint in the cluster's
322    /// endpoint list. Set during load balancing; used
323    /// for passive health recording in the logging hook.
324    pub selected_endpoint_index: Option<usize>,
325
326    /// Endpoints already attempted (for alternate-host retry).
327    pub attempted_endpoints: Vec<Arc<str>>,
328
329    /// Snapshot of the resolved retry policy for this request.
330    pub retry_policy: Option<Arc<praxis_core::config::RetryPolicy>>,
331
332    /// Optional route-level retry policy override (merged by the load balancer).
333    pub route_retry_policy: Option<Arc<praxis_core::config::RetryPolicy>>,
334
335    /// Shared cluster retry state (budget + active requests).
336    pub cluster_retry_state: Option<Arc<praxis_core::retry::ClusterRetryState>>,
337
338    /// Whether `cluster_retry_state.leave()` has already been called.
339    pub cluster_retry_state_released: bool,
340
341    /// Reselector for alternate-host selection on retry.
342    pub endpoint_reselector: Option<Arc<praxis_filter::EndpointReselector>>,
343
344    /// Pending backoff delay to apply before the next `upstream_peer` call.
345    pub pending_backoff: Option<std::time::Duration>,
346
347    /// Whether the next upstream attempt should re-select (alternate host).
348    pub reselect_on_retry: bool,
349
350    /// Rewritten URI path for the upstream request.
351    ///
352    /// Set by the `path_rewrite` filter via [`HttpFilterContext`] and
353    /// applied in `upstream_request_filter`.
354    ///
355    /// [`HttpFilterContext`]: praxis_filter::HttpFilterContext
356    pub rewritten_path: Option<String>,
357
358    /// Upstream endpoint selected by the load balancer filter.
359    pub upstream: Option<Upstream>,
360
361    /// Saved upstream for retry (cloned before first use).
362    pub upstream_for_retry: Option<Upstream>,
363
364    /// Whether the upstream was contacted (a peer was resolved for at least
365    /// one attempt) during this request. Unlike `upstream_for_retry`, which a
366    /// retry decision clears to force reselection, this stays set once the
367    /// upstream has been reached, so response-phase health accounting (passive
368    /// health, circuit breaker) can tell a genuine connect/read failure from a
369    /// request that never reached the cluster. Reset per request.
370    pub upstream_contacted: bool,
371}
372
373/// Build an [`HttpFilterContext`] from a `PingoraRequestCtx`.
374///
375/// Macro (not a function) so Rust's disjoint field borrowing
376/// works: `filter_context_for` borrows `self.request_snapshot`
377/// immutably while `cluster`, `upstream`, and `rewritten_path`
378/// are taken mutably. A function call with `&mut self` would
379/// collapse these into a single mutable borrow.
380///
381/// [`HttpFilterContext`]: praxis_filter::HttpFilterContext
382macro_rules! filter_context {
383    ($ctx:expr, $pipeline:expr, $request:expr, $response_header:expr) => {{
384        praxis_filter::HttpFilterContext {
385            buffered_request_body: $ctx
386                .pre_read_body
387                .as_ref()
388                .map(|chunks| chunks.front().cloned().unwrap_or_default()),
389            body_done_indices: std::mem::take(&mut $ctx.cached_body_done_indices),
390            branch_iterations: std::collections::HashMap::new(),
391            client_addr: $ctx.client_addr,
392            cluster: $ctx.cluster.take(),
393            current_filter_id: None,
394            downstream_tls: $ctx.downstream_tls,
395            metrics_route: $ctx.metrics_route.clone(),
396            peer_identity: $ctx.peer_identity.clone(),
397            extensions: std::mem::take(&mut $ctx.extensions),
398            executed_filter_indices: std::mem::take(&mut $ctx.cached_executed_filter_indices),
399            extra_request_headers: Vec::new(),
400            request_headers_to_remove: Vec::new(),
401            request_headers_to_set: Vec::new(),
402            filter_metadata: std::mem::take(&mut $ctx.filter_metadata),
403            // Seeded per pre-read pass by `pre_read_body`; empty otherwise.
404            prior_pre_read_mutations: Vec::new(),
405            grpc_completion: $ctx.grpc_completion.clone(),
406            pre_read_mutations: std::mem::take(&mut $ctx.pre_read_mutations),
407            structured_metadata: std::mem::take(&mut $ctx.structured_metadata),
408            filter_results: std::mem::take(&mut $ctx.filter_results),
409            filter_state: std::mem::take(&mut $ctx.filter_state),
410            health_registry: $pipeline.health_registry(),
411            id_generator: $pipeline.id_generator(),
412            kv_stores: $pipeline.kv_stores(),
413            session_stores: $pipeline.session_stores(),
414            subrequest_client: $pipeline.subrequest_client(),
415            subrequest_response_mode: praxis_filter::SubRequestResponseMode::Buffered,
416            request: $request,
417            request_body_bytes: $ctx.request_body_bytes,
418            request_body_mode: $ctx.request_body_mode,
419            request_start: $ctx.request_start,
420            response_body_bytes: $ctx.response_body_bytes,
421            response_body_mode: $ctx.response_body_mode,
422            response_header: $response_header,
423            response_headers_modified: false,
424            upstream_reached: $ctx.upstream_contacted,
425            rewritten_path: $ctx.rewritten_path.take(),
426            selected_endpoint_index: $ctx.selected_endpoint_index,
427            attempted_endpoints: std::mem::take(&mut $ctx.attempted_endpoints),
428            retry_policy: $ctx.retry_policy.clone(),
429            route_retry_policy: $ctx.route_retry_policy.clone(),
430            cluster_retry_state: $ctx.cluster_retry_state.clone(),
431            cluster_retry_state_released: $ctx.cluster_retry_state_released,
432            endpoint_reselector: $ctx.endpoint_reselector.clone(),
433            pinned_endpoint_address: None,
434            time_source: $pipeline.time_source(),
435            upstream: $ctx.upstream.take().or_else(|| $ctx.upstream_for_retry.clone()),
436        }
437    }};
438}
439
440impl PingoraRequestCtx {
441    /// Build an [`HttpFilterContext`] using an external request reference.
442    ///
443    /// Takes `cluster` and `upstream` from `self` (leaving `None`
444    /// behind) so that filters can reassign them. The caller must
445    /// write those fields back after filter execution.
446    ///
447    /// ```
448    /// use praxis_filter::{FilterPipeline, FilterRegistry, Request};
449    /// use praxis_protocol::http::pingora::context::PingoraRequestCtx;
450    ///
451    /// let registry = FilterRegistry::with_builtins();
452    /// let pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
453    /// let request = Request {
454    ///     method: http::Method::GET,
455    ///     uri: http::Uri::from_static("/"),
456    ///     headers: http::HeaderMap::new(),
457    /// };
458    /// let mut ctx = PingoraRequestCtx::default();
459    /// let filter_ctx = ctx.build_filter_context(&pipeline, &request, None);
460    /// assert!(filter_ctx.cluster.is_none());
461    /// ```
462    ///
463    /// [`HttpFilterContext`]: praxis_filter::HttpFilterContext
464    pub fn build_filter_context<'a>(
465        &mut self,
466        pipeline: &'a FilterPipeline,
467        request: &'a Request,
468        response_header: Option<&'a mut Response>,
469    ) -> praxis_filter::HttpFilterContext<'a> {
470        filter_context!(self, pipeline, request, response_header)
471    }
472
473    /// Build an [`HttpFilterContext`] from the stored [`request_snapshot`].
474    ///
475    /// Returns `None` when `request_snapshot` is not set.
476    ///
477    /// ```
478    /// use praxis_filter::{FilterPipeline, FilterRegistry, Request};
479    /// use praxis_protocol::http::pingora::context::PingoraRequestCtx;
480    ///
481    /// let registry = FilterRegistry::with_builtins();
482    /// let pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
483    /// let mut ctx = PingoraRequestCtx::default();
484    /// ctx.request_snapshot = Some(Request {
485    ///     method: http::Method::GET,
486    ///     uri: http::Uri::from_static("/"),
487    ///     headers: http::HeaderMap::new(),
488    /// });
489    /// let filter_ctx = ctx.filter_context_for(&pipeline, None);
490    /// assert!(filter_ctx.is_some());
491    /// ```
492    ///
493    /// [`HttpFilterContext`]: praxis_filter::HttpFilterContext
494    /// [`request_snapshot`]: PingoraRequestCtx::request_snapshot
495    pub fn filter_context_for<'a>(
496        &'a mut self,
497        pipeline: &'a FilterPipeline,
498        response_header: Option<&'a mut Response>,
499    ) -> Option<praxis_filter::HttpFilterContext<'a>> {
500        let request = self.request_snapshot.as_ref()?;
501        Some(filter_context!(self, pipeline, request, response_header))
502    }
503
504    /// Build an [`HttpFilterContext`] plus the saved response header for body conditions.
505    ///
506    /// Response body hooks do not receive mutable headers, but their
507    /// `response_conditions` still need the response status and headers
508    /// captured during the response phase.
509    ///
510    /// [`HttpFilterContext`]: praxis_filter::HttpFilterContext
511    pub fn response_body_context_for<'a>(
512        &'a mut self,
513        pipeline: &'a FilterPipeline,
514    ) -> Option<(praxis_filter::HttpFilterContext<'a>, Option<&'a Response>)> {
515        let request = self.request_snapshot.as_ref()?;
516        let response_header = self.response_header_snapshot.as_ref();
517        Some((filter_context!(self, pipeline, request, None), response_header))
518    }
519
520    /// Pin the current pipeline for this request's entire lifecycle.
521    ///
522    /// Clones the [`Arc`] from the [`ArcSwap`] and stores it in
523    /// [`pinned_pipeline`]. All subsequent hooks should call
524    /// [`pipeline`] instead of re-loading from the [`ArcSwap`].
525    ///
526    /// Called by `early_request_filter` before compression negotiation;
527    /// repeated calls in later hooks reuse the same generation.
528    ///
529    /// [`ArcSwap`]: arc_swap::ArcSwap
530    /// [`pinned_pipeline`]: Self::pinned_pipeline
531    /// [`pipeline`]: Self::pipeline
532    pub fn pin_pipeline(&mut self, swap: &arc_swap::ArcSwap<FilterPipeline>) -> Arc<FilterPipeline> {
533        if let Some(existing) = &self.pinned_pipeline {
534            return Arc::clone(existing);
535        }
536        let pipeline = swap.load_full();
537        // Extension preparation is once per request, when the pipeline is
538        // pinned — not per filter context (body chunks build one each).
539        pipeline.prepare_extensions(&mut self.extensions);
540        self.pinned_pipeline = Some(Arc::clone(&pipeline));
541        pipeline
542    }
543
544    /// Record the first proxy error cause seen for this request.
545    ///
546    /// Later causes are ignored: a single failing request can pass through
547    /// several classification sites (connect failure, retry, final proxy
548    /// error), and the counter must reflect the root cause once.
549    pub(crate) fn stamp_error_type(&mut self, error_type: &'static str) {
550        self.error_type.get_or_insert(error_type);
551    }
552
553    /// Return the pinned pipeline, falling back to a fresh
554    /// [`ArcSwap`] load when no pipeline was pinned.
555    ///
556    /// The fallback covers early-failure paths where a lifecycle
557    /// hook runs before `request_filter` (e.g. after
558    /// `early_request_filter` rejection triggers `logging`).
559    ///
560    /// [`ArcSwap`]: arc_swap::ArcSwap
561    pub fn pipeline(&self, swap: &arc_swap::ArcSwap<FilterPipeline>) -> Arc<FilterPipeline> {
562        self.pinned_pipeline
563            .as_ref()
564            .map_or_else(|| swap.load_full(), Arc::clone)
565    }
566}
567
568impl Default for PingoraRequestCtx {
569    #[expect(
570        clippy::too_many_lines,
571        reason = "context default enumerates all lifecycle fields explicitly"
572    )]
573    fn default() -> Self {
574        Self {
575            _connection_permit: None,
576            _global_connection_permit: None,
577            cached_body_done_indices: Vec::new(),
578            cached_executed_filter_indices: Vec::new(),
579            client_addr: None,
580            client_http_version: None,
581            cluster: None,
582            connection_upgraded: false,
583            downstream_tls: false,
584            peer_identity: None,
585            extensions: praxis_filter::RequestExtensions::new(),
586            filter_metadata: std::collections::HashMap::new(),
587            pre_read_mutations: Vec::new(),
588            structured_metadata: std::collections::HashMap::new(),
589            mutated_request_body_len: None,
590            pinned_pipeline: None,
591            filter_results: std::collections::HashMap::new(),
592            filter_state: std::collections::HashMap::new(),
593            metrics_cluster: None,
594            metrics_cluster_shared: None,
595            metrics_route: None,
596            error_type: None,
597            _active_request: None,
598            upstream_connect_start: None,
599            pre_read_body: None,
600            retained_pre_read_body: None,
601            adapted_request_body: None,
602            retained_adapted_request_body: None,
603            adapted_request_body_len: None,
604            request_body_buffer: None,
605            request_body_bytes: 0,
606            request_body_mode: BodyMode::Stream,
607            request_body_released: false,
608            request_is_idempotent: false,
609            request_snapshot: None,
610            request_span: Span::none(),
611            request_start: Instant::now(),
612            upstream_exchange_span: Span::none(),
613            response_body_buffer: None,
614            response_body_bytes: 0,
615            response_body_mode: BodyMode::Stream,
616            response_body_released: false,
617            response_header_snapshot: None,
618            upstream_response_status: None,
619            grpc_completion: None,
620            response_phase_done: false,
621            response_delivery_complete: false,
622            pending_rejection: None,
623            retries: 0,
624            rewritten_path: None,
625            selected_endpoint_index: None,
626            attempted_endpoints: Vec::new(),
627            retry_policy: None,
628            route_retry_policy: None,
629            cluster_retry_state: None,
630            cluster_retry_state_released: false,
631            endpoint_reselector: None,
632            pending_backoff: None,
633            reselect_on_retry: false,
634            upstream: None,
635            upstream_for_retry: None,
636            upstream_contacted: false,
637        }
638    }
639}
640
641// -----------------------------------------------------------------------------
642// Tests
643// -----------------------------------------------------------------------------
644
645#[cfg(test)]
646#[expect(clippy::allow_attributes, reason = "blanket test suppressions")]
647#[allow(
648    clippy::unwrap_used,
649    clippy::expect_used,
650    clippy::indexing_slicing,
651    clippy::significant_drop_tightening,
652    clippy::too_many_lines,
653    reason = "tests"
654)]
655mod tests {
656    use std::net::Ipv4Addr;
657
658    use http::{HeaderMap, Method, Uri};
659    use praxis_filter::FilterRegistry;
660
661    use super::*;
662
663    #[test]
664    fn default_state_has_no_client_addr() {
665        let ctx = default_ctx();
666        assert!(ctx.client_addr.is_none(), "default client_addr should be None");
667    }
668
669    #[test]
670    fn default_state_has_no_cluster() {
671        let ctx = default_ctx();
672        assert!(ctx.cluster.is_none(), "default cluster should be None");
673    }
674
675    #[test]
676    fn default_state_has_zero_retries() {
677        let ctx = default_ctx();
678        assert_eq!(ctx.retries, 0, "default retries should be zero");
679    }
680
681    #[test]
682    fn default_state_flags_are_false() {
683        let ctx = default_ctx();
684        assert!(
685            !ctx.request_body_released,
686            "default request_body_released should be false"
687        );
688        assert!(
689            !ctx.response_body_released,
690            "default response_body_released should be false"
691        );
692        assert!(
693            !ctx.request_is_idempotent,
694            "default request_is_idempotent should be false"
695        );
696        assert!(!ctx.response_phase_done, "default response_phase_done should be false");
697    }
698
699    #[test]
700    fn default_state_buffers_are_none() {
701        let ctx = default_ctx();
702        assert!(
703            ctx.request_body_buffer.is_none(),
704            "default request_body_buffer should be None"
705        );
706        assert!(
707            ctx.response_body_buffer.is_none(),
708            "default response_body_buffer should be None"
709        );
710        assert!(ctx.pre_read_body.is_none(), "default pre_read_body should be None");
711    }
712
713    #[test]
714    fn default_state_request_span_is_disabled() {
715        let ctx = default_ctx();
716        assert!(
717            ctx.request_span.is_disabled(),
718            "default request_span should be a disabled (none) span"
719        );
720    }
721
722    #[test]
723    fn default_state_upstream_exchange_span_is_disabled() {
724        let ctx = default_ctx();
725        assert!(
726            ctx.upstream_exchange_span.is_disabled(),
727            "default upstream_exchange_span should be a disabled (none) span"
728        );
729    }
730
731    #[test]
732    fn default_state_snapshots_are_none() {
733        let ctx = default_ctx();
734        assert!(
735            ctx.request_snapshot.is_none(),
736            "default request_snapshot should be None"
737        );
738        assert!(ctx.upstream.is_none(), "default upstream should be None");
739        assert!(
740            ctx.upstream_for_retry.is_none(),
741            "default upstream_for_retry should be None"
742        );
743    }
744
745    #[test]
746    fn set_client_addr() {
747        let mut ctx = default_ctx();
748        let addr = IpAddr::V4(Ipv4Addr::new(192, 168, 1, 100));
749        ctx.client_addr = Some(addr);
750        assert_eq!(
751            ctx.client_addr.unwrap(),
752            addr,
753            "client_addr should match assigned value"
754        );
755    }
756
757    #[test]
758    fn set_cluster() {
759        let mut ctx = default_ctx();
760        ctx.cluster = Some(Arc::from("api-cluster"));
761        assert_eq!(
762            ctx.cluster.as_deref(),
763            Some("api-cluster"),
764            "cluster should match assigned value"
765        );
766    }
767
768    #[test]
769    fn set_upstream() {
770        let mut ctx = default_ctx();
771        let upstream = Upstream {
772            address: Arc::from("10.0.0.1:80"),
773            authority: None,
774            tls: None,
775            connection: Arc::new(praxis_core::connectivity::ConnectionOptions::default()),
776        };
777        ctx.upstream = Some(upstream.clone());
778        assert_eq!(
779            &*ctx.upstream.as_ref().unwrap().address,
780            "10.0.0.1:80",
781            "upstream address should match assigned value"
782        );
783    }
784
785    #[test]
786    fn increment_retries() {
787        let mut ctx = default_ctx();
788        ctx.retries += 1;
789        ctx.retries += 1;
790        assert_eq!(ctx.retries, 2, "retries should be 2 after two increments");
791    }
792
793    #[test]
794    fn release_request_body_flag() {
795        let mut ctx = default_ctx();
796        assert!(!ctx.request_body_released, "request_body_released should start false");
797        ctx.request_body_released = true;
798        assert!(
799            ctx.request_body_released,
800            "request_body_released should be true after setting"
801        );
802    }
803
804    #[test]
805    fn release_response_body_flag() {
806        let mut ctx = default_ctx();
807        assert!(!ctx.response_body_released, "response_body_released should start false");
808        ctx.response_body_released = true;
809        assert!(
810            ctx.response_body_released,
811            "response_body_released should be true after setting"
812        );
813    }
814
815    #[test]
816    fn response_phase_done_flag() {
817        let mut ctx = default_ctx();
818        assert!(!ctx.response_phase_done, "response_phase_done should start false");
819        ctx.response_phase_done = true;
820        assert!(
821            ctx.response_phase_done,
822            "response_phase_done should be true after setting"
823        );
824    }
825
826    #[test]
827    fn set_pre_read_body() {
828        let mut ctx = default_ctx();
829        let chunks = VecDeque::from([Bytes::from_static(b"chunk1"), Bytes::from_static(b"chunk2")]);
830        ctx.pre_read_body = Some(chunks);
831        let body = ctx.pre_read_body.as_ref().unwrap();
832        assert_eq!(body.len(), 2, "pre_read_body should contain 2 chunks");
833        assert_eq!(body[0], Bytes::from_static(b"chunk1"), "first chunk should be 'chunk1'");
834        assert_eq!(
835            body[1],
836            Bytes::from_static(b"chunk2"),
837            "second chunk should be 'chunk2'"
838        );
839    }
840
841    #[test]
842    fn set_request_snapshot() {
843        let mut ctx = default_ctx();
844        let snapshot = Request {
845            method: Method::POST,
846            uri: "/api/data".parse::<Uri>().unwrap(),
847            headers: HeaderMap::new(),
848        };
849        ctx.request_snapshot = Some(snapshot);
850        let snap = ctx.request_snapshot.as_ref().unwrap();
851        assert_eq!(snap.method, Method::POST, "snapshot method should be POST");
852        assert_eq!(snap.uri.path(), "/api/data", "snapshot URI path should be /api/data");
853    }
854
855    #[test]
856    fn request_body_buffer_lifecycle() {
857        let mut ctx = default_ctx();
858        let mut buf = BodyBuffer::new(100);
859        buf.push(Bytes::from_static(b"data")).unwrap();
860        ctx.request_body_buffer = Some(buf);
861
862        assert!(
863            ctx.request_body_buffer.is_some(),
864            "buffer should be present after assignment"
865        );
866        let taken = ctx.request_body_buffer.take().unwrap();
867        assert_eq!(
868            taken.freeze(),
869            Bytes::from_static(b"data"),
870            "frozen buffer should contain pushed data"
871        );
872        assert!(ctx.request_body_buffer.is_none(), "buffer should be None after take");
873    }
874
875    #[test]
876    fn default_request_body_mode_is_stream() {
877        let ctx = default_ctx();
878        assert_eq!(
879            ctx.request_body_mode,
880            BodyMode::Stream,
881            "default request_body_mode should be Stream"
882        );
883    }
884
885    #[test]
886    fn default_response_body_mode_is_stream() {
887        let ctx = default_ctx();
888        assert_eq!(
889            ctx.response_body_mode,
890            BodyMode::Stream,
891            "default response_body_mode should be Stream"
892        );
893    }
894
895    #[test]
896    fn set_request_body_mode() {
897        let mut ctx = default_ctx();
898        ctx.request_body_mode = BodyMode::StreamBuffer { max_bytes: Some(4096) };
899        assert_eq!(
900            ctx.request_body_mode,
901            BodyMode::StreamBuffer { max_bytes: Some(4096) },
902            "request_body_mode should match assigned value"
903        );
904    }
905
906    #[test]
907    fn set_response_body_mode() {
908        let mut ctx = default_ctx();
909        ctx.response_body_mode = BodyMode::StreamBuffer { max_bytes: Some(8192) };
910        assert_eq!(
911            ctx.response_body_mode,
912            BodyMode::StreamBuffer { max_bytes: Some(8192) },
913            "response_body_mode should match assigned value"
914        );
915    }
916
917    // -------------------------------------------------------------------------
918    // Metadata Roundtrip Tests
919    // -------------------------------------------------------------------------
920
921    #[test]
922    fn metadata_roundtrip_through_filter_context() {
923        let registry = FilterRegistry::with_builtins();
924        let pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
925        let request = Request {
926            method: Method::GET,
927            uri: "/".parse::<Uri>().unwrap(),
928            headers: HeaderMap::new(),
929        };
930
931        let mut ctx = default_ctx();
932        ctx.filter_metadata.insert("rpc.method".to_owned(), "echo".to_owned());
933        ctx.filter_metadata.insert("rpc.status".to_owned(), "ok".to_owned());
934
935        let fctx = ctx.build_filter_context(&pipeline, &request, None);
936        assert_eq!(
937            fctx.get_metadata("rpc.method"),
938            Some("echo"),
939            "metadata written before build should survive into filter context"
940        );
941        assert_eq!(
942            fctx.get_metadata("rpc.status"),
943            Some("ok"),
944            "multiple metadata keys should round-trip"
945        );
946    }
947
948    #[test]
949    fn metadata_written_in_filter_context_persists_back() {
950        let registry = FilterRegistry::with_builtins();
951        let pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
952        let request = Request {
953            method: Method::GET,
954            uri: "/".parse::<Uri>().unwrap(),
955            headers: HeaderMap::new(),
956        };
957
958        let mut ctx = default_ctx();
959        let mut fctx = ctx.build_filter_context(&pipeline, &request, None);
960        fctx.set_metadata("trace.id", "abc-123");
961        fctx.set_metadata("trace.span", "42");
962
963        ctx.filter_metadata = fctx.filter_metadata;
964        assert_eq!(
965            ctx.filter_metadata.get("trace.id").map(String::as_str),
966            Some("abc-123"),
967            "metadata set in filter context should persist back to protocol context"
968        );
969        assert_eq!(
970            ctx.filter_metadata.get("trace.span").map(String::as_str),
971            Some("42"),
972            "multiple metadata keys should persist back"
973        );
974    }
975
976    // -------------------------------------------------------------------------
977    // Hot-Reload Pipeline Pinning (via production utilities)
978    // -------------------------------------------------------------------------
979
980    #[test]
981    fn pin_pipeline_captures_current_arc() {
982        let registry = FilterRegistry::with_builtins();
983        let pipeline_a = Arc::new(FilterPipeline::build(&mut [], &registry).unwrap());
984        let swap = arc_swap::ArcSwap::from(Arc::clone(&pipeline_a));
985
986        let mut ctx = default_ctx();
987        let pinned = ctx.pin_pipeline(&swap);
988        assert!(
989            Arc::ptr_eq(&pinned, &pipeline_a),
990            "pin_pipeline should return the current pipeline"
991        );
992        assert!(
993            Arc::ptr_eq(ctx.pinned_pipeline.as_ref().unwrap(), &pipeline_a),
994            "pinned_pipeline should be stored in ctx"
995        );
996    }
997
998    #[test]
999    fn stamp_error_type_is_first_write_wins() {
1000        let mut ctx = PingoraRequestCtx::default();
1001        ctx.stamp_error_type(crate::http::pingora::metrics::ERROR_TYPE_FILTER_REJECT);
1002        ctx.stamp_error_type(crate::http::pingora::metrics::ERROR_TYPE_INTERNAL);
1003        assert_eq!(
1004            ctx.error_type,
1005            Some(crate::http::pingora::metrics::ERROR_TYPE_FILTER_REJECT),
1006            "the first classification wins; a later site must not overwrite it"
1007        );
1008    }
1009
1010    #[test]
1011    fn error_type_is_unset_until_stamped() {
1012        let ctx = PingoraRequestCtx::default();
1013        assert_eq!(ctx.error_type, None, "a fresh context carries no error cause");
1014    }
1015
1016    #[test]
1017    fn pipeline_returns_pinned_after_reload() {
1018        let registry = FilterRegistry::with_builtins();
1019        let pipeline_a = Arc::new(FilterPipeline::build(&mut [], &registry).unwrap());
1020        let swap = arc_swap::ArcSwap::from(Arc::clone(&pipeline_a));
1021
1022        let mut ctx = default_ctx();
1023        ctx.pin_pipeline(&swap);
1024
1025        let pipeline_b = Arc::new(FilterPipeline::build(&mut [], &registry).unwrap());
1026        swap.store(pipeline_b);
1027
1028        let later = ctx.pipeline(&swap);
1029        assert!(
1030            Arc::ptr_eq(&later, &pipeline_a),
1031            "later hooks should still return pipeline A after reload"
1032        );
1033    }
1034
1035    #[test]
1036    fn new_request_after_reload_pins_new_pipeline() {
1037        let registry = FilterRegistry::with_builtins();
1038        let pipeline_a = Arc::new(FilterPipeline::build(&mut [], &registry).unwrap());
1039        let swap = arc_swap::ArcSwap::from(Arc::clone(&pipeline_a));
1040
1041        let mut ctx_a = default_ctx();
1042        ctx_a.pin_pipeline(&swap);
1043
1044        let pipeline_b = Arc::new(FilterPipeline::build(&mut [], &registry).unwrap());
1045        swap.store(Arc::clone(&pipeline_b));
1046
1047        let mut ctx_b = default_ctx();
1048        ctx_b.pin_pipeline(&swap);
1049
1050        assert!(
1051            Arc::ptr_eq(&ctx_a.pipeline(&swap), &pipeline_a),
1052            "request A should use old pipeline"
1053        );
1054        assert!(
1055            Arc::ptr_eq(&ctx_b.pipeline(&swap), &pipeline_b),
1056            "request B should use new pipeline"
1057        );
1058    }
1059
1060    #[test]
1061    fn old_pipeline_drops_after_ctx_drops() {
1062        let registry = FilterRegistry::with_builtins();
1063        let pipeline_a = Arc::new(FilterPipeline::build(&mut [], &registry).unwrap());
1064        let weak_a = Arc::downgrade(&pipeline_a);
1065        let swap = arc_swap::ArcSwap::from(pipeline_a);
1066
1067        let mut ctx = default_ctx();
1068        ctx.pin_pipeline(&swap);
1069
1070        swap.store(Arc::new(FilterPipeline::build(&mut [], &registry).unwrap()));
1071
1072        assert!(weak_a.upgrade().is_some(), "old pipeline alive while ctx exists");
1073
1074        drop(ctx);
1075
1076        assert!(weak_a.upgrade().is_none(), "old pipeline drops with ctx");
1077    }
1078
1079    #[test]
1080    fn pipeline_utility_returns_pinned_for_every_phase() {
1081        let registry = FilterRegistry::with_builtins();
1082        let pipeline_a = Arc::new(FilterPipeline::build(&mut [], &registry).unwrap());
1083        let swap = arc_swap::ArcSwap::from(Arc::clone(&pipeline_a));
1084
1085        let mut ctx = default_ctx();
1086        ctx.pin_pipeline(&swap);
1087
1088        swap.store(Arc::new(FilterPipeline::build(&mut [], &registry).unwrap()));
1089
1090        for phase in ["request_body", "response", "response_body", "logging"] {
1091            let p = ctx.pipeline(&swap);
1092            assert!(
1093                Arc::ptr_eq(&p, &pipeline_a),
1094                "{phase}: should still return pinned pipeline A"
1095            );
1096        }
1097    }
1098
1099    #[test]
1100    fn pipeline_fallback_when_not_pinned() {
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 ctx = default_ctx();
1106
1107        let loaded = ctx.pipeline(&swap);
1108        assert!(
1109            Arc::ptr_eq(&loaded, &pipeline_a),
1110            "unpinned ctx should fall back to current ArcSwap value"
1111        );
1112    }
1113
1114    #[test]
1115    fn pin_pipeline_is_idempotent_after_reload() {
1116        let registry = FilterRegistry::with_builtins();
1117        let pipeline_a = Arc::new(FilterPipeline::build(&mut [], &registry).unwrap());
1118        let swap = arc_swap::ArcSwap::from(Arc::clone(&pipeline_a));
1119
1120        let mut ctx = default_ctx();
1121        ctx.pin_pipeline(&swap);
1122
1123        let pipeline_b = Arc::new(FilterPipeline::build(&mut [], &registry).unwrap());
1124        swap.store(pipeline_b);
1125
1126        let second_pin = ctx.pin_pipeline(&swap);
1127        assert!(
1128            Arc::ptr_eq(&second_pin, &pipeline_a),
1129            "repeated pin_pipeline after reload should return the original pin"
1130        );
1131    }
1132
1133    #[test]
1134    fn filter_state_isolated_across_pipelines_with_same_ids() {
1135        let registry = FilterRegistry::with_builtins();
1136        let pipeline_a = Arc::new(FilterPipeline::build(&mut [], &registry).unwrap());
1137        let pipeline_b = Arc::new(FilterPipeline::build(&mut [], &registry).unwrap());
1138
1139        let request = Request {
1140            method: Method::GET,
1141            uri: "/".parse::<Uri>().unwrap(),
1142            headers: HeaderMap::new(),
1143        };
1144
1145        let mut ctx_a = default_ctx();
1146        ctx_a.pinned_pipeline = Some(Arc::clone(&pipeline_a));
1147        ctx_a.request_snapshot = Some(request.clone());
1148        let mut fctx_a = ctx_a.build_filter_context(&pipeline_a, &request, None);
1149        fctx_a.current_filter_id = Some(0);
1150        fctx_a.insert_filter_state(String::from("from_pipeline_a"));
1151        ctx_a.filter_state = fctx_a.filter_state;
1152
1153        let mut ctx_b = default_ctx();
1154        ctx_b.pinned_pipeline = Some(Arc::clone(&pipeline_b));
1155        ctx_b.request_snapshot = Some(request.clone());
1156        let fctx_b = ctx_b.build_filter_context(&pipeline_b, &request, None);
1157
1158        assert!(
1159            fctx_b.filter_state.is_empty(),
1160            "request B should have its own empty state map"
1161        );
1162
1163        let fctx_a2 = ctx_a.filter_context_for(&pipeline_a, None).unwrap();
1164        assert_eq!(
1165            fctx_a2.filter_state.get(&0).and_then(|v| v.downcast_ref::<String>()),
1166            Some(&String::from("from_pipeline_a")),
1167            "request A should still see its own state in a later phase"
1168        );
1169    }
1170
1171    #[test]
1172    fn default_ctx_has_no_adapted_body() {
1173        let ctx = PingoraRequestCtx::default();
1174        assert!(ctx.adapted_request_body.is_none(), "adapted body defaults to None");
1175        assert!(
1176            ctx.retained_adapted_request_body.is_none(),
1177            "retained adapted body defaults to None"
1178        );
1179        assert!(
1180            ctx.adapted_request_body_len.is_none(),
1181            "adapted body length defaults to None"
1182        );
1183    }
1184
1185    // -------------------------------------------------------------------------
1186    // Test Utilities
1187    // -------------------------------------------------------------------------
1188
1189    /// Create a default request context for tests.
1190    fn default_ctx() -> PingoraRequestCtx {
1191        PingoraRequestCtx::default()
1192    }
1193}