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