praxis_protocol/http/pingora/handler/
with_body.rs1use 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::{ProxyHttp, Session};
25use praxis_filter::{CompressionConfig, FilterPipeline};
26use tokio::sync::Semaphore;
27use tracing::debug;
28
29use super::{
30 adjust_compression, emit_request_metrics, handle_connect_failure, logging_cleanup, record_passive_health,
31 request_body_filter, request_filter, response_body_filter, response_filter, upstream_peer, upstream_request, via,
32};
33use crate::http::pingora::context::PingoraRequestCtx;
34
35pub struct PingoraHttpHandler {
65 compression: Option<CompressionConfig>,
79
80 connection_semaphore: Option<Arc<Semaphore>>,
82
83 downstream_read_timeout: Option<Duration>,
85
86 pipeline: Arc<ArcSwap<FilterPipeline>>,
88}
89
90impl PingoraHttpHandler {
91 pub(super) fn new(
93 pipeline: Arc<ArcSwap<FilterPipeline>>,
94 downstream_read_timeout: Option<Duration>,
95 connection_semaphore: Option<Arc<Semaphore>>,
96 ) -> Self {
97 let compression = pipeline.load().compression_config().cloned();
98 Self {
99 compression,
100 connection_semaphore,
101 downstream_read_timeout,
102 pipeline,
103 }
104 }
105}
106
107#[async_trait]
108impl ProxyHttp for PingoraHttpHandler {
109 type CTX = PingoraRequestCtx;
110
111 fn new_ctx(&self) -> Self::CTX {
112 PingoraRequestCtx::default()
113 }
114
115 fn init_downstream_modules(&self, modules: &mut HttpModules) {
119 if let Some(cfg) = &self.compression {
120 debug!(level = cfg.default_level, "registering compression module");
121 modules.add_module(ResponseCompressionBuilder::enable(cfg.default_level));
122 }
123 }
124
125 #[expect(clippy::cast_possible_truncation, reason = "millis fit u64")]
126 async fn early_request_filter(&self, session: &mut Session, ctx: &mut Self::CTX) -> Result<()>
127 where
128 Self::CTX: Send + Sync,
129 {
130 if praxis_core::memory::is_exceeded() {
131 return reject_503(session, "5", "memory pressure exceeded").await;
132 }
133
134 let (exceeded, permit) = crate::connections::try_acquire_global();
135 ctx._global_connection_permit = permit;
136 if exceeded {
137 return reject_503(session, "1", "global max connections exceeded").await;
138 }
139
140 if let Some(sem) = &self.connection_semaphore {
141 if let Ok(permit) = Arc::clone(sem).try_acquire_owned() {
142 ctx._connection_permit = Some(permit);
143 } else {
144 return reject_503(session, "1", "max connections exceeded").await;
145 }
146 }
147
148 if let Some(timeout) = self.downstream_read_timeout {
149 debug!(
150 timeout_ms = timeout.as_millis() as u64,
151 "applying downstream read timeout"
152 );
153 session.set_read_timeout(Some(timeout));
154 }
155 Ok(())
156 }
157
158 async fn request_filter(&self, session: &mut Session, ctx: &mut Self::CTX) -> Result<bool> {
159 let pipeline = ctx.pin_pipeline(&self.pipeline);
160 request_filter::execute(&pipeline, session, ctx).await
161 }
162
163 async fn request_body_filter(
164 &self,
165 session: &mut Session,
166 body: &mut Option<Bytes>,
167 end_of_stream: bool,
168 ctx: &mut Self::CTX,
169 ) -> Result<()>
170 where
171 Self::CTX: Send + Sync,
172 {
173 let pipeline = ctx.pipeline(&self.pipeline);
174 request_body_filter::execute(&pipeline, session, body, end_of_stream, ctx).await
175 }
176
177 fn response_body_filter(
178 &self,
179 _session: &mut Session,
180 body: &mut Option<Bytes>,
181 end_of_stream: bool,
182 ctx: &mut Self::CTX,
183 ) -> Result<Option<Duration>>
184 where
185 Self::CTX: Send + Sync,
186 {
187 let pipeline = ctx.pipeline(&self.pipeline);
188 response_body_filter::execute(&pipeline, body, end_of_stream, ctx)
189 }
190
191 fn fail_to_connect(
192 &self,
193 _session: &mut Session,
194 _peer: &HttpPeer,
195 ctx: &mut Self::CTX,
196 e: Box<pingora_core::Error>,
197 ) -> Box<pingora_core::Error> {
198 handle_connect_failure(ctx, e)
199 }
200
201 async fn upstream_request_filter(
202 &self,
203 session: &mut Session,
204 upstream_request: &mut pingora_http::RequestHeader,
205 ctx: &mut Self::CTX,
206 ) -> Result<()>
207 where
208 Self::CTX: Send + Sync,
209 {
210 let is_upgrade = session.is_upgrade_req();
211 upstream_request::strip_hop_by_hop(upstream_request, is_upgrade);
212 upstream_request::strip_reserved_internal(upstream_request);
213 upstream_request::apply_rewritten_path(upstream_request, ctx)?;
214 upstream_request::apply_mutated_content_length(upstream_request, ctx);
215 let client_ver = ctx.client_http_version.unwrap_or(http::Version::HTTP_11);
216 via::append_request_via(upstream_request, client_ver);
217 Ok(())
218 }
219
220 async fn response_filter(
221 &self,
222 session: &mut Session,
223 upstream_response: &mut pingora_http::ResponseHeader,
224 ctx: &mut Self::CTX,
225 ) -> Result<()>
226 where
227 Self::CTX: Send + Sync,
228 {
229 let pipeline = ctx.pipeline(&self.pipeline);
230 let result = response_filter::execute(&pipeline, upstream_response, ctx).await;
231 if result.is_ok() {
232 let client_ver = ctx.client_http_version.unwrap_or(http::Version::HTTP_11);
233 via::append_response_via(upstream_response, client_ver);
234 adjust_compression(session, upstream_response, pipeline.compression_config());
235 }
236 result
237 }
238
239 async fn upstream_peer(&self, _session: &mut Session, ctx: &mut Self::CTX) -> Result<Box<HttpPeer>> {
240 upstream_peer::execute(ctx).await
241 }
242
243 async fn logging(&self, session: &mut Session, e: Option<&pingora_core::Error>, ctx: &mut Self::CTX) {
244 let pipeline = ctx.pipeline(&self.pipeline);
245 emit_request_metrics(session, ctx);
246 record_passive_health(&pipeline, e, ctx);
247 logging_cleanup(&pipeline, ctx).await;
248 }
249}
250
251async fn reject_503(session: &mut Session, retry_after: &'static str, reason: &'static str) -> Result<()> {
257 tracing::warn!(reason, "rejecting request");
258 let mut header = pingora_http::ResponseHeader::build(503, None)?;
259 header.append_header("Retry-After", retry_after)?;
260 session.write_response_header(Box::new(header), true).await?;
261 Err(pingora_core::Error::explain(
262 pingora_core::ErrorType::HTTPStatus(503),
263 reason,
264 ))
265}