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, 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
44pub struct PingoraHttpHandler {
75 compression: Option<CompressionConfig>,
89
90 connection_semaphore: Option<Arc<Semaphore>>,
92
93 downstream_read_timeout: Option<Duration>,
95
96 listener_name: ::metrics::SharedString,
98
99 pipeline: Arc<ArcSwap<FilterPipeline>>,
101}
102
103impl PingoraHttpHandler {
104 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
122fn reused_only_replay_safe(client_reused: bool, retry_buffer_truncated: bool) -> bool {
133 client_reused && !retry_buffer_truncated
134}
135
136fn 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 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 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 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 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 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 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 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 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 let truncated = session.as_mut().retry_buffer_truncated();
368 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 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 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 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 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 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 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
519async 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#[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}