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