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