1use std::{sync::Arc, time::Duration};
17
18use arc_swap::ArcSwap;
19use pingora_core::{Result, server::Server, services::listening::Service};
20use pingora_proxy::http_proxy;
21use praxis_filter::FilterPipeline;
22use tokio::sync::Semaphore;
23use tracing::debug;
24
25use super::metrics;
26
27mod compression;
29mod connected_to_upstream;
31mod fail_to_proxy;
33mod hop_by_hop;
35mod normalize;
37mod request_body_filter;
39mod request_filter;
41mod reserved_headers;
43mod response_body_filter;
45mod response_filter;
47mod response_trailer_filter;
49mod response_trailers;
51mod retry;
53mod upstream_peer;
55mod upstream_request;
57mod upstream_response;
59mod via;
61mod with_body;
63
64mod body_util;
66mod health_util;
68mod logging_util;
70mod metrics_util;
72mod retry_util;
74mod server_options;
76mod span_util;
78
79pub use upstream_peer::{UpstreamRetryGateRelease, arm_upstream_retry_gate, lock_upstream_retry_gate_tests};
80pub use with_body::PingoraHttpHandler;
81
82pub fn load_http_handler(
128 server: &mut Server,
129 listener: &praxis_core::config::Listener,
130 pipeline: Arc<ArcSwap<FilterPipeline>>,
131 cert_watcher_shutdowns: &mut Vec<tokio::sync::watch::Sender<bool>>,
132) -> Result<(), praxis_core::ProxyError> {
133 let downstream_read_timeout = listener.downstream_read_timeout_ms.map(Duration::from_millis);
134 let connection_semaphore = listener
135 .max_connections
136 .map(|max| Arc::new(Semaphore::new(max as usize)));
137
138 debug!(listener = %listener.name, "loading HTTP handler with body filters");
141 let handler = PingoraHttpHandler::new(
142 pipeline,
143 downstream_read_timeout,
144 connection_semaphore,
145 ::metrics::SharedString::from_shared(Arc::from(listener.name.as_str())),
149 );
150 wire_service(server, listener, handler, cert_watcher_shutdowns)?;
151 Ok(())
152}
153
154fn wire_service<H>(
156 server: &mut Server,
157 listener: &praxis_core::config::Listener,
158 handler: H,
159 cert_watcher_shutdowns: &mut Vec<tokio::sync::watch::Sender<bool>>,
160) -> Result<(), praxis_core::ProxyError>
161where
162 H: pingora_proxy::ProxyHttp + Send + Sync + 'static,
163 H::CTX: Send + Sync,
164{
165 let service_name = format!("http-proxy:{name}", name = listener.name);
166 let mut proxy = http_proxy(&server.configuration, handler);
167 proxy.server_options = Some(server_options::h2c_server_options());
168 proxy.h2_options = Some(server_options::h2_server_options());
169 let mut service = Service::new(service_name, proxy);
170 if let Some(tx) = super::listener::add_listener(&mut service, listener)? {
171 cert_watcher_shutdowns.push(tx);
172 }
173 server.add_service(service);
174 Ok(())
175}
176
177#[cfg(test)]
182#[expect(clippy::allow_attributes, reason = "blanket test suppressions")]
183#[allow(
184 clippy::unwrap_used,
185 clippy::expect_used,
186 clippy::indexing_slicing,
187 clippy::field_reassign_with_default,
188 clippy::too_many_lines,
189 clippy::cast_possible_truncation,
190 clippy::significant_drop_tightening,
191 reason = "tests"
192)]
193mod tests {
194 use std::collections::HashMap;
195
196 use bytes::Bytes;
197 use praxis_core::{
198 connectivity::{ConnectionOptions, Upstream},
199 health::{ClusterHealthEntry, EndpointHealth},
200 };
201 use praxis_filter::{BodyBuffer, BodyMode, RequestExtensions};
202
203 use super::*;
204 use crate::http::pingora::context::PingoraRequestCtx;
205
206 const MAX_RETRIES: usize = praxis_core::config::DEFAULT_MAX_RETRIES as usize;
208
209 const RETRY_BODY_LIMIT: u64 = praxis_core::config::DEFAULT_RETRY_BODY_LIMIT_BYTES;
211
212 #[test]
213 fn first_failure_idempotent_sets_retry() {
214 let mut ctx = PingoraRequestCtx::default();
215 ctx.request_is_idempotent = true;
216 let e = retry_util::handle_connect_failure(&mut ctx, make_error());
217 assert!(e.retry(), "first failure should set retry flag");
218 assert_eq!(ctx.retries, 1);
219 }
220
221 #[test]
222 fn large_body_skips_retry() {
223 let mut ctx = PingoraRequestCtx::default();
224 ctx.request_is_idempotent = true;
225 ctx.request_body_bytes = RETRY_BODY_LIMIT + 1;
226 let e = retry_util::handle_connect_failure(&mut ctx, make_error());
227 assert!(!e.retry(), "should not retry when body exceeds retry buffer limit");
228 assert_eq!(ctx.retries, 0, "retry counter should not increment");
229 }
230
231 #[test]
232 fn mutated_body_exceeding_limit_skips_retry() {
233 let mut ctx = PingoraRequestCtx::default();
234 ctx.request_is_idempotent = true;
235 ctx.request_body_bytes = 1024;
236 ctx.mutated_request_body_len = Some((RETRY_BODY_LIMIT + 1) as usize);
237 let e = retry_util::handle_connect_failure(&mut ctx, make_error());
238 assert!(
239 !e.retry(),
240 "should not retry when mutated body exceeds retry buffer limit"
241 );
242 assert_eq!(ctx.retries, 0);
243 }
244
245 #[test]
246 fn body_at_limit_allows_retry() {
247 let mut ctx = PingoraRequestCtx::default();
248 ctx.request_is_idempotent = true;
249 ctx.request_body_bytes = RETRY_BODY_LIMIT;
250 let e = retry_util::handle_connect_failure(&mut ctx, make_error());
251 assert!(e.retry(), "body exactly at limit should allow retry");
252 assert_eq!(ctx.retries, 1);
253 }
254
255 #[test]
256 fn zero_body_allows_retry() {
257 let mut ctx = PingoraRequestCtx::default();
258 ctx.request_is_idempotent = true;
259 ctx.request_body_bytes = 0;
260 let e = retry_util::handle_connect_failure(&mut ctx, make_error());
261 assert!(e.retry(), "zero-length body should allow retry");
262 assert_eq!(ctx.retries, 1);
263 }
264
265 #[test]
266 fn max_retries_exhausted_does_not_retry() {
267 let mut ctx = PingoraRequestCtx::default();
268 ctx.request_is_idempotent = true;
269 ctx.retries = MAX_RETRIES as u32;
270 let e = retry_util::handle_connect_failure(&mut ctx, make_error());
271 assert!(!e.retry(), "should not retry after MAX_RETRIES");
272 assert_eq!(ctx.retries as usize, MAX_RETRIES);
273 }
274
275 #[test]
276 fn counter_increments_across_calls() {
277 let mut ctx = PingoraRequestCtx::default();
278 ctx.request_is_idempotent = true;
279 for expected in 1..=MAX_RETRIES {
280 let _result = retry_util::handle_connect_failure(&mut ctx, make_error());
281 assert_eq!(ctx.retries as usize, expected);
282 }
283 let e = retry_util::handle_connect_failure(&mut ctx, make_error());
284 assert!(!e.retry(), "should not retry after reaching MAX_RETRIES");
285 assert_eq!(ctx.retries as usize, MAX_RETRIES);
286 }
287
288 #[test]
289 fn non_idempotent_request_never_retries() {
290 let mut ctx = PingoraRequestCtx::default();
291 ctx.request_is_idempotent = false;
292 let e = retry_util::handle_connect_failure(&mut ctx, make_error());
293 assert!(!e.retry(), "non-idempotent request should never retry");
294 assert_eq!(ctx.retries, 0);
295 }
296
297 #[test]
298 fn connect_failure_clears_upstream_connect_start() {
299 let mut ctx = PingoraRequestCtx::default();
300 ctx.upstream_connect_start = Some(std::time::Instant::now());
301 let _e = retry_util::handle_connect_failure(&mut ctx, make_error());
302 assert!(
303 ctx.upstream_connect_start.is_none(),
304 "failed connect should consume upstream_connect_start for duration recording"
305 );
306 }
307
308 #[test]
309 fn response_503_retries_when_status5xx_enabled() {
310 let mut ctx = PingoraRequestCtx::default();
311 ctx.request_is_idempotent = true;
312 ctx.retry_policy = Some(Arc::new(praxis_core::config::RetryPolicy {
313 configured: true,
314 retriable_conditions: vec![praxis_core::config::RetriableCondition::Status5xx],
315 ..praxis_core::config::RetryPolicy::legacy_default()
316 }));
317 let e = retry_util::maybe_retry_response(&mut ctx, 503).expect("503 should be retriable");
318 assert!(e.retry(), "503 under Status5xx should set retry");
319 assert_eq!(ctx.retries, 1);
320 assert!(ctx.reselect_on_retry);
321 assert!(ctx.pending_backoff.is_some());
322 }
323
324 #[test]
325 fn response_404_does_not_retry() {
326 let mut ctx = PingoraRequestCtx::default();
327 ctx.request_is_idempotent = true;
328 ctx.retry_policy = Some(Arc::new(praxis_core::config::RetryPolicy {
329 retriable_conditions: vec![praxis_core::config::RetriableCondition::Status5xx],
330 ..praxis_core::config::RetryPolicy::legacy_default()
331 }));
332 assert!(
333 retry_util::maybe_retry_response(&mut ctx, 404).is_none(),
334 "404 must never trigger status-based retry"
335 );
336 assert_eq!(ctx.retries, 0);
337 }
338
339 #[test]
340 fn response_502_does_not_retry_under_legacy_default() {
341 let mut ctx = PingoraRequestCtx::default();
342 ctx.request_is_idempotent = true;
343 assert!(
344 retry_util::maybe_retry_response(&mut ctx, 502).is_none(),
345 "legacy default must forward 5xx without retry"
346 );
347 }
348
349 #[test]
350 fn max_retries_zero_disables_connect_retry() {
351 let mut ctx = PingoraRequestCtx::default();
352 ctx.request_is_idempotent = true;
353 ctx.retry_policy = Some(Arc::new(praxis_core::config::RetryPolicy {
354 max_retries: Some(0),
355 ..praxis_core::config::RetryPolicy::legacy_default()
356 }));
357 let e = retry_util::handle_connect_failure(&mut ctx, make_error());
358 assert!(!e.retry(), "max_retries: 0 must disable retries");
359 assert_eq!(ctx.retries, 0);
360 }
361
362 #[test]
363 fn non_idempotent_clears_pingora_default_retry_flag() {
364 let mut ctx = PingoraRequestCtx::default();
365 ctx.request_is_idempotent = false;
366 let mut e = make_error();
367 e.set_retry(true);
368 let e = retry_util::handle_connect_failure(&mut ctx, e);
369 assert!(!e.retry(), "policy denial must clear Pingora's default retry flag");
370 assert_eq!(ctx.retries, 0);
371 }
372
373 #[tokio::test]
374 async fn logging_cleanup_noop_when_response_phase_done() {
375 let registry = praxis_filter::FilterRegistry::with_builtins();
376 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
377 let mut ctx = PingoraRequestCtx::default();
378 ctx.response_phase_done = true;
379 ctx.request_snapshot = Some(praxis_filter::Request {
380 method: http::Method::GET,
381 uri: "/".parse().unwrap(),
382 headers: http::HeaderMap::new(),
383 });
384 logging_util::logging_cleanup(&pipeline, &mut ctx).await;
385 }
386
387 #[tokio::test]
388 async fn logging_cleanup_noop_when_no_snapshot() {
389 let registry = praxis_filter::FilterRegistry::with_builtins();
390 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
391 let mut ctx = PingoraRequestCtx::default();
392 ctx.response_phase_done = false;
393 ctx.request_snapshot = None;
394 logging_util::logging_cleanup(&pipeline, &mut ctx).await;
395 }
396
397 #[tokio::test]
398 async fn logging_cleanup_runs_response_pipeline_when_needed() {
399 let registry = praxis_filter::FilterRegistry::with_builtins();
400 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
401 let mut ctx = PingoraRequestCtx::default();
402 ctx.response_phase_done = false;
403 ctx.cluster = Some(Arc::from("test-cluster"));
404 ctx.request_snapshot = Some(praxis_filter::Request {
405 method: http::Method::GET,
406 uri: "/test".parse().unwrap(),
407 headers: http::HeaderMap::new(),
408 });
409 logging_util::logging_cleanup(&pipeline, &mut ctx).await;
410 assert_eq!(
411 ctx.cluster.as_deref(),
412 Some("test-cluster"),
413 "cluster must be restored so the fallback access record can attribute the failure"
414 );
415 }
416
417 #[tokio::test]
418 async fn logging_cleanup_preserves_filter_metadata() {
419 let registry = praxis_filter::FilterRegistry::with_builtins();
420 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
421 let mut ctx = PingoraRequestCtx::default();
422 ctx.response_phase_done = false;
423 ctx.filter_metadata
424 .insert("json_rpc.method".to_owned(), "service/invoke".to_owned());
425 ctx.request_snapshot = Some(praxis_filter::Request {
426 method: http::Method::POST,
427 uri: "/api".parse().unwrap(),
428 headers: http::HeaderMap::new(),
429 });
430 logging_util::logging_cleanup(&pipeline, &mut ctx).await;
431 assert_eq!(
432 ctx.filter_metadata.get("json_rpc.method").map(String::as_str),
433 Some("service/invoke"),
434 "filter_metadata should survive logging_cleanup"
435 );
436 }
437
438 #[tokio::test]
439 async fn logging_cleanup_preserves_extensions() {
440 let registry = praxis_filter::FilterRegistry::with_builtins();
441 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
442 let mut ctx = PingoraRequestCtx::default();
443 ctx.response_phase_done = false;
444 ctx.extensions.insert(42_u32);
445 ctx.request_snapshot = Some(praxis_filter::Request {
446 method: http::Method::POST,
447 uri: "/test".parse().unwrap(),
448 headers: http::HeaderMap::new(),
449 });
450 logging_util::logging_cleanup(&pipeline, &mut ctx).await;
451 assert_eq!(
452 ctx.extensions.get::<u32>(),
453 Some(&42),
454 "extensions should survive logging_cleanup"
455 );
456 }
457
458 #[test]
459 fn passive_health_error_is_failure() {
460 let (pipeline, ctx) = make_passive_scenario(Some(3), Some(2));
461 let error = make_error();
462 health_util::record_passive_health(&pipeline, Some(&error), &ctx);
463
464 let registry = pipeline.health_registry().unwrap();
465 let entry = registry.get("test-cluster").unwrap();
466 assert!(
467 entry.endpoints()[0].is_healthy(),
468 "single failure should not yet mark unhealthy (threshold=3)"
469 );
470 }
471
472 #[test]
473 fn passive_health_downstream_error_is_not_failure() {
474 let (pipeline, ctx) = make_passive_scenario(Some(1), Some(1));
475 let error = make_error().into_down();
476 health_util::record_passive_health(&pipeline, Some(&error), &ctx);
477
478 let registry = pipeline.health_registry().unwrap();
479 let entry = registry.get("test-cluster").unwrap();
480 assert!(
481 entry.endpoints()[0].is_healthy(),
482 "a downstream/client error must not mark the endpoint unhealthy"
483 );
484 }
485
486 #[test]
487 fn passive_health_downstream_error_with_5xx_still_counts_as_failure() {
488 let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
489 ctx.upstream_response_status = Some(503);
490 let error = make_error().into_down();
491 health_util::record_passive_health(&pipeline, Some(&error), &ctx);
492
493 let registry = pipeline.health_registry().unwrap();
494 let entry = registry.get("test-cluster").unwrap();
495 assert!(
496 !entry.endpoints()[0].is_healthy(),
497 "a downstream error with an upstream 503 must still mark the endpoint unhealthy"
498 );
499 }
500
501 #[test]
502 fn passive_health_downstream_error_does_not_reset_failure_streak() {
503 let (pipeline, ctx) = make_passive_scenario(Some(2), Some(1));
504 let mut upstream_err = make_error();
505 upstream_err.as_up();
506 let downstream_err = make_error().into_down();
507
508 health_util::record_passive_health(&pipeline, Some(&upstream_err), &ctx);
509 health_util::record_passive_health(&pipeline, Some(&downstream_err), &ctx);
510 health_util::record_passive_health(&pipeline, Some(&upstream_err), &ctx);
511
512 let registry = pipeline.health_registry().unwrap();
513 let entry = registry.get("test-cluster").unwrap();
514 assert!(
515 !entry.endpoints()[0].is_healthy(),
516 "two upstream failures must eject the endpoint even with an interleaved client error"
517 );
518 }
519
520 #[test]
521 fn passive_health_upstream_error_is_failure() {
522 let (pipeline, ctx) = make_passive_scenario(Some(1), Some(1));
523 let mut error = make_error();
524 error.as_up();
525 health_util::record_passive_health(&pipeline, Some(&error), &ctx);
526
527 let registry = pipeline.health_registry().unwrap();
528 let entry = registry.get("test-cluster").unwrap();
529 assert!(
530 !entry.endpoints()[0].is_healthy(),
531 "an upstream error at unhealthy-threshold 1 must mark the endpoint unhealthy"
532 );
533 }
534
535 #[test]
536 fn passive_health_skips_observations_without_upstream_contact() {
537 let (pipeline, mut ctx) = make_passive_scenario(Some(2), Some(1));
538 let mut upstream_err = make_error();
539 upstream_err.as_up();
540
541 health_util::record_passive_health(&pipeline, Some(&upstream_err), &ctx);
542
543 ctx.upstream_contacted = false;
544 health_util::record_passive_health(&pipeline, None, &ctx);
545
546 ctx.upstream_response_status = Some(200);
547 health_util::record_passive_health(&pipeline, None, &ctx);
548 ctx.upstream_response_status = None;
549
550 ctx.upstream_contacted = true;
551 health_util::record_passive_health(&pipeline, Some(&upstream_err), &ctx);
552
553 let registry = pipeline.health_registry().unwrap();
554 let entry = registry.get("test-cluster").unwrap();
555 assert!(
556 !entry.endpoints()[0].is_healthy(),
557 "observations without upstream contact must not reset the failure streak"
558 );
559 }
560
561 #[test]
562 fn passive_health_records_connect_failure_after_reselect_clears_upstream() {
563 let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
564 ctx.upstream_for_retry = None;
565 ctx.upstream_contacted = true;
566 let mut error = make_error();
567 error.as_up();
568 health_util::record_passive_health(&pipeline, Some(&error), &ctx);
569
570 let registry = pipeline.health_registry().unwrap();
571 let entry = registry.get("test-cluster").unwrap();
572 assert!(
573 !entry.endpoints()[0].is_healthy(),
574 "a connect failure after reselection cleared upstream_for_retry must still count"
575 );
576 }
577
578 #[test]
579 fn passive_health_status_500_is_failure() {
580 let (pipeline, mut ctx) = make_passive_scenario(Some(3), Some(2));
581 ctx.upstream_response_status = Some(500);
582 health_util::record_passive_health(&pipeline, None, &ctx);
583
584 let registry = pipeline.health_registry().unwrap();
585 let entry = registry.get("test-cluster").unwrap();
586 assert!(
587 entry.endpoints()[0].is_healthy(),
588 "single 500 should not yet mark unhealthy (threshold=3)"
589 );
590 }
591
592 #[test]
593 fn passive_health_status_below_500_is_success() {
594 let (pipeline, mut ctx) = make_passive_scenario(Some(2), Some(1));
595 ctx.upstream_response_status = Some(499);
596 health_util::record_passive_health(&pipeline, None, &ctx);
597
598 let registry = pipeline.health_registry().unwrap();
599 let entry = registry.get("test-cluster").unwrap();
600 assert!(entry.endpoints()[0].is_healthy(), "status 499 should count as success");
601 }
602
603 #[test]
604 fn passive_unhealthy_threshold_transition() {
605 let (pipeline, ctx) = make_passive_scenario(Some(2), Some(1));
606 let error = make_error();
607 health_util::record_passive_health(&pipeline, Some(&error), &ctx);
608 health_util::record_passive_health(&pipeline, Some(&error), &ctx);
609
610 let registry = pipeline.health_registry().unwrap();
611 let entry = registry.get("test-cluster").unwrap();
612 assert!(
613 !entry.endpoints()[0].is_healthy(),
614 "2 consecutive failures should mark unhealthy (threshold=2)"
615 );
616 }
617
618 #[test]
619 fn passive_healthy_threshold_recovery() {
620 let (pipeline, ctx) = make_passive_scenario(Some(1), Some(2));
621 let error = make_error();
622 health_util::record_passive_health(&pipeline, Some(&error), &ctx);
623
624 let registry = pipeline.health_registry().unwrap();
625 let entry = registry.get("test-cluster").unwrap();
626 assert!(
627 !entry.endpoints()[0].is_healthy(),
628 "should be unhealthy after 1 failure"
629 );
630
631 let ctx_ok = make_passive_ctx("test-cluster", 0, Some(200));
632 health_util::record_passive_health(&pipeline, None, &ctx_ok);
633 assert!(
634 !entry.endpoints()[0].is_healthy(),
635 "one success should not recover (threshold=2)"
636 );
637
638 health_util::record_passive_health(&pipeline, None, &ctx_ok);
639 assert!(
640 entry.endpoints()[0].is_healthy(),
641 "2 consecutive successes should recover (threshold=2)"
642 );
643 }
644
645 #[test]
646 fn passive_health_no_thresholds_is_noop() {
647 let (pipeline, ctx) = make_passive_scenario(None, None);
648 let error = make_error();
649 health_util::record_passive_health(&pipeline, Some(&error), &ctx);
650
651 let registry = pipeline.health_registry().unwrap();
652 let entry = registry.get("test-cluster").unwrap();
653 assert!(
654 entry.endpoints()[0].is_healthy(),
655 "no passive thresholds means failures are no-op"
656 );
657 }
658
659 #[test]
660 fn passive_health_endpoint_index_out_of_bounds() {
661 let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
662 ctx.selected_endpoint_index = Some(999);
663 let error = make_error();
664 health_util::record_passive_health(&pipeline, Some(&error), &ctx);
665
666 let registry = pipeline.health_registry().unwrap();
667 let entry = registry.get("test-cluster").unwrap();
668 assert!(entry.endpoints()[0].is_healthy(), "out-of-bounds index should be no-op");
669 }
670
671 #[test]
672 fn passive_health_missing_cluster_is_noop() {
673 let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
674 ctx.cluster = None;
675 ctx.metrics_cluster = None;
676 let error = make_error();
677 health_util::record_passive_health(&pipeline, Some(&error), &ctx);
678 }
679
680 #[test]
681 fn passive_health_falls_back_to_metrics_cluster() {
682 let (pipeline, mut ctx) = make_passive_scenario(Some(2), Some(1));
683 ctx.cluster = None;
684 ctx.metrics_cluster = Some(Arc::from("test-cluster"));
685 let error = make_error();
686 health_util::record_passive_health(&pipeline, Some(&error), &ctx);
687 health_util::record_passive_health(&pipeline, Some(&error), &ctx);
688
689 let registry = pipeline.health_registry().unwrap();
690 let entry = registry.get("test-cluster").unwrap();
691 assert!(
692 !entry.endpoints()[0].is_healthy(),
693 "fallback to metrics_cluster should still record passive health"
694 );
695 }
696
697 #[test]
698 fn passive_health_missing_endpoint_index_is_noop() {
699 let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
700 ctx.selected_endpoint_index = None;
701 let error = make_error();
702 health_util::record_passive_health(&pipeline, Some(&error), &ctx);
703 }
704
705 #[test]
706 fn passive_health_missing_registry_is_noop() {
707 let registry = praxis_filter::FilterRegistry::with_builtins();
708 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
709 let mut ctx = PingoraRequestCtx::default();
710 ctx.cluster = Some(Arc::from("test-cluster"));
711 ctx.selected_endpoint_index = Some(0);
712 let error = make_error();
713 health_util::record_passive_health(&pipeline, Some(&error), &ctx);
714 }
715
716 #[test]
717 fn passive_health_unknown_cluster_is_noop() {
718 let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
719 ctx.cluster = Some(Arc::from("nonexistent"));
720 let error = make_error();
721 health_util::record_passive_health(&pipeline, Some(&error), &ctx);
722 }
723
724 #[test]
725 fn size_limit_none_body_returns_false() {
726 let mut bytes = 0_u64;
727 assert!(!body_util::check_body_size_limit(None, &mut bytes, 100));
728 assert_eq!(bytes, 0, "accumulated bytes unchanged for None body");
729 }
730
731 #[test]
732 fn size_limit_within_limit() {
733 let mut bytes = 0_u64;
734 let body = Some(Bytes::from_static(b"hello"));
735 assert!(!body_util::check_body_size_limit(body.as_ref(), &mut bytes, 10));
736 assert_eq!(bytes, 5);
737 }
738
739 #[test]
740 fn size_limit_at_exact_limit() {
741 let mut bytes = 0_u64;
742 let body = Some(Bytes::from_static(b"exact"));
743 assert!(!body_util::check_body_size_limit(body.as_ref(), &mut bytes, 5));
744 assert_eq!(bytes, 5);
745 }
746
747 #[test]
748 fn size_limit_exceeds_limit() {
749 let mut bytes = 0_u64;
750 let body = Some(Bytes::from_static(b"toolong"));
751 assert!(body_util::check_body_size_limit(body.as_ref(), &mut bytes, 3));
752 }
753
754 #[test]
755 fn size_limit_cumulative_overflow() {
756 let mut bytes = 0_u64;
757 let first = Some(Bytes::from_static(b"aaa"));
758 assert!(!body_util::check_body_size_limit(first.as_ref(), &mut bytes, 5));
759
760 let second = Some(Bytes::from_static(b"bbb"));
761 assert!(body_util::check_body_size_limit(second.as_ref(), &mut bytes, 5));
762 assert_eq!(bytes, 6);
763 }
764
765 #[test]
766 fn stream_buffer_accumulates_chunks() {
767 let mut body = Some(Bytes::from_static(b"hello "));
768 let mut buf: Option<BodyBuffer> = None;
769 assert!(!body_util::accumulate_stream_buffer(
770 &mut body,
771 &mut buf,
772 false,
773 Some(100)
774 ));
775 assert!(buf.is_some());
776
777 body = Some(Bytes::from_static(b"world"));
778 assert!(!body_util::accumulate_stream_buffer(
779 &mut body,
780 &mut buf,
781 false,
782 Some(100)
783 ));
784
785 let frozen = buf.take().unwrap().freeze();
786 assert_eq!(frozen, Bytes::from_static(b"hello world"));
787 }
788
789 #[test]
790 fn stream_buffer_freezes_at_eos() {
791 let mut body = Some(Bytes::from_static(b"data"));
792 let mut buf: Option<BodyBuffer> = None;
793 assert!(!body_util::accumulate_stream_buffer(
794 &mut body,
795 &mut buf,
796 false,
797 Some(100)
798 ));
799
800 body = Some(Bytes::from_static(b" end"));
801 assert!(!body_util::accumulate_stream_buffer(
802 &mut body,
803 &mut buf,
804 true,
805 Some(100)
806 ));
807 assert!(buf.is_none(), "buffer should be taken at EOS");
808 assert_eq!(body.unwrap(), Bytes::from_static(b"data end"));
809 }
810
811 #[test]
812 fn stream_buffer_overflow() {
813 let mut body = Some(Bytes::from_static(b"too long"));
814 let mut buf: Option<BodyBuffer> = None;
815 assert!(body_util::accumulate_stream_buffer(&mut body, &mut buf, false, Some(5)));
816 }
817
818 #[test]
819 fn stream_buffer_none_body() {
820 let mut body: Option<Bytes> = None;
821 let mut buf: Option<BodyBuffer> = None;
822 assert!(!body_util::accumulate_stream_buffer(
823 &mut body,
824 &mut buf,
825 false,
826 Some(100)
827 ));
828 assert!(buf.is_none());
829 }
830
831 #[test]
832 fn stream_buffer_uses_absolute_max_when_none() {
833 let mut body = Some(Bytes::from_static(b"data"));
834 let mut buf: Option<BodyBuffer> = None;
835 assert!(!body_util::accumulate_stream_buffer(&mut body, &mut buf, false, None));
836 assert!(buf.is_some(), "should create buffer with absolute max");
837 }
838
839 #[test]
840 fn suppress_clears_body_when_buffering() {
841 let mut body = Some(Bytes::from_static(b"data"));
842 body_util::suppress_stream_buffer_chunk(&mut body, true, false, false);
843 assert!(body.is_none());
844 }
845
846 #[test]
847 fn suppress_noop_when_not_stream_buffer() {
848 let mut body = Some(Bytes::from_static(b"data"));
849 body_util::suppress_stream_buffer_chunk(&mut body, false, false, false);
850 assert!(body.is_some());
851 }
852
853 #[test]
854 fn suppress_noop_when_released() {
855 let mut body = Some(Bytes::from_static(b"data"));
856 body_util::suppress_stream_buffer_chunk(&mut body, true, true, false);
857 assert!(body.is_some());
858 }
859
860 #[test]
861 fn suppress_noop_at_eos() {
862 let mut body = Some(Bytes::from_static(b"data"));
863 body_util::suppress_stream_buffer_chunk(&mut body, true, false, true);
864 assert!(body.is_some());
865 }
866
867 #[test]
868 fn release_sets_flag_and_flushes_buffer() {
869 let mut body: Option<Bytes> = None;
870 let mut released = false;
871 let mut buf = Some(BodyBuffer::new(100));
872 buf.as_mut().unwrap().push(Bytes::from_static(b"buffered")).unwrap();
873
874 body_util::release_stream_buffer(&mut body, true, &mut released, &mut buf, false);
875 assert!(released);
876 assert_eq!(body.unwrap(), Bytes::from_static(b"buffered"));
877 assert!(buf.is_none());
878 }
879
880 #[test]
881 fn release_noop_when_already_released() {
882 let mut body: Option<Bytes> = None;
883 let mut released = true;
884 let mut buf: Option<BodyBuffer> = None;
885
886 body_util::release_stream_buffer(&mut body, true, &mut released, &mut buf, false);
887 assert!(body.is_none(), "body should be unchanged when already released");
888 }
889
890 #[test]
891 fn release_noop_when_not_stream_buffer() {
892 let mut body: Option<Bytes> = None;
893 let mut released = false;
894 let mut buf: Option<BodyBuffer> = None;
895
896 body_util::release_stream_buffer(&mut body, false, &mut released, &mut buf, false);
897 assert!(!released, "released flag should be unchanged for non-stream-buffer");
898 }
899
900 #[test]
901 fn release_at_eos_sets_flag_but_no_flush() {
902 let mut body: Option<Bytes> = None;
903 let mut released = false;
904 let mut buf = Some(BodyBuffer::new(100));
905 buf.as_mut().unwrap().push(Bytes::from_static(b"data")).unwrap();
906
907 body_util::release_stream_buffer(&mut body, true, &mut released, &mut buf, true);
908 assert!(released);
909 assert!(body.is_none(), "body should not be overwritten at EOS");
910 assert!(buf.is_some(), "buffer should not be taken at EOS");
911 }
912
913 #[test]
914 fn write_back_transfers_fields() {
915 let mut ctx = PingoraRequestCtx::default();
916
917 let mut extensions = RequestExtensions::new();
918 extensions.insert(42_u32);
919
920 let state_val: Box<dyn std::any::Any + Send + Sync> = Box::new(99_i32);
921 let filter_state = HashMap::from([(0_usize, state_val)]);
922
923 let output = body_util::BodyFilterOutput {
924 cluster: Some(Arc::from("test-cluster")),
925 upstream: Some(Upstream {
926 address: Arc::from("10.0.0.1:80"),
927 authority: None,
928 connection: Arc::new(ConnectionOptions::default()),
929 tls: None,
930 }),
931 extensions,
932 attempted_endpoints: Vec::new(),
933 filter_metadata: HashMap::from([("key".to_owned(), "val".to_owned())]),
934 filter_state,
935 executed_filter_indices: vec![true, false],
936 body_done_indices: vec![false, true],
937 };
938 output.write_back(&mut ctx);
939
940 assert_eq!(ctx.cluster.as_deref(), Some("test-cluster"));
941 assert!(ctx.upstream.is_some(), "upstream should transfer");
942 assert_eq!(ctx.upstream.as_ref().unwrap().address.as_ref(), "10.0.0.1:80");
943 assert_eq!(ctx.extensions.get::<u32>(), Some(&42));
944 assert_eq!(ctx.filter_metadata.get("key").map(String::as_str), Some("val"));
945 assert_eq!(ctx.filter_state.len(), 1, "filter_state should transfer");
946 assert_eq!(
947 ctx.filter_state.get(&0).and_then(|v| v.downcast_ref::<i32>()),
948 Some(&99)
949 );
950 assert_eq!(ctx.cached_executed_filter_indices, vec![true, false]);
951 assert_eq!(ctx.cached_body_done_indices, vec![false, true]);
952 }
953
954 #[test]
959 fn fallback_access_log_emits_for_incomplete_request() {
960 let pipeline = access_log_pipeline();
961 let mut ctx = make_fallback_ctx();
962
963 let events = capture_access_events(|| logging_util::maybe_emit_fallback_access_log(&pipeline, 502, &mut ctx));
964 assert_eq!(
965 events.len(),
966 1,
967 "incomplete request must produce a fallback access record"
968 );
969 }
970
971 #[test]
972 fn fallback_access_log_skips_completed_delivery() {
973 let pipeline = access_log_pipeline();
974 let mut ctx = make_fallback_ctx();
975 ctx.response_delivery_complete = true;
976
977 let events = capture_access_events(|| logging_util::maybe_emit_fallback_access_log(&pipeline, 200, &mut ctx));
978 assert!(events.is_empty(), "completed delivery already logged via the filter");
979 }
980
981 #[test]
982 fn fallback_access_log_skips_upgraded_connections() {
983 let pipeline = access_log_pipeline();
984 let mut ctx = make_fallback_ctx();
985 ctx.connection_upgraded = true;
986
987 let events = capture_access_events(|| logging_util::maybe_emit_fallback_access_log(&pipeline, 101, &mut ctx));
988 assert!(events.is_empty(), "upgraded connections have no body completion");
989 }
990
991 #[test]
992 fn fallback_access_log_skips_without_access_log_filter() {
993 let registry = praxis_filter::FilterRegistry::with_builtins();
994 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
995 let mut ctx = make_fallback_ctx();
996
997 let events = capture_access_events(|| logging_util::maybe_emit_fallback_access_log(&pipeline, 502, &mut ctx));
998 assert!(events.is_empty(), "no access_log filter means no fallback record");
999 }
1000
1001 #[test]
1002 fn fallback_access_log_honors_entry_conditions() {
1003 let registry = praxis_filter::FilterRegistry::with_builtins();
1004 let mut entries = vec![praxis_filter::FilterEntry {
1005 branch_chains: None,
1006 conditions: vec![serde_yaml::from_str("when:\n path_prefix: /api\n").unwrap()],
1007 failure_mode: praxis_filter::FailureMode::default(),
1008 filter_type: "access_log".to_owned(),
1009 config: serde_yaml::Value::Null,
1010 name: None,
1011 response_conditions: vec![],
1012 }];
1013 let pipeline = FilterPipeline::build(&mut entries, ®istry).unwrap();
1014
1015 let mut excluded = make_fallback_ctx();
1016 let events =
1017 capture_access_events(|| logging_util::maybe_emit_fallback_access_log(&pipeline, 502, &mut excluded));
1018 assert!(
1019 events.is_empty(),
1020 "requests the operator scoped out must not gain fallback records"
1021 );
1022
1023 let mut included = make_fallback_ctx();
1024 if let Some(snapshot) = included.request_snapshot.as_mut() {
1025 snapshot.uri = "/api/users".parse().unwrap();
1026 }
1027 let events =
1028 capture_access_events(|| logging_util::maybe_emit_fallback_access_log(&pipeline, 502, &mut included));
1029 assert_eq!(events.len(), 1, "in-scope incomplete requests still get a record");
1030 }
1031
1032 #[cfg(feature = "upstream-binding")]
1033 #[test]
1034 fn fallback_access_log_honors_a_bound_upstream_condition() {
1035 for (axis, bound_records, unbound_records) in [("when", 1, 0), ("unless", 0, 1)] {
1036 let pipeline = access_log_pipeline_gated_on_binding(axis);
1037 let mut bound = make_fallback_ctx();
1038 bind_through_the_pipeline(&pipeline, &mut bound);
1039 let mut unbound = make_fallback_ctx();
1040
1041 let bound_events =
1042 capture_access_events(|| logging_util::maybe_emit_fallback_access_log(&pipeline, 502, &mut bound));
1043 let unbound_events =
1044 capture_access_events(|| logging_util::maybe_emit_fallback_access_log(&pipeline, 502, &mut unbound));
1045
1046 assert_eq!(
1047 bound_events.len(),
1048 bound_records,
1049 "{axis}: a request the router bound to openai is matched against its real binding"
1050 );
1051 assert_eq!(
1052 unbound_events.len(),
1053 unbound_records,
1054 "{axis}: a request that failed before routing has no binding to match"
1055 );
1056 }
1057 }
1058
1059 #[test]
1060 fn aborted_response_body_at_eos_is_not_marked_delivered() {
1061 let pipeline = access_log_pipeline();
1062 let mut ctx = make_fallback_ctx();
1063 ctx.response_body_mode = BodyMode::SizeLimit { max_bytes: 4 };
1064 let mut body = Some(Bytes::from_static(b"exceeds the limit"));
1065
1066 let result = response_body_filter::execute(&pipeline, &mut body, true, &mut ctx);
1067 assert!(result.is_err(), "over-limit body must abort");
1068 assert!(
1069 !ctx.response_delivery_complete,
1070 "a response aborted at end-of-stream was not delivered; the fallback record must fire"
1071 );
1072 }
1073
1074 #[test]
1075 fn response_body_eos_marks_delivery_complete() {
1076 let registry = praxis_filter::FilterRegistry::with_builtins();
1077 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
1078 let mut ctx = PingoraRequestCtx::default();
1079 let mut body: Option<Bytes> = None;
1080
1081 let _timeout = response_body_filter::execute(&pipeline, &mut body, false, &mut ctx).unwrap();
1082 assert!(
1083 !ctx.response_delivery_complete,
1084 "mid-stream chunks must not mark delivery complete"
1085 );
1086
1087 let _timeout = response_body_filter::execute(&pipeline, &mut body, true, &mut ctx).unwrap();
1088 assert!(ctx.response_delivery_complete, "end-of-stream marks delivery complete");
1089 }
1090
1091 #[test]
1096 fn http_version_label_http_09() {
1097 assert_eq!(
1098 span_util::http_version_label(http::Version::HTTP_09),
1099 "0.9",
1100 "HTTP/0.9 should map to '0.9'"
1101 );
1102 }
1103
1104 #[test]
1105 fn http_version_label_http_10() {
1106 assert_eq!(
1107 span_util::http_version_label(http::Version::HTTP_10),
1108 "1.0",
1109 "HTTP/1.0 should map to '1.0'"
1110 );
1111 }
1112
1113 #[test]
1114 fn http_version_label_http_11() {
1115 assert_eq!(
1116 span_util::http_version_label(http::Version::HTTP_11),
1117 "1.1",
1118 "HTTP/1.1 should map to '1.1'"
1119 );
1120 }
1121
1122 #[test]
1123 fn http_version_label_http_2() {
1124 assert_eq!(
1125 span_util::http_version_label(http::Version::HTTP_2),
1126 "2",
1127 "HTTP/2 should map to '2'"
1128 );
1129 }
1130
1131 #[test]
1132 fn http_version_label_http_3() {
1133 assert_eq!(
1134 span_util::http_version_label(http::Version::HTTP_3),
1135 "3",
1136 "HTTP/3 should map to '3'"
1137 );
1138 }
1139
1140 #[test]
1141 fn record_response_span_attributes_noop_for_disabled_span() {
1142 let ctx = PingoraRequestCtx::default();
1143 assert!(ctx.request_span.is_disabled(), "default span should be disabled");
1144 }
1145
1146 #[derive(Clone, Default)]
1148 struct RecordCapture(Arc<std::sync::Mutex<Vec<(String, String)>>>);
1149
1150 impl<S> tracing_subscriber::Layer<S> for RecordCapture
1151 where
1152 S: tracing::Subscriber + for<'a> tracing_subscriber::registry::LookupSpan<'a>,
1153 {
1154 fn on_record(
1155 &self,
1156 _id: &tracing::span::Id,
1157 values: &tracing::span::Record<'_>,
1158 _ctx: tracing_subscriber::layer::Context<'_, S>,
1159 ) {
1160 struct Visitor<'a>(&'a mut Vec<(String, String)>);
1161 impl tracing::field::Visit for Visitor<'_> {
1162 fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
1163 self.0.push((field.name().to_owned(), format!("{value:?}")));
1164 }
1165 }
1166 let mut captured = self.0.lock().expect("capture lock");
1167 values.record(&mut Visitor(&mut captured));
1168 }
1169 }
1170
1171 #[test]
1172 fn record_response_span_fields_records_status_upstream_and_cluster() {
1173 use tracing_subscriber::layer::SubscriberExt as _;
1174
1175 let capture = RecordCapture::default();
1176 let subscriber = tracing_subscriber::registry().with(capture.clone());
1177 let _guard = tracing::subscriber::set_default(subscriber);
1178
1179 let mut ctx = PingoraRequestCtx::default();
1180 ctx.metrics_cluster = Some(Arc::from("api-cluster"));
1181 ctx.upstream_for_retry = Some(Upstream {
1182 address: Arc::from("10.0.0.1:80"),
1183 authority: None,
1184 connection: Arc::new(ConnectionOptions::default()),
1185 tls: None,
1186 });
1187 ctx.request_span = tracing::info_span!(
1188 "test_span",
1189 "http.response.status_code" = tracing::field::Empty,
1190 "otel.status_code" = tracing::field::Empty,
1191 "upstream.address" = tracing::field::Empty,
1192 "upstream.cluster" = tracing::field::Empty,
1193 );
1194
1195 span_util::record_response_span_fields(Some(http::StatusCode::SERVICE_UNAVAILABLE), "GET", None, &ctx);
1196
1197 let captured = capture.0.lock().expect("capture lock");
1198 let get = |name: &str| {
1199 let value = captured.iter().find(|(f, _)| f == name).map(|(_, v)| v.clone());
1200 assert!(value.is_some(), "field {name} not recorded; got {captured:?}");
1201 value.unwrap_or_default()
1202 };
1203 assert_eq!(get("http.response.status_code"), "503");
1204 assert_eq!(get("otel.status_code"), "\"ERROR\"", "5xx must set otel error status");
1205 assert_eq!(get("upstream.address"), "\"10.0.0.1:80\"");
1206 assert_eq!(get("upstream.cluster"), "\"api-cluster\"");
1207 }
1208
1209 #[test]
1210 fn record_response_span_fields_success_has_no_error_status() {
1211 use tracing_subscriber::layer::SubscriberExt as _;
1212
1213 let capture = RecordCapture::default();
1214 let subscriber = tracing_subscriber::registry().with(capture.clone());
1215 let _guard = tracing::subscriber::set_default(subscriber);
1216
1217 let mut ctx = PingoraRequestCtx::default();
1218 ctx.request_span = tracing::info_span!(
1219 "test_span",
1220 "http.response.status_code" = tracing::field::Empty,
1221 "otel.status_code" = tracing::field::Empty,
1222 );
1223
1224 span_util::record_response_span_fields(Some(http::StatusCode::OK), "GET", None, &ctx);
1225
1226 let captured = capture.0.lock().expect("capture lock");
1227 assert!(
1228 captured
1229 .iter()
1230 .any(|(f, v)| f == "http.response.status_code" && v == "200"),
1231 "status should be recorded: {captured:?}"
1232 );
1233 assert!(
1234 !captured.iter().any(|(f, _)| f == "otel.status_code"),
1235 "2xx must not set otel error status: {captured:?}"
1236 );
1237 }
1238
1239 #[test]
1240 fn record_response_span_attributes_records_exchange_span_fields() {
1241 let mut ctx = PingoraRequestCtx::default();
1242 ctx.request_span = tracing::info_span!(
1243 "test_request",
1244 "http.response.status_code" = tracing::field::Empty,
1245 "server.address" = tracing::field::Empty,
1246 "upstream.cluster" = tracing::field::Empty,
1247 );
1248 ctx.upstream_exchange_span = tracing::info_span!(
1249 parent: &ctx.request_span,
1250 "upstream_exchange",
1251 "http.response.status_code" = tracing::field::Empty,
1252 "http.response.body.size" = tracing::field::Empty,
1253 );
1254 ctx.response_body_bytes = 4096;
1255
1256 ctx.upstream_exchange_span.record("http.response.status_code", 200_u16);
1257 ctx.upstream_exchange_span.record("http.response.body.size", 4096_u64);
1258 }
1259
1260 #[test]
1261 fn record_response_span_attributes_skips_exchange_when_disabled() {
1262 let mut ctx = PingoraRequestCtx::default();
1263 ctx.request_span = tracing::info_span!(
1264 "test_request",
1265 "http.response.status_code" = tracing::field::Empty,
1266 "server.address" = tracing::field::Empty,
1267 "upstream.cluster" = tracing::field::Empty,
1268 );
1269 assert!(
1270 ctx.upstream_exchange_span.is_disabled(),
1271 "exchange span should be disabled by default"
1272 );
1273 }
1274
1275 #[test]
1280 fn retry_with_upstream_address_sets_retry_flag() {
1281 let mut ctx = PingoraRequestCtx::default();
1282 ctx.request_is_idempotent = true;
1283 ctx.upstream_for_retry = Some(Upstream {
1284 address: Arc::from("10.0.0.1:8080"),
1285 connection: Arc::new(ConnectionOptions::default()),
1286 tls: None,
1287 authority: None,
1288 });
1289 let e = retry_util::handle_connect_failure(&mut ctx, make_error());
1290 assert!(e.retry(), "should retry with upstream address present");
1291 assert_eq!(ctx.retries, 1, "retry counter should increment to 1");
1292 }
1293
1294 #[test]
1295 fn retry_without_upstream_address_uses_fallback() {
1296 let mut ctx = PingoraRequestCtx::default();
1297 ctx.request_is_idempotent = true;
1298 ctx.upstream_for_retry = None;
1299 let e = retry_util::handle_connect_failure(&mut ctx, make_error());
1300 assert!(
1301 e.retry(),
1302 "should retry even when upstream_for_retry is None (address defaults to unknown)"
1303 );
1304 assert_eq!(ctx.retries, 1, "retry counter should increment to 1");
1305 }
1306
1307 #[test]
1308 fn retry_exhausted_with_upstream_address_does_not_retry() {
1309 let mut ctx = PingoraRequestCtx::default();
1310 ctx.request_is_idempotent = true;
1311 ctx.retries = MAX_RETRIES as u32;
1312 ctx.upstream_for_retry = Some(Upstream {
1313 address: Arc::from("10.0.0.2:443"),
1314 connection: Arc::new(ConnectionOptions::default()),
1315 tls: None,
1316 authority: None,
1317 });
1318 let e = retry_util::handle_connect_failure(&mut ctx, make_error());
1319 assert!(
1320 !e.retry(),
1321 "should not retry after MAX_RETRIES even with upstream address"
1322 );
1323 }
1324
1325 #[test]
1326 fn large_body_skip_with_upstream_address() {
1327 let mut ctx = PingoraRequestCtx::default();
1328 ctx.request_is_idempotent = true;
1329 ctx.request_body_bytes = RETRY_BODY_LIMIT + 1;
1330 ctx.upstream_for_retry = Some(Upstream {
1331 address: Arc::from("10.0.0.3:8080"),
1332 connection: Arc::new(ConnectionOptions::default()),
1333 tls: None,
1334 authority: None,
1335 });
1336 let e = retry_util::handle_connect_failure(&mut ctx, make_error());
1337 assert!(!e.retry(), "should not retry large body even with upstream address");
1338 assert_eq!(ctx.retries, 0, "retry counter should not increment");
1339 }
1340
1341 fn make_error() -> Box<pingora_core::Error> {
1347 pingora_core::Error::explain(pingora_core::ErrorType::ConnectError, "test connect failure")
1348 }
1349
1350 fn access_log_pipeline() -> FilterPipeline {
1352 let registry = praxis_filter::FilterRegistry::with_builtins();
1353 let mut entries = vec![praxis_filter::FilterEntry {
1354 branch_chains: None,
1355 conditions: vec![],
1356 failure_mode: praxis_filter::FailureMode::default(),
1357 filter_type: "access_log".to_owned(),
1358 config: serde_yaml::Value::Null,
1359 name: None,
1360 response_conditions: vec![],
1361 }];
1362 FilterPipeline::build(&mut entries, ®istry).unwrap()
1363 }
1364
1365 #[cfg(feature = "upstream-binding")]
1369 fn access_log_pipeline_gated_on_binding(axis: &str) -> FilterPipeline {
1370 let registry = praxis_filter::FilterRegistry::with_builtins();
1371 let mut entries: Vec<praxis_filter::FilterEntry> = serde_yaml::from_str(&format!(
1372 r#"
1373- filter: router
1374 routes: [{{path_prefix: "/", cluster: backend}}]
1375- filter: access_log
1376 conditions: [{{{axis}: {{bound_upstream: {{application_provider: openai}}}}}}]
1377- filter: load_balancer
1378 clusters: [{{name: backend, http: {{application_provider: openai}}, endpoints: ["127.0.0.1:9"]}}]
1379"#
1380 ))
1381 .unwrap();
1382 FilterPipeline::build(&mut entries, ®istry).unwrap()
1383 }
1384
1385 #[cfg(feature = "upstream-binding")]
1388 fn bind_through_the_pipeline(pipeline: &FilterPipeline, ctx: &mut PingoraRequestCtx) {
1389 let runtime = tokio::runtime::Builder::new_current_thread().build().unwrap();
1390 let extensions = {
1391 let mut filter_ctx = ctx.filter_context_for(pipeline, None).expect("the snapshot is present");
1392 drop(
1393 runtime
1394 .block_on(pipeline.execute_http_request(&mut filter_ctx))
1395 .expect("the router binds the catch-all route"),
1396 );
1397 std::mem::take(&mut filter_ctx.extensions)
1398 };
1399 ctx.extensions = extensions;
1400 }
1401
1402 fn make_fallback_ctx() -> PingoraRequestCtx {
1403 let mut ctx = PingoraRequestCtx::default();
1404 ctx.request_snapshot = Some(praxis_filter::Request {
1405 method: http::Method::GET,
1406 uri: "/incomplete".parse().unwrap(),
1407 headers: http::HeaderMap::new(),
1408 });
1409 ctx
1410 }
1411
1412 fn capture_access_events<F: FnOnce()>(f: F) -> Vec<String> {
1414 use tracing_subscriber::layer::SubscriberExt as _;
1415
1416 let messages = Arc::new(std::sync::Mutex::new(Vec::<String>::new()));
1417 let capture = AccessCapture(Arc::clone(&messages));
1418 let subscriber = tracing_subscriber::registry().with(capture);
1419 tracing::subscriber::with_default(subscriber, f);
1420 let mut guard = messages.lock().unwrap();
1421 std::mem::take(&mut *guard)
1422 }
1423
1424 struct AccessCapture(Arc<std::sync::Mutex<Vec<String>>>);
1426
1427 impl<S: tracing::Subscriber> tracing_subscriber::Layer<S> for AccessCapture {
1428 fn on_event(&self, event: &tracing::Event<'_>, _ctx: tracing_subscriber::layer::Context<'_, S>) {
1429 let mut visitor = AccessMessageVisitor(String::new());
1430 event.record(&mut visitor);
1431 if visitor.0.contains("access") {
1432 self.0.lock().unwrap().push(visitor.0);
1433 }
1434 }
1435 }
1436
1437 struct AccessMessageVisitor(String);
1439
1440 impl tracing::field::Visit for AccessMessageVisitor {
1441 fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
1442 if field.name() == "message" {
1443 self.0 = format!("{value:?}");
1444 }
1445 }
1446 }
1447
1448 fn make_passive_ctx(cluster: &str, endpoint_idx: usize, status: Option<u16>) -> PingoraRequestCtx {
1450 let mut ctx = PingoraRequestCtx::default();
1451 ctx.cluster = Some(Arc::from(cluster));
1452 ctx.selected_endpoint_index = Some(endpoint_idx);
1453 ctx.upstream_response_status = status;
1454 ctx.upstream_contacted = true;
1456 ctx
1457 }
1458
1459 fn make_passive_scenario(
1462 passive_unhealthy: Option<u32>,
1463 passive_healthy: Option<u32>,
1464 ) -> (FilterPipeline, PingoraRequestCtx) {
1465 let entry = ClusterHealthEntry::new(
1466 vec![EndpointHealth::new()],
1467 vec![Arc::from("10.0.0.1:80")],
1468 passive_unhealthy,
1469 passive_healthy,
1470 );
1471 let mut map = HashMap::new();
1472 map.insert(Arc::from("test-cluster"), Arc::new(entry));
1473 let health_registry = Arc::new(map);
1474
1475 let registry = praxis_filter::FilterRegistry::with_builtins();
1476 let mut pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
1477 pipeline.set_health_registry(health_registry);
1478
1479 let ctx = make_passive_ctx("test-cluster", 0, None);
1480
1481 (pipeline, ctx)
1482 }
1483}