1use 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, upstream_peer, upstream_request, via,
35};
36use crate::http::pingora::{context::PingoraRequestCtx, metrics};
37
38pub struct PingoraHttpHandler {
69 compression: Option<CompressionConfig>,
83
84 connection_semaphore: Option<Arc<Semaphore>>,
86
87 downstream_read_timeout: Option<Duration>,
89
90 listener_name: ::metrics::SharedString,
92
93 pipeline: Arc<ArcSwap<FilterPipeline>>,
95}
96
97impl PingoraHttpHandler {
98 pub(super) fn new(
100 pipeline: Arc<ArcSwap<FilterPipeline>>,
101 downstream_read_timeout: Option<Duration>,
102 connection_semaphore: Option<Arc<Semaphore>>,
103 listener_name: ::metrics::SharedString,
104 ) -> Self {
105 let compression = pipeline.load().compression_config().cloned();
106 Self {
107 compression,
108 connection_semaphore,
109 downstream_read_timeout,
110 listener_name,
111 pipeline,
112 }
113 }
114}
115
116fn resolve_reused_only_retry(
122 ctx: &PingoraRequestCtx,
123 session: &Session,
124 client_reused: bool,
125 mut e: Box<pingora_core::Error>,
126) -> Box<pingora_core::Error> {
127 let policy = ctx.retry_policy.clone().unwrap_or_else(super::legacy_default_policy);
128 let replay_safe = client_reused
129 && !session.as_ref().retry_buffer_truncated()
130 && (ctx.request_is_idempotent || policy.allow_non_idempotent());
131 if !replay_safe {
132 debug!("clearing reused-connection retry: replay is not safe for this request");
133 }
134 e.set_retry(replay_safe);
135 e
136}
137
138#[async_trait]
139impl ProxyHttp for PingoraHttpHandler {
140 type CTX = PingoraRequestCtx;
141
142 fn new_ctx(&self) -> Self::CTX {
143 PingoraRequestCtx::default()
144 }
145
146 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(0));
154 }
155 }
156
157 #[expect(
158 clippy::too_many_lines,
159 reason = "admission guards and compression setup precede all request module hooks"
160 )]
161 async fn early_request_filter(&self, session: &mut Session, ctx: &mut Self::CTX) -> Result<()>
162 where
163 Self::CTX: Send + Sync,
164 {
165 ctx.upstream_for_retry = None;
172 ctx.upstream_contacted = false;
173
174 if praxis_core::memory::is_exceeded() {
175 metrics::record_overload_reject(metrics::OVERLOAD_REASON_MEMORY);
176 return reject_503(session, "5", "memory pressure exceeded").await;
177 }
178
179 let (exceeded, permit) = crate::connections::try_acquire_global();
180 ctx._global_connection_permit = permit;
181 if exceeded {
182 metrics::record_overload_reject(metrics::OVERLOAD_REASON_GLOBAL_CONNECTIONS);
183 return reject_503(session, "1", "global max connections exceeded").await;
184 }
185
186 if let Some(sem) = &self.connection_semaphore {
187 if let Ok(permit) = Arc::clone(sem).try_acquire_owned() {
188 ctx._connection_permit = Some(permit);
189 } else {
190 metrics::record_overload_reject(metrics::OVERLOAD_REASON_LISTENER_CONNECTIONS);
191 return reject_503(session, "1", "max connections exceeded").await;
192 }
193 }
194
195 ctx._active_request = Some(metrics::ActiveRequestGuard::acquire(self.listener_name.clone()));
196
197 if let Some(timeout) = self.downstream_read_timeout {
198 debug!(
199 timeout_ms = u64::try_from(timeout.as_millis()).unwrap_or(u64::MAX),
200 "applying downstream read timeout"
201 );
202 session.set_read_timeout(Some(timeout));
203 }
204
205 let pipeline = ctx.pin_pipeline(&self.pipeline);
209 configure_compression(&mut session.downstream_modules_ctx, pipeline.compression_config());
210 Ok(())
211 }
212
213 async fn request_filter(&self, session: &mut Session, ctx: &mut Self::CTX) -> Result<bool> {
214 let pipeline = ctx.pin_pipeline(&self.pipeline);
215 request_filter::execute(&pipeline, session, ctx).await
216 }
217
218 async fn request_body_filter(
219 &self,
220 session: &mut Session,
221 body: &mut Option<Bytes>,
222 end_of_stream: bool,
223 ctx: &mut Self::CTX,
224 ) -> Result<()>
225 where
226 Self::CTX: Send + Sync,
227 {
228 if ctx.connection_upgraded {
233 return Ok(());
234 }
235 if ctx.pre_read_body.is_none()
236 && let Some(pinned) = &ctx.pinned_pipeline
237 && !pinned.body_capabilities().needs_request_body
238 {
239 return Ok(());
240 }
241 let pipeline = ctx.pipeline(&self.pipeline);
242 let span = ctx.request_span.clone();
243 request_body_filter::execute(&pipeline, session, body, end_of_stream, ctx)
244 .instrument(span)
245 .await
246 }
247
248 fn response_body_filter(
249 &self,
250 _session: &mut Session,
251 body: &mut Option<Bytes>,
252 end_of_stream: bool,
253 ctx: &mut Self::CTX,
254 ) -> Result<Option<Duration>>
255 where
256 Self::CTX: Send + Sync,
257 {
258 if ctx.connection_upgraded {
261 return Ok(None);
262 }
263 if let Some(pinned) = &ctx.pinned_pipeline
264 && !pinned.body_capabilities().needs_response_body
265 {
266 if end_of_stream {
267 ctx.response_delivery_complete = true;
268 }
269 return Ok(None);
270 }
271 let span = ctx.request_span.clone();
272 let _entered = span.enter();
273 let pipeline = ctx.pipeline(&self.pipeline);
274 response_body_filter::execute(&pipeline, body, end_of_stream, ctx)
275 }
276
277 fn fail_to_connect(
278 &self,
279 session: &mut Session,
280 _peer: &HttpPeer,
281 ctx: &mut Self::CTX,
282 e: Box<pingora_core::Error>,
283 ) -> Box<pingora_core::Error> {
284 let span = ctx.request_span.clone();
285 let _entered = span.enter();
286 if session.as_mut().retry_buffer_truncated() {
289 let mut e = e;
290 e.set_retry(false);
291 return e;
292 }
293 handle_connect_failure(ctx, e)
294 }
295
296 fn error_while_proxy(
297 &self,
298 peer: &HttpPeer,
299 session: &mut Session,
300 e: Box<pingora_core::Error>,
301 ctx: &mut Self::CTX,
302 client_reused: bool,
303 ) -> Box<pingora_core::Error> {
304 if session.response_written().is_some_and(fail_to_proxy::is_final_response) {
310 let mut e = e;
311 e.set_retry(false);
312 return e;
313 }
314 let truncated = session.as_mut().retry_buffer_truncated();
317 if matches!(e.retry, pingora_core::RetryType::Decided(true)) {
325 if truncated {
326 let mut e = e;
327 e.set_retry(false);
328 return e;
329 }
330 return e;
331 }
332 if matches!(e.retry, pingora_core::RetryType::ReusedOnly) {
336 return resolve_reused_only_retry(ctx, session, client_reused, e);
337 }
338 let e = e.more_context(format!("Peer: {peer}"));
339 if truncated {
340 let mut e = e;
341 e.set_retry(false);
342 return e;
343 }
344 handle_connect_failure(ctx, e)
348 }
349
350 async fn fail_to_proxy(&self, session: &mut Session, e: &pingora_core::Error, ctx: &mut Self::CTX) -> FailToProxy
351 where
352 Self::CTX: Send + Sync,
353 {
354 let span = ctx.request_span.clone();
355 fail_to_proxy::execute(session, e, ctx).instrument(span).await
356 }
357
358 async fn upstream_request_filter(
359 &self,
360 session: &mut Session,
361 upstream_request: &mut pingora_http::RequestHeader,
362 ctx: &mut Self::CTX,
363 ) -> Result<()>
364 where
365 Self::CTX: Send + Sync,
366 {
367 let span = ctx.request_span.clone();
368 let _entered = span.enter();
369 let pipeline = ctx.pipeline(&self.pipeline);
371 pipeline.clear_request_body_done(&mut ctx.cached_body_done_indices);
372
373 let is_upgrade = session.is_upgrade_req();
374 upstream_request::strip_hop_by_hop(upstream_request, is_upgrade);
375 upstream_request.strip_reserved_internal();
376 upstream_request::apply_authority_override(upstream_request, ctx)?;
377 upstream_request::apply_rewritten_path(upstream_request, ctx)?;
378 upstream_request::apply_mutated_content_length(upstream_request, ctx);
379 upstream_request::reseed_retry_body(ctx);
383 let client_ver = ctx.client_http_version.unwrap_or(http::Version::HTTP_11);
384 via::append_request_via(upstream_request, client_ver);
385 Ok(())
386 }
387
388 async fn response_filter(
389 &self,
390 session: &mut Session,
391 upstream_response: &mut pingora_http::ResponseHeader,
392 ctx: &mut Self::CTX,
393 ) -> Result<()>
394 where
395 Self::CTX: Send + Sync,
396 {
397 let pipeline = ctx.pipeline(&self.pipeline);
398 let span = ctx.request_span.clone();
399 let exchange_span = ctx.upstream_exchange_span.clone();
400 let result = response_filter::execute(&pipeline, upstream_response, ctx)
401 .instrument(exchange_span)
402 .instrument(span)
403 .await;
404 if result.is_ok() {
405 let upstream_ver = upstream_response.version;
409 via::append_response_via(upstream_response, upstream_ver);
410 adjust_compression(session, upstream_response, pipeline.compression_config());
411 }
412 result
413 }
414
415 async fn upstream_peer(&self, _session: &mut Session, ctx: &mut Self::CTX) -> Result<Box<HttpPeer>> {
416 let span = ctx.request_span.clone();
417 upstream_peer::execute(ctx).instrument(span).await
418 }
419
420 async fn connected_to_upstream(
421 &self,
422 _session: &mut Session,
423 reused: bool,
424 peer: &HttpPeer,
425 #[cfg(unix)] _fd: std::os::unix::io::RawFd,
426 #[cfg(windows)] _sock: std::os::windows::io::RawSocket,
427 digest: Option<&pingora_core::protocols::Digest>,
428 ctx: &mut Self::CTX,
429 ) -> Result<()>
430 where
431 Self::CTX: Send + Sync,
432 {
433 let span = ctx.request_span.clone();
434 let _entered = span.enter();
435 let cluster = ctx.metrics_cluster_shared.clone().unwrap_or_else(metrics::cluster_none);
436 if !reused && let Some(start) = ctx.upstream_connect_start.take() {
437 metrics::record_upstream_connect_duration(cluster.clone(), start.elapsed().as_secs_f64());
438 }
439 if ctx.retries > 0 {
440 metrics::record_upstream_retry(cluster, metrics::RETRY_RESULT_SUCCESS);
441 }
442 connected_to_upstream::execute(reused, peer, digest, ctx);
443 Ok(())
444 }
445
446 async fn logging(&self, session: &mut Session, e: Option<&pingora_core::Error>, ctx: &mut Self::CTX) {
447 record_response_span_attributes(session, ctx);
448 let _exchange_span = std::mem::replace(&mut ctx.upstream_exchange_span, tracing::Span::none());
451 drop(_exchange_span);
452 let span = std::mem::replace(&mut ctx.request_span, tracing::Span::none());
453 let written_status = session.response_written().map_or(0, |resp| resp.status.as_u16());
454 async {
455 let pipeline = ctx.pipeline(&self.pipeline);
456 emit_request_metrics(session, ctx);
457 record_passive_health(&pipeline, e, ctx);
458 release_retry_state(ctx);
459 logging_cleanup(&pipeline, ctx).await;
460 super::maybe_emit_fallback_access_log(&pipeline, written_status, ctx);
461 }
462 .instrument(span)
463 .await;
464 }
465}
466
467async fn reject_503(session: &mut Session, retry_after: &'static str, reason: &'static str) -> Result<()> {
473 tracing::warn!(reason, "rejecting request");
474 let mut header = pingora_http::ResponseHeader::build(503, None)?;
475 header.append_header("Retry-After", retry_after)?;
476 session.write_response_header(Box::new(header), true).await?;
477 Err(pingora_core::Error::explain(
478 pingora_core::ErrorType::HTTPStatus(503),
479 reason,
480 ))
481}