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