Skip to main content

praxis_protocol/http/pingora/handler/
mod.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2024 Praxis Contributors
3
4//! Pingora `ProxyHttp` implementation: the main HTTP reverse-proxy
5//! handler.
6//!
7//! Bridges Pingora's hook-based lifecycle (`request_filter`,
8//! `upstream_peer`, `upstream_request_filter`, etc.) to the Praxis
9//! filter pipeline via `PingoraHttpHandler` (body-capable). Body
10//! hooks are always available so hot reload can add body filters and
11//! Pingora compression init remains one-shot.
12//!
13//! Each submodule implements one Pingora hook. The pipeline is held
14//! behind `Arc<ArcSwap<FilterPipeline>>` for lock-free hot reload.
15
16use 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
27/// Safe per-request compression configuration.
28mod compression;
29/// Upstream connection established hook.
30mod connected_to_upstream;
31/// Structured error responses for fatal proxy errors.
32mod fail_to_proxy;
33/// Shared hop-by-hop header stripping logic.
34mod hop_by_hop;
35/// Request header normalization (duplicate headers, obs-fold).
36mod normalize;
37/// Request body filter hook.
38mod request_body_filter;
39/// Request filter hook.
40mod request_filter;
41/// Reserved internal header utilities.
42mod reserved_headers;
43/// Response body filter hook.
44mod response_body_filter;
45/// Response filter hook.
46mod response_filter;
47/// Response trailer hook: filter-driven trailer rewriting.
48mod response_trailer_filter;
49/// Response trailer hook: gRPC completion capture.
50mod response_trailers;
51/// Policy-aware retry decision engine.
52mod retry;
53/// Upstream peer selection hook.
54mod upstream_peer;
55/// Upstream request transformation hook.
56mod upstream_request;
57/// Upstream response hop-by-hop stripping hook.
58mod upstream_response;
59/// Via header injection hook.
60mod via;
61/// HTTP handler with body filter hooks.
62mod with_body;
63
64/// Body mode clamping and stream buffer utilities.
65mod body_util;
66/// Passive health checking utilities.
67mod health_util;
68/// Fallback access log emission and response filter cleanup.
69mod logging_util;
70/// Request metrics emission utilities.
71mod metrics_util;
72/// Retry decision logic and state management.
73mod retry_util;
74/// Pingora HTTP/2 server option builders.
75mod server_options;
76/// Span attribute recording for request tracing.
77mod span_util;
78
79pub use upstream_peer::{UpstreamRetryGateRelease, arm_upstream_retry_gate, lock_upstream_retry_gate_tests};
80pub use with_body::PingoraHttpHandler;
81
82// -----------------------------------------------------------------------------
83// Load Handler
84// -----------------------------------------------------------------------------
85
86/// Load an HTTP handler for a single listener.
87///
88/// Any TLS certificate watcher shutdown senders are appended to
89/// `cert_watcher_shutdowns`. The watcher tasks run for the process
90/// lifetime; the caller keeps this `Vec` to stop them early via
91/// `send(true)` (dropping the senders does not stop them).
92///
93/// ```ignore
94/// use std::sync::Arc;
95///
96/// use pingora_core::server::Server;
97/// use praxis_core::config::Listener;
98/// use praxis_filter::{FilterPipeline, FilterRegistry};
99/// use praxis_protocol::http::pingora::handler::load_http_handler;
100///
101/// let mut server = Server::new(None).unwrap();
102/// server.bootstrap();
103/// let registry = FilterRegistry::with_builtins();
104/// let pipeline = Arc::new(FilterPipeline::build(&mut [], &registry).unwrap());
105/// let listener = Listener {
106///     name: "http".into(),
107///     address: "127.0.0.1:8080".into(),
108///     cluster: None,
109///     downstream_read_timeout_ms: None,
110///     filter_chains: vec![],
111///     max_connections: None,
112///     protocol: Default::default(),
113///     tcp_session_timeout_ms: None,
114///     tcp_max_duration_secs: None,
115///     tls: None,
116///     upstream: None,
117/// };
118/// let mut shutdowns = Vec::new();
119/// load_http_handler(&mut server, &listener, pipeline, &mut shutdowns).unwrap();
120/// ```
121///
122/// # Errors
123///
124/// Returns [`ProxyError`] if the listener fails to bind.
125///
126/// [`ProxyError`]: praxis_core::ProxyError
127pub 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    // Always use the body-capable handler: a reload may add body
139    // filters, and compression init is one-shot in Pingora.
140    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        // `from_shared` keeps the label as a refcounted `Arc<str>`: the
146        // handler clones it per connection, and an owned `String` label
147        // would deep-copy on every clone.
148        ::metrics::SharedString::from_shared(Arc::from(listener.name.as_str())),
149    );
150    wire_service(server, listener, handler, cert_watcher_shutdowns)?;
151    Ok(())
152}
153
154/// Create a Pingora HTTP proxy service, bind the listener, and add it to the server.
155fn 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// -----------------------------------------------------------------------------
178// Tests
179// -----------------------------------------------------------------------------
180
181#[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    /// Maximum number of upstream connection retries for the legacy default policy.
207    const MAX_RETRIES: usize = praxis_core::config::DEFAULT_MAX_RETRIES as usize;
208
209    /// Default Pingora retry body buffer limit (64 `KiB`).
210    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 [], &registry).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 [], &registry).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 [], &registry).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 [], &registry).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 [], &registry).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 [], &registry).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    // -------------------------------------------------------------------------
955    // Fallback Access Log
956    // -------------------------------------------------------------------------
957
958    #[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 [], &registry).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, &registry).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 [], &registry).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    // -------------------------------------------------------------------------
1092    // Span Attribute Utilities
1093    // -------------------------------------------------------------------------
1094
1095    #[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    /// Layer that captures every `Span::record` call as `(field, value)` pairs.
1147    #[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    // -------------------------------------------------------------------------
1276    // Span Event Tests
1277    // -------------------------------------------------------------------------
1278
1279    #[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    // -------------------------------------------------------------------------
1342    // Test Utilities
1343    // -------------------------------------------------------------------------
1344
1345    /// Create a connect error for tests.
1346    fn make_error() -> Box<pingora_core::Error> {
1347        pingora_core::Error::explain(pingora_core::ErrorType::ConnectError, "test connect failure")
1348    }
1349
1350    /// Build a pipeline containing an `access_log` filter.
1351    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, &registry).unwrap()
1363    }
1364
1365    /// Build a context with a request snapshot for fallback logging tests.
1366    /// A router, an `access_log` gated on the openai binding with `axis`, and a
1367    /// load balancer whose cluster carries the openai tag.
1368    #[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, &registry).unwrap()
1383    }
1384
1385    /// Run the request phase so the router publishes the binding into `ctx`,
1386    /// the way the handler leaves it for a request that later fails.
1387    #[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    /// Capture `access` info events emitted while running `f`.
1413    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    /// Layer capturing `access` records for assertions.
1425    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    /// Visitor extracting the `message` field from an event.
1438    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    /// Build a [`PingoraRequestCtx`] for passive health testing.
1449    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        // A passive-health observation implies the upstream was contacted.
1455        ctx.upstream_contacted = true;
1456        ctx
1457    }
1458
1459    /// Build a pipeline with a health registry and a matching context
1460    /// for passive health testing.
1461    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 [], &registry).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}