Skip to main content

praxis_protocol/http/pingora/handler/
with_body.rs

1// SPDX-License-Identifier: MIT
2// Copyright (c) 2024 Praxis Contributors
3
4//! Pingora HTTP handler with body filter hooks enabled.
5//!
6//! [`PingoraHttpHandler`] is the full-featured `ProxyHttp` implementation
7//! used when the pipeline's [`BodyCapabilities`] declare request or
8//! response body access. It delegates each Pingora lifecycle hook to
9//! the corresponding submodule and enables Pingora's compression
10//! module when a compression filter is configured.
11//!
12//! [`BodyCapabilities`]: praxis_filter::body::BodyCapabilities
13
14use std::{sync::Arc, time::Duration};
15
16use arc_swap::ArcSwap;
17use async_trait::async_trait;
18use bytes::Bytes;
19use pingora_core::{
20    Result,
21    modules::http::{HttpModules, compression::ResponseCompressionBuilder},
22    upstreams::peer::HttpPeer,
23};
24use pingora_proxy::{ProxyHttp, Session};
25use praxis_filter::{CompressionConfig, FilterPipeline};
26use tokio::sync::Semaphore;
27use tracing::debug;
28
29use super::{
30    adjust_compression, emit_request_metrics, handle_connect_failure, logging_cleanup, record_passive_health,
31    request_body_filter, request_filter, response_body_filter, response_filter, upstream_peer, upstream_request, via,
32};
33use crate::http::pingora::context::PingoraRequestCtx;
34
35// -----------------------------------------------------------------------------
36// PingoraHttpHandler
37// -----------------------------------------------------------------------------
38
39/// Pingora HTTP handler that overrides body filter hooks.
40///
41/// Used when the pipeline contains filters that declare
42/// body access via [`BodyAccess`].
43///
44/// The pipeline is held behind [`ArcSwap`] so it can be
45/// atomically replaced by hot config reload without
46/// disrupting in-flight requests.
47///
48/// ```ignore
49/// // Requires a `FilterPipeline` and Pingora server runtime.
50/// use std::sync::Arc;
51///
52/// use arc_swap::ArcSwap;
53/// use praxis_protocol::http::pingora::handler::PingoraHttpHandler;
54///
55/// let handler = PingoraHttpHandler::new(
56///     Arc::new(ArcSwap::from_pointee(pipeline)),
57///     None,
58///     None,
59/// );
60/// ```
61///
62/// [`BodyAccess`]: praxis_filter::BodyAccess
63/// [`ArcSwap`]: arc_swap::ArcSwap
64pub struct PingoraHttpHandler {
65    /// Compression configuration snapshot for module registration.
66    ///
67    /// Used only by [`init_downstream_modules`] to register the
68    /// compression module at startup. Per-request compression
69    /// levels are read from the live pipeline via [`ArcSwap`]
70    /// so that hot-reload updates take effect immediately.
71    ///
72    /// Module registration itself is one-shot in Pingora;
73    /// adding compression to a listener that had none at
74    /// startup requires a restart.
75    ///
76    /// [`init_downstream_modules`]: Self::init_downstream_modules
77    /// [`ArcSwap`]: arc_swap::ArcSwap
78    compression: Option<CompressionConfig>,
79
80    /// Per-listener connection semaphore for max connections.
81    connection_semaphore: Option<Arc<Semaphore>>,
82
83    /// Per-listener downstream read timeout.
84    downstream_read_timeout: Option<Duration>,
85
86    /// Swappable filter pipeline.
87    pipeline: Arc<ArcSwap<FilterPipeline>>,
88}
89
90impl PingoraHttpHandler {
91    /// Create a handler with body filter support.
92    pub(super) fn new(
93        pipeline: Arc<ArcSwap<FilterPipeline>>,
94        downstream_read_timeout: Option<Duration>,
95        connection_semaphore: Option<Arc<Semaphore>>,
96    ) -> Self {
97        let compression = pipeline.load().compression_config().cloned();
98        Self {
99            compression,
100            connection_semaphore,
101            downstream_read_timeout,
102            pipeline,
103        }
104    }
105}
106
107#[async_trait]
108impl ProxyHttp for PingoraHttpHandler {
109    type CTX = PingoraRequestCtx;
110
111    fn new_ctx(&self) -> Self::CTX {
112        PingoraRequestCtx::default()
113    }
114
115    /// Registers Pingora's compression module when compression is
116    /// configured. Otherwise skips module registration to avoid
117    /// per-request `Box` allocation overhead.
118    fn init_downstream_modules(&self, modules: &mut HttpModules) {
119        if let Some(cfg) = &self.compression {
120            debug!(level = cfg.default_level, "registering compression module");
121            modules.add_module(ResponseCompressionBuilder::enable(cfg.default_level));
122        }
123    }
124
125    #[expect(clippy::cast_possible_truncation, reason = "millis fit u64")]
126    async fn early_request_filter(&self, session: &mut Session, ctx: &mut Self::CTX) -> Result<()>
127    where
128        Self::CTX: Send + Sync,
129    {
130        if praxis_core::memory::is_exceeded() {
131            return reject_503(session, "5", "memory pressure exceeded").await;
132        }
133
134        let (exceeded, permit) = crate::connections::try_acquire_global();
135        ctx._global_connection_permit = permit;
136        if exceeded {
137            return reject_503(session, "1", "global max connections exceeded").await;
138        }
139
140        if let Some(sem) = &self.connection_semaphore {
141            if let Ok(permit) = Arc::clone(sem).try_acquire_owned() {
142                ctx._connection_permit = Some(permit);
143            } else {
144                return reject_503(session, "1", "max connections exceeded").await;
145            }
146        }
147
148        if let Some(timeout) = self.downstream_read_timeout {
149            debug!(
150                timeout_ms = timeout.as_millis() as u64,
151                "applying downstream read timeout"
152            );
153            session.set_read_timeout(Some(timeout));
154        }
155        Ok(())
156    }
157
158    async fn request_filter(&self, session: &mut Session, ctx: &mut Self::CTX) -> Result<bool> {
159        let pipeline = ctx.pin_pipeline(&self.pipeline);
160        request_filter::execute(&pipeline, session, ctx).await
161    }
162
163    async fn request_body_filter(
164        &self,
165        session: &mut Session,
166        body: &mut Option<Bytes>,
167        end_of_stream: bool,
168        ctx: &mut Self::CTX,
169    ) -> Result<()>
170    where
171        Self::CTX: Send + Sync,
172    {
173        let pipeline = ctx.pipeline(&self.pipeline);
174        request_body_filter::execute(&pipeline, session, body, end_of_stream, ctx).await
175    }
176
177    fn response_body_filter(
178        &self,
179        _session: &mut Session,
180        body: &mut Option<Bytes>,
181        end_of_stream: bool,
182        ctx: &mut Self::CTX,
183    ) -> Result<Option<Duration>>
184    where
185        Self::CTX: Send + Sync,
186    {
187        let pipeline = ctx.pipeline(&self.pipeline);
188        response_body_filter::execute(&pipeline, body, end_of_stream, ctx)
189    }
190
191    fn fail_to_connect(
192        &self,
193        _session: &mut Session,
194        _peer: &HttpPeer,
195        ctx: &mut Self::CTX,
196        e: Box<pingora_core::Error>,
197    ) -> Box<pingora_core::Error> {
198        handle_connect_failure(ctx, e)
199    }
200
201    async fn upstream_request_filter(
202        &self,
203        session: &mut Session,
204        upstream_request: &mut pingora_http::RequestHeader,
205        ctx: &mut Self::CTX,
206    ) -> Result<()>
207    where
208        Self::CTX: Send + Sync,
209    {
210        let is_upgrade = session.is_upgrade_req();
211        upstream_request::strip_hop_by_hop(upstream_request, is_upgrade);
212        upstream_request::strip_reserved_internal(upstream_request);
213        upstream_request::apply_rewritten_path(upstream_request, ctx)?;
214        upstream_request::apply_mutated_content_length(upstream_request, ctx);
215        let client_ver = ctx.client_http_version.unwrap_or(http::Version::HTTP_11);
216        via::append_request_via(upstream_request, client_ver);
217        Ok(())
218    }
219
220    async fn response_filter(
221        &self,
222        session: &mut Session,
223        upstream_response: &mut pingora_http::ResponseHeader,
224        ctx: &mut Self::CTX,
225    ) -> Result<()>
226    where
227        Self::CTX: Send + Sync,
228    {
229        let pipeline = ctx.pipeline(&self.pipeline);
230        let result = response_filter::execute(&pipeline, upstream_response, ctx).await;
231        if result.is_ok() {
232            let client_ver = ctx.client_http_version.unwrap_or(http::Version::HTTP_11);
233            via::append_response_via(upstream_response, client_ver);
234            adjust_compression(session, upstream_response, pipeline.compression_config());
235        }
236        result
237    }
238
239    async fn upstream_peer(&self, _session: &mut Session, ctx: &mut Self::CTX) -> Result<Box<HttpPeer>> {
240        upstream_peer::execute(ctx).await
241    }
242
243    async fn logging(&self, session: &mut Session, e: Option<&pingora_core::Error>, ctx: &mut Self::CTX) {
244        let pipeline = ctx.pipeline(&self.pipeline);
245        emit_request_metrics(session, ctx);
246        record_passive_health(&pipeline, e, ctx);
247        logging_cleanup(&pipeline, ctx).await;
248    }
249}
250
251// -----------------------------------------------------------------------------
252// Utilities
253// -----------------------------------------------------------------------------
254
255/// Write a 503 response with `Retry-After` and return the corresponding error.
256async fn reject_503(session: &mut Session, retry_after: &'static str, reason: &'static str) -> Result<()> {
257    tracing::warn!(reason, "rejecting request");
258    let mut header = pingora_http::ResponseHeader::build(503, None)?;
259    header.append_header("Retry-After", retry_after)?;
260    session.write_response_header(Box::new(header), true).await?;
261    Err(pingora_core::Error::explain(
262        pingora_core::ErrorType::HTTPStatus(503),
263        reason,
264    ))
265}