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, response_trailer_filter, response_trailers, upstream_peer,
35 upstream_request, via,
36};
37use crate::http::pingora::{context::PingoraRequestCtx, metrics};
38
39pub struct PingoraHttpHandler {
70 compression: Option<CompressionConfig>,
84
85 connection_semaphore: Option<Arc<Semaphore>>,
87
88 downstream_read_timeout: Option<Duration>,
90
91 listener_name: ::metrics::SharedString,
93
94 pipeline: Arc<ArcSwap<FilterPipeline>>,
96}
97
98impl PingoraHttpHandler {
99 pub(super) fn new(
101 pipeline: Arc<ArcSwap<FilterPipeline>>,
102 downstream_read_timeout: Option<Duration>,
103 connection_semaphore: Option<Arc<Semaphore>>,
104 listener_name: ::metrics::SharedString,
105 ) -> Self {
106 let compression = pipeline.load().compression_config().cloned();
107 Self {
108 compression,
109 connection_semaphore,
110 downstream_read_timeout,
111 listener_name,
112 pipeline,
113 }
114 }
115}
116
117fn reused_only_replay_safe(client_reused: bool, retry_buffer_truncated: bool) -> bool {
128 client_reused && !retry_buffer_truncated
129}
130
131fn resolve_reused_only_retry(
136 session: &Session,
137 client_reused: bool,
138 mut e: Box<pingora_core::Error>,
139) -> Box<pingora_core::Error> {
140 let replay_safe = reused_only_replay_safe(client_reused, session.as_ref().retry_buffer_truncated());
141 if !replay_safe {
142 debug!(
143 client_reused,
144 "clearing reused-connection retry: fresh connection or truncated replay buffer"
145 );
146 }
147 e.set_retry(replay_safe);
148 e
149}
150
151#[async_trait]
152impl ProxyHttp for PingoraHttpHandler {
153 type CTX = PingoraRequestCtx;
154
155 fn new_ctx(&self) -> Self::CTX {
156 PingoraRequestCtx::default()
157 }
158
159 fn init_downstream_modules(&self, modules: &mut HttpModules) {
162 if let Some(cfg) = &self.compression {
163 debug!(level = cfg.default_level, "registering compression module");
164 modules.add_module(ResponseCompressionBuilder::enable(0));
167 }
168 }
169
170 #[expect(
171 clippy::too_many_lines,
172 reason = "admission guards and compression setup precede all request module hooks"
173 )]
174 async fn early_request_filter(&self, session: &mut Session, ctx: &mut Self::CTX) -> Result<()>
175 where
176 Self::CTX: Send + Sync,
177 {
178 ctx.upstream_for_retry = None;
185 ctx.upstream_contacted = false;
186
187 if praxis_core::memory::is_exceeded() {
188 metrics::record_overload_reject(metrics::OVERLOAD_REASON_MEMORY);
189 return reject_503(session, "5", "memory pressure exceeded").await;
190 }
191
192 let (exceeded, permit) = crate::connections::try_acquire_global();
193 ctx._global_connection_permit = permit;
194 if exceeded {
195 metrics::record_overload_reject(metrics::OVERLOAD_REASON_GLOBAL_CONNECTIONS);
196 return reject_503(session, "1", "global max connections exceeded").await;
197 }
198
199 if let Some(sem) = &self.connection_semaphore {
200 if let Ok(permit) = Arc::clone(sem).try_acquire_owned() {
201 ctx._connection_permit = Some(permit);
202 } else {
203 metrics::record_overload_reject(metrics::OVERLOAD_REASON_LISTENER_CONNECTIONS);
204 return reject_503(session, "1", "max connections exceeded").await;
205 }
206 }
207
208 ctx._active_request = Some(metrics::ActiveRequestGuard::acquire(self.listener_name.clone()));
209
210 if let Some(timeout) = self.downstream_read_timeout {
211 debug!(
212 timeout_ms = u64::try_from(timeout.as_millis()).unwrap_or(u64::MAX),
213 "applying downstream read timeout"
214 );
215 session.set_read_timeout(Some(timeout));
216 }
217
218 let pipeline = ctx.pin_pipeline(&self.pipeline);
222 configure_compression(&mut session.downstream_modules_ctx, pipeline.compression_config());
223 Ok(())
224 }
225
226 async fn request_filter(&self, session: &mut Session, ctx: &mut Self::CTX) -> Result<bool> {
227 let pipeline = ctx.pin_pipeline(&self.pipeline);
228 request_filter::execute(&pipeline, session, ctx).await
229 }
230
231 async fn request_body_filter(
232 &self,
233 session: &mut Session,
234 body: &mut Option<Bytes>,
235 end_of_stream: bool,
236 ctx: &mut Self::CTX,
237 ) -> Result<()>
238 where
239 Self::CTX: Send + Sync,
240 {
241 if ctx.connection_upgraded {
246 return Ok(());
247 }
248 if ctx.pre_read_body.is_none()
249 && let Some(pinned) = &ctx.pinned_pipeline
250 && !pinned.body_capabilities().needs_request_body
251 {
252 return Ok(());
253 }
254 let pipeline = ctx.pipeline(&self.pipeline);
255 let span = ctx.request_span.clone();
256 request_body_filter::execute(&pipeline, session, body, end_of_stream, ctx)
257 .instrument(span)
258 .await
259 }
260
261 fn upstream_response_trailer_filter(
262 &self,
263 _session: &mut Session,
264 upstream_trailers: &mut http::HeaderMap,
265 ctx: &mut Self::CTX,
266 ) -> Result<()>
267 where
268 Self::CTX: Send + Sync,
269 {
270 let span = ctx.request_span.clone();
271 let _entered = span.enter();
272 response_trailers::capture(upstream_trailers, ctx);
273 Ok(())
274 }
275
276 async fn response_trailer_filter(
277 &self,
278 _session: &mut Session,
279 upstream_trailers: &mut http::HeaderMap,
280 ctx: &mut Self::CTX,
281 ) -> Result<Option<Bytes>>
282 where
283 Self::CTX: Send + Sync,
284 {
285 let pipeline = ctx.pipeline(&self.pipeline);
286 if !pipeline.needs_response_trailers() {
287 return Ok(None);
288 }
289 let span = ctx.request_span.clone();
290 let _entered = span.enter();
291 Ok(response_trailer_filter::execute(&pipeline, upstream_trailers, ctx))
292 }
293
294 fn response_body_filter(
295 &self,
296 _session: &mut Session,
297 body: &mut Option<Bytes>,
298 end_of_stream: bool,
299 ctx: &mut Self::CTX,
300 ) -> Result<Option<Duration>>
301 where
302 Self::CTX: Send + Sync,
303 {
304 if ctx.connection_upgraded {
307 return Ok(None);
308 }
309 if let Some(pinned) = &ctx.pinned_pipeline
310 && !pinned.body_capabilities().needs_response_body
311 {
312 if end_of_stream {
313 ctx.response_delivery_complete = true;
314 }
315 return Ok(None);
316 }
317 let span = ctx.request_span.clone();
318 let _entered = span.enter();
319 let pipeline = ctx.pipeline(&self.pipeline);
320 response_body_filter::execute(&pipeline, body, end_of_stream, ctx)
321 }
322
323 fn fail_to_connect(
324 &self,
325 session: &mut Session,
326 _peer: &HttpPeer,
327 ctx: &mut Self::CTX,
328 e: Box<pingora_core::Error>,
329 ) -> Box<pingora_core::Error> {
330 let span = ctx.request_span.clone();
331 let _entered = span.enter();
332 if session.as_mut().retry_buffer_truncated() {
335 let mut e = e;
336 e.set_retry(false);
337 return e;
338 }
339 handle_connect_failure(ctx, e)
340 }
341
342 fn error_while_proxy(
343 &self,
344 peer: &HttpPeer,
345 session: &mut Session,
346 e: Box<pingora_core::Error>,
347 ctx: &mut Self::CTX,
348 client_reused: bool,
349 ) -> Box<pingora_core::Error> {
350 if session.response_written().is_some_and(fail_to_proxy::is_final_response) {
356 let mut e = e;
357 e.set_retry(false);
358 return e;
359 }
360 let truncated = session.as_mut().retry_buffer_truncated();
363 if matches!(e.retry, pingora_core::RetryType::Decided(true)) {
371 if truncated {
372 let mut e = e;
373 e.set_retry(false);
374 return e;
375 }
376 return e;
377 }
378 if matches!(e.retry, pingora_core::RetryType::ReusedOnly) {
382 return resolve_reused_only_retry(session, client_reused, e);
383 }
384 let e = e.more_context(format!("Peer: {peer}"));
385 if truncated {
386 let mut e = e;
387 e.set_retry(false);
388 return e;
389 }
390 handle_connect_failure(ctx, e)
394 }
395
396 async fn fail_to_proxy(&self, session: &mut Session, e: &pingora_core::Error, ctx: &mut Self::CTX) -> FailToProxy
397 where
398 Self::CTX: Send + Sync,
399 {
400 let span = ctx.request_span.clone();
401 fail_to_proxy::execute(session, e, ctx).instrument(span).await
402 }
403
404 async fn upstream_request_filter(
405 &self,
406 session: &mut Session,
407 upstream_request: &mut pingora_http::RequestHeader,
408 ctx: &mut Self::CTX,
409 ) -> Result<()>
410 where
411 Self::CTX: Send + Sync,
412 {
413 let span = ctx.request_span.clone();
414 let _entered = span.enter();
415 let pipeline = ctx.pipeline(&self.pipeline);
417 pipeline.clear_request_body_done(&mut ctx.cached_body_done_indices);
418
419 let is_upgrade = session.is_upgrade_req();
420 upstream_request::strip_hop_by_hop(upstream_request, is_upgrade);
421 upstream_request.strip_reserved_internal();
422 upstream_request::apply_authority_override(upstream_request, ctx)?;
423 upstream_request::apply_rewritten_path(upstream_request, ctx)?;
424 upstream_request::apply_mutated_content_length(upstream_request, ctx);
425 upstream_request::reseed_retry_body(ctx);
429 upstream_request::apply_grpc_deadline_header(upstream_request, ctx);
430 let client_ver = ctx.client_http_version.unwrap_or(http::Version::HTTP_11);
431 via::append_request_via(upstream_request, client_ver);
432 Ok(())
433 }
434
435 async fn response_filter(
436 &self,
437 session: &mut Session,
438 upstream_response: &mut pingora_http::ResponseHeader,
439 ctx: &mut Self::CTX,
440 ) -> Result<()>
441 where
442 Self::CTX: Send + Sync,
443 {
444 let pipeline = ctx.pipeline(&self.pipeline);
445 let span = ctx.request_span.clone();
446 let exchange_span = ctx.upstream_exchange_span.clone();
447 let result = response_filter::execute(&pipeline, upstream_response, ctx)
448 .instrument(exchange_span)
449 .instrument(span)
450 .await;
451 if result.is_ok() {
452 let upstream_ver = upstream_response.version;
456 via::append_response_via(upstream_response, upstream_ver);
457 adjust_compression(session, upstream_response, pipeline.compression_config());
458 }
459 result
460 }
461
462 async fn upstream_peer(&self, _session: &mut Session, ctx: &mut Self::CTX) -> Result<Box<HttpPeer>> {
463 let span = ctx.request_span.clone();
464 upstream_peer::execute(ctx).instrument(span).await
465 }
466
467 async fn connected_to_upstream(
468 &self,
469 _session: &mut Session,
470 reused: bool,
471 peer: &HttpPeer,
472 #[cfg(unix)] _fd: std::os::unix::io::RawFd,
473 #[cfg(windows)] _sock: std::os::windows::io::RawSocket,
474 digest: Option<&pingora_core::protocols::Digest>,
475 ctx: &mut Self::CTX,
476 ) -> Result<()>
477 where
478 Self::CTX: Send + Sync,
479 {
480 let span = ctx.request_span.clone();
481 let _entered = span.enter();
482 let cluster = ctx.metrics_cluster_shared.clone().unwrap_or_else(metrics::cluster_none);
483 if !reused && let Some(start) = ctx.upstream_connect_start.take() {
484 metrics::record_upstream_connect_duration(cluster.clone(), start.elapsed().as_secs_f64());
485 }
486 if ctx.retries > 0 {
487 metrics::record_upstream_retry(cluster, metrics::RETRY_RESULT_SUCCESS);
488 }
489 connected_to_upstream::execute(reused, peer, digest, ctx);
490 Ok(())
491 }
492
493 async fn logging(&self, session: &mut Session, e: Option<&pingora_core::Error>, ctx: &mut Self::CTX) {
494 record_response_span_attributes(session, ctx);
495 let _exchange_span = std::mem::replace(&mut ctx.upstream_exchange_span, tracing::Span::none());
498 drop(_exchange_span);
499 let span = std::mem::replace(&mut ctx.request_span, tracing::Span::none());
500 let written_status = session.response_written().map_or(0, |resp| resp.status.as_u16());
501 async {
502 let pipeline = ctx.pipeline(&self.pipeline);
503 emit_request_metrics(session, ctx);
504 record_passive_health(&pipeline, e, ctx);
505 release_retry_state(ctx);
506 logging_cleanup(&pipeline, ctx).await;
507 super::maybe_emit_fallback_access_log(&pipeline, written_status, ctx);
508 }
509 .instrument(span)
510 .await;
511 }
512}
513
514async fn reject_503(session: &mut Session, retry_after: &'static str, reason: &'static str) -> Result<()> {
520 tracing::warn!(reason, "rejecting request");
521 let mut header = pingora_http::ResponseHeader::build(503, None)?;
522 header.append_header("Retry-After", retry_after)?;
523 session.write_response_header(Box::new(header), true).await?;
524 Err(pingora_core::Error::explain(
525 pingora_core::ErrorType::HTTPStatus(503),
526 reason,
527 ))
528}
529
530#[cfg(test)]
535mod tests {
536 use super::*;
537
538 #[test]
539 fn reused_only_replays_when_reused_and_buffer_intact() {
540 assert!(
541 reused_only_replay_safe(true, false),
542 "a stale reused connection with an intact replay buffer must retry, regardless of method"
543 );
544 }
545
546 #[test]
547 fn reused_only_does_not_replay_on_fresh_connection() {
548 assert!(
549 !reused_only_replay_safe(false, false),
550 "a fresh-connection failure may already have been processed and must not replay"
551 );
552 assert!(
553 !reused_only_replay_safe(false, true),
554 "a fresh connection with a truncated buffer must not replay"
555 );
556 }
557
558 #[test]
559 fn reused_only_does_not_replay_when_buffer_truncated() {
560 assert!(
561 !reused_only_replay_safe(true, true),
562 "a truncated replay buffer cannot resend the full body, so no retry"
563 );
564 }
565}