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