Skip to main content

praxis_protocol/http/pingora/handler/
mod.rs

1// SPDX-License-Identifier: MIT
2// Copyright (c) 2024 Praxis Contributors
3
4//! Pingora `ProxyHttp` implementation: the main HTTP reverse-proxy handler.
5
6use std::{sync::Arc, time::Duration};
7
8use arc_swap::ArcSwap;
9use pingora_core::{Result, apps::HttpServerOptions, server::Server, services::listening::Service};
10use pingora_proxy::{Session, http_proxy};
11use praxis_filter::{BodyMode, CompressionConfig, FilterPipeline};
12use tokio::sync::Semaphore;
13use tracing::{debug, warn};
14
15use super::{context::PingoraRequestCtx, metrics};
16
17/// Shared hop-by-hop header stripping logic.
18mod hop_by_hop;
19/// HTTP handler without body filter hooks.
20mod no_body;
21/// Request header normalization (duplicate headers, obs-fold).
22mod normalize;
23/// Request body filter hook.
24mod request_body_filter;
25/// Request filter hook.
26mod request_filter;
27/// Reserved internal header helpers.
28mod reserved_headers;
29/// Response body filter hook.
30mod response_body_filter;
31/// Response filter hook.
32mod response_filter;
33/// Upstream peer selection hook.
34mod upstream_peer;
35/// Upstream request transformation hook.
36mod upstream_request;
37/// Upstream response hop-by-hop stripping hook.
38mod upstream_response;
39/// Via header injection hook.
40mod via;
41/// HTTP handler with body filter hooks.
42mod with_body;
43
44pub use no_body::PingoraHttpHandlerNoBody;
45pub use with_body::PingoraHttpHandler;
46
47// -----------------------------------------------------------------------------
48// Constants
49// -----------------------------------------------------------------------------
50
51/// Maximum number of upstream connection retries for idempotent requests.
52const MAX_RETRIES: usize = 3;
53
54/// Pingora still replays retry attempts from a fixed-size internal buffer.
55///
56/// `StreamBuffer` initial forwarding uses Praxis-owned `pre_read_body`, but
57/// retry replay cannot safely cover bodies larger than Pingora's retry
58/// buffer.
59const RETRY_BODY_LIMIT: u64 = 65_536; // 64 KiB
60
61// -----------------------------------------------------------------------------
62// Load Handler
63// -----------------------------------------------------------------------------
64
65/// Load an HTTP handler for a single listener.
66///
67/// Any TLS certificate watcher shutdown senders are appended to
68/// `cert_watcher_shutdowns`. The caller must keep this `Vec` alive
69/// until server shutdown; dropping the senders signals the watcher
70/// tasks to stop.
71///
72/// ```ignore
73/// use std::sync::Arc;
74///
75/// use pingora_core::server::Server;
76/// use praxis_core::config::Listener;
77/// use praxis_filter::{FilterPipeline, FilterRegistry};
78/// use praxis_protocol::http::pingora::handler::load_http_handler;
79///
80/// let mut server = Server::new(None).unwrap();
81/// server.bootstrap();
82/// let registry = FilterRegistry::with_builtins();
83/// let pipeline = Arc::new(FilterPipeline::build(&mut [], &registry).unwrap());
84/// let listener = Listener {
85///     name: "http".into(),
86///     address: "127.0.0.1:8080".into(),
87///     cluster: None,
88///     downstream_read_timeout_ms: None,
89///     filter_chains: vec![],
90///     max_connections: None,
91///     protocol: Default::default(),
92///     tcp_session_timeout_ms: None,
93///     tcp_max_duration_secs: None,
94///     tls: None,
95///     upstream: None,
96/// };
97/// let mut shutdowns = Vec::new();
98/// load_http_handler(&mut server, &listener, pipeline, &mut shutdowns).unwrap();
99/// ```
100///
101/// # Errors
102///
103/// Returns [`ProxyError`] if the listener fails to bind.
104///
105/// [`ProxyError`]: praxis_core::ProxyError
106pub fn load_http_handler(
107    server: &mut Server,
108    listener: &praxis_core::config::Listener,
109    pipeline: Arc<ArcSwap<FilterPipeline>>,
110    cert_watcher_shutdowns: &mut Vec<tokio::sync::watch::Sender<bool>>,
111) -> Result<(), praxis_core::ProxyError> {
112    let downstream_read_timeout = listener.downstream_read_timeout_ms.map(Duration::from_millis);
113    let connection_semaphore = listener
114        .max_connections
115        .map(|max| Arc::new(Semaphore::new(max as usize)));
116
117    // Always use the body-capable handler: a reload may add body
118    // filters, and compression init is one-shot in Pingora.
119    debug!(listener = %listener.name, "loading HTTP handler with body filters");
120    let handler = PingoraHttpHandler::new(pipeline, downstream_read_timeout, connection_semaphore);
121    wire_service(server, listener, handler, cert_watcher_shutdowns)?;
122    Ok(())
123}
124
125/// Create a Pingora HTTP proxy service, bind the listener, and add it to the server.
126fn wire_service<H>(
127    server: &mut Server,
128    listener: &praxis_core::config::Listener,
129    handler: H,
130    cert_watcher_shutdowns: &mut Vec<tokio::sync::watch::Sender<bool>>,
131) -> Result<(), praxis_core::ProxyError>
132where
133    H: pingora_proxy::ProxyHttp + Send + Sync + 'static,
134    H::CTX: Send + Sync,
135{
136    let service_name = format!("http-proxy:{name}", name = listener.name);
137    let mut proxy = http_proxy(&server.configuration, handler);
138    proxy.server_options = Some(h2c_server_options());
139    let mut service = Service::new(service_name, proxy);
140    if let Some(tx) = super::listener::add_listener(&mut service, listener)? {
141        cert_watcher_shutdowns.push(tx);
142    }
143    server.add_service(service);
144    Ok(())
145}
146
147// -----------------------------------------------------------------------------
148// Shared Utilities
149// -----------------------------------------------------------------------------
150
151/// Clamp a runtime-selected body mode to the byte ceiling implied by `baseline`.
152///
153/// `baseline` is the mode established before request/response-phase filter hooks
154/// run (typically from pipeline capabilities + global body limits). Runtime
155/// `set_*_body_mode` calls may widen limits; this helper preserves the original
156/// ceiling while still allowing upgrades between body mode variants.
157///
158/// `Stream` mode passes through unconditionally because it delivers chunks
159/// as they arrive without accumulating them — there is no buffer to cap.
160/// A filter that downgrades from `StreamBuffer` to `Stream` at runtime is
161/// opting out of buffering entirely, which is always safe from a memory
162/// perspective. The pipeline-level body size limit (enforced separately
163/// via `SizeLimit`) remains the backstop for oversized payloads.
164fn clamp_body_mode_to_ceiling(mode: BodyMode, baseline: BodyMode) -> BodyMode {
165    let ceiling = match baseline {
166        BodyMode::StreamBuffer { max_bytes: Some(v) } | BodyMode::SizeLimit { max_bytes: v } => Some(v),
167        _ => None,
168    };
169
170    match (mode, ceiling) {
171        (BodyMode::StreamBuffer { max_bytes }, Some(limit)) => BodyMode::StreamBuffer {
172            max_bytes: Some(max_bytes.map_or(limit, |v| v.min(limit))),
173        },
174        (BodyMode::SizeLimit { max_bytes }, Some(limit)) => BodyMode::SizeLimit {
175            max_bytes: max_bytes.min(limit),
176        },
177        // Stream has no buffer to clamp; other modes pass through when the
178        // baseline imposes no ceiling (e.g. unbounded StreamBuffer).
179        (m, None | Some(_)) => m,
180    }
181}
182
183/// Apply compression settings from the pipeline config to the Pingora response.
184fn adjust_compression(
185    session: &mut Session,
186    upstream_response: &pingora_http::ResponseHeader,
187    compression: Option<&CompressionConfig>,
188) {
189    use pingora_core::{modules::http::compression::ResponseCompression, protocols::http::compression::Algorithm};
190
191    let Some(cfg) = compression else {
192        return;
193    };
194
195    let Some(module) = session.downstream_modules_ctx.get_mut::<ResponseCompression>() else {
196        return;
197    };
198
199    let headers = &upstream_response.headers;
200
201    if !cfg.should_compress(headers) {
202        debug!("disabling compression: response does not qualify");
203        module.adjust_level(0);
204        return;
205    }
206
207    for (enabled, level, algo) in [
208        (cfg.gzip_enabled, cfg.gzip_level, Algorithm::Gzip),
209        (cfg.brotli_enabled, cfg.brotli_level, Algorithm::Brotli),
210        (cfg.zstd_enabled, cfg.zstd_level, Algorithm::Zstd),
211    ] {
212        if !enabled {
213            module.adjust_algorithm_level(algo, 0);
214        } else if let Some(lvl) = level {
215            module.adjust_algorithm_level(algo, lvl);
216        }
217    }
218}
219
220/// Handle upstream connect failures with retry logic.
221///
222/// Retries are skipped when the effective forwarded body size exceeds
223/// Pingora's retry buffer limit.
224///
225/// Body-mutating filters can change payload length after `request_filter`,
226/// so the retry guard uses the larger of the original and mutated lengths.
227fn handle_connect_failure(ctx: &mut PingoraRequestCtx, e: Box<pingora_core::Error>) -> Box<pingora_core::Error> {
228    if ctx.request_is_idempotent {
229        let mutated_len = ctx.mutated_request_body_len.unwrap_or(0) as u64;
230        let effective_body_size = std::cmp::max(ctx.request_body_bytes, mutated_len);
231        if effective_body_size > RETRY_BODY_LIMIT {
232            warn!(
233                body_bytes = ctx.request_body_bytes,
234                mutated_len = ?ctx.mutated_request_body_len,
235                limit = RETRY_BODY_LIMIT,
236                "skipping retry: request body exceeds Pingora retry buffer limit"
237            );
238            return e;
239        }
240        if (ctx.retries as usize) < MAX_RETRIES {
241            ctx.retries += 1;
242            debug!(
243                retries = ctx.retries,
244                max = MAX_RETRIES,
245                "retrying idempotent request after connect failure"
246            );
247            let mut e = e;
248            e.set_retry(true);
249            return e;
250        }
251        warn!(
252            retries = ctx.retries,
253            max = MAX_RETRIES,
254            "retry limit reached for idempotent request"
255        );
256    }
257    e
258}
259
260/// Run response filters during the logging phase if the
261/// response phase never executed (upstream error, filter
262/// rejection, etc.).
263async fn logging_cleanup(pipeline: &FilterPipeline, ctx: &mut PingoraRequestCtx) {
264    if !ctx.response_phase_done
265        && let Some(mut filter_ctx) = ctx.filter_context_for(pipeline, None)
266    {
267        let _result = pipeline.execute_http_response(&mut filter_ctx).await;
268        let extensions = filter_ctx.extensions;
269        let metadata = filter_ctx.filter_metadata;
270        let state = filter_ctx.filter_state;
271        let exec_idx = filter_ctx.executed_filter_indices;
272        let body_idx = filter_ctx.body_done_indices;
273        ctx.extensions = extensions;
274        ctx.filter_metadata = metadata;
275        ctx.filter_state = state;
276        ctx.cached_executed_filter_indices = exec_idx;
277        ctx.cached_body_done_indices = body_idx;
278    }
279}
280
281/// Emit Prometheus metrics for a completed HTTP request.
282///
283/// No-op when the Prometheus recorder has not been installed.
284fn emit_request_metrics(session: &Session, ctx: &PingoraRequestCtx) {
285    if !metrics::is_recorder_installed() {
286        return;
287    }
288
289    let status_code = session.response_written().map_or(0, |resp| resp.status.as_u16());
290    let status_class = metrics::status_class(status_code);
291
292    let request_method = session.req_header().method.as_str();
293    let raw_method = if request_method.is_empty() {
294        ctx.request_snapshot.as_ref().map_or("UNKNOWN", |r| r.method.as_str())
295    } else {
296        request_method
297    };
298    let method = metrics::method_label(raw_method);
299
300    let cluster = ctx
301        .metrics_cluster_shared
302        .clone()
303        .unwrap_or_else(|| ::metrics::SharedString::const_str("none"));
304
305    let labels = metrics::RequestMetricLabels {
306        cluster,
307        method,
308        route: "unknown",
309        status_class,
310    };
311
312    let duration_secs = ctx.request_start.elapsed().as_secs_f64();
313    metrics::record_request_metrics(labels, duration_secs);
314}
315
316/// Record a passive health observation for the selected upstream endpoint.
317///
318/// Called from the `logging` hook on every completed request. Determines
319/// success/failure from the error argument and the stashed upstream
320/// response status code.
321///
322/// No-op when no upstream was selected, no health registry is available,
323/// or passive checking is not configured for the cluster.
324fn record_passive_health(pipeline: &FilterPipeline, error: Option<&pingora_core::Error>, ctx: &PingoraRequestCtx) {
325    let cluster_name = ctx.cluster.as_ref().or(ctx.metrics_cluster.as_ref());
326    let Some(cluster_name) = cluster_name else {
327        return;
328    };
329    let Some(idx) = ctx.selected_endpoint_index else {
330        return;
331    };
332    let Some(registry) = pipeline.health_registry() else {
333        return;
334    };
335    let Some(health) = registry.get(cluster_name) else {
336        return;
337    };
338
339    let is_failure = error.is_some() || ctx.upstream_response_status.is_some_and(|s| s >= 500);
340    apply_passive_threshold(health, idx, cluster_name, is_failure);
341}
342
343/// Apply passive health threshold for a single endpoint observation.
344fn apply_passive_threshold(
345    health: &praxis_core::health::ClusterHealthEntry,
346    idx: usize,
347    cluster_name: &Arc<str>,
348    is_failure: bool,
349) {
350    if is_failure {
351        if let Some(threshold) = health.passive_unhealthy_threshold()
352            && health
353                .endpoints()
354                .get(idx)
355                .is_some_and(|ep| ep.record_failure(threshold))
356        {
357            tracing::warn!(
358                cluster = %cluster_name,
359                endpoint_index = idx,
360                threshold,
361                "passive health: endpoint marked unhealthy"
362            );
363        }
364    } else if let Some(threshold) = health.passive_healthy_threshold()
365        && health
366            .endpoints()
367            .get(idx)
368            .is_some_and(|ep| ep.record_success(threshold))
369    {
370        tracing::info!(
371            cluster = %cluster_name,
372            endpoint_index = idx,
373            threshold,
374            "passive health: endpoint recovered"
375        );
376    }
377}
378
379/// Build [`HttpServerOptions`] with h2c enabled.
380///
381/// [`HttpServerOptions`]: pingora_core::apps::HttpServerOptions
382fn h2c_server_options() -> HttpServerOptions {
383    let mut opts = HttpServerOptions::default();
384    opts.h2c = true;
385    opts
386}
387
388// -----------------------------------------------------------------------------
389// Tests
390// -----------------------------------------------------------------------------
391
392#[cfg(test)]
393#[expect(clippy::allow_attributes, reason = "blanket test suppressions")]
394#[allow(
395    clippy::unwrap_used,
396    clippy::expect_used,
397    clippy::indexing_slicing,
398    clippy::field_reassign_with_default,
399    clippy::too_many_lines,
400    clippy::cast_possible_truncation,
401    clippy::significant_drop_tightening,
402    reason = "tests"
403)]
404mod tests {
405    use std::sync::Arc;
406
407    use super::*;
408
409    #[test]
410    fn first_failure_idempotent_sets_retry() {
411        let mut ctx = PingoraRequestCtx::default();
412        ctx.request_is_idempotent = true;
413        let e = handle_connect_failure(&mut ctx, make_error());
414        assert!(e.retry(), "first failure should set retry flag");
415        assert_eq!(ctx.retries, 1);
416    }
417
418    #[test]
419    fn large_body_skips_retry() {
420        let mut ctx = PingoraRequestCtx::default();
421        ctx.request_is_idempotent = true;
422        ctx.request_body_bytes = RETRY_BODY_LIMIT + 1;
423        let e = handle_connect_failure(&mut ctx, make_error());
424        assert!(!e.retry(), "should not retry when body exceeds retry buffer limit");
425        assert_eq!(ctx.retries, 0, "retry counter should not increment");
426    }
427
428    #[test]
429    fn mutated_body_exceeding_limit_skips_retry() {
430        let mut ctx = PingoraRequestCtx::default();
431        ctx.request_is_idempotent = true;
432        ctx.request_body_bytes = 1024;
433        ctx.mutated_request_body_len = Some((RETRY_BODY_LIMIT + 1) as usize);
434        let e = handle_connect_failure(&mut ctx, make_error());
435        assert!(
436            !e.retry(),
437            "should not retry when mutated body exceeds retry buffer limit"
438        );
439        assert_eq!(ctx.retries, 0);
440    }
441
442    #[test]
443    fn body_at_limit_allows_retry() {
444        let mut ctx = PingoraRequestCtx::default();
445        ctx.request_is_idempotent = true;
446        ctx.request_body_bytes = RETRY_BODY_LIMIT;
447        let e = handle_connect_failure(&mut ctx, make_error());
448        assert!(e.retry(), "body exactly at limit should allow retry");
449        assert_eq!(ctx.retries, 1);
450    }
451
452    #[test]
453    fn zero_body_allows_retry() {
454        let mut ctx = PingoraRequestCtx::default();
455        ctx.request_is_idempotent = true;
456        ctx.request_body_bytes = 0;
457        let e = handle_connect_failure(&mut ctx, make_error());
458        assert!(e.retry(), "zero-length body should allow retry");
459        assert_eq!(ctx.retries, 1);
460    }
461
462    #[test]
463    fn max_retries_exhausted_does_not_retry() {
464        let mut ctx = PingoraRequestCtx::default();
465        ctx.request_is_idempotent = true;
466        ctx.retries = MAX_RETRIES as u32;
467        let e = handle_connect_failure(&mut ctx, make_error());
468        assert!(!e.retry(), "should not retry after MAX_RETRIES");
469        assert_eq!(ctx.retries as usize, MAX_RETRIES);
470    }
471
472    #[test]
473    fn counter_increments_across_calls() {
474        let mut ctx = PingoraRequestCtx::default();
475        ctx.request_is_idempotent = true;
476        for expected in 1..=MAX_RETRIES {
477            let _result = handle_connect_failure(&mut ctx, make_error());
478            assert_eq!(ctx.retries as usize, expected);
479        }
480        let e = handle_connect_failure(&mut ctx, make_error());
481        assert!(!e.retry(), "should not retry after reaching MAX_RETRIES");
482        assert_eq!(ctx.retries as usize, MAX_RETRIES);
483    }
484
485    #[test]
486    fn non_idempotent_request_never_retries() {
487        let mut ctx = PingoraRequestCtx::default();
488        ctx.request_is_idempotent = false;
489        let e = handle_connect_failure(&mut ctx, make_error());
490        assert!(!e.retry(), "non-idempotent request should never retry");
491        assert_eq!(ctx.retries, 0);
492    }
493
494    #[tokio::test]
495    async fn logging_cleanup_noop_when_response_phase_done() {
496        let registry = praxis_filter::FilterRegistry::with_builtins();
497        let pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
498        let mut ctx = PingoraRequestCtx::default();
499        ctx.response_phase_done = true;
500        ctx.request_snapshot = Some(praxis_filter::Request {
501            method: http::Method::GET,
502            uri: "/".parse().unwrap(),
503            headers: http::HeaderMap::new(),
504        });
505        logging_cleanup(&pipeline, &mut ctx).await;
506    }
507
508    #[tokio::test]
509    async fn logging_cleanup_noop_when_no_snapshot() {
510        let registry = praxis_filter::FilterRegistry::with_builtins();
511        let pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
512        let mut ctx = PingoraRequestCtx::default();
513        ctx.response_phase_done = false;
514        ctx.request_snapshot = None;
515        logging_cleanup(&pipeline, &mut ctx).await;
516    }
517
518    #[tokio::test]
519    async fn logging_cleanup_runs_response_pipeline_when_needed() {
520        let registry = praxis_filter::FilterRegistry::with_builtins();
521        let pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
522        let mut ctx = PingoraRequestCtx::default();
523        ctx.response_phase_done = false;
524        ctx.cluster = Some(Arc::from("test-cluster"));
525        ctx.request_snapshot = Some(praxis_filter::Request {
526            method: http::Method::GET,
527            uri: "/test".parse().unwrap(),
528            headers: http::HeaderMap::new(),
529        });
530        logging_cleanup(&pipeline, &mut ctx).await;
531        assert!(ctx.cluster.is_none(), "cluster should be taken by logging_cleanup");
532        assert!(ctx.upstream.is_none(), "upstream should be taken by logging_cleanup");
533    }
534
535    #[tokio::test]
536    async fn logging_cleanup_preserves_filter_metadata() {
537        let registry = praxis_filter::FilterRegistry::with_builtins();
538        let pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
539        let mut ctx = PingoraRequestCtx::default();
540        ctx.response_phase_done = false;
541        ctx.filter_metadata
542            .insert("json_rpc.method".to_owned(), "service/invoke".to_owned());
543        ctx.request_snapshot = Some(praxis_filter::Request {
544            method: http::Method::POST,
545            uri: "/api".parse().unwrap(),
546            headers: http::HeaderMap::new(),
547        });
548        logging_cleanup(&pipeline, &mut ctx).await;
549        assert_eq!(
550            ctx.filter_metadata.get("json_rpc.method").map(String::as_str),
551            Some("service/invoke"),
552            "filter_metadata should survive logging_cleanup"
553        );
554    }
555
556    #[tokio::test]
557    async fn logging_cleanup_preserves_extensions() {
558        let registry = praxis_filter::FilterRegistry::with_builtins();
559        let pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
560        let mut ctx = PingoraRequestCtx::default();
561        ctx.response_phase_done = false;
562        ctx.extensions.insert(42_u32);
563        ctx.request_snapshot = Some(praxis_filter::Request {
564            method: http::Method::POST,
565            uri: "/test".parse().unwrap(),
566            headers: http::HeaderMap::new(),
567        });
568        logging_cleanup(&pipeline, &mut ctx).await;
569        assert_eq!(
570            ctx.extensions.get::<u32>(),
571            Some(&42),
572            "extensions should survive logging_cleanup"
573        );
574    }
575
576    #[test]
577    fn passive_health_error_is_failure() {
578        let (pipeline, ctx) = make_passive_scenario(Some(3), Some(2));
579        let error = make_error();
580        record_passive_health(&pipeline, Some(&error), &ctx);
581
582        let registry = pipeline.health_registry().unwrap();
583        let entry = registry.get("test-cluster").unwrap();
584        assert!(
585            entry.endpoints()[0].is_healthy(),
586            "single failure should not yet mark unhealthy (threshold=3)"
587        );
588    }
589
590    #[test]
591    fn passive_health_status_500_is_failure() {
592        let (pipeline, mut ctx) = make_passive_scenario(Some(3), Some(2));
593        ctx.upstream_response_status = Some(500);
594        record_passive_health(&pipeline, None, &ctx);
595
596        let registry = pipeline.health_registry().unwrap();
597        let entry = registry.get("test-cluster").unwrap();
598        assert!(
599            entry.endpoints()[0].is_healthy(),
600            "single 500 should not yet mark unhealthy (threshold=3)"
601        );
602    }
603
604    #[test]
605    fn passive_health_status_below_500_is_success() {
606        let (pipeline, mut ctx) = make_passive_scenario(Some(2), Some(1));
607        ctx.upstream_response_status = Some(499);
608        record_passive_health(&pipeline, None, &ctx);
609
610        let registry = pipeline.health_registry().unwrap();
611        let entry = registry.get("test-cluster").unwrap();
612        assert!(entry.endpoints()[0].is_healthy(), "status 499 should count as success");
613    }
614
615    #[test]
616    fn passive_unhealthy_threshold_transition() {
617        let (pipeline, ctx) = make_passive_scenario(Some(2), Some(1));
618        let error = make_error();
619        record_passive_health(&pipeline, Some(&error), &ctx);
620        record_passive_health(&pipeline, Some(&error), &ctx);
621
622        let registry = pipeline.health_registry().unwrap();
623        let entry = registry.get("test-cluster").unwrap();
624        assert!(
625            !entry.endpoints()[0].is_healthy(),
626            "2 consecutive failures should mark unhealthy (threshold=2)"
627        );
628    }
629
630    #[test]
631    fn passive_healthy_threshold_recovery() {
632        let (pipeline, ctx) = make_passive_scenario(Some(1), Some(2));
633        let error = make_error();
634        record_passive_health(&pipeline, Some(&error), &ctx);
635
636        let registry = pipeline.health_registry().unwrap();
637        let entry = registry.get("test-cluster").unwrap();
638        assert!(
639            !entry.endpoints()[0].is_healthy(),
640            "should be unhealthy after 1 failure"
641        );
642
643        let ctx_ok = make_passive_ctx("test-cluster", 0, Some(200));
644        record_passive_health(&pipeline, None, &ctx_ok);
645        assert!(
646            !entry.endpoints()[0].is_healthy(),
647            "one success should not recover (threshold=2)"
648        );
649
650        record_passive_health(&pipeline, None, &ctx_ok);
651        assert!(
652            entry.endpoints()[0].is_healthy(),
653            "2 consecutive successes should recover (threshold=2)"
654        );
655    }
656
657    #[test]
658    fn passive_health_no_thresholds_is_noop() {
659        let (pipeline, ctx) = make_passive_scenario(None, None);
660        let error = make_error();
661        record_passive_health(&pipeline, Some(&error), &ctx);
662
663        let registry = pipeline.health_registry().unwrap();
664        let entry = registry.get("test-cluster").unwrap();
665        assert!(
666            entry.endpoints()[0].is_healthy(),
667            "no passive thresholds means failures are no-op"
668        );
669    }
670
671    #[test]
672    fn passive_health_endpoint_index_out_of_bounds() {
673        let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
674        ctx.selected_endpoint_index = Some(999);
675        let error = make_error();
676        record_passive_health(&pipeline, Some(&error), &ctx);
677
678        let registry = pipeline.health_registry().unwrap();
679        let entry = registry.get("test-cluster").unwrap();
680        assert!(entry.endpoints()[0].is_healthy(), "out-of-bounds index should be no-op");
681    }
682
683    #[test]
684    fn passive_health_missing_cluster_is_noop() {
685        let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
686        ctx.cluster = None;
687        ctx.metrics_cluster = None;
688        let error = make_error();
689        record_passive_health(&pipeline, Some(&error), &ctx);
690    }
691
692    #[test]
693    fn passive_health_falls_back_to_metrics_cluster() {
694        let (pipeline, mut ctx) = make_passive_scenario(Some(2), Some(1));
695        ctx.cluster = None;
696        ctx.metrics_cluster = Some(Arc::from("test-cluster"));
697        let error = make_error();
698        record_passive_health(&pipeline, Some(&error), &ctx);
699        record_passive_health(&pipeline, Some(&error), &ctx);
700
701        let registry = pipeline.health_registry().unwrap();
702        let entry = registry.get("test-cluster").unwrap();
703        assert!(
704            !entry.endpoints()[0].is_healthy(),
705            "fallback to metrics_cluster should still record passive health"
706        );
707    }
708
709    #[test]
710    fn passive_health_missing_endpoint_index_is_noop() {
711        let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
712        ctx.selected_endpoint_index = None;
713        let error = make_error();
714        record_passive_health(&pipeline, Some(&error), &ctx);
715    }
716
717    #[test]
718    fn passive_health_missing_registry_is_noop() {
719        let registry = praxis_filter::FilterRegistry::with_builtins();
720        let pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
721        let mut ctx = PingoraRequestCtx::default();
722        ctx.cluster = Some(Arc::from("test-cluster"));
723        ctx.selected_endpoint_index = Some(0);
724        let error = make_error();
725        record_passive_health(&pipeline, Some(&error), &ctx);
726    }
727
728    #[test]
729    fn passive_health_unknown_cluster_is_noop() {
730        let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
731        ctx.cluster = Some(Arc::from("nonexistent"));
732        let error = make_error();
733        record_passive_health(&pipeline, Some(&error), &ctx);
734    }
735
736    // -------------------------------------------------------------------------
737    // Test Utilities
738    // -------------------------------------------------------------------------
739
740    /// Create a connect error for tests.
741    fn make_error() -> Box<pingora_core::Error> {
742        pingora_core::Error::explain(pingora_core::ErrorType::ConnectError, "test connect failure")
743    }
744
745    /// Build a [`PingoraRequestCtx`] for passive health testing.
746    fn make_passive_ctx(cluster: &str, endpoint_idx: usize, status: Option<u16>) -> PingoraRequestCtx {
747        let mut ctx = PingoraRequestCtx::default();
748        ctx.cluster = Some(Arc::from(cluster));
749        ctx.selected_endpoint_index = Some(endpoint_idx);
750        ctx.upstream_response_status = status;
751        ctx
752    }
753
754    /// Build a pipeline with a health registry and a matching context
755    /// for passive health testing.
756    fn make_passive_scenario(
757        passive_unhealthy: Option<u32>,
758        passive_healthy: Option<u32>,
759    ) -> (FilterPipeline, PingoraRequestCtx) {
760        use std::collections::HashMap;
761
762        use praxis_core::health::{ClusterHealthEntry, EndpointHealth};
763
764        let entry = ClusterHealthEntry::new(
765            vec![EndpointHealth::new()],
766            vec![Arc::from("10.0.0.1:80")],
767            passive_unhealthy,
768            passive_healthy,
769        );
770        let mut map = HashMap::new();
771        map.insert(Arc::from("test-cluster"), Arc::new(entry));
772        let health_registry = Arc::new(map);
773
774        let registry = praxis_filter::FilterRegistry::with_builtins();
775        let mut pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
776        pipeline.set_health_registry(health_registry);
777
778        let ctx = make_passive_ctx("test-cluster", 0, None);
779
780        (pipeline, ctx)
781    }
782}