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 if praxis_core::memory::is_exceeded() {
160 metrics::record_overload_reject(metrics::OVERLOAD_REASON_MEMORY);
161 return reject_503(session, "5", "memory pressure exceeded").await;
162 }
163
164 let (exceeded, permit) = crate::connections::try_acquire_global();
165 ctx._global_connection_permit = permit;
166 if exceeded {
167 metrics::record_overload_reject(metrics::OVERLOAD_REASON_GLOBAL_CONNECTIONS);
168 return reject_503(session, "1", "global max connections exceeded").await;
169 }
170
171 if let Some(sem) = &self.connection_semaphore {
172 if let Ok(permit) = Arc::clone(sem).try_acquire_owned() {
173 ctx._connection_permit = Some(permit);
174 } else {
175 metrics::record_overload_reject(metrics::OVERLOAD_REASON_LISTENER_CONNECTIONS);
176 return reject_503(session, "1", "max connections exceeded").await;
177 }
178 }
179
180 ctx._active_connection = Some(metrics::ActiveConnectionGuard::acquire(self.listener_name.clone()));
181
182 if let Some(timeout) = self.downstream_read_timeout {
183 debug!(
184 timeout_ms = u64::try_from(timeout.as_millis()).unwrap_or(u64::MAX),
185 "applying downstream read timeout"
186 );
187 session.set_read_timeout(Some(timeout));
188 }
189 Ok(())
190 }
191
192 async fn request_filter(&self, session: &mut Session, ctx: &mut Self::CTX) -> Result<bool> {
193 let pipeline = ctx.pin_pipeline(&self.pipeline);
194 request_filter::execute(&pipeline, session, ctx).await
195 }
196
197 async fn request_body_filter(
198 &self,
199 session: &mut Session,
200 body: &mut Option<Bytes>,
201 end_of_stream: bool,
202 ctx: &mut Self::CTX,
203 ) -> Result<()>
204 where
205 Self::CTX: Send + Sync,
206 {
207 if ctx.connection_upgraded {
212 return Ok(());
213 }
214 if ctx.pre_read_body.is_none()
215 && let Some(pinned) = &ctx.pinned_pipeline
216 && !pinned.body_capabilities().needs_request_body
217 {
218 return Ok(());
219 }
220 let pipeline = ctx.pipeline(&self.pipeline);
221 let span = ctx.request_span.clone();
222 request_body_filter::execute(&pipeline, session, body, end_of_stream, ctx)
223 .instrument(span)
224 .await
225 }
226
227 fn response_body_filter(
228 &self,
229 _session: &mut Session,
230 body: &mut Option<Bytes>,
231 end_of_stream: bool,
232 ctx: &mut Self::CTX,
233 ) -> Result<Option<Duration>>
234 where
235 Self::CTX: Send + Sync,
236 {
237 if ctx.connection_upgraded {
240 return Ok(None);
241 }
242 if let Some(pinned) = &ctx.pinned_pipeline
243 && !pinned.body_capabilities().needs_response_body
244 {
245 if end_of_stream {
246 ctx.response_delivery_complete = true;
247 }
248 return Ok(None);
249 }
250 let span = ctx.request_span.clone();
251 let _entered = span.enter();
252 let pipeline = ctx.pipeline(&self.pipeline);
253 response_body_filter::execute(&pipeline, body, end_of_stream, ctx)
254 }
255
256 fn fail_to_connect(
257 &self,
258 session: &mut Session,
259 _peer: &HttpPeer,
260 ctx: &mut Self::CTX,
261 e: Box<pingora_core::Error>,
262 ) -> Box<pingora_core::Error> {
263 let span = ctx.request_span.clone();
264 let _entered = span.enter();
265 if session.as_mut().retry_buffer_truncated() {
268 let mut e = e;
269 e.set_retry(false);
270 return e;
271 }
272 handle_connect_failure(ctx, e)
273 }
274
275 fn error_while_proxy(
276 &self,
277 peer: &HttpPeer,
278 session: &mut Session,
279 e: Box<pingora_core::Error>,
280 ctx: &mut Self::CTX,
281 client_reused: bool,
282 ) -> Box<pingora_core::Error> {
283 if session.response_written().is_some_and(fail_to_proxy::is_final_response) {
289 let mut e = e;
290 e.set_retry(false);
291 return e;
292 }
293 let truncated = session.as_mut().retry_buffer_truncated();
296 if matches!(e.retry, pingora_core::RetryType::Decided(true)) {
304 if truncated {
305 let mut e = e;
306 e.set_retry(false);
307 return e;
308 }
309 return e;
310 }
311 if matches!(e.retry, pingora_core::RetryType::ReusedOnly) {
315 return resolve_reused_only_retry(ctx, session, client_reused, e);
316 }
317 let e = e.more_context(format!("Peer: {peer}"));
318 if truncated {
319 let mut e = e;
320 e.set_retry(false);
321 return e;
322 }
323 handle_connect_failure(ctx, e)
327 }
328
329 async fn fail_to_proxy(&self, session: &mut Session, e: &pingora_core::Error, ctx: &mut Self::CTX) -> FailToProxy
330 where
331 Self::CTX: Send + Sync,
332 {
333 let span = ctx.request_span.clone();
334 fail_to_proxy::execute(session, e, ctx).instrument(span).await
335 }
336
337 async fn upstream_request_filter(
338 &self,
339 session: &mut Session,
340 upstream_request: &mut pingora_http::RequestHeader,
341 ctx: &mut Self::CTX,
342 ) -> Result<()>
343 where
344 Self::CTX: Send + Sync,
345 {
346 let span = ctx.request_span.clone();
347 let _entered = span.enter();
348 let is_upgrade = session.is_upgrade_req();
349 upstream_request::strip_hop_by_hop(upstream_request, is_upgrade);
350 upstream_request.strip_reserved_internal();
351 upstream_request::apply_authority_override(upstream_request, ctx)?;
352 upstream_request::apply_rewritten_path(upstream_request, ctx)?;
353 upstream_request::apply_mutated_content_length(upstream_request, ctx);
354 let client_ver = ctx.client_http_version.unwrap_or(http::Version::HTTP_11);
355 via::append_request_via(upstream_request, client_ver);
356 Ok(())
357 }
358
359 async fn response_filter(
360 &self,
361 session: &mut Session,
362 upstream_response: &mut pingora_http::ResponseHeader,
363 ctx: &mut Self::CTX,
364 ) -> Result<()>
365 where
366 Self::CTX: Send + Sync,
367 {
368 let pipeline = ctx.pipeline(&self.pipeline);
369 let span = ctx.request_span.clone();
370 let exchange_span = ctx.upstream_exchange_span.clone();
371 let result = response_filter::execute(&pipeline, upstream_response, ctx)
372 .instrument(exchange_span)
373 .instrument(span)
374 .await;
375 if result.is_ok() {
376 let upstream_ver = upstream_response.version;
380 via::append_response_via(upstream_response, upstream_ver);
381 adjust_compression(session, upstream_response, pipeline.compression_config());
382 }
383 result
384 }
385
386 async fn upstream_peer(&self, _session: &mut Session, ctx: &mut Self::CTX) -> Result<Box<HttpPeer>> {
387 let span = ctx.request_span.clone();
388 upstream_peer::execute(ctx).instrument(span).await
389 }
390
391 async fn connected_to_upstream(
392 &self,
393 _session: &mut Session,
394 reused: bool,
395 peer: &HttpPeer,
396 #[cfg(unix)] _fd: std::os::unix::io::RawFd,
397 #[cfg(windows)] _sock: std::os::windows::io::RawSocket,
398 digest: Option<&pingora_core::protocols::Digest>,
399 ctx: &mut Self::CTX,
400 ) -> Result<()>
401 where
402 Self::CTX: Send + Sync,
403 {
404 let span = ctx.request_span.clone();
405 let _entered = span.enter();
406 let cluster = ctx.metrics_cluster_shared.clone().unwrap_or_else(metrics::cluster_none);
407 if !reused && let Some(start) = ctx.upstream_connect_start.take() {
408 metrics::record_upstream_connect_duration(cluster.clone(), start.elapsed().as_secs_f64());
409 }
410 if ctx.retries > 0 {
411 metrics::record_upstream_retry(cluster, metrics::RETRY_RESULT_SUCCESS);
412 }
413 connected_to_upstream::execute(reused, peer, digest, ctx);
414 Ok(())
415 }
416
417 async fn logging(&self, session: &mut Session, e: Option<&pingora_core::Error>, ctx: &mut Self::CTX) {
418 record_response_span_attributes(session, ctx);
419 let _exchange_span = std::mem::replace(&mut ctx.upstream_exchange_span, tracing::Span::none());
422 drop(_exchange_span);
423 let span = std::mem::replace(&mut ctx.request_span, tracing::Span::none());
424 let written_status = session.response_written().map_or(0, |resp| resp.status.as_u16());
425 async {
426 let pipeline = ctx.pipeline(&self.pipeline);
427 emit_request_metrics(session, ctx);
428 record_passive_health(&pipeline, e, ctx);
429 release_retry_state(ctx);
430 logging_cleanup(&pipeline, ctx).await;
431 super::maybe_emit_fallback_access_log(&pipeline, written_status, ctx);
432 }
433 .instrument(span)
434 .await;
435 }
436}
437
438async fn reject_503(session: &mut Session, retry_after: &'static str, reason: &'static str) -> Result<()> {
444 tracing::warn!(reason, "rejecting request");
445 let mut header = pingora_http::ResponseHeader::build(503, None)?;
446 header.append_header("Retry-After", retry_after)?;
447 session.write_response_header(Box::new(header), true).await?;
448 Err(pingora_core::Error::explain(
449 pingora_core::ErrorType::HTTPStatus(503),
450 reason,
451 ))
452}