praxis_protocol/http/pingora/handler/
no_body.rs1use std::{sync::Arc, time::Duration};
11
12use arc_swap::ArcSwap;
13use async_trait::async_trait;
14use pingora_core::{
15 Result,
16 modules::http::{HttpModules, compression::ResponseCompressionBuilder},
17 upstreams::peer::HttpPeer,
18};
19use pingora_proxy::{ProxyHttp, Session};
20use praxis_filter::{CompressionConfig, FilterPipeline};
21use tokio::sync::Semaphore;
22use tracing::debug;
23
24use super::{
25 adjust_compression, emit_request_metrics, handle_connect_failure, logging_cleanup, record_passive_health,
26 request_filter, response_filter, upstream_peer, upstream_request, via,
27};
28use crate::http::pingora::context::PingoraRequestCtx;
29
30pub struct PingoraHttpHandlerNoBody {
47 compression: Option<CompressionConfig>,
53
54 connection_semaphore: Option<Arc<Semaphore>>,
56
57 downstream_read_timeout: Option<Duration>,
59
60 pipeline: Arc<ArcSwap<FilterPipeline>>,
62}
63
64impl PingoraHttpHandlerNoBody {
65 #[expect(dead_code, reason = "reserved for non-reload paths")]
67 pub(super) fn new(
68 pipeline: Arc<ArcSwap<FilterPipeline>>,
69 downstream_read_timeout: Option<Duration>,
70 connection_semaphore: Option<Arc<Semaphore>>,
71 ) -> Self {
72 let compression = pipeline.load().compression_config().cloned();
73 Self {
74 compression,
75 connection_semaphore,
76 downstream_read_timeout,
77 pipeline,
78 }
79 }
80}
81
82#[async_trait]
83impl ProxyHttp for PingoraHttpHandlerNoBody {
84 type CTX = PingoraRequestCtx;
85
86 fn new_ctx(&self) -> Self::CTX {
87 PingoraRequestCtx::default()
88 }
89
90 fn init_downstream_modules(&self, modules: &mut HttpModules) {
93 if let Some(cfg) = &self.compression {
94 debug!(level = cfg.default_level, "registering compression module");
95 modules.add_module(ResponseCompressionBuilder::enable(cfg.default_level));
96 }
97 }
98
99 #[expect(clippy::cast_possible_truncation, reason = "millis fit u64")]
100 async fn early_request_filter(&self, session: &mut Session, ctx: &mut Self::CTX) -> Result<()>
101 where
102 Self::CTX: Send + Sync,
103 {
104 if praxis_core::memory::is_exceeded() {
105 return reject_503(session, "5", "memory pressure exceeded").await;
106 }
107
108 let (exceeded, permit) = crate::connections::try_acquire_global();
109 ctx._global_connection_permit = permit;
110 if exceeded {
111 return reject_503(session, "1", "global max connections exceeded").await;
112 }
113
114 if let Some(sem) = &self.connection_semaphore {
115 if let Ok(permit) = Arc::clone(sem).try_acquire_owned() {
116 ctx._connection_permit = Some(permit);
117 } else {
118 return reject_503(session, "1", "max connections exceeded").await;
119 }
120 }
121
122 if let Some(timeout) = self.downstream_read_timeout {
123 debug!(
124 timeout_ms = timeout.as_millis() as u64,
125 "applying downstream read timeout"
126 );
127 session.set_read_timeout(Some(timeout));
128 }
129 Ok(())
130 }
131
132 async fn request_filter(&self, session: &mut Session, ctx: &mut Self::CTX) -> Result<bool> {
133 let pipeline = ctx.pin_pipeline(&self.pipeline);
134 request_filter::execute(&pipeline, session, ctx).await
135 }
136
137 async fn response_filter(
138 &self,
139 session: &mut Session,
140 upstream_response: &mut pingora_http::ResponseHeader,
141 ctx: &mut Self::CTX,
142 ) -> Result<()>
143 where
144 Self::CTX: Send + Sync,
145 {
146 let pipeline = ctx.pipeline(&self.pipeline);
147 let result = response_filter::execute(&pipeline, upstream_response, ctx).await;
148 if result.is_ok() {
149 let client_ver = ctx.client_http_version.unwrap_or(http::Version::HTTP_11);
150 via::append_response_via(upstream_response, client_ver);
151 adjust_compression(session, upstream_response, pipeline.compression_config());
152 }
153 result
154 }
155
156 fn fail_to_connect(
157 &self,
158 _session: &mut Session,
159 _peer: &HttpPeer,
160 ctx: &mut Self::CTX,
161 e: Box<pingora_core::Error>,
162 ) -> Box<pingora_core::Error> {
163 handle_connect_failure(ctx, e)
164 }
165
166 async fn upstream_request_filter(
167 &self,
168 session: &mut Session,
169 upstream_request: &mut pingora_http::RequestHeader,
170 ctx: &mut Self::CTX,
171 ) -> Result<()>
172 where
173 Self::CTX: Send + Sync,
174 {
175 let is_upgrade = session.is_upgrade_req();
176 upstream_request::strip_hop_by_hop(upstream_request, is_upgrade);
177 upstream_request::strip_reserved_internal(upstream_request);
178 upstream_request::apply_rewritten_path(upstream_request, ctx)?;
179 upstream_request::apply_mutated_content_length(upstream_request, ctx);
180 let client_ver = ctx.client_http_version.unwrap_or(http::Version::HTTP_11);
181 via::append_request_via(upstream_request, client_ver);
182 Ok(())
183 }
184
185 async fn upstream_peer(&self, _session: &mut Session, ctx: &mut Self::CTX) -> Result<Box<HttpPeer>> {
186 upstream_peer::execute(ctx).await
187 }
188
189 async fn logging(&self, session: &mut Session, e: Option<&pingora_core::Error>, ctx: &mut Self::CTX) {
190 let pipeline = ctx.pipeline(&self.pipeline);
191 emit_request_metrics(session, ctx);
192 record_passive_health(&pipeline, e, ctx);
193 logging_cleanup(&pipeline, ctx).await;
194 }
195}
196
197async fn reject_503(session: &mut Session, retry_after: &'static str, reason: &'static str) -> Result<()> {
203 tracing::warn!(reason, "rejecting request");
204 let mut header = pingora_http::ResponseHeader::build(503, None)?;
205 header.append_header("Retry-After", retry_after)?;
206 session.write_response_header(Box::new(header), true).await?;
207 Err(pingora_core::Error::explain(
208 pingora_core::ErrorType::HTTPStatus(503),
209 reason,
210 ))
211}