Skip to main content

praxis_protocol/http/pingora/handler/
with_body.rs

1// SPDX-License-Identifier: Apache-2.0
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::{FailToProxy, ProxyHttp, Session};
25use praxis_filter::{CompressionConfig, FilterPipeline};
26use tokio::sync::Semaphore;
27use tracing::{Instrument as _, debug};
28
29use super::{
30    adjust_compression, connected_to_upstream, emit_request_metrics, fail_to_proxy, handle_connect_failure,
31    hop_by_hop::RemoveHeader as _, logging_cleanup, record_passive_health, record_response_span_attributes,
32    release_retry_state, request_body_filter, request_filter, response_body_filter, response_filter, upstream_peer,
33    upstream_request, via,
34};
35use crate::http::pingora::{context::PingoraRequestCtx, metrics};
36
37// -----------------------------------------------------------------------------
38// PingoraHttpHandler
39// -----------------------------------------------------------------------------
40
41/// Pingora HTTP handler that overrides body filter hooks.
42///
43/// Used when the pipeline contains filters that declare
44/// body access via [`BodyAccess`].
45///
46/// The pipeline is held behind [`ArcSwap`] so it can be
47/// atomically replaced by hot config reload without
48/// disrupting in-flight requests.
49///
50/// ```ignore
51/// // Requires a `FilterPipeline` and Pingora server runtime.
52/// use std::sync::Arc;
53///
54/// use arc_swap::ArcSwap;
55/// use praxis_protocol::http::pingora::handler::PingoraHttpHandler;
56///
57/// let handler = PingoraHttpHandler::new(
58///     Arc::new(ArcSwap::from_pointee(pipeline)),
59///     None,
60///     None,
61///     ::metrics::SharedString::const_str("http"),
62/// );
63/// ```
64///
65/// [`BodyAccess`]: praxis_filter::BodyAccess
66/// [`ArcSwap`]: arc_swap::ArcSwap
67pub struct PingoraHttpHandler {
68    /// Compression configuration snapshot for module registration.
69    ///
70    /// Used only by [`init_downstream_modules`] to register the
71    /// compression module at startup. Per-request compression
72    /// levels are read from the live pipeline via [`ArcSwap`]
73    /// so that hot-reload updates take effect immediately.
74    ///
75    /// Module registration itself is one-shot in Pingora;
76    /// adding compression to a listener that had none at
77    /// startup requires a restart.
78    ///
79    /// [`init_downstream_modules`]: Self::init_downstream_modules
80    /// [`ArcSwap`]: arc_swap::ArcSwap
81    compression: Option<CompressionConfig>,
82
83    /// Per-listener connection semaphore for max connections.
84    connection_semaphore: Option<Arc<Semaphore>>,
85
86    /// Per-listener downstream read timeout.
87    downstream_read_timeout: Option<Duration>,
88
89    /// Listener name for connection metrics.
90    listener_name: ::metrics::SharedString,
91
92    /// Swappable filter pipeline.
93    pipeline: Arc<ArcSwap<FilterPipeline>>,
94}
95
96impl PingoraHttpHandler {
97    /// Create a handler with body filter support.
98    pub(super) fn new(
99        pipeline: Arc<ArcSwap<FilterPipeline>>,
100        downstream_read_timeout: Option<Duration>,
101        connection_semaphore: Option<Arc<Semaphore>>,
102        listener_name: ::metrics::SharedString,
103    ) -> Self {
104        let compression = pipeline.load().compression_config().cloned();
105        Self {
106            compression,
107            connection_semaphore,
108            downstream_read_timeout,
109            listener_name,
110            pipeline,
111        }
112    }
113}
114
115/// Resolve retry safety for a stale (`ReusedOnly`) upstream connection.
116///
117/// A pooled connection that closed while idle is not a real attempt, but the
118/// request bytes were already written upstream, so replay must be safe: an
119/// idempotent method (or explicit opt-in) and an intact buffered body.
120fn resolve_reused_only_retry(
121    ctx: &PingoraRequestCtx,
122    session: &Session,
123    client_reused: bool,
124    mut e: Box<pingora_core::Error>,
125) -> Box<pingora_core::Error> {
126    let policy = ctx.retry_policy.clone().unwrap_or_else(super::legacy_default_policy);
127    let replay_safe = client_reused
128        && !session.as_ref().retry_buffer_truncated()
129        && (ctx.request_is_idempotent || policy.allow_non_idempotent());
130    if !replay_safe {
131        debug!("clearing reused-connection retry: replay is not safe for this request");
132    }
133    e.set_retry(replay_safe);
134    e
135}
136
137#[async_trait]
138impl ProxyHttp for PingoraHttpHandler {
139    type CTX = PingoraRequestCtx;
140
141    fn new_ctx(&self) -> Self::CTX {
142        PingoraRequestCtx::default()
143    }
144
145    /// Registers Pingora's compression module when compression is
146    /// configured. Otherwise skips module registration to avoid
147    /// per-request `Box` allocation overhead.
148    fn init_downstream_modules(&self, modules: &mut HttpModules) {
149        if let Some(cfg) = &self.compression {
150            debug!(level = cfg.default_level, "registering compression module");
151            modules.add_module(ResponseCompressionBuilder::enable(cfg.default_level));
152        }
153    }
154
155    async fn early_request_filter(&self, session: &mut Session, ctx: &mut Self::CTX) -> Result<()>
156    where
157        Self::CTX: Send + Sync,
158    {
159        if praxis_core::memory::is_exceeded() {
160            metrics::record_overload_reject(metrics::OVERLOAD_REASON_MEMORY);
161            return reject_503(session, "5", "memory pressure exceeded").await;
162        }
163
164        let (exceeded, permit) = crate::connections::try_acquire_global();
165        ctx._global_connection_permit = permit;
166        if exceeded {
167            metrics::record_overload_reject(metrics::OVERLOAD_REASON_GLOBAL_CONNECTIONS);
168            return reject_503(session, "1", "global max connections exceeded").await;
169        }
170
171        if let Some(sem) = &self.connection_semaphore {
172            if let Ok(permit) = Arc::clone(sem).try_acquire_owned() {
173                ctx._connection_permit = Some(permit);
174            } else {
175                metrics::record_overload_reject(metrics::OVERLOAD_REASON_LISTENER_CONNECTIONS);
176                return reject_503(session, "1", "max connections exceeded").await;
177            }
178        }
179
180        ctx._active_connection = Some(metrics::ActiveConnectionGuard::acquire(self.listener_name.clone()));
181
182        if let Some(timeout) = self.downstream_read_timeout {
183            debug!(
184                timeout_ms = u64::try_from(timeout.as_millis()).unwrap_or(u64::MAX),
185                "applying downstream read timeout"
186            );
187            session.set_read_timeout(Some(timeout));
188        }
189        Ok(())
190    }
191
192    async fn request_filter(&self, session: &mut Session, ctx: &mut Self::CTX) -> Result<bool> {
193        let pipeline = ctx.pin_pipeline(&self.pipeline);
194        request_filter::execute(&pipeline, session, ctx).await
195    }
196
197    async fn request_body_filter(
198        &self,
199        session: &mut Session,
200        body: &mut Option<Bytes>,
201        end_of_stream: bool,
202        ctx: &mut Self::CTX,
203    ) -> Result<()>
204    where
205        Self::CTX: Send + Sync,
206    {
207        // Pure no-op fast path (mirroring execute's own early returns,
208        // which emit nothing): skip the pipeline Arc clone, span clone,
209        // and future instrumentation per chunk when no body filter can
210        // run — the default configuration for proxied bodies.
211        if ctx.connection_upgraded {
212            return Ok(());
213        }
214        if ctx.pre_read_body.is_none()
215            && let Some(pinned) = &ctx.pinned_pipeline
216            && !pinned.body_capabilities().needs_request_body
217        {
218            return Ok(());
219        }
220        let pipeline = ctx.pipeline(&self.pipeline);
221        let span = ctx.request_span.clone();
222        request_body_filter::execute(&pipeline, session, body, end_of_stream, ctx)
223            .instrument(span)
224            .await
225    }
226
227    fn response_body_filter(
228        &self,
229        _session: &mut Session,
230        body: &mut Option<Bytes>,
231        end_of_stream: bool,
232        ctx: &mut Self::CTX,
233    ) -> Result<Option<Duration>>
234    where
235        Self::CTX: Send + Sync,
236    {
237        // Same silent fast path as the request side, keeping execute's
238        // delivery-complete bookkeeping.
239        if ctx.connection_upgraded {
240            return Ok(None);
241        }
242        if let Some(pinned) = &ctx.pinned_pipeline
243            && !pinned.body_capabilities().needs_response_body
244        {
245            if end_of_stream {
246                ctx.response_delivery_complete = true;
247            }
248            return Ok(None);
249        }
250        let span = ctx.request_span.clone();
251        let _entered = span.enter();
252        let pipeline = ctx.pipeline(&self.pipeline);
253        response_body_filter::execute(&pipeline, body, end_of_stream, ctx)
254    }
255
256    fn fail_to_connect(
257        &self,
258        session: &mut Session,
259        _peer: &HttpPeer,
260        ctx: &mut Self::CTX,
261        e: Box<pingora_core::Error>,
262    ) -> Box<pingora_core::Error> {
263        let span = ctx.request_span.clone();
264        let _entered = span.enter();
265        // A truncated replay buffer means a retry would resend a partial
266        // request body — refuse rather than corrupt the request upstream.
267        if session.as_mut().retry_buffer_truncated() {
268            let mut e = e;
269            e.set_retry(false);
270            return e;
271        }
272        handle_connect_failure(ctx, e)
273    }
274
275    fn error_while_proxy(
276        &self,
277        peer: &HttpPeer,
278        session: &mut Session,
279        e: Box<pingora_core::Error>,
280        ctx: &mut Self::CTX,
281        client_reused: bool,
282    ) -> Box<pingora_core::Error> {
283        // Never retry once a final response reached the client: a second
284        // attempt cannot rewrite the response line, so its body would be
285        // spliced after whatever was already sent. A non-final 1xx (e.g.
286        // 100 Continue) does not commit the final response, so it must not
287        // block a retry; mirror fail_to_proxy's final-response predicate.
288        if session.response_written().is_some_and(fail_to_proxy::is_final_response) {
289            let mut e = e;
290            e.set_retry(false);
291            return e;
292        }
293        // A truncated replay buffer means a retry would resend a partial
294        // request body — silent request corruption.
295        let truncated = session.as_mut().retry_buffer_truncated();
296        // Preserve an explicit retry decision from the response-status path
297        // (already validated by the policy engine) — but still refuse it if the
298        // replay buffer was truncated. should_retry's body-size guard makes
299        // this unreachable while retry_body_limit_bytes stays capped at
300        // Pingora's replay-buffer size, but that invariant lives in config
301        // validation, not here; guarding locally keeps the two other retry
302        // paths' replay-safety property from resting on an external cap.
303        if matches!(e.retry, pingora_core::RetryType::Decided(true)) {
304            if truncated {
305                let mut e = e;
306                e.set_retry(false);
307                return e;
308            }
309            return e;
310        }
311        // Stale-connection (ReusedOnly) errors skip the retry budget but still
312        // require replay safety; see resolve_reused_only_retry.
313        // See docs/architecture/http-correctness.md.
314        if matches!(e.retry, pingora_core::RetryType::ReusedOnly) {
315            return resolve_reused_only_retry(ctx, session, client_reused, e);
316        }
317        let e = e.more_context(format!("Peer: {peer}"));
318        if truncated {
319            let mut e = e;
320            e.set_retry(false);
321            return e;
322        }
323        // Mid-proxy errors (reset, refused stream, etc.) go through the same
324        // policy engine as connect failures — Pingora's decide_reuse must not
325        // bypass idempotency / budget / max_retries guards.
326        handle_connect_failure(ctx, e)
327    }
328
329    async fn fail_to_proxy(&self, session: &mut Session, e: &pingora_core::Error, ctx: &mut Self::CTX) -> FailToProxy
330    where
331        Self::CTX: Send + Sync,
332    {
333        let span = ctx.request_span.clone();
334        fail_to_proxy::execute(session, e, ctx).instrument(span).await
335    }
336
337    async fn upstream_request_filter(
338        &self,
339        session: &mut Session,
340        upstream_request: &mut pingora_http::RequestHeader,
341        ctx: &mut Self::CTX,
342    ) -> Result<()>
343    where
344        Self::CTX: Send + Sync,
345    {
346        let span = ctx.request_span.clone();
347        let _entered = span.enter();
348        let is_upgrade = session.is_upgrade_req();
349        upstream_request::strip_hop_by_hop(upstream_request, is_upgrade);
350        upstream_request.strip_reserved_internal();
351        upstream_request::apply_authority_override(upstream_request, ctx)?;
352        upstream_request::apply_rewritten_path(upstream_request, ctx)?;
353        upstream_request::apply_mutated_content_length(upstream_request, ctx);
354        let client_ver = ctx.client_http_version.unwrap_or(http::Version::HTTP_11);
355        via::append_request_via(upstream_request, client_ver);
356        Ok(())
357    }
358
359    async fn response_filter(
360        &self,
361        session: &mut Session,
362        upstream_response: &mut pingora_http::ResponseHeader,
363        ctx: &mut Self::CTX,
364    ) -> Result<()>
365    where
366        Self::CTX: Send + Sync,
367    {
368        let pipeline = ctx.pipeline(&self.pipeline);
369        let span = ctx.request_span.clone();
370        let exchange_span = ctx.upstream_exchange_span.clone();
371        let result = response_filter::execute(&pipeline, upstream_response, ctx)
372            .instrument(exchange_span)
373            .instrument(span)
374            .await;
375        if result.is_ok() {
376            // RFC 9110 §7.6.3: the response Via received-protocol is the leg this
377            // proxy received the response on — the upstream connection — not the
378            // downstream client's version.
379            let upstream_ver = upstream_response.version;
380            via::append_response_via(upstream_response, upstream_ver);
381            adjust_compression(session, upstream_response, pipeline.compression_config());
382        }
383        result
384    }
385
386    async fn upstream_peer(&self, _session: &mut Session, ctx: &mut Self::CTX) -> Result<Box<HttpPeer>> {
387        let span = ctx.request_span.clone();
388        upstream_peer::execute(ctx).instrument(span).await
389    }
390
391    async fn connected_to_upstream(
392        &self,
393        _session: &mut Session,
394        reused: bool,
395        peer: &HttpPeer,
396        #[cfg(unix)] _fd: std::os::unix::io::RawFd,
397        #[cfg(windows)] _sock: std::os::windows::io::RawSocket,
398        digest: Option<&pingora_core::protocols::Digest>,
399        ctx: &mut Self::CTX,
400    ) -> Result<()>
401    where
402        Self::CTX: Send + Sync,
403    {
404        let span = ctx.request_span.clone();
405        let _entered = span.enter();
406        let cluster = ctx.metrics_cluster_shared.clone().unwrap_or_else(metrics::cluster_none);
407        if !reused && let Some(start) = ctx.upstream_connect_start.take() {
408            metrics::record_upstream_connect_duration(cluster.clone(), start.elapsed().as_secs_f64());
409        }
410        if ctx.retries > 0 {
411            metrics::record_upstream_retry(cluster, metrics::RETRY_RESULT_SUCCESS);
412        }
413        connected_to_upstream::execute(reused, peer, digest, ctx);
414        Ok(())
415    }
416
417    async fn logging(&self, session: &mut Session, e: Option<&pingora_core::Error>, ctx: &mut Self::CTX) {
418        record_response_span_attributes(session, ctx);
419        // Drop the exchange span before the request span so child
420        // ends before parent in tracing output.
421        let _exchange_span = std::mem::replace(&mut ctx.upstream_exchange_span, tracing::Span::none());
422        drop(_exchange_span);
423        let span = std::mem::replace(&mut ctx.request_span, tracing::Span::none());
424        let written_status = session.response_written().map_or(0, |resp| resp.status.as_u16());
425        async {
426            let pipeline = ctx.pipeline(&self.pipeline);
427            emit_request_metrics(session, ctx);
428            record_passive_health(&pipeline, e, ctx);
429            release_retry_state(ctx);
430            logging_cleanup(&pipeline, ctx).await;
431            super::maybe_emit_fallback_access_log(&pipeline, written_status, ctx);
432        }
433        .instrument(span)
434        .await;
435    }
436}
437
438// -----------------------------------------------------------------------------
439// Utilities
440// -----------------------------------------------------------------------------
441
442/// Write a 503 response with `Retry-After` and return the corresponding error.
443async fn reject_503(session: &mut Session, retry_after: &'static str, reason: &'static str) -> Result<()> {
444    tracing::warn!(reason, "rejecting request");
445    let mut header = pingora_http::ResponseHeader::build(503, None)?;
446    header.append_header("Retry-After", retry_after)?;
447    session.write_response_header(Box::new(header), true).await?;
448    Err(pingora_core::Error::explain(
449        pingora_core::ErrorType::HTTPStatus(503),
450        reason,
451    ))
452}