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_branch_filters: vec![false, true, true],
936 executed_filter_indices: vec![true, false],
937 body_done_indices: vec![false, true],
938 };
939 output.write_back(&mut ctx);
940
941 assert_eq!(ctx.cluster.as_deref(), Some("test-cluster"));
942 assert!(ctx.upstream.is_some(), "upstream should transfer");
943 assert_eq!(ctx.upstream.as_ref().unwrap().address.as_ref(), "10.0.0.1:80");
944 assert_eq!(ctx.extensions.get::<u32>(), Some(&42));
945 assert_eq!(ctx.filter_metadata.get("key").map(String::as_str), Some("val"));
946 assert_eq!(ctx.filter_state.len(), 1, "filter_state should transfer");
947 assert_eq!(
948 ctx.filter_state.get(&0).and_then(|v| v.downcast_ref::<i32>()),
949 Some(&99)
950 );
951 assert_eq!(ctx.cached_executed_branch_filters, vec![false, true, true]);
952 assert_eq!(ctx.cached_executed_filter_indices, vec![true, false]);
953 assert_eq!(ctx.cached_body_done_indices, vec![false, true]);
954 }
955
956 #[test]
961 fn fallback_access_log_emits_for_incomplete_request() {
962 let pipeline = access_log_pipeline();
963 let mut ctx = make_fallback_ctx();
964
965 let events = capture_access_events(|| logging_util::maybe_emit_fallback_access_log(&pipeline, 502, &mut ctx));
966 assert_eq!(
967 events.len(),
968 1,
969 "incomplete request must produce a fallback access record"
970 );
971 }
972
973 #[test]
974 fn fallback_access_log_skips_completed_delivery() {
975 let pipeline = access_log_pipeline();
976 let mut ctx = make_fallback_ctx();
977 ctx.response_delivery_complete = true;
978
979 let events = capture_access_events(|| logging_util::maybe_emit_fallback_access_log(&pipeline, 200, &mut ctx));
980 assert!(events.is_empty(), "completed delivery already logged via the filter");
981 }
982
983 #[test]
984 fn fallback_access_log_skips_upgraded_connections() {
985 let pipeline = access_log_pipeline();
986 let mut ctx = make_fallback_ctx();
987 ctx.connection_upgraded = true;
988
989 let events = capture_access_events(|| logging_util::maybe_emit_fallback_access_log(&pipeline, 101, &mut ctx));
990 assert!(events.is_empty(), "upgraded connections have no body completion");
991 }
992
993 #[test]
994 fn fallback_access_log_skips_without_access_log_filter() {
995 let registry = praxis_filter::FilterRegistry::with_builtins();
996 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
997 let mut ctx = make_fallback_ctx();
998
999 let events = capture_access_events(|| logging_util::maybe_emit_fallback_access_log(&pipeline, 502, &mut ctx));
1000 assert!(events.is_empty(), "no access_log filter means no fallback record");
1001 }
1002
1003 #[test]
1004 fn fallback_access_log_honors_entry_conditions() {
1005 let registry = praxis_filter::FilterRegistry::with_builtins();
1006 let mut entries = vec![praxis_filter::FilterEntry {
1007 branch_chains: None,
1008 conditions: vec![serde_yaml::from_str("when:\n path_prefix: /api\n").unwrap()],
1009 failure_mode: praxis_filter::FailureMode::default(),
1010 filter_type: "access_log".to_owned(),
1011 config: serde_yaml::Value::Null,
1012 name: None,
1013 response_conditions: vec![],
1014 }];
1015 let pipeline = FilterPipeline::build(&mut entries, ®istry).unwrap();
1016
1017 let mut excluded = make_fallback_ctx();
1018 let events =
1019 capture_access_events(|| logging_util::maybe_emit_fallback_access_log(&pipeline, 502, &mut excluded));
1020 assert!(
1021 events.is_empty(),
1022 "requests the operator scoped out must not gain fallback records"
1023 );
1024
1025 let mut included = make_fallback_ctx();
1026 if let Some(snapshot) = included.request_snapshot.as_mut() {
1027 snapshot.uri = "/api/users".parse().unwrap();
1028 }
1029 let events =
1030 capture_access_events(|| logging_util::maybe_emit_fallback_access_log(&pipeline, 502, &mut included));
1031 assert_eq!(events.len(), 1, "in-scope incomplete requests still get a record");
1032 }
1033
1034 #[cfg(feature = "upstream-binding")]
1035 #[test]
1036 fn fallback_access_log_honors_a_bound_upstream_condition() {
1037 for (axis, bound_records, unbound_records) in [("when", 1, 0), ("unless", 0, 1)] {
1038 let pipeline = access_log_pipeline_gated_on_binding(axis);
1039 let mut bound = make_fallback_ctx();
1040 bind_through_the_pipeline(&pipeline, &mut bound);
1041 let mut unbound = make_fallback_ctx();
1042
1043 let bound_events =
1044 capture_access_events(|| logging_util::maybe_emit_fallback_access_log(&pipeline, 502, &mut bound));
1045 let unbound_events =
1046 capture_access_events(|| logging_util::maybe_emit_fallback_access_log(&pipeline, 502, &mut unbound));
1047
1048 assert_eq!(
1049 bound_events.len(),
1050 bound_records,
1051 "{axis}: a request the router bound to openai is matched against its real binding"
1052 );
1053 assert_eq!(
1054 unbound_events.len(),
1055 unbound_records,
1056 "{axis}: a request that failed before routing has no binding to match"
1057 );
1058 }
1059 }
1060
1061 #[test]
1062 fn aborted_response_body_at_eos_is_not_marked_delivered() {
1063 let pipeline = access_log_pipeline();
1064 let mut ctx = make_fallback_ctx();
1065 ctx.response_body_mode = BodyMode::SizeLimit { max_bytes: 4 };
1066 let mut body = Some(Bytes::from_static(b"exceeds the limit"));
1067
1068 let result = response_body_filter::execute(&pipeline, &mut body, true, &mut ctx);
1069 assert!(result.is_err(), "over-limit body must abort");
1070 assert!(
1071 !ctx.response_delivery_complete,
1072 "a response aborted at end-of-stream was not delivered; the fallback record must fire"
1073 );
1074 }
1075
1076 #[test]
1077 fn response_body_eos_marks_delivery_complete() {
1078 let registry = praxis_filter::FilterRegistry::with_builtins();
1079 let pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
1080 let mut ctx = PingoraRequestCtx::default();
1081 let mut body: Option<Bytes> = None;
1082
1083 let _timeout = response_body_filter::execute(&pipeline, &mut body, false, &mut ctx).unwrap();
1084 assert!(
1085 !ctx.response_delivery_complete,
1086 "mid-stream chunks must not mark delivery complete"
1087 );
1088
1089 let _timeout = response_body_filter::execute(&pipeline, &mut body, true, &mut ctx).unwrap();
1090 assert!(ctx.response_delivery_complete, "end-of-stream marks delivery complete");
1091 }
1092
1093 #[test]
1098 fn http_version_label_http_09() {
1099 assert_eq!(
1100 span_util::http_version_label(http::Version::HTTP_09),
1101 "0.9",
1102 "HTTP/0.9 should map to '0.9'"
1103 );
1104 }
1105
1106 #[test]
1107 fn http_version_label_http_10() {
1108 assert_eq!(
1109 span_util::http_version_label(http::Version::HTTP_10),
1110 "1.0",
1111 "HTTP/1.0 should map to '1.0'"
1112 );
1113 }
1114
1115 #[test]
1116 fn http_version_label_http_11() {
1117 assert_eq!(
1118 span_util::http_version_label(http::Version::HTTP_11),
1119 "1.1",
1120 "HTTP/1.1 should map to '1.1'"
1121 );
1122 }
1123
1124 #[test]
1125 fn http_version_label_http_2() {
1126 assert_eq!(
1127 span_util::http_version_label(http::Version::HTTP_2),
1128 "2",
1129 "HTTP/2 should map to '2'"
1130 );
1131 }
1132
1133 #[test]
1134 fn http_version_label_http_3() {
1135 assert_eq!(
1136 span_util::http_version_label(http::Version::HTTP_3),
1137 "3",
1138 "HTTP/3 should map to '3'"
1139 );
1140 }
1141
1142 #[test]
1143 fn record_response_span_attributes_noop_for_disabled_span() {
1144 let ctx = PingoraRequestCtx::default();
1145 assert!(ctx.request_span.is_disabled(), "default span should be disabled");
1146 }
1147
1148 #[derive(Clone, Default)]
1150 struct RecordCapture(Arc<std::sync::Mutex<Vec<(String, String)>>>);
1151
1152 impl<S> tracing_subscriber::Layer<S> for RecordCapture
1153 where
1154 S: tracing::Subscriber + for<'a> tracing_subscriber::registry::LookupSpan<'a>,
1155 {
1156 fn on_record(
1157 &self,
1158 _id: &tracing::span::Id,
1159 values: &tracing::span::Record<'_>,
1160 _ctx: tracing_subscriber::layer::Context<'_, S>,
1161 ) {
1162 struct Visitor<'a>(&'a mut Vec<(String, String)>);
1163 impl tracing::field::Visit for Visitor<'_> {
1164 fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
1165 self.0.push((field.name().to_owned(), format!("{value:?}")));
1166 }
1167 }
1168 let mut captured = self.0.lock().expect("capture lock");
1169 values.record(&mut Visitor(&mut captured));
1170 }
1171 }
1172
1173 #[test]
1174 fn record_response_span_fields_records_status_upstream_and_cluster() {
1175 use tracing_subscriber::layer::SubscriberExt as _;
1176
1177 let capture = RecordCapture::default();
1178 let subscriber = tracing_subscriber::registry().with(capture.clone());
1179 let _guard = tracing::subscriber::set_default(subscriber);
1180
1181 let mut ctx = PingoraRequestCtx::default();
1182 ctx.metrics_cluster = Some(Arc::from("api-cluster"));
1183 ctx.upstream_for_retry = Some(Upstream {
1184 address: Arc::from("10.0.0.1:80"),
1185 authority: None,
1186 connection: Arc::new(ConnectionOptions::default()),
1187 tls: None,
1188 });
1189 ctx.request_span = tracing::info_span!(
1190 "test_span",
1191 "http.response.status_code" = tracing::field::Empty,
1192 "otel.status_code" = tracing::field::Empty,
1193 "upstream.address" = tracing::field::Empty,
1194 "upstream.cluster" = tracing::field::Empty,
1195 );
1196
1197 span_util::record_response_span_fields(Some(http::StatusCode::SERVICE_UNAVAILABLE), "GET", None, &ctx);
1198
1199 let captured = capture.0.lock().expect("capture lock");
1200 let get = |name: &str| {
1201 let value = captured.iter().find(|(f, _)| f == name).map(|(_, v)| v.clone());
1202 assert!(value.is_some(), "field {name} not recorded; got {captured:?}");
1203 value.unwrap_or_default()
1204 };
1205 assert_eq!(get("http.response.status_code"), "503");
1206 assert_eq!(get("otel.status_code"), "\"ERROR\"", "5xx must set otel error status");
1207 assert_eq!(get("upstream.address"), "\"10.0.0.1:80\"");
1208 assert_eq!(get("upstream.cluster"), "\"api-cluster\"");
1209 }
1210
1211 #[test]
1212 fn record_response_span_fields_success_has_no_error_status() {
1213 use tracing_subscriber::layer::SubscriberExt as _;
1214
1215 let capture = RecordCapture::default();
1216 let subscriber = tracing_subscriber::registry().with(capture.clone());
1217 let _guard = tracing::subscriber::set_default(subscriber);
1218
1219 let mut ctx = PingoraRequestCtx::default();
1220 ctx.request_span = tracing::info_span!(
1221 "test_span",
1222 "http.response.status_code" = tracing::field::Empty,
1223 "otel.status_code" = tracing::field::Empty,
1224 );
1225
1226 span_util::record_response_span_fields(Some(http::StatusCode::OK), "GET", None, &ctx);
1227
1228 let captured = capture.0.lock().expect("capture lock");
1229 assert!(
1230 captured
1231 .iter()
1232 .any(|(f, v)| f == "http.response.status_code" && v == "200"),
1233 "status should be recorded: {captured:?}"
1234 );
1235 assert!(
1236 !captured.iter().any(|(f, _)| f == "otel.status_code"),
1237 "2xx must not set otel error status: {captured:?}"
1238 );
1239 }
1240
1241 #[test]
1242 fn record_response_span_attributes_records_exchange_span_fields() {
1243 let mut ctx = PingoraRequestCtx::default();
1244 ctx.request_span = tracing::info_span!(
1245 "test_request",
1246 "http.response.status_code" = tracing::field::Empty,
1247 "server.address" = tracing::field::Empty,
1248 "upstream.cluster" = tracing::field::Empty,
1249 );
1250 ctx.upstream_exchange_span = tracing::info_span!(
1251 parent: &ctx.request_span,
1252 "upstream_exchange",
1253 "http.response.status_code" = tracing::field::Empty,
1254 "http.response.body.size" = tracing::field::Empty,
1255 );
1256 ctx.response_body_bytes = 4096;
1257
1258 ctx.upstream_exchange_span.record("http.response.status_code", 200_u16);
1259 ctx.upstream_exchange_span.record("http.response.body.size", 4096_u64);
1260 }
1261
1262 #[test]
1263 fn record_response_span_attributes_skips_exchange_when_disabled() {
1264 let mut ctx = PingoraRequestCtx::default();
1265 ctx.request_span = tracing::info_span!(
1266 "test_request",
1267 "http.response.status_code" = tracing::field::Empty,
1268 "server.address" = tracing::field::Empty,
1269 "upstream.cluster" = tracing::field::Empty,
1270 );
1271 assert!(
1272 ctx.upstream_exchange_span.is_disabled(),
1273 "exchange span should be disabled by default"
1274 );
1275 }
1276
1277 #[test]
1282 fn retry_with_upstream_address_sets_retry_flag() {
1283 let mut ctx = PingoraRequestCtx::default();
1284 ctx.request_is_idempotent = true;
1285 ctx.upstream_for_retry = Some(Upstream {
1286 address: Arc::from("10.0.0.1:8080"),
1287 connection: Arc::new(ConnectionOptions::default()),
1288 tls: None,
1289 authority: None,
1290 });
1291 let e = retry_util::handle_connect_failure(&mut ctx, make_error());
1292 assert!(e.retry(), "should retry with upstream address present");
1293 assert_eq!(ctx.retries, 1, "retry counter should increment to 1");
1294 }
1295
1296 #[test]
1297 fn retry_without_upstream_address_uses_fallback() {
1298 let mut ctx = PingoraRequestCtx::default();
1299 ctx.request_is_idempotent = true;
1300 ctx.upstream_for_retry = None;
1301 let e = retry_util::handle_connect_failure(&mut ctx, make_error());
1302 assert!(
1303 e.retry(),
1304 "should retry even when upstream_for_retry is None (address defaults to unknown)"
1305 );
1306 assert_eq!(ctx.retries, 1, "retry counter should increment to 1");
1307 }
1308
1309 #[test]
1310 fn retry_exhausted_with_upstream_address_does_not_retry() {
1311 let mut ctx = PingoraRequestCtx::default();
1312 ctx.request_is_idempotent = true;
1313 ctx.retries = MAX_RETRIES as u32;
1314 ctx.upstream_for_retry = Some(Upstream {
1315 address: Arc::from("10.0.0.2:443"),
1316 connection: Arc::new(ConnectionOptions::default()),
1317 tls: None,
1318 authority: None,
1319 });
1320 let e = retry_util::handle_connect_failure(&mut ctx, make_error());
1321 assert!(
1322 !e.retry(),
1323 "should not retry after MAX_RETRIES even with upstream address"
1324 );
1325 }
1326
1327 #[test]
1328 fn large_body_skip_with_upstream_address() {
1329 let mut ctx = PingoraRequestCtx::default();
1330 ctx.request_is_idempotent = true;
1331 ctx.request_body_bytes = RETRY_BODY_LIMIT + 1;
1332 ctx.upstream_for_retry = Some(Upstream {
1333 address: Arc::from("10.0.0.3:8080"),
1334 connection: Arc::new(ConnectionOptions::default()),
1335 tls: None,
1336 authority: None,
1337 });
1338 let e = retry_util::handle_connect_failure(&mut ctx, make_error());
1339 assert!(!e.retry(), "should not retry large body even with upstream address");
1340 assert_eq!(ctx.retries, 0, "retry counter should not increment");
1341 }
1342
1343 fn make_error() -> Box<pingora_core::Error> {
1349 pingora_core::Error::explain(pingora_core::ErrorType::ConnectError, "test connect failure")
1350 }
1351
1352 fn access_log_pipeline() -> FilterPipeline {
1354 let registry = praxis_filter::FilterRegistry::with_builtins();
1355 let mut entries = vec![praxis_filter::FilterEntry {
1356 branch_chains: None,
1357 conditions: vec![],
1358 failure_mode: praxis_filter::FailureMode::default(),
1359 filter_type: "access_log".to_owned(),
1360 config: serde_yaml::Value::Null,
1361 name: None,
1362 response_conditions: vec![],
1363 }];
1364 FilterPipeline::build(&mut entries, ®istry).unwrap()
1365 }
1366
1367 #[cfg(feature = "upstream-binding")]
1371 fn access_log_pipeline_gated_on_binding(axis: &str) -> FilterPipeline {
1372 let registry = praxis_filter::FilterRegistry::with_builtins();
1373 let mut entries: Vec<praxis_filter::FilterEntry> = serde_yaml::from_str(&format!(
1374 r#"
1375- filter: router
1376 routes: [{{path_prefix: "/", cluster: backend}}]
1377- filter: access_log
1378 conditions: [{{{axis}: {{bound_upstream: {{application_provider: openai}}}}}}]
1379- filter: load_balancer
1380 clusters: [{{name: backend, http: {{application_provider: openai}}, endpoints: ["127.0.0.1:9"]}}]
1381"#
1382 ))
1383 .unwrap();
1384 FilterPipeline::build(&mut entries, ®istry).unwrap()
1385 }
1386
1387 #[cfg(feature = "upstream-binding")]
1390 fn bind_through_the_pipeline(pipeline: &FilterPipeline, ctx: &mut PingoraRequestCtx) {
1391 let runtime = tokio::runtime::Builder::new_current_thread().build().unwrap();
1392 let extensions = {
1393 let mut filter_ctx = ctx.filter_context_for(pipeline, None).expect("the snapshot is present");
1394 drop(
1395 runtime
1396 .block_on(pipeline.execute_http_request(&mut filter_ctx))
1397 .expect("the router binds the catch-all route"),
1398 );
1399 std::mem::take(&mut filter_ctx.extensions)
1400 };
1401 ctx.extensions = extensions;
1402 }
1403
1404 fn make_fallback_ctx() -> PingoraRequestCtx {
1405 let mut ctx = PingoraRequestCtx::default();
1406 ctx.request_snapshot = Some(praxis_filter::Request {
1407 method: http::Method::GET,
1408 uri: "/incomplete".parse().unwrap(),
1409 headers: http::HeaderMap::new(),
1410 });
1411 ctx
1412 }
1413
1414 fn capture_access_events<F: FnOnce()>(f: F) -> Vec<String> {
1416 use tracing_subscriber::layer::SubscriberExt as _;
1417
1418 let messages = Arc::new(std::sync::Mutex::new(Vec::<String>::new()));
1419 let capture = AccessCapture(Arc::clone(&messages));
1420 let subscriber = tracing_subscriber::registry().with(capture);
1421 tracing::subscriber::with_default(subscriber, f);
1422 let mut guard = messages.lock().unwrap();
1423 std::mem::take(&mut *guard)
1424 }
1425
1426 struct AccessCapture(Arc<std::sync::Mutex<Vec<String>>>);
1428
1429 impl<S: tracing::Subscriber> tracing_subscriber::Layer<S> for AccessCapture {
1430 fn on_event(&self, event: &tracing::Event<'_>, _ctx: tracing_subscriber::layer::Context<'_, S>) {
1431 let mut visitor = AccessMessageVisitor(String::new());
1432 event.record(&mut visitor);
1433 if visitor.0.contains("access") {
1434 self.0.lock().unwrap().push(visitor.0);
1435 }
1436 }
1437 }
1438
1439 struct AccessMessageVisitor(String);
1441
1442 impl tracing::field::Visit for AccessMessageVisitor {
1443 fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
1444 if field.name() == "message" {
1445 self.0 = format!("{value:?}");
1446 }
1447 }
1448 }
1449
1450 fn make_passive_ctx(cluster: &str, endpoint_idx: usize, status: Option<u16>) -> PingoraRequestCtx {
1452 let mut ctx = PingoraRequestCtx::default();
1453 ctx.cluster = Some(Arc::from(cluster));
1454 ctx.selected_endpoint_index = Some(endpoint_idx);
1455 ctx.upstream_response_status = status;
1456 ctx.upstream_contacted = true;
1458 ctx
1459 }
1460
1461 fn make_passive_scenario(
1464 passive_unhealthy: Option<u32>,
1465 passive_healthy: Option<u32>,
1466 ) -> (FilterPipeline, PingoraRequestCtx) {
1467 let entry = ClusterHealthEntry::new(
1468 vec![EndpointHealth::new()],
1469 vec![Arc::from("10.0.0.1:80")],
1470 passive_unhealthy,
1471 passive_healthy,
1472 );
1473 let mut map = HashMap::new();
1474 map.insert(Arc::from("test-cluster"), Arc::new(entry));
1475 let health_registry = Arc::new(map);
1476
1477 let registry = praxis_filter::FilterRegistry::with_builtins();
1478 let mut pipeline = FilterPipeline::build(&mut [], ®istry).unwrap();
1479 pipeline.set_health_registry(health_registry);
1480
1481 let ctx = make_passive_ctx("test-cluster", 0, None);
1482
1483 (pipeline, ctx)
1484 }
1485}