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
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::{collections::HashMap, sync::Arc, time::Duration};
17
18use arc_swap::ArcSwap;
19use bytes::Bytes;
20use pingora_core::{
21    Result, apps::HttpServerOptions, protocols::http::v2::server::H2Options, server::Server,
22    services::listening::Service,
23};
24use pingora_proxy::{Session, http_proxy};
25use praxis_core::{config::ABSOLUTE_MAX_BODY_BYTES, connectivity::Upstream};
26use praxis_filter::{BodyBuffer, BodyMode, CompressionConfig, FilterPipeline, HttpFilterContext, RequestExtensions};
27use tokio::sync::Semaphore;
28use tracing::{debug, warn};
29
30use super::{context::PingoraRequestCtx, metrics};
31
32/// Structured error responses for fatal proxy errors.
33mod fail_to_proxy;
34/// Shared hop-by-hop header stripping logic.
35mod hop_by_hop;
36/// Request header normalization (duplicate headers, obs-fold).
37mod normalize;
38/// Request body filter hook.
39mod request_body_filter;
40/// Request filter hook.
41mod request_filter;
42/// Reserved internal header helpers.
43mod reserved_headers;
44/// Response body filter hook.
45mod response_body_filter;
46/// Response filter hook.
47mod response_filter;
48/// Upstream peer selection hook.
49mod upstream_peer;
50/// Upstream request transformation hook.
51mod upstream_request;
52/// Upstream response hop-by-hop stripping hook.
53mod upstream_response;
54/// Via header injection hook.
55mod via;
56/// HTTP handler with body filter hooks.
57mod with_body;
58
59pub use upstream_peer::{UpstreamRetryGateRelease, arm_upstream_retry_gate, lock_upstream_retry_gate_tests};
60pub use with_body::PingoraHttpHandler;
61
62// -----------------------------------------------------------------------------
63// Constants
64// -----------------------------------------------------------------------------
65
66/// Maximum number of upstream connection retries for idempotent requests.
67const MAX_RETRIES: usize = 3;
68
69/// Pingora still replays retry attempts from a fixed-size internal buffer.
70///
71/// `StreamBuffer` initial forwarding uses Praxis-owned `pre_read_body`, but
72/// retry replay cannot safely cover bodies larger than Pingora's retry
73/// buffer.
74const RETRY_BODY_LIMIT: u64 = 65_536; // 64 KiB
75
76// -----------------------------------------------------------------------------
77// Load Handler
78// -----------------------------------------------------------------------------
79
80/// Load an HTTP handler for a single listener.
81///
82/// Any TLS certificate watcher shutdown senders are appended to
83/// `cert_watcher_shutdowns`. The caller must keep this `Vec` alive
84/// until server shutdown; dropping the senders signals the watcher
85/// tasks to stop.
86///
87/// ```ignore
88/// use std::sync::Arc;
89///
90/// use pingora_core::server::Server;
91/// use praxis_core::config::Listener;
92/// use praxis_filter::{FilterPipeline, FilterRegistry};
93/// use praxis_protocol::http::pingora::handler::load_http_handler;
94///
95/// let mut server = Server::new(None).unwrap();
96/// server.bootstrap();
97/// let registry = FilterRegistry::with_builtins();
98/// let pipeline = Arc::new(FilterPipeline::build(&mut [], &registry).unwrap());
99/// let listener = Listener {
100///     name: "http".into(),
101///     address: "127.0.0.1:8080".into(),
102///     cluster: None,
103///     downstream_read_timeout_ms: None,
104///     filter_chains: vec![],
105///     max_connections: None,
106///     protocol: Default::default(),
107///     tcp_session_timeout_ms: None,
108///     tcp_max_duration_secs: None,
109///     tls: None,
110///     upstream: None,
111/// };
112/// let mut shutdowns = Vec::new();
113/// load_http_handler(&mut server, &listener, pipeline, &mut shutdowns).unwrap();
114/// ```
115///
116/// # Errors
117///
118/// Returns [`ProxyError`] if the listener fails to bind.
119///
120/// [`ProxyError`]: praxis_core::ProxyError
121pub fn load_http_handler(
122    server: &mut Server,
123    listener: &praxis_core::config::Listener,
124    pipeline: Arc<ArcSwap<FilterPipeline>>,
125    cert_watcher_shutdowns: &mut Vec<tokio::sync::watch::Sender<bool>>,
126) -> Result<(), praxis_core::ProxyError> {
127    let downstream_read_timeout = listener.downstream_read_timeout_ms.map(Duration::from_millis);
128    let connection_semaphore = listener
129        .max_connections
130        .map(|max| Arc::new(Semaphore::new(max as usize)));
131
132    // Always use the body-capable handler: a reload may add body
133    // filters, and compression init is one-shot in Pingora.
134    debug!(listener = %listener.name, "loading HTTP handler with body filters");
135    let handler = PingoraHttpHandler::new(
136        pipeline,
137        downstream_read_timeout,
138        connection_semaphore,
139        ::metrics::SharedString::from(listener.name.clone()),
140    );
141    wire_service(server, listener, handler, cert_watcher_shutdowns)?;
142    Ok(())
143}
144
145/// Create a Pingora HTTP proxy service, bind the listener, and add it to the server.
146fn wire_service<H>(
147    server: &mut Server,
148    listener: &praxis_core::config::Listener,
149    handler: H,
150    cert_watcher_shutdowns: &mut Vec<tokio::sync::watch::Sender<bool>>,
151) -> Result<(), praxis_core::ProxyError>
152where
153    H: pingora_proxy::ProxyHttp + Send + Sync + 'static,
154    H::CTX: Send + Sync,
155{
156    let service_name = format!("http-proxy:{name}", name = listener.name);
157    let mut proxy = http_proxy(&server.configuration, handler);
158    proxy.server_options = Some(h2c_server_options());
159    proxy.h2_options = Some(h2_server_options());
160    let mut service = Service::new(service_name, proxy);
161    if let Some(tx) = super::listener::add_listener(&mut service, listener)? {
162        cert_watcher_shutdowns.push(tx);
163    }
164    server.add_service(service);
165    Ok(())
166}
167
168// -----------------------------------------------------------------------------
169// Shared Utilities
170// -----------------------------------------------------------------------------
171
172/// Clamp a runtime-selected body mode to the byte ceiling implied by `baseline`.
173///
174/// `baseline` is the mode established before request/response-phase filter hooks
175/// run (typically from pipeline capabilities + global body limits). Runtime
176/// `set_*_body_mode` calls may widen limits; this helper preserves the original
177/// ceiling while still allowing upgrades between body mode variants.
178///
179/// `Stream` mode passes through unconditionally because it delivers chunks
180/// as they arrive without accumulating them — there is no buffer to cap.
181/// A filter that downgrades from `StreamBuffer` to `Stream` at runtime is
182/// opting out of buffering entirely, which is always safe from a memory
183/// perspective. The pipeline-level body size limit (enforced separately
184/// via `SizeLimit`) remains the backstop for oversized payloads.
185fn clamp_body_mode_to_ceiling(mode: BodyMode, baseline: BodyMode) -> BodyMode {
186    let ceiling = match baseline {
187        BodyMode::StreamBuffer { max_bytes: Some(v) } | BodyMode::SizeLimit { max_bytes: v } => Some(v),
188        _ => None,
189    };
190
191    match (mode, ceiling) {
192        (BodyMode::StreamBuffer { max_bytes }, Some(limit)) => BodyMode::StreamBuffer {
193            max_bytes: Some(max_bytes.map_or(limit, |v| v.min(limit))),
194        },
195        (BodyMode::SizeLimit { max_bytes }, Some(limit)) => BodyMode::SizeLimit {
196            max_bytes: max_bytes.min(limit),
197        },
198        // Stream has no buffer to clamp; other modes pass through when the
199        // baseline imposes no ceiling (e.g. unbounded StreamBuffer).
200        (m, None | Some(_)) => m,
201    }
202}
203
204/// Apply compression settings from the pipeline config to the Pingora response.
205fn adjust_compression(
206    session: &mut Session,
207    upstream_response: &pingora_http::ResponseHeader,
208    compression: Option<&CompressionConfig>,
209) {
210    use pingora_core::{modules::http::compression::ResponseCompression, protocols::http::compression::Algorithm};
211
212    let Some(cfg) = compression else {
213        return;
214    };
215
216    let Some(module) = session.downstream_modules_ctx.get_mut::<ResponseCompression>() else {
217        return;
218    };
219
220    let headers = &upstream_response.headers;
221
222    if !cfg.should_compress(headers) {
223        debug!("disabling compression: response does not qualify");
224        module.adjust_level(0);
225        return;
226    }
227
228    for (enabled, level, algo) in [
229        (cfg.gzip_enabled, cfg.gzip_level, Algorithm::Gzip),
230        (cfg.brotli_enabled, cfg.brotli_level, Algorithm::Brotli),
231        (cfg.zstd_enabled, cfg.zstd_level, Algorithm::Zstd),
232    ] {
233        if !enabled {
234            module.adjust_algorithm_level(algo, 0);
235        } else if let Some(lvl) = level {
236            module.adjust_algorithm_level(algo, lvl);
237        }
238    }
239}
240
241/// Handle upstream connect failures with retry logic.
242///
243/// Retries are skipped when the effective forwarded body size exceeds
244/// Pingora's retry buffer limit.
245///
246/// Body-mutating filters can change payload length after `request_filter`,
247/// so the retry guard uses the larger of the original and mutated lengths.
248fn handle_connect_failure(ctx: &mut PingoraRequestCtx, e: Box<pingora_core::Error>) -> Box<pingora_core::Error> {
249    let cluster = ctx.metrics_cluster_shared.clone().unwrap_or_else(metrics::cluster_none);
250    if let Some(start) = ctx.upstream_connect_start.take() {
251        metrics::record_upstream_connect_duration(cluster.clone(), start.elapsed().as_secs_f64());
252    }
253    metrics::record_upstream_connect_failure(cluster.clone());
254    if ctx.request_is_idempotent {
255        maybe_retry_idempotent_connect(ctx, cluster, e)
256    } else {
257        e
258    }
259}
260
261/// Decide whether an idempotent request may retry after a connect failure.
262fn maybe_retry_idempotent_connect(
263    ctx: &mut PingoraRequestCtx,
264    cluster: ::metrics::SharedString,
265    e: Box<pingora_core::Error>,
266) -> Box<pingora_core::Error> {
267    let mutated_len = ctx.mutated_request_body_len.unwrap_or(0) as u64;
268    if std::cmp::max(ctx.request_body_bytes, mutated_len) > RETRY_BODY_LIMIT {
269        warn!(
270            body_bytes = ctx.request_body_bytes,
271            mutated_len = ?ctx.mutated_request_body_len,
272            limit = RETRY_BODY_LIMIT,
273            "skipping retry: request body exceeds Pingora retry buffer limit"
274        );
275        record_retry_exhausted_if_attempted(ctx, cluster);
276        return e;
277    }
278    if (ctx.retries as usize) < MAX_RETRIES {
279        ctx.retries += 1;
280        debug!(
281            retries = ctx.retries,
282            max = MAX_RETRIES,
283            "retrying idempotent request after connect failure"
284        );
285        let mut e = e;
286        e.set_retry(true);
287        return e;
288    }
289    warn!(
290        retries = ctx.retries,
291        max = MAX_RETRIES,
292        "retry limit reached for idempotent request"
293    );
294    metrics::record_upstream_retry(cluster, metrics::RETRY_RESULT_EXHAUSTED);
295    e
296}
297
298/// Record `result=exhausted` only when at least one retry was already attempted.
299fn record_retry_exhausted_if_attempted(ctx: &PingoraRequestCtx, cluster: ::metrics::SharedString) {
300    if ctx.retries > 0 {
301        metrics::record_upstream_retry(cluster, metrics::RETRY_RESULT_EXHAUSTED);
302    }
303}
304
305/// Run response filters during the logging phase if the
306/// response phase never executed (upstream error, filter
307/// rejection, etc.).
308async fn logging_cleanup(pipeline: &FilterPipeline, ctx: &mut PingoraRequestCtx) {
309    if !ctx.response_phase_done
310        && let Some(mut filter_ctx) = ctx.filter_context_for(pipeline, None)
311    {
312        let _result = pipeline.execute_http_response(&mut filter_ctx).await;
313        let extensions = filter_ctx.extensions;
314        let metadata = filter_ctx.filter_metadata;
315        let state = filter_ctx.filter_state;
316        let exec_idx = filter_ctx.executed_filter_indices;
317        let body_idx = filter_ctx.body_done_indices;
318        ctx.extensions = extensions;
319        ctx.filter_metadata = metadata;
320        ctx.filter_state = state;
321        ctx.cached_executed_filter_indices = exec_idx;
322        ctx.cached_body_done_indices = body_idx;
323    }
324}
325
326/// Emit Prometheus metrics for a completed HTTP request.
327///
328/// No-op when the Prometheus recorder has not been installed.
329fn emit_request_metrics(session: &Session, ctx: &PingoraRequestCtx) {
330    if !metrics::is_recorder_installed() {
331        return;
332    }
333
334    let status_code = session.response_written().map_or(0, |resp| resp.status.as_u16());
335    let status_class = metrics::status_class(status_code);
336
337    let request_method = session.req_header().method.as_str();
338    let raw_method = if request_method.is_empty() {
339        ctx.request_snapshot.as_ref().map_or("UNKNOWN", |r| r.method.as_str())
340    } else {
341        request_method
342    };
343    let method = metrics::method_label(raw_method);
344
345    let cluster = ctx.metrics_cluster_shared.clone().unwrap_or_else(metrics::cluster_none);
346
347    let route = ctx.metrics_route.clone().unwrap_or_else(metrics::route_unknown);
348
349    let labels = metrics::RequestMetricLabels {
350        cluster: cluster.clone(),
351        method,
352        route,
353        status_class,
354    };
355
356    let duration_secs = ctx.request_start.elapsed().as_secs_f64();
357    metrics::record_request_metrics(labels, duration_secs);
358    metrics::record_body_size_metrics(
359        method,
360        status_class,
361        cluster,
362        ctx.request_body_bytes,
363        ctx.response_body_bytes,
364    );
365}
366
367/// Record a passive health observation for the selected upstream endpoint.
368///
369/// Called from the `logging` hook on every completed request. Determines
370/// success/failure from the error argument and the stashed upstream
371/// response status code.
372///
373/// No-op when no upstream was selected, no health registry is available,
374/// or passive checking is not configured for the cluster.
375fn record_passive_health(pipeline: &FilterPipeline, error: Option<&pingora_core::Error>, ctx: &PingoraRequestCtx) {
376    let cluster_name = ctx.cluster.as_ref().or(ctx.metrics_cluster.as_ref());
377    let Some(cluster_name) = cluster_name else {
378        return;
379    };
380    let Some(idx) = ctx.selected_endpoint_index else {
381        return;
382    };
383    let Some(registry) = pipeline.health_registry() else {
384        return;
385    };
386    let Some(health) = registry.get(cluster_name) else {
387        return;
388    };
389
390    let is_failure = error.is_some() || ctx.upstream_response_status.is_some_and(|s| s >= 500);
391    apply_passive_threshold(health, idx, cluster_name, is_failure);
392}
393
394/// Apply passive health threshold for a single endpoint observation.
395fn apply_passive_threshold(
396    health: &praxis_core::health::ClusterHealthEntry,
397    idx: usize,
398    cluster_name: &Arc<str>,
399    is_failure: bool,
400) {
401    if is_failure {
402        if let Some(threshold) = health.passive_unhealthy_threshold()
403            && health
404                .endpoints()
405                .get(idx)
406                .is_some_and(|ep| ep.record_failure(threshold))
407        {
408            tracing::warn!(
409                cluster = %cluster_name,
410                endpoint_index = idx,
411                threshold,
412                "passive health: endpoint marked unhealthy"
413            );
414            emit_passive_health_transition(health, cluster_name, metrics::HEALTH_RESULT_UNHEALTHY);
415        }
416    } else if let Some(threshold) = health.passive_healthy_threshold()
417        && health
418            .endpoints()
419            .get(idx)
420            .is_some_and(|ep| ep.record_success(threshold))
421    {
422        tracing::info!(
423            cluster = %cluster_name,
424            endpoint_index = idx,
425            threshold,
426            "passive health: endpoint recovered"
427        );
428        emit_passive_health_transition(health, cluster_name, metrics::HEALTH_RESULT_HEALTHY);
429    }
430}
431
432/// Refresh health gauges and increment the transition counter after a passive flip.
433fn emit_passive_health_transition(
434    health: &praxis_core::health::ClusterHealthEntry,
435    cluster_name: &Arc<str>,
436    result: &'static str,
437) {
438    let (healthy, total) = metrics::count_healthy_endpoints(health);
439    metrics::record_health_transition(
440        ::metrics::SharedString::from(Arc::clone(cluster_name)),
441        result,
442        healthy,
443        total,
444    );
445}
446
447/// Map an [`http::Version`] to the [OTel `network.protocol.version`] value.
448///
449/// [OTel `network.protocol.version`]: https://opentelemetry.io/docs/specs/semconv/attributes-registry/network/
450pub(super) fn http_version_label(version: http::Version) -> &'static str {
451    match version {
452        http::Version::HTTP_09 => "0.9",
453        http::Version::HTTP_10 => "1.0",
454        http::Version::HTTP_11 => "1.1",
455        http::Version::HTTP_2 => "2",
456        http::Version::HTTP_3 => "3",
457        _ => "unknown",
458    }
459}
460
461/// Record response-phase span attributes that are only available after
462/// the upstream exchange.
463///
464/// Called from the `logging` hook to fill in `http.response.status_code`,
465/// `upstream.address`, and `upstream.cluster` on the root request span.
466fn record_response_span_attributes(session: &Session, ctx: &PingoraRequestCtx) {
467    if ctx.request_span.is_disabled() {
468        return;
469    }
470
471    if let Some(resp) = session.response_written() {
472        let status = resp.status.as_u16();
473        if status > 0 {
474            ctx.request_span.record("http.response.status_code", status);
475        }
476        if resp.status.is_server_error() {
477            ctx.request_span.record("otel.status_code", "ERROR");
478        }
479    }
480
481    if let Some(upstream) = &ctx.upstream_for_retry {
482        ctx.request_span.record("upstream.address", upstream.address.as_ref());
483    }
484
485    if let Some(cluster) = &ctx.metrics_cluster {
486        ctx.request_span.record("upstream.cluster", cluster.as_ref());
487    }
488}
489
490/// Build [`HttpServerOptions`] with h2c enabled.
491///
492/// [`HttpServerOptions`]: pingora_core::apps::HttpServerOptions
493fn h2c_server_options() -> HttpServerOptions {
494    let mut opts = HttpServerOptions::default();
495    opts.h2c = true;
496    opts
497}
498
499/// Build [`H2Options`] with limits to mitigate HPACK amplification attacks
500/// (CWE-409).
501///
502/// Without explicit limits the `h2` crate defaults allow unbounded header
503/// list sizes and concurrent streams, enabling a small compressed request
504/// to allocate hundreds of megabytes on the server.
505///
506/// [`H2Options`]: pingora_core::protocols::http::v2::server::H2Options
507fn h2_server_options() -> H2Options {
508    let mut opts = H2Options::new();
509    opts.max_header_list_size(65_536); // 64 KiB
510    opts.max_concurrent_streams(128);
511    opts
512}
513
514/// Accumulate `chunk.len()` into `accumulated_bytes` and return `true` when
515/// the total exceeds `max_bytes`. Returns `false` when the body is `None`.
516fn check_body_size_limit(body: Option<&Bytes>, accumulated_bytes: &mut u64, max_bytes: usize) -> bool {
517    if let Some(chunk) = body {
518        #[expect(clippy::allow_attributes, reason = "cast lint is platform-dependent")]
519        #[allow(clippy::cast_possible_truncation, reason = "chunk length fits u64")]
520        let chunk_len = chunk.len() as u64;
521        *accumulated_bytes += chunk_len;
522
523        #[expect(clippy::allow_attributes, reason = "cast lint is platform-dependent")]
524        #[allow(clippy::cast_possible_truncation, reason = "max_bytes fits u64")]
525        let limit = max_bytes as u64;
526        return *accumulated_bytes > limit;
527    }
528    false
529}
530
531/// Push `chunk` into the stream buffer, creating it if absent. At end-of-stream
532/// the buffer is frozen into `body`. Returns `true` when the push overflows.
533fn accumulate_stream_buffer(
534    body: &mut Option<Bytes>,
535    body_buffer: &mut Option<BodyBuffer>,
536    end_of_stream: bool,
537    max_bytes: Option<usize>,
538) -> bool {
539    if let Some(chunk) = &*body {
540        let limit = max_bytes.unwrap_or(ABSOLUTE_MAX_BODY_BYTES);
541        let buf = body_buffer.get_or_insert_with(|| BodyBuffer::new(limit));
542
543        if buf.push(chunk.clone()).is_err() {
544            return true;
545        }
546    }
547
548    if end_of_stream {
549        tracing::trace!("stream buffer: freezing accumulated body before pipeline at EOS");
550        *body = body_buffer.take().map(BodyBuffer::freeze);
551    } else {
552        tracing::trace!("stream buffer: filters see the original chunk");
553    }
554    false
555}
556
557/// Suppress the body chunk while the stream buffer is still accumulating
558/// (i.e. `Continue`/`BodyDone` before release).
559#[expect(
560    clippy::fn_params_excessive_bools,
561    reason = "mirrors the caller's existing condition flags"
562)]
563fn suppress_stream_buffer_chunk(body: &mut Option<Bytes>, is_stream_buffer: bool, released: bool, end_of_stream: bool) {
564    if is_stream_buffer && !released && !end_of_stream {
565        *body = None;
566    }
567}
568
569/// Release the accumulated stream buffer on `FilterAction::Release`.
570fn release_stream_buffer(
571    body: &mut Option<Bytes>,
572    is_stream_buffer: bool,
573    released: &mut bool,
574    body_buffer: &mut Option<BodyBuffer>,
575    end_of_stream: bool,
576) {
577    if is_stream_buffer && !*released {
578        *released = true;
579        if !end_of_stream {
580            *body = body_buffer.take().map(BodyBuffer::freeze);
581        }
582    }
583}
584
585/// Shared fields extracted from an `HttpFilterContext` after body filter
586/// execution. Written back to `PingoraRequestCtx` via [`write_back`].
587///
588/// [`write_back`]: BodyFilterOutput::write_back
589struct BodyFilterOutput {
590    /// Cluster selected by the filter pipeline.
591    cluster: Option<Arc<str>>,
592    /// Upstream endpoint selected by the load balancer.
593    upstream: Option<Upstream>,
594    /// Type-safe request-scoped extension container.
595    extensions: RequestExtensions,
596    /// Durable per-request metadata that persists across phases.
597    filter_metadata: HashMap<String, String>,
598    /// Typed per-filter state keyed by stable filter invocation ID.
599    filter_state: HashMap<usize, Box<dyn std::any::Any + Send + Sync>>,
600    /// Per-filter execution tracking indices.
601    executed_filter_indices: Vec<bool>,
602    /// Per-filter body-done tracking indices.
603    body_done_indices: Vec<bool>,
604}
605
606impl BodyFilterOutput {
607    /// Move the shared fields out of the filter context, replacing each
608    /// with its `Default` value (zero-allocation no-ops for the types involved).
609    fn take_from(fctx: &mut HttpFilterContext<'_>) -> Self {
610        Self {
611            cluster: fctx.cluster.take(),
612            upstream: fctx.upstream.take(),
613            extensions: std::mem::take(&mut fctx.extensions),
614            filter_metadata: std::mem::take(&mut fctx.filter_metadata),
615            filter_state: std::mem::take(&mut fctx.filter_state),
616            executed_filter_indices: std::mem::take(&mut fctx.executed_filter_indices),
617            body_done_indices: std::mem::take(&mut fctx.body_done_indices),
618        }
619    }
620
621    /// Write the shared fields back to the protocol context.
622    fn write_back(self, ctx: &mut PingoraRequestCtx) {
623        ctx.cluster = self.cluster;
624        ctx.upstream = self.upstream;
625        ctx.extensions = self.extensions;
626        ctx.filter_metadata = self.filter_metadata;
627        ctx.filter_state = self.filter_state;
628        ctx.cached_executed_filter_indices = self.executed_filter_indices;
629        ctx.cached_body_done_indices = self.body_done_indices;
630    }
631}
632
633// -----------------------------------------------------------------------------
634// Tests
635// -----------------------------------------------------------------------------
636
637#[cfg(test)]
638#[expect(clippy::allow_attributes, reason = "blanket test suppressions")]
639#[allow(
640    clippy::unwrap_used,
641    clippy::expect_used,
642    clippy::indexing_slicing,
643    clippy::field_reassign_with_default,
644    clippy::too_many_lines,
645    clippy::cast_possible_truncation,
646    clippy::significant_drop_tightening,
647    reason = "tests"
648)]
649mod tests {
650    use praxis_core::connectivity::ConnectionOptions;
651
652    use super::*;
653
654    #[test]
655    fn first_failure_idempotent_sets_retry() {
656        let mut ctx = PingoraRequestCtx::default();
657        ctx.request_is_idempotent = true;
658        let e = handle_connect_failure(&mut ctx, make_error());
659        assert!(e.retry(), "first failure should set retry flag");
660        assert_eq!(ctx.retries, 1);
661    }
662
663    #[test]
664    fn large_body_skips_retry() {
665        let mut ctx = PingoraRequestCtx::default();
666        ctx.request_is_idempotent = true;
667        ctx.request_body_bytes = RETRY_BODY_LIMIT + 1;
668        let e = handle_connect_failure(&mut ctx, make_error());
669        assert!(!e.retry(), "should not retry when body exceeds retry buffer limit");
670        assert_eq!(ctx.retries, 0, "retry counter should not increment");
671    }
672
673    #[test]
674    fn mutated_body_exceeding_limit_skips_retry() {
675        let mut ctx = PingoraRequestCtx::default();
676        ctx.request_is_idempotent = true;
677        ctx.request_body_bytes = 1024;
678        ctx.mutated_request_body_len = Some((RETRY_BODY_LIMIT + 1) as usize);
679        let e = handle_connect_failure(&mut ctx, make_error());
680        assert!(
681            !e.retry(),
682            "should not retry when mutated body exceeds retry buffer limit"
683        );
684        assert_eq!(ctx.retries, 0);
685    }
686
687    #[test]
688    fn body_at_limit_allows_retry() {
689        let mut ctx = PingoraRequestCtx::default();
690        ctx.request_is_idempotent = true;
691        ctx.request_body_bytes = RETRY_BODY_LIMIT;
692        let e = handle_connect_failure(&mut ctx, make_error());
693        assert!(e.retry(), "body exactly at limit should allow retry");
694        assert_eq!(ctx.retries, 1);
695    }
696
697    #[test]
698    fn zero_body_allows_retry() {
699        let mut ctx = PingoraRequestCtx::default();
700        ctx.request_is_idempotent = true;
701        ctx.request_body_bytes = 0;
702        let e = handle_connect_failure(&mut ctx, make_error());
703        assert!(e.retry(), "zero-length body should allow retry");
704        assert_eq!(ctx.retries, 1);
705    }
706
707    #[test]
708    fn max_retries_exhausted_does_not_retry() {
709        let mut ctx = PingoraRequestCtx::default();
710        ctx.request_is_idempotent = true;
711        ctx.retries = MAX_RETRIES as u32;
712        let e = handle_connect_failure(&mut ctx, make_error());
713        assert!(!e.retry(), "should not retry after MAX_RETRIES");
714        assert_eq!(ctx.retries as usize, MAX_RETRIES);
715    }
716
717    #[test]
718    fn counter_increments_across_calls() {
719        let mut ctx = PingoraRequestCtx::default();
720        ctx.request_is_idempotent = true;
721        for expected in 1..=MAX_RETRIES {
722            let _result = handle_connect_failure(&mut ctx, make_error());
723            assert_eq!(ctx.retries as usize, expected);
724        }
725        let e = handle_connect_failure(&mut ctx, make_error());
726        assert!(!e.retry(), "should not retry after reaching MAX_RETRIES");
727        assert_eq!(ctx.retries as usize, MAX_RETRIES);
728    }
729
730    #[test]
731    fn non_idempotent_request_never_retries() {
732        let mut ctx = PingoraRequestCtx::default();
733        ctx.request_is_idempotent = false;
734        let e = handle_connect_failure(&mut ctx, make_error());
735        assert!(!e.retry(), "non-idempotent request should never retry");
736        assert_eq!(ctx.retries, 0);
737    }
738
739    #[test]
740    fn connect_failure_clears_upstream_connect_start() {
741        let mut ctx = PingoraRequestCtx::default();
742        ctx.upstream_connect_start = Some(std::time::Instant::now());
743        let _e = handle_connect_failure(&mut ctx, make_error());
744        assert!(
745            ctx.upstream_connect_start.is_none(),
746            "failed connect should consume upstream_connect_start for duration recording"
747        );
748    }
749
750    #[tokio::test]
751    async fn logging_cleanup_noop_when_response_phase_done() {
752        let registry = praxis_filter::FilterRegistry::with_builtins();
753        let pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
754        let mut ctx = PingoraRequestCtx::default();
755        ctx.response_phase_done = true;
756        ctx.request_snapshot = Some(praxis_filter::Request {
757            method: http::Method::GET,
758            uri: "/".parse().unwrap(),
759            headers: http::HeaderMap::new(),
760        });
761        logging_cleanup(&pipeline, &mut ctx).await;
762    }
763
764    #[tokio::test]
765    async fn logging_cleanup_noop_when_no_snapshot() {
766        let registry = praxis_filter::FilterRegistry::with_builtins();
767        let pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
768        let mut ctx = PingoraRequestCtx::default();
769        ctx.response_phase_done = false;
770        ctx.request_snapshot = None;
771        logging_cleanup(&pipeline, &mut ctx).await;
772    }
773
774    #[tokio::test]
775    async fn logging_cleanup_runs_response_pipeline_when_needed() {
776        let registry = praxis_filter::FilterRegistry::with_builtins();
777        let pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
778        let mut ctx = PingoraRequestCtx::default();
779        ctx.response_phase_done = false;
780        ctx.cluster = Some(Arc::from("test-cluster"));
781        ctx.request_snapshot = Some(praxis_filter::Request {
782            method: http::Method::GET,
783            uri: "/test".parse().unwrap(),
784            headers: http::HeaderMap::new(),
785        });
786        logging_cleanup(&pipeline, &mut ctx).await;
787        assert!(ctx.cluster.is_none(), "cluster should be taken by logging_cleanup");
788        assert!(ctx.upstream.is_none(), "upstream should be taken by logging_cleanup");
789    }
790
791    #[tokio::test]
792    async fn logging_cleanup_preserves_filter_metadata() {
793        let registry = praxis_filter::FilterRegistry::with_builtins();
794        let pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
795        let mut ctx = PingoraRequestCtx::default();
796        ctx.response_phase_done = false;
797        ctx.filter_metadata
798            .insert("json_rpc.method".to_owned(), "service/invoke".to_owned());
799        ctx.request_snapshot = Some(praxis_filter::Request {
800            method: http::Method::POST,
801            uri: "/api".parse().unwrap(),
802            headers: http::HeaderMap::new(),
803        });
804        logging_cleanup(&pipeline, &mut ctx).await;
805        assert_eq!(
806            ctx.filter_metadata.get("json_rpc.method").map(String::as_str),
807            Some("service/invoke"),
808            "filter_metadata should survive logging_cleanup"
809        );
810    }
811
812    #[tokio::test]
813    async fn logging_cleanup_preserves_extensions() {
814        let registry = praxis_filter::FilterRegistry::with_builtins();
815        let pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
816        let mut ctx = PingoraRequestCtx::default();
817        ctx.response_phase_done = false;
818        ctx.extensions.insert(42_u32);
819        ctx.request_snapshot = Some(praxis_filter::Request {
820            method: http::Method::POST,
821            uri: "/test".parse().unwrap(),
822            headers: http::HeaderMap::new(),
823        });
824        logging_cleanup(&pipeline, &mut ctx).await;
825        assert_eq!(
826            ctx.extensions.get::<u32>(),
827            Some(&42),
828            "extensions should survive logging_cleanup"
829        );
830    }
831
832    #[test]
833    fn passive_health_error_is_failure() {
834        let (pipeline, ctx) = make_passive_scenario(Some(3), Some(2));
835        let error = make_error();
836        record_passive_health(&pipeline, Some(&error), &ctx);
837
838        let registry = pipeline.health_registry().unwrap();
839        let entry = registry.get("test-cluster").unwrap();
840        assert!(
841            entry.endpoints()[0].is_healthy(),
842            "single failure should not yet mark unhealthy (threshold=3)"
843        );
844    }
845
846    #[test]
847    fn passive_health_status_500_is_failure() {
848        let (pipeline, mut ctx) = make_passive_scenario(Some(3), Some(2));
849        ctx.upstream_response_status = Some(500);
850        record_passive_health(&pipeline, None, &ctx);
851
852        let registry = pipeline.health_registry().unwrap();
853        let entry = registry.get("test-cluster").unwrap();
854        assert!(
855            entry.endpoints()[0].is_healthy(),
856            "single 500 should not yet mark unhealthy (threshold=3)"
857        );
858    }
859
860    #[test]
861    fn passive_health_status_below_500_is_success() {
862        let (pipeline, mut ctx) = make_passive_scenario(Some(2), Some(1));
863        ctx.upstream_response_status = Some(499);
864        record_passive_health(&pipeline, None, &ctx);
865
866        let registry = pipeline.health_registry().unwrap();
867        let entry = registry.get("test-cluster").unwrap();
868        assert!(entry.endpoints()[0].is_healthy(), "status 499 should count as success");
869    }
870
871    #[test]
872    fn passive_unhealthy_threshold_transition() {
873        let (pipeline, ctx) = make_passive_scenario(Some(2), Some(1));
874        let error = make_error();
875        record_passive_health(&pipeline, Some(&error), &ctx);
876        record_passive_health(&pipeline, Some(&error), &ctx);
877
878        let registry = pipeline.health_registry().unwrap();
879        let entry = registry.get("test-cluster").unwrap();
880        assert!(
881            !entry.endpoints()[0].is_healthy(),
882            "2 consecutive failures should mark unhealthy (threshold=2)"
883        );
884    }
885
886    #[test]
887    fn passive_healthy_threshold_recovery() {
888        let (pipeline, ctx) = make_passive_scenario(Some(1), Some(2));
889        let error = make_error();
890        record_passive_health(&pipeline, Some(&error), &ctx);
891
892        let registry = pipeline.health_registry().unwrap();
893        let entry = registry.get("test-cluster").unwrap();
894        assert!(
895            !entry.endpoints()[0].is_healthy(),
896            "should be unhealthy after 1 failure"
897        );
898
899        let ctx_ok = make_passive_ctx("test-cluster", 0, Some(200));
900        record_passive_health(&pipeline, None, &ctx_ok);
901        assert!(
902            !entry.endpoints()[0].is_healthy(),
903            "one success should not recover (threshold=2)"
904        );
905
906        record_passive_health(&pipeline, None, &ctx_ok);
907        assert!(
908            entry.endpoints()[0].is_healthy(),
909            "2 consecutive successes should recover (threshold=2)"
910        );
911    }
912
913    #[test]
914    fn passive_health_no_thresholds_is_noop() {
915        let (pipeline, ctx) = make_passive_scenario(None, None);
916        let error = make_error();
917        record_passive_health(&pipeline, Some(&error), &ctx);
918
919        let registry = pipeline.health_registry().unwrap();
920        let entry = registry.get("test-cluster").unwrap();
921        assert!(
922            entry.endpoints()[0].is_healthy(),
923            "no passive thresholds means failures are no-op"
924        );
925    }
926
927    #[test]
928    fn passive_health_endpoint_index_out_of_bounds() {
929        let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
930        ctx.selected_endpoint_index = Some(999);
931        let error = make_error();
932        record_passive_health(&pipeline, Some(&error), &ctx);
933
934        let registry = pipeline.health_registry().unwrap();
935        let entry = registry.get("test-cluster").unwrap();
936        assert!(entry.endpoints()[0].is_healthy(), "out-of-bounds index should be no-op");
937    }
938
939    #[test]
940    fn passive_health_missing_cluster_is_noop() {
941        let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
942        ctx.cluster = None;
943        ctx.metrics_cluster = None;
944        let error = make_error();
945        record_passive_health(&pipeline, Some(&error), &ctx);
946    }
947
948    #[test]
949    fn passive_health_falls_back_to_metrics_cluster() {
950        let (pipeline, mut ctx) = make_passive_scenario(Some(2), Some(1));
951        ctx.cluster = None;
952        ctx.metrics_cluster = Some(Arc::from("test-cluster"));
953        let error = make_error();
954        record_passive_health(&pipeline, Some(&error), &ctx);
955        record_passive_health(&pipeline, Some(&error), &ctx);
956
957        let registry = pipeline.health_registry().unwrap();
958        let entry = registry.get("test-cluster").unwrap();
959        assert!(
960            !entry.endpoints()[0].is_healthy(),
961            "fallback to metrics_cluster should still record passive health"
962        );
963    }
964
965    #[test]
966    fn passive_health_missing_endpoint_index_is_noop() {
967        let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
968        ctx.selected_endpoint_index = None;
969        let error = make_error();
970        record_passive_health(&pipeline, Some(&error), &ctx);
971    }
972
973    #[test]
974    fn passive_health_missing_registry_is_noop() {
975        let registry = praxis_filter::FilterRegistry::with_builtins();
976        let pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
977        let mut ctx = PingoraRequestCtx::default();
978        ctx.cluster = Some(Arc::from("test-cluster"));
979        ctx.selected_endpoint_index = Some(0);
980        let error = make_error();
981        record_passive_health(&pipeline, Some(&error), &ctx);
982    }
983
984    #[test]
985    fn passive_health_unknown_cluster_is_noop() {
986        let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
987        ctx.cluster = Some(Arc::from("nonexistent"));
988        let error = make_error();
989        record_passive_health(&pipeline, Some(&error), &ctx);
990    }
991
992    #[test]
993    fn size_limit_none_body_returns_false() {
994        let mut bytes = 0_u64;
995        assert!(!check_body_size_limit(None, &mut bytes, 100));
996        assert_eq!(bytes, 0, "accumulated bytes unchanged for None body");
997    }
998
999    #[test]
1000    fn size_limit_within_limit() {
1001        let mut bytes = 0_u64;
1002        let body = Some(Bytes::from_static(b"hello"));
1003        assert!(!check_body_size_limit(body.as_ref(), &mut bytes, 10));
1004        assert_eq!(bytes, 5);
1005    }
1006
1007    #[test]
1008    fn size_limit_at_exact_limit() {
1009        let mut bytes = 0_u64;
1010        let body = Some(Bytes::from_static(b"exact"));
1011        assert!(!check_body_size_limit(body.as_ref(), &mut bytes, 5));
1012        assert_eq!(bytes, 5);
1013    }
1014
1015    #[test]
1016    fn size_limit_exceeds_limit() {
1017        let mut bytes = 0_u64;
1018        let body = Some(Bytes::from_static(b"toolong"));
1019        assert!(check_body_size_limit(body.as_ref(), &mut bytes, 3));
1020    }
1021
1022    #[test]
1023    fn size_limit_cumulative_overflow() {
1024        let mut bytes = 0_u64;
1025        let first = Some(Bytes::from_static(b"aaa"));
1026        assert!(!check_body_size_limit(first.as_ref(), &mut bytes, 5));
1027
1028        let second = Some(Bytes::from_static(b"bbb"));
1029        assert!(check_body_size_limit(second.as_ref(), &mut bytes, 5));
1030        assert_eq!(bytes, 6);
1031    }
1032
1033    #[test]
1034    fn stream_buffer_accumulates_chunks() {
1035        let mut body = Some(Bytes::from_static(b"hello "));
1036        let mut buf: Option<BodyBuffer> = None;
1037        assert!(!accumulate_stream_buffer(&mut body, &mut buf, false, Some(100)));
1038        assert!(buf.is_some());
1039
1040        body = Some(Bytes::from_static(b"world"));
1041        assert!(!accumulate_stream_buffer(&mut body, &mut buf, false, Some(100)));
1042
1043        let frozen = buf.take().unwrap().freeze();
1044        assert_eq!(frozen, Bytes::from_static(b"hello world"));
1045    }
1046
1047    #[test]
1048    fn stream_buffer_freezes_at_eos() {
1049        let mut body = Some(Bytes::from_static(b"data"));
1050        let mut buf: Option<BodyBuffer> = None;
1051        assert!(!accumulate_stream_buffer(&mut body, &mut buf, false, Some(100)));
1052
1053        body = Some(Bytes::from_static(b" end"));
1054        assert!(!accumulate_stream_buffer(&mut body, &mut buf, true, Some(100)));
1055        assert!(buf.is_none(), "buffer should be taken at EOS");
1056        assert_eq!(body.unwrap(), Bytes::from_static(b"data end"));
1057    }
1058
1059    #[test]
1060    fn stream_buffer_overflow() {
1061        let mut body = Some(Bytes::from_static(b"too long"));
1062        let mut buf: Option<BodyBuffer> = None;
1063        assert!(accumulate_stream_buffer(&mut body, &mut buf, false, Some(5)));
1064    }
1065
1066    #[test]
1067    fn stream_buffer_none_body() {
1068        let mut body: Option<Bytes> = None;
1069        let mut buf: Option<BodyBuffer> = None;
1070        assert!(!accumulate_stream_buffer(&mut body, &mut buf, false, Some(100)));
1071        assert!(buf.is_none());
1072    }
1073
1074    #[test]
1075    fn stream_buffer_uses_absolute_max_when_none() {
1076        let mut body = Some(Bytes::from_static(b"data"));
1077        let mut buf: Option<BodyBuffer> = None;
1078        assert!(!accumulate_stream_buffer(&mut body, &mut buf, false, None));
1079        assert!(buf.is_some(), "should create buffer with absolute max");
1080    }
1081
1082    #[test]
1083    fn suppress_clears_body_when_buffering() {
1084        let mut body = Some(Bytes::from_static(b"data"));
1085        suppress_stream_buffer_chunk(&mut body, true, false, false);
1086        assert!(body.is_none());
1087    }
1088
1089    #[test]
1090    fn suppress_noop_when_not_stream_buffer() {
1091        let mut body = Some(Bytes::from_static(b"data"));
1092        suppress_stream_buffer_chunk(&mut body, false, false, false);
1093        assert!(body.is_some());
1094    }
1095
1096    #[test]
1097    fn suppress_noop_when_released() {
1098        let mut body = Some(Bytes::from_static(b"data"));
1099        suppress_stream_buffer_chunk(&mut body, true, true, false);
1100        assert!(body.is_some());
1101    }
1102
1103    #[test]
1104    fn suppress_noop_at_eos() {
1105        let mut body = Some(Bytes::from_static(b"data"));
1106        suppress_stream_buffer_chunk(&mut body, true, false, true);
1107        assert!(body.is_some());
1108    }
1109
1110    #[test]
1111    fn release_sets_flag_and_flushes_buffer() {
1112        let mut body: Option<Bytes> = None;
1113        let mut released = false;
1114        let mut buf = Some(BodyBuffer::new(100));
1115        buf.as_mut().unwrap().push(Bytes::from_static(b"buffered")).unwrap();
1116
1117        release_stream_buffer(&mut body, true, &mut released, &mut buf, false);
1118        assert!(released);
1119        assert_eq!(body.unwrap(), Bytes::from_static(b"buffered"));
1120        assert!(buf.is_none());
1121    }
1122
1123    #[test]
1124    fn release_noop_when_already_released() {
1125        let mut body: Option<Bytes> = None;
1126        let mut released = true;
1127        let mut buf: Option<BodyBuffer> = None;
1128
1129        release_stream_buffer(&mut body, true, &mut released, &mut buf, false);
1130        assert!(body.is_none(), "body should be unchanged when already released");
1131    }
1132
1133    #[test]
1134    fn release_noop_when_not_stream_buffer() {
1135        let mut body: Option<Bytes> = None;
1136        let mut released = false;
1137        let mut buf: Option<BodyBuffer> = None;
1138
1139        release_stream_buffer(&mut body, false, &mut released, &mut buf, false);
1140        assert!(!released, "released flag should be unchanged for non-stream-buffer");
1141    }
1142
1143    #[test]
1144    fn release_at_eos_sets_flag_but_no_flush() {
1145        let mut body: Option<Bytes> = None;
1146        let mut released = false;
1147        let mut buf = Some(BodyBuffer::new(100));
1148        buf.as_mut().unwrap().push(Bytes::from_static(b"data")).unwrap();
1149
1150        release_stream_buffer(&mut body, true, &mut released, &mut buf, true);
1151        assert!(released);
1152        assert!(body.is_none(), "body should not be overwritten at EOS");
1153        assert!(buf.is_some(), "buffer should not be taken at EOS");
1154    }
1155
1156    #[test]
1157    fn write_back_transfers_fields() {
1158        let mut ctx = PingoraRequestCtx::default();
1159
1160        let mut extensions = RequestExtensions::new();
1161        extensions.insert(42_u32);
1162
1163        let state_val: Box<dyn std::any::Any + Send + Sync> = Box::new(99_i32);
1164        let filter_state = HashMap::from([(0_usize, state_val)]);
1165
1166        let output = BodyFilterOutput {
1167            cluster: Some(Arc::from("test-cluster")),
1168            upstream: Some(Upstream {
1169                address: Arc::from("10.0.0.1:80"),
1170                connection: Arc::new(ConnectionOptions::default()),
1171                tls: None,
1172            }),
1173            extensions,
1174            filter_metadata: HashMap::from([("key".to_owned(), "val".to_owned())]),
1175            filter_state,
1176            executed_filter_indices: vec![true, false],
1177            body_done_indices: vec![false, true],
1178        };
1179        output.write_back(&mut ctx);
1180
1181        assert_eq!(ctx.cluster.as_deref(), Some("test-cluster"));
1182        assert!(ctx.upstream.is_some(), "upstream should transfer");
1183        assert_eq!(ctx.upstream.as_ref().unwrap().address.as_ref(), "10.0.0.1:80");
1184        assert_eq!(ctx.extensions.get::<u32>(), Some(&42));
1185        assert_eq!(ctx.filter_metadata.get("key").map(String::as_str), Some("val"));
1186        assert_eq!(ctx.filter_state.len(), 1, "filter_state should transfer");
1187        assert_eq!(
1188            ctx.filter_state.get(&0).and_then(|v| v.downcast_ref::<i32>()),
1189            Some(&99)
1190        );
1191        assert_eq!(ctx.cached_executed_filter_indices, vec![true, false]);
1192        assert_eq!(ctx.cached_body_done_indices, vec![false, true]);
1193    }
1194
1195    // -------------------------------------------------------------------------
1196    // Span Attribute Helpers
1197    // -------------------------------------------------------------------------
1198
1199    #[test]
1200    fn http_version_label_http_09() {
1201        assert_eq!(
1202            http_version_label(http::Version::HTTP_09),
1203            "0.9",
1204            "HTTP/0.9 should map to '0.9'"
1205        );
1206    }
1207
1208    #[test]
1209    fn http_version_label_http_10() {
1210        assert_eq!(
1211            http_version_label(http::Version::HTTP_10),
1212            "1.0",
1213            "HTTP/1.0 should map to '1.0'"
1214        );
1215    }
1216
1217    #[test]
1218    fn http_version_label_http_11() {
1219        assert_eq!(
1220            http_version_label(http::Version::HTTP_11),
1221            "1.1",
1222            "HTTP/1.1 should map to '1.1'"
1223        );
1224    }
1225
1226    #[test]
1227    fn http_version_label_http_2() {
1228        assert_eq!(
1229            http_version_label(http::Version::HTTP_2),
1230            "2",
1231            "HTTP/2 should map to '2'"
1232        );
1233    }
1234
1235    #[test]
1236    fn http_version_label_http_3() {
1237        assert_eq!(
1238            http_version_label(http::Version::HTTP_3),
1239            "3",
1240            "HTTP/3 should map to '3'"
1241        );
1242    }
1243
1244    #[test]
1245    fn record_response_span_attributes_noop_for_disabled_span() {
1246        let ctx = PingoraRequestCtx::default();
1247        assert!(ctx.request_span.is_disabled(), "default span should be disabled");
1248    }
1249
1250    #[test]
1251    fn record_response_span_attributes_records_upstream_cluster() {
1252        let mut ctx = PingoraRequestCtx::default();
1253        ctx.metrics_cluster = Some(Arc::from("api-cluster"));
1254        ctx.upstream_for_retry = Some(Upstream {
1255            address: Arc::from("10.0.0.1:80"),
1256            connection: Arc::new(ConnectionOptions::default()),
1257            tls: None,
1258        });
1259        ctx.request_span = tracing::info_span!(
1260            "test_span",
1261            "http.response.status_code" = tracing::field::Empty,
1262            "upstream.address" = tracing::field::Empty,
1263            "upstream.cluster" = tracing::field::Empty,
1264        );
1265        ctx.request_span.record("upstream.cluster", "api-cluster");
1266        ctx.request_span.record("upstream.address", "10.0.0.1:80");
1267    }
1268
1269    // -------------------------------------------------------------------------
1270    // Test Utilities
1271    // -------------------------------------------------------------------------
1272
1273    /// Create a connect error for tests.
1274    fn make_error() -> Box<pingora_core::Error> {
1275        pingora_core::Error::explain(pingora_core::ErrorType::ConnectError, "test connect failure")
1276    }
1277
1278    /// Build a [`PingoraRequestCtx`] for passive health testing.
1279    fn make_passive_ctx(cluster: &str, endpoint_idx: usize, status: Option<u16>) -> PingoraRequestCtx {
1280        let mut ctx = PingoraRequestCtx::default();
1281        ctx.cluster = Some(Arc::from(cluster));
1282        ctx.selected_endpoint_index = Some(endpoint_idx);
1283        ctx.upstream_response_status = status;
1284        ctx
1285    }
1286
1287    /// Build a pipeline with a health registry and a matching context
1288    /// for passive health testing.
1289    fn make_passive_scenario(
1290        passive_unhealthy: Option<u32>,
1291        passive_healthy: Option<u32>,
1292    ) -> (FilterPipeline, PingoraRequestCtx) {
1293        use praxis_core::health::{ClusterHealthEntry, EndpointHealth};
1294
1295        let entry = ClusterHealthEntry::new(
1296            vec![EndpointHealth::new()],
1297            vec![Arc::from("10.0.0.1:80")],
1298            passive_unhealthy,
1299            passive_healthy,
1300        );
1301        let mut map = HashMap::new();
1302        map.insert(Arc::from("test-cluster"), Arc::new(entry));
1303        let health_registry = Arc::new(map);
1304
1305        let registry = praxis_filter::FilterRegistry::with_builtins();
1306        let mut pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
1307        pipeline.set_health_registry(health_registry);
1308
1309        let ctx = make_passive_ctx("test-cluster", 0, None);
1310
1311        (pipeline, ctx)
1312    }
1313}