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        // Defensively clear upstream-contact state at the first per-request hook.
160        // The current Pingora fork builds a fresh context per request
161        // (persist_connection_context is off), so nothing leaks across keep-alive
162        // requests today; this is belt-and-suspenders that keeps the invariant
163        // true before every rejection path (the 503s below and the early exits in
164        // request_filter) should context reuse ever be enabled.
165        ctx.upstream_for_retry = None;
166        ctx.upstream_contacted = false;
167
168        if praxis_core::memory::is_exceeded() {
169            metrics::record_overload_reject(metrics::OVERLOAD_REASON_MEMORY);
170            return reject_503(session, "5", "memory pressure exceeded").await;
171        }
172
173        let (exceeded, permit) = crate::connections::try_acquire_global();
174        ctx._global_connection_permit = permit;
175        if exceeded {
176            metrics::record_overload_reject(metrics::OVERLOAD_REASON_GLOBAL_CONNECTIONS);
177            return reject_503(session, "1", "global max connections exceeded").await;
178        }
179
180        if let Some(sem) = &self.connection_semaphore {
181            if let Ok(permit) = Arc::clone(sem).try_acquire_owned() {
182                ctx._connection_permit = Some(permit);
183            } else {
184                metrics::record_overload_reject(metrics::OVERLOAD_REASON_LISTENER_CONNECTIONS);
185                return reject_503(session, "1", "max connections exceeded").await;
186            }
187        }
188
189        ctx._active_request = Some(metrics::ActiveRequestGuard::acquire(self.listener_name.clone()));
190
191        if let Some(timeout) = self.downstream_read_timeout {
192            debug!(
193                timeout_ms = u64::try_from(timeout.as_millis()).unwrap_or(u64::MAX),
194                "applying downstream read timeout"
195            );
196            session.set_read_timeout(Some(timeout));
197        }
198        Ok(())
199    }
200
201    async fn request_filter(&self, session: &mut Session, ctx: &mut Self::CTX) -> Result<bool> {
202        let pipeline = ctx.pin_pipeline(&self.pipeline);
203        request_filter::execute(&pipeline, session, ctx).await
204    }
205
206    async fn request_body_filter(
207        &self,
208        session: &mut Session,
209        body: &mut Option<Bytes>,
210        end_of_stream: bool,
211        ctx: &mut Self::CTX,
212    ) -> Result<()>
213    where
214        Self::CTX: Send + Sync,
215    {
216        // Pure no-op fast path (mirroring execute's own early returns,
217        // which emit nothing): skip the pipeline Arc clone, span clone,
218        // and future instrumentation per chunk when no body filter can
219        // run — the default configuration for proxied bodies.
220        if ctx.connection_upgraded {
221            return Ok(());
222        }
223        if ctx.pre_read_body.is_none()
224            && let Some(pinned) = &ctx.pinned_pipeline
225            && !pinned.body_capabilities().needs_request_body
226        {
227            return Ok(());
228        }
229        let pipeline = ctx.pipeline(&self.pipeline);
230        let span = ctx.request_span.clone();
231        request_body_filter::execute(&pipeline, session, body, end_of_stream, ctx)
232            .instrument(span)
233            .await
234    }
235
236    fn response_body_filter(
237        &self,
238        _session: &mut Session,
239        body: &mut Option<Bytes>,
240        end_of_stream: bool,
241        ctx: &mut Self::CTX,
242    ) -> Result<Option<Duration>>
243    where
244        Self::CTX: Send + Sync,
245    {
246        // Same silent fast path as the request side, keeping execute's
247        // delivery-complete bookkeeping.
248        if ctx.connection_upgraded {
249            return Ok(None);
250        }
251        if let Some(pinned) = &ctx.pinned_pipeline
252            && !pinned.body_capabilities().needs_response_body
253        {
254            if end_of_stream {
255                ctx.response_delivery_complete = true;
256            }
257            return Ok(None);
258        }
259        let span = ctx.request_span.clone();
260        let _entered = span.enter();
261        let pipeline = ctx.pipeline(&self.pipeline);
262        response_body_filter::execute(&pipeline, body, end_of_stream, ctx)
263    }
264
265    fn fail_to_connect(
266        &self,
267        session: &mut Session,
268        _peer: &HttpPeer,
269        ctx: &mut Self::CTX,
270        e: Box<pingora_core::Error>,
271    ) -> Box<pingora_core::Error> {
272        let span = ctx.request_span.clone();
273        let _entered = span.enter();
274        // A truncated replay buffer means a retry would resend a partial
275        // request body — refuse rather than corrupt the request upstream.
276        if session.as_mut().retry_buffer_truncated() {
277            let mut e = e;
278            e.set_retry(false);
279            return e;
280        }
281        handle_connect_failure(ctx, e)
282    }
283
284    fn error_while_proxy(
285        &self,
286        peer: &HttpPeer,
287        session: &mut Session,
288        e: Box<pingora_core::Error>,
289        ctx: &mut Self::CTX,
290        client_reused: bool,
291    ) -> Box<pingora_core::Error> {
292        // Never retry once a final response reached the client: a second
293        // attempt cannot rewrite the response line, so its body would be
294        // spliced after whatever was already sent. A non-final 1xx (e.g.
295        // 100 Continue) does not commit the final response, so it must not
296        // block a retry; mirror fail_to_proxy's final-response predicate.
297        if session.response_written().is_some_and(fail_to_proxy::is_final_response) {
298            let mut e = e;
299            e.set_retry(false);
300            return e;
301        }
302        // A truncated replay buffer means a retry would resend a partial
303        // request body — silent request corruption.
304        let truncated = session.as_mut().retry_buffer_truncated();
305        // Preserve an explicit retry decision from the response-status path
306        // (already validated by the policy engine) — but still refuse it if the
307        // replay buffer was truncated. should_retry's body-size guard makes
308        // this unreachable while retry_body_limit_bytes stays capped at
309        // Pingora's replay-buffer size, but that invariant lives in config
310        // validation, not here; guarding locally keeps the two other retry
311        // paths' replay-safety property from resting on an external cap.
312        if matches!(e.retry, pingora_core::RetryType::Decided(true)) {
313            if truncated {
314                let mut e = e;
315                e.set_retry(false);
316                return e;
317            }
318            return e;
319        }
320        // Stale-connection (ReusedOnly) errors skip the retry budget but still
321        // require replay safety; see resolve_reused_only_retry.
322        // See docs/architecture/http-correctness.md.
323        if matches!(e.retry, pingora_core::RetryType::ReusedOnly) {
324            return resolve_reused_only_retry(ctx, session, client_reused, e);
325        }
326        let e = e.more_context(format!("Peer: {peer}"));
327        if truncated {
328            let mut e = e;
329            e.set_retry(false);
330            return e;
331        }
332        // Mid-proxy errors (reset, refused stream, etc.) go through the same
333        // policy engine as connect failures — Pingora's decide_reuse must not
334        // bypass idempotency / budget / max_retries guards.
335        handle_connect_failure(ctx, e)
336    }
337
338    async fn fail_to_proxy(&self, session: &mut Session, e: &pingora_core::Error, ctx: &mut Self::CTX) -> FailToProxy
339    where
340        Self::CTX: Send + Sync,
341    {
342        let span = ctx.request_span.clone();
343        fail_to_proxy::execute(session, e, ctx).instrument(span).await
344    }
345
346    async fn upstream_request_filter(
347        &self,
348        session: &mut Session,
349        upstream_request: &mut pingora_http::RequestHeader,
350        ctx: &mut Self::CTX,
351    ) -> Result<()>
352    where
353        Self::CTX: Send + Sync,
354    {
355        let span = ctx.request_span.clone();
356        let _entered = span.enter();
357        let is_upgrade = session.is_upgrade_req();
358        upstream_request::strip_hop_by_hop(upstream_request, is_upgrade);
359        upstream_request.strip_reserved_internal();
360        upstream_request::apply_authority_override(upstream_request, ctx)?;
361        upstream_request::apply_rewritten_path(upstream_request, ctx)?;
362        upstream_request::apply_mutated_content_length(upstream_request, ctx);
363        // Runs once per attempt: on a retry, re-seed the retained mutated body
364        // so the replayed bytes match the re-stamped Content-Length above (a
365        // no-op on the first attempt and when no body writer ran).
366        upstream_request::reseed_retry_body(ctx);
367        let client_ver = ctx.client_http_version.unwrap_or(http::Version::HTTP_11);
368        via::append_request_via(upstream_request, client_ver);
369        Ok(())
370    }
371
372    async fn response_filter(
373        &self,
374        session: &mut Session,
375        upstream_response: &mut pingora_http::ResponseHeader,
376        ctx: &mut Self::CTX,
377    ) -> Result<()>
378    where
379        Self::CTX: Send + Sync,
380    {
381        let pipeline = ctx.pipeline(&self.pipeline);
382        let span = ctx.request_span.clone();
383        let exchange_span = ctx.upstream_exchange_span.clone();
384        let result = response_filter::execute(&pipeline, upstream_response, ctx)
385            .instrument(exchange_span)
386            .instrument(span)
387            .await;
388        if result.is_ok() {
389            // RFC 9110 §7.6.3: the response Via received-protocol is the leg this
390            // proxy received the response on — the upstream connection — not the
391            // downstream client's version.
392            let upstream_ver = upstream_response.version;
393            via::append_response_via(upstream_response, upstream_ver);
394            adjust_compression(session, upstream_response, pipeline.compression_config());
395        }
396        result
397    }
398
399    async fn upstream_peer(&self, _session: &mut Session, ctx: &mut Self::CTX) -> Result<Box<HttpPeer>> {
400        let span = ctx.request_span.clone();
401        upstream_peer::execute(ctx).instrument(span).await
402    }
403
404    async fn connected_to_upstream(
405        &self,
406        _session: &mut Session,
407        reused: bool,
408        peer: &HttpPeer,
409        #[cfg(unix)] _fd: std::os::unix::io::RawFd,
410        #[cfg(windows)] _sock: std::os::windows::io::RawSocket,
411        digest: Option<&pingora_core::protocols::Digest>,
412        ctx: &mut Self::CTX,
413    ) -> Result<()>
414    where
415        Self::CTX: Send + Sync,
416    {
417        let span = ctx.request_span.clone();
418        let _entered = span.enter();
419        let cluster = ctx.metrics_cluster_shared.clone().unwrap_or_else(metrics::cluster_none);
420        if !reused && let Some(start) = ctx.upstream_connect_start.take() {
421            metrics::record_upstream_connect_duration(cluster.clone(), start.elapsed().as_secs_f64());
422        }
423        if ctx.retries > 0 {
424            metrics::record_upstream_retry(cluster, metrics::RETRY_RESULT_SUCCESS);
425        }
426        connected_to_upstream::execute(reused, peer, digest, ctx);
427        Ok(())
428    }
429
430    async fn logging(&self, session: &mut Session, e: Option<&pingora_core::Error>, ctx: &mut Self::CTX) {
431        record_response_span_attributes(session, ctx);
432        // Drop the exchange span before the request span so child
433        // ends before parent in tracing output.
434        let _exchange_span = std::mem::replace(&mut ctx.upstream_exchange_span, tracing::Span::none());
435        drop(_exchange_span);
436        let span = std::mem::replace(&mut ctx.request_span, tracing::Span::none());
437        let written_status = session.response_written().map_or(0, |resp| resp.status.as_u16());
438        async {
439            let pipeline = ctx.pipeline(&self.pipeline);
440            emit_request_metrics(session, ctx);
441            record_passive_health(&pipeline, e, ctx);
442            release_retry_state(ctx);
443            logging_cleanup(&pipeline, ctx).await;
444            super::maybe_emit_fallback_access_log(&pipeline, written_status, ctx);
445        }
446        .instrument(span)
447        .await;
448    }
449}
450
451// -----------------------------------------------------------------------------
452// Utilities
453// -----------------------------------------------------------------------------
454
455/// Write a 503 response with `Retry-After` and return the corresponding error.
456async fn reject_503(session: &mut Session, retry_after: &'static str, reason: &'static str) -> Result<()> {
457    tracing::warn!(reason, "rejecting request");
458    let mut header = pingora_http::ResponseHeader::build(503, None)?;
459    header.append_header("Retry-After", retry_after)?;
460    session.write_response_header(Box::new(header), true).await?;
461    Err(pingora_core::Error::explain(
462        pingora_core::ErrorType::HTTPStatus(503),
463        reason,
464    ))
465}