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::{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/// Upstream connection established hook.
33mod connected_to_upstream;
34/// Structured error responses for fatal proxy errors.
35mod fail_to_proxy;
36/// Shared hop-by-hop header stripping logic.
37mod hop_by_hop;
38/// Request header normalization (duplicate headers, obs-fold).
39mod normalize;
40/// Request body filter hook.
41mod request_body_filter;
42/// Request filter hook.
43mod request_filter;
44/// Reserved internal header helpers.
45mod reserved_headers;
46/// Response body filter hook.
47mod response_body_filter;
48/// Response filter hook.
49mod response_filter;
50/// Policy-aware retry decision engine.
51mod retry;
52/// Upstream peer selection hook.
53mod upstream_peer;
54/// Upstream request transformation hook.
55mod upstream_request;
56/// Upstream response hop-by-hop stripping hook.
57mod upstream_response;
58/// Via header injection hook.
59mod via;
60/// HTTP handler with body filter hooks.
61mod with_body;
62
63pub use upstream_peer::{UpstreamRetryGateRelease, arm_upstream_retry_gate, lock_upstream_retry_gate_tests};
64pub use with_body::PingoraHttpHandler;
65
66// -----------------------------------------------------------------------------
67// Load Handler
68// -----------------------------------------------------------------------------
69
70/// Load an HTTP handler for a single listener.
71///
72/// Any TLS certificate watcher shutdown senders are appended to
73/// `cert_watcher_shutdowns`. The caller must keep this `Vec` alive
74/// until server shutdown; dropping the senders signals the watcher
75/// tasks to stop.
76///
77/// ```ignore
78/// use std::sync::Arc;
79///
80/// use pingora_core::server::Server;
81/// use praxis_core::config::Listener;
82/// use praxis_filter::{FilterPipeline, FilterRegistry};
83/// use praxis_protocol::http::pingora::handler::load_http_handler;
84///
85/// let mut server = Server::new(None).unwrap();
86/// server.bootstrap();
87/// let registry = FilterRegistry::with_builtins();
88/// let pipeline = Arc::new(FilterPipeline::build(&mut [], &registry).unwrap());
89/// let listener = Listener {
90///     name: "http".into(),
91///     address: "127.0.0.1:8080".into(),
92///     cluster: None,
93///     downstream_read_timeout_ms: None,
94///     filter_chains: vec![],
95///     max_connections: None,
96///     protocol: Default::default(),
97///     tcp_session_timeout_ms: None,
98///     tcp_max_duration_secs: None,
99///     tls: None,
100///     upstream: None,
101/// };
102/// let mut shutdowns = Vec::new();
103/// load_http_handler(&mut server, &listener, pipeline, &mut shutdowns).unwrap();
104/// ```
105///
106/// # Errors
107///
108/// Returns [`ProxyError`] if the listener fails to bind.
109///
110/// [`ProxyError`]: praxis_core::ProxyError
111pub fn load_http_handler(
112    server: &mut Server,
113    listener: &praxis_core::config::Listener,
114    pipeline: Arc<ArcSwap<FilterPipeline>>,
115    cert_watcher_shutdowns: &mut Vec<tokio::sync::watch::Sender<bool>>,
116) -> Result<(), praxis_core::ProxyError> {
117    let downstream_read_timeout = listener.downstream_read_timeout_ms.map(Duration::from_millis);
118    let connection_semaphore = listener
119        .max_connections
120        .map(|max| Arc::new(Semaphore::new(max as usize)));
121
122    // Always use the body-capable handler: a reload may add body
123    // filters, and compression init is one-shot in Pingora.
124    debug!(listener = %listener.name, "loading HTTP handler with body filters");
125    let handler = PingoraHttpHandler::new(
126        pipeline,
127        downstream_read_timeout,
128        connection_semaphore,
129        // `from_shared` keeps the label as a refcounted `Arc<str>`: the
130        // handler clones it per connection, and an owned `String` label
131        // would deep-copy on every clone.
132        ::metrics::SharedString::from_shared(Arc::from(listener.name.as_str())),
133    );
134    wire_service(server, listener, handler, cert_watcher_shutdowns)?;
135    Ok(())
136}
137
138/// Create a Pingora HTTP proxy service, bind the listener, and add it to the server.
139fn wire_service<H>(
140    server: &mut Server,
141    listener: &praxis_core::config::Listener,
142    handler: H,
143    cert_watcher_shutdowns: &mut Vec<tokio::sync::watch::Sender<bool>>,
144) -> Result<(), praxis_core::ProxyError>
145where
146    H: pingora_proxy::ProxyHttp + Send + Sync + 'static,
147    H::CTX: Send + Sync,
148{
149    let service_name = format!("http-proxy:{name}", name = listener.name);
150    let mut proxy = http_proxy(&server.configuration, handler);
151    proxy.server_options = Some(h2c_server_options());
152    proxy.h2_options = Some(h2_server_options());
153    let mut service = Service::new(service_name, proxy);
154    if let Some(tx) = super::listener::add_listener(&mut service, listener)? {
155        cert_watcher_shutdowns.push(tx);
156    }
157    server.add_service(service);
158    Ok(())
159}
160
161// -----------------------------------------------------------------------------
162// Shared Utilities
163// -----------------------------------------------------------------------------
164
165/// Clamp a runtime-selected body mode to the byte ceiling implied by `baseline`.
166///
167/// `baseline` is the mode established before request/response-phase filter hooks
168/// run (typically from pipeline capabilities + global body limits). Runtime
169/// `set_*_body_mode` calls may widen limits; this helper preserves the original
170/// ceiling while still allowing upgrades between body mode variants.
171///
172/// `Stream` mode passes through unconditionally because it delivers chunks
173/// as they arrive without accumulating them — there is no buffer to cap.
174/// A filter that downgrades from `StreamBuffer` to `Stream` at runtime is
175/// opting out of buffering entirely, which is always safe from a memory
176/// perspective. The pipeline-level body size limit (enforced separately
177/// via `SizeLimit`) remains the backstop for oversized payloads.
178fn clamp_body_mode_to_ceiling(mode: BodyMode, baseline: BodyMode) -> BodyMode {
179    let ceiling = match baseline {
180        BodyMode::StreamBuffer { max_bytes: Some(v) } | BodyMode::SizeLimit { max_bytes: v } => Some(v),
181        _ => None,
182    };
183
184    match (mode, ceiling) {
185        (BodyMode::StreamBuffer { max_bytes }, Some(limit)) => BodyMode::StreamBuffer {
186            max_bytes: Some(max_bytes.map_or(limit, |v| v.min(limit))),
187        },
188        (BodyMode::SizeLimit { max_bytes }, Some(limit)) => BodyMode::SizeLimit {
189            max_bytes: max_bytes.min(limit),
190        },
191        // Stream has no buffer to clamp; other modes pass through when the
192        // baseline imposes no ceiling (e.g. unbounded StreamBuffer).
193        (m, None | Some(_)) => m,
194    }
195}
196
197/// Apply compression settings from the pipeline config to the Pingora response.
198fn adjust_compression(
199    session: &mut Session,
200    upstream_response: &pingora_http::ResponseHeader,
201    compression: Option<&CompressionConfig>,
202) {
203    use pingora_core::{modules::http::compression::ResponseCompression, protocols::http::compression::Algorithm};
204
205    let Some(cfg) = compression else {
206        return;
207    };
208
209    let Some(module) = session.downstream_modules_ctx.get_mut::<ResponseCompression>() else {
210        return;
211    };
212
213    let headers = &upstream_response.headers;
214
215    if !cfg.should_compress(headers) {
216        debug!("disabling compression: response does not qualify");
217        module.adjust_level(0);
218        return;
219    }
220
221    for (enabled, level, algo) in [
222        (cfg.gzip_enabled, cfg.gzip_level, Algorithm::Gzip),
223        (cfg.brotli_enabled, cfg.brotli_level, Algorithm::Brotli),
224        (cfg.zstd_enabled, cfg.zstd_level, Algorithm::Zstd),
225    ] {
226        if !enabled {
227            module.adjust_algorithm_level(algo, 0);
228        } else if let Some(lvl) = level {
229            module.adjust_algorithm_level(algo, lvl);
230        }
231    }
232}
233
234/// Shared legacy-default retry policy for requests that carry none.
235///
236/// The retry hooks run on every upstream response and connect failure;
237/// building a fresh `Arc<RetryPolicy>` there would heap-allocate per
238/// event for a value that never changes.
239fn legacy_default_policy() -> Arc<praxis_core::config::RetryPolicy> {
240    static LEGACY_DEFAULT: std::sync::LazyLock<Arc<praxis_core::config::RetryPolicy>> =
241        std::sync::LazyLock::new(|| Arc::new(praxis_core::config::RetryPolicy::legacy_default()));
242    Arc::clone(&LEGACY_DEFAULT)
243}
244
245/// Handle upstream connect failures with the policy-aware retry engine.
246///
247/// Retries are skipped when the effective forwarded body size exceeds
248/// the configured replay limit, the method is non-idempotent without
249/// opt-in, the budget is exhausted, or the overall deadline has passed.
250#[expect(clippy::too_many_lines, reason = "sequential guard checks")]
251fn handle_connect_failure(ctx: &mut PingoraRequestCtx, e: Box<pingora_core::Error>) -> Box<pingora_core::Error> {
252    let cluster = ctx.metrics_cluster_shared.clone().unwrap_or_else(metrics::cluster_none);
253    if let Some(start) = ctx.upstream_connect_start.take() {
254        metrics::record_upstream_connect_duration(cluster.clone(), start.elapsed().as_secs_f64());
255    }
256    metrics::record_upstream_connect_failure(cluster.clone());
257
258    let policy = ctx.retry_policy.clone().unwrap_or_else(legacy_default_policy);
259    let outcome = retry::classify_error(&e);
260    let decision = retry::should_retry(ctx, &policy, outcome, ctx.cluster_retry_state.as_deref());
261
262    match decision {
263        retry::RetryDecision::Retry { backoff } => {
264            ctx.retries += 1;
265            ctx.pending_backoff = Some(backoff);
266            // Legacy (unconfigured) policies keep the historical
267            // retry-same-endpoint behavior; only operator-configured
268            // policies opt into endpoint reselection.
269            ctx.reselect_on_retry = policy.configured;
270            if let Some(upstream) = ctx.upstream_for_retry.as_ref() {
271                let addr = Arc::clone(&upstream.address);
272                if !ctx.attempted_endpoints.iter().any(|e| e.as_ref() == addr.as_ref()) {
273                    ctx.attempted_endpoints.push(addr);
274                }
275            }
276            // Under reselection, release the failed endpoint's in-flight
277            // counter and clear the saved upstream so upstream_peer picks an
278            // alternate host. Legacy same-endpoint retries keep the saved
279            // upstream (and its counter) for the next attempt.
280            if policy.configured {
281                if let Some(upstream) = ctx.upstream_for_retry.as_ref()
282                    && let Some(reselector) = ctx.endpoint_reselector.as_ref()
283                {
284                    reselector.release(&upstream.address);
285                }
286                ctx.upstream_for_retry = None;
287            }
288            let upstream_address = ctx
289                .upstream_for_retry
290                .as_ref()
291                .map_or("unknown", |u| u.address.as_ref());
292            debug!(
293                retries = ctx.retries,
294                max = policy.effective_max_retries(),
295                ?backoff,
296                upstream_address,
297                "retrying after connect failure"
298            );
299            let mut e = e;
300            e.set_retry(true);
301            e
302        },
303        retry::RetryDecision::DoNotRetry => {
304            if ctx.retries > 0 {
305                warn!(
306                    retries = ctx.retries,
307                    max = policy.effective_max_retries(),
308                    upstream_address = ctx
309                        .upstream_for_retry
310                        .as_ref()
311                        .map_or("unknown", |u| u.address.as_ref()),
312                    "retry limit exhausted"
313                );
314            }
315            record_retry_exhausted_if_attempted(ctx, cluster);
316            // Pingora may mark some errors retriable by default; clear the
317            // flag so the policy decision is authoritative.
318            let mut e = e;
319            e.set_retry(false);
320            e
321        },
322    }
323}
324
325/// Decide whether an HTTP response status should trigger a retry.
326///
327/// Returns `Some(error)` marked retriable when the status is retriable
328/// and all guards pass; `None` when the response should be forwarded.
329#[expect(clippy::too_many_lines, reason = "sequential guard checks")]
330fn maybe_retry_response(ctx: &mut PingoraRequestCtx, status: u16) -> Option<Box<pingora_core::Error>> {
331    let policy = ctx.retry_policy.clone().unwrap_or_else(legacy_default_policy);
332    let outcome = retry::RetryOutcome::StatusCode(status);
333    let decision = retry::should_retry(ctx, &policy, outcome, ctx.cluster_retry_state.as_deref());
334    match decision {
335        retry::RetryDecision::Retry { backoff } => {
336            ctx.retries += 1;
337            ctx.pending_backoff = Some(backoff);
338            // Legacy (unconfigured) policies keep the historical
339            // retry-same-endpoint behavior; only operator-configured
340            // policies opt into endpoint reselection.
341            ctx.reselect_on_retry = policy.configured;
342            if let Some(upstream) = ctx.upstream_for_retry.as_ref() {
343                let addr = Arc::clone(&upstream.address);
344                if !ctx.attempted_endpoints.iter().any(|e| e.as_ref() == addr.as_ref()) {
345                    ctx.attempted_endpoints.push(addr);
346                }
347            }
348            // Under reselection, release the failed endpoint's in-flight
349            // counter and clear the saved upstream so upstream_peer picks an
350            // alternate host. Legacy same-endpoint retries keep the saved
351            // upstream (and its counter) for the next attempt.
352            if policy.configured {
353                if let Some(upstream) = ctx.upstream_for_retry.as_ref()
354                    && let Some(reselector) = ctx.endpoint_reselector.as_ref()
355                {
356                    reselector.release(&upstream.address);
357                }
358                ctx.upstream_for_retry = None;
359            }
360            debug!(
361                status,
362                retries = ctx.retries,
363                max = policy.effective_max_retries(),
364                ?backoff,
365                "retrying after retriable response status"
366            );
367            let mut e =
368                pingora_core::Error::explain(pingora_core::ErrorType::HTTPStatus(status), "retriable upstream status");
369            e.set_retry(true);
370            Some(e)
371        },
372        retry::RetryDecision::DoNotRetry => None,
373    }
374}
375
376/// Release the active-request counter if it has not already been released.
377fn release_retry_state(ctx: &mut PingoraRequestCtx) {
378    if !ctx.cluster_retry_state_released
379        && let Some(state) = ctx.cluster_retry_state.take()
380    {
381        state.leave();
382        ctx.cluster_retry_state_released = true;
383    }
384}
385
386/// Record `result=exhausted` only when at least one retry was already attempted.
387fn record_retry_exhausted_if_attempted(ctx: &PingoraRequestCtx, cluster: ::metrics::SharedString) {
388    if ctx.retries > 0 {
389        metrics::record_upstream_retry(cluster, metrics::RETRY_RESULT_EXHAUSTED);
390    }
391}
392
393/// Emit a fallback access record for requests whose lifecycle ended
394/// before the access log filter's completion hooks could run.
395///
396/// Covers pre-upstream rejections, upstream connect and read failures,
397/// and streamed responses aborted mid-body: none of these reach the
398/// bodyless response phase or body end-of-stream where the filter
399/// emits. Only fires when the pipeline configures an `access_log`
400/// filter; these records bypass the filter's sampling because
401/// incomplete requests are always worth a record.
402fn maybe_emit_fallback_access_log(pipeline: &FilterPipeline, status: u16, ctx: &mut PingoraRequestCtx) {
403    if ctx.response_delivery_complete || ctx.connection_upgraded || !pipeline.contains_filter("access_log") {
404        return;
405    }
406    if let Some(filter_ctx) = ctx.filter_context_for(pipeline, None) {
407        // The access_log filter already logged this request (e.g. a bodyless
408        // response whose on_response emitted before a later filter rejected):
409        // no fallback record, or it would duplicate.
410        if praxis_filter::access_record_already_emitted(&filter_ctx) {
411            return;
412        }
413        // Honor the entry's request conditions: a scoped access_log (e.g.
414        // only /api paths) must not gain fallback records for requests the
415        // operator excluded. Sampling is still deliberately bypassed.
416        if !pipeline.filter_request_conditions_match("access_log", filter_ctx.request) {
417            return;
418        }
419        praxis_filter::emit_access_record(&filter_ctx, status);
420    }
421}
422
423/// Run response filters during the logging phase if the
424/// response phase never executed (upstream error, filter
425/// rejection, etc.).
426async fn logging_cleanup(pipeline: &FilterPipeline, ctx: &mut PingoraRequestCtx) {
427    if !ctx.response_phase_done
428        && let Some(mut filter_ctx) = ctx.filter_context_for(pipeline, None)
429    {
430        let _result = pipeline.execute_http_response(&mut filter_ctx).await;
431        let extensions = filter_ctx.extensions;
432        let metadata = filter_ctx.filter_metadata;
433        let state = filter_ctx.filter_state;
434        let exec_idx = filter_ctx.executed_filter_indices;
435        let body_idx = filter_ctx.body_done_indices;
436        // The context macro takes cluster/upstream out of ctx; restore them
437        // so the fallback access record that follows can attribute the
438        // failure to the routed cluster and selected endpoint.
439        let cluster = filter_ctx.cluster;
440        let upstream = filter_ctx.upstream;
441        ctx.extensions = extensions;
442        ctx.filter_metadata = metadata;
443        ctx.filter_state = state;
444        ctx.cached_executed_filter_indices = exec_idx;
445        ctx.cached_body_done_indices = body_idx;
446        ctx.cluster = cluster;
447        ctx.upstream = upstream;
448    }
449}
450
451/// Emit Prometheus metrics for a completed HTTP request.
452///
453/// No-op when the Prometheus recorder has not been installed.
454fn emit_request_metrics(session: &Session, ctx: &PingoraRequestCtx) {
455    if !metrics::is_recorder_installed() {
456        return;
457    }
458
459    let status_code = session.response_written().map_or(0, |resp| resp.status.as_u16());
460    let status_class = metrics::status_class(status_code);
461
462    let request_method = session.req_header().method.as_str();
463    let raw_method = if request_method.is_empty() {
464        ctx.request_snapshot.as_ref().map_or("UNKNOWN", |r| r.method.as_str())
465    } else {
466        request_method
467    };
468    let method = metrics::method_label(raw_method);
469
470    let cluster = ctx.metrics_cluster_shared.clone().unwrap_or_else(metrics::cluster_none);
471
472    let route = ctx.metrics_route.clone().unwrap_or_else(metrics::route_unknown);
473
474    let labels = metrics::RequestMetricLabels {
475        cluster: cluster.clone(),
476        method,
477        route,
478        status_class,
479    };
480
481    let duration_secs = ctx.request_start.elapsed().as_secs_f64();
482    metrics::record_request_metrics(labels, duration_secs);
483    metrics::record_body_size_metrics(
484        method,
485        status_class,
486        cluster,
487        ctx.request_body_bytes,
488        ctx.response_body_bytes,
489    );
490}
491
492/// Record a passive health observation for the selected upstream endpoint.
493///
494/// Called from the `logging` hook on every completed request. Determines
495/// success/failure from the error argument and the stashed upstream
496/// response status code.
497///
498/// No-op when no upstream was selected, no health registry is available,
499/// or passive checking is not configured for the cluster.
500fn record_passive_health(pipeline: &FilterPipeline, error: Option<&pingora_core::Error>, ctx: &PingoraRequestCtx) {
501    let cluster_name = ctx.cluster.as_ref().or(ctx.metrics_cluster.as_ref());
502    let Some(cluster_name) = cluster_name else {
503        return;
504    };
505    let Some(idx) = ctx.selected_endpoint_index else {
506        return;
507    };
508    let Some(registry) = pipeline.health_registry() else {
509        return;
510    };
511    let Some(health) = registry.get(cluster_name) else {
512        return;
513    };
514
515    let is_failure = error.is_some() || ctx.upstream_response_status.is_some_and(|s| s >= 500);
516    apply_passive_threshold(health, idx, cluster_name, is_failure);
517}
518
519/// Apply passive health threshold for a single endpoint observation.
520fn apply_passive_threshold(
521    health: &praxis_core::health::ClusterHealthEntry,
522    idx: usize,
523    cluster_name: &Arc<str>,
524    is_failure: bool,
525) {
526    if is_failure {
527        if let Some(threshold) = health.passive_unhealthy_threshold()
528            && health
529                .endpoints()
530                .get(idx)
531                .is_some_and(|ep| ep.record_failure(threshold))
532        {
533            tracing::warn!(
534                cluster = %cluster_name,
535                endpoint_index = idx,
536                threshold,
537                "passive health: endpoint marked unhealthy"
538            );
539            emit_passive_health_transition(health, cluster_name, metrics::HEALTH_RESULT_UNHEALTHY);
540        }
541    } else if let Some(threshold) = health.passive_healthy_threshold()
542        && health
543            .endpoints()
544            .get(idx)
545            .is_some_and(|ep| ep.record_success(threshold))
546    {
547        tracing::info!(
548            cluster = %cluster_name,
549            endpoint_index = idx,
550            threshold,
551            "passive health: endpoint recovered"
552        );
553        emit_passive_health_transition(health, cluster_name, metrics::HEALTH_RESULT_HEALTHY);
554    }
555}
556
557/// Refresh health gauges and increment the transition counter after a passive flip.
558fn emit_passive_health_transition(
559    health: &praxis_core::health::ClusterHealthEntry,
560    cluster_name: &Arc<str>,
561    result: &'static str,
562) {
563    let (healthy, total) = metrics::count_healthy_endpoints(health);
564    metrics::record_health_transition(
565        ::metrics::SharedString::from(Arc::clone(cluster_name)),
566        result,
567        healthy,
568        total,
569    );
570}
571
572/// Map an [`http::Version`] to the [OTel `network.protocol.version`] value.
573///
574/// [OTel `network.protocol.version`]: https://opentelemetry.io/docs/specs/semconv/attributes-registry/network/
575pub(super) fn http_version_label(version: http::Version) -> &'static str {
576    match version {
577        http::Version::HTTP_09 => "0.9",
578        http::Version::HTTP_10 => "1.0",
579        http::Version::HTTP_11 => "1.1",
580        http::Version::HTTP_2 => "2",
581        http::Version::HTTP_3 => "3",
582        _ => "unknown",
583    }
584}
585
586/// Record response-phase span attributes that are only available after
587/// the upstream exchange.
588///
589/// Called from the `logging` hook to fill in `http.response.status_code`,
590/// `otel.status_code` and `error.type` (5xx only), `http.route` and the
591/// `otel.name` upgrade to `{method} {route}` (when a route matched),
592/// `upstream.address`, and `upstream.cluster` on the root request span,
593/// and response attributes on the upstream exchange span.
594fn record_response_span_attributes(session: &Session, ctx: &PingoraRequestCtx) {
595    if ctx.request_span.is_disabled() {
596        return;
597    }
598    let response = session.response_written();
599    let status = response.map(|resp| resp.status);
600    let method = session.req_header().method.as_str();
601    record_response_span_fields(status, method, response, ctx);
602}
603
604/// Record the response-phase fields once the status and method have been extracted.
605///
606/// Split from [`record_response_span_attributes`] so the recording logic is
607/// unit-testable without constructing a live Pingora session.
608fn record_response_span_fields(
609    status: Option<http::StatusCode>,
610    method: &str,
611    response: Option<&pingora_http::ResponseHeader>,
612    ctx: &PingoraRequestCtx,
613) {
614    if let Some(status) = status {
615        let code = status.as_u16();
616        if code > 0 {
617            ctx.request_span.record("http.response.status_code", code);
618        }
619        if status.is_server_error() {
620            ctx.request_span.record("otel.status_code", "ERROR");
621            // OTel semconv: error.type for an HTTP status is the numeric code
622            // as a string, not StatusCode's "{code} {reason}" Display form.
623            ctx.request_span.record("error.type", code.to_string().as_str());
624        }
625    }
626
627    if let Some(route) = &ctx.metrics_route {
628        ctx.request_span.record("http.route", route.as_ref());
629        ctx.request_span
630            .record("otel.name", format!("{method} {route}").as_str());
631    }
632
633    if let Some(upstream) = &ctx.upstream_for_retry {
634        ctx.request_span.record("upstream.address", upstream.address.as_ref());
635    }
636
637    if let Some(cluster) = &ctx.metrics_cluster {
638        ctx.request_span.record("upstream.cluster", cluster.as_ref());
639    }
640
641    record_upstream_exchange_span(ctx, response);
642}
643
644/// Record the upstream-exchange child span's response fields.
645fn record_upstream_exchange_span(ctx: &PingoraRequestCtx, response: Option<&pingora_http::ResponseHeader>) {
646    if ctx.upstream_exchange_span.is_disabled() {
647        return;
648    }
649    // Prefer the upstream's own status (captured before any response-phase
650    // rewrite); fall back to the written response when it was not captured.
651    if let Some(status) = ctx
652        .upstream_response_status
653        .or_else(|| response.map(|resp| resp.status.as_u16()))
654    {
655        ctx.upstream_exchange_span.record("http.response.status_code", status);
656    }
657    ctx.upstream_exchange_span
658        .record("http.response.body.size", ctx.response_body_bytes);
659}
660
661/// Build [`HttpServerOptions`] with h2c enabled.
662///
663/// [`HttpServerOptions`]: pingora_core::apps::HttpServerOptions
664fn h2c_server_options() -> HttpServerOptions {
665    let mut opts = HttpServerOptions::default();
666    opts.h2c = true;
667    opts
668}
669
670/// Build [`H2Options`] with limits to mitigate HPACK amplification attacks
671/// (CWE-409).
672///
673/// Without explicit limits the `h2` crate defaults allow unbounded header
674/// list sizes and concurrent streams, enabling a small compressed request
675/// to allocate hundreds of megabytes on the server.
676///
677/// [`H2Options`]: pingora_core::protocols::http::v2::server::H2Options
678fn h2_server_options() -> H2Options {
679    let mut opts = H2Options::new();
680    opts.max_header_list_size(65_536); // 64 KiB
681    opts.max_concurrent_streams(128);
682    opts
683}
684
685/// Accumulate `chunk.len()` into `accumulated_bytes` and return `true` when
686/// the total exceeds `max_bytes`. Returns `false` when the body is `None`.
687fn check_body_size_limit(body: Option<&Bytes>, accumulated_bytes: &mut u64, max_bytes: usize) -> bool {
688    if let Some(chunk) = body {
689        let chunk_len = chunk.len() as u64;
690        *accumulated_bytes += chunk_len;
691
692        let limit = max_bytes as u64;
693        return *accumulated_bytes > limit;
694    }
695    false
696}
697
698/// Push `chunk` into the stream buffer, creating it if absent. At end-of-stream
699/// the buffer is frozen into `body`. Returns `true` when the push overflows.
700fn accumulate_stream_buffer(
701    body: &mut Option<Bytes>,
702    body_buffer: &mut Option<BodyBuffer>,
703    end_of_stream: bool,
704    max_bytes: Option<usize>,
705) -> bool {
706    if let Some(chunk) = &*body {
707        let limit = max_bytes.unwrap_or(ABSOLUTE_MAX_BODY_BYTES);
708        let buf = body_buffer.get_or_insert_with(|| BodyBuffer::new(limit));
709
710        if buf.push(chunk.clone()).is_err() {
711            return true;
712        }
713    }
714
715    if end_of_stream {
716        tracing::trace!("stream buffer: freezing accumulated body before pipeline at EOS");
717        *body = body_buffer.take().map(BodyBuffer::freeze);
718    } else {
719        tracing::trace!("stream buffer: filters see the original chunk");
720    }
721    false
722}
723
724/// Suppress the body chunk while the stream buffer is still accumulating
725/// (i.e. `Continue`/`BodyDone` before release).
726#[expect(
727    clippy::fn_params_excessive_bools,
728    reason = "mirrors the caller's existing condition flags"
729)]
730fn suppress_stream_buffer_chunk(body: &mut Option<Bytes>, is_stream_buffer: bool, released: bool, end_of_stream: bool) {
731    if is_stream_buffer && !released && !end_of_stream {
732        *body = None;
733    }
734}
735
736/// Release the accumulated stream buffer on `FilterAction::Release`.
737fn release_stream_buffer(
738    body: &mut Option<Bytes>,
739    is_stream_buffer: bool,
740    released: &mut bool,
741    body_buffer: &mut Option<BodyBuffer>,
742    end_of_stream: bool,
743) {
744    if is_stream_buffer && !*released {
745        *released = true;
746        if !end_of_stream {
747            *body = body_buffer.take().map(BodyBuffer::freeze);
748        }
749    }
750}
751
752/// Shared fields extracted from an `HttpFilterContext` after body filter
753/// execution. Written back to `PingoraRequestCtx` via [`write_back`].
754///
755/// [`write_back`]: BodyFilterOutput::write_back
756struct BodyFilterOutput {
757    /// Cluster selected by the filter pipeline.
758    cluster: Option<Arc<str>>,
759    /// Upstream endpoint selected by the load balancer.
760    upstream: Option<Upstream>,
761    /// Type-safe request-scoped extension container.
762    extensions: RequestExtensions,
763    /// Durable per-request metadata that persists across phases.
764    filter_metadata: HashMap<String, String>,
765    /// Typed per-filter state keyed by stable filter invocation ID.
766    filter_state: HashMap<usize, Box<dyn std::any::Any + Send + Sync>>,
767    /// Per-filter execution tracking indices.
768    executed_filter_indices: Vec<bool>,
769    /// Per-filter body-done tracking indices.
770    body_done_indices: Vec<bool>,
771    /// Endpoints already attempted for this request (retry exclusion set).
772    attempted_endpoints: Vec<Arc<str>>,
773}
774
775impl BodyFilterOutput {
776    /// Move the shared fields out of the filter context, replacing each
777    /// with its `Default` value (zero-allocation no-ops for the types involved).
778    fn take_from(fctx: &mut HttpFilterContext<'_>) -> Self {
779        Self {
780            cluster: fctx.cluster.take(),
781            upstream: fctx.upstream.take(),
782            extensions: std::mem::take(&mut fctx.extensions),
783            filter_metadata: std::mem::take(&mut fctx.filter_metadata),
784            filter_state: std::mem::take(&mut fctx.filter_state),
785            executed_filter_indices: std::mem::take(&mut fctx.executed_filter_indices),
786            body_done_indices: std::mem::take(&mut fctx.body_done_indices),
787            attempted_endpoints: std::mem::take(&mut fctx.attempted_endpoints),
788        }
789    }
790
791    /// Write the shared fields back to the protocol context.
792    fn write_back(self, ctx: &mut PingoraRequestCtx) {
793        ctx.cluster = self.cluster;
794        ctx.upstream = self.upstream;
795        ctx.extensions = self.extensions;
796        ctx.filter_metadata = self.filter_metadata;
797        ctx.filter_state = self.filter_state;
798        ctx.cached_executed_filter_indices = self.executed_filter_indices;
799        ctx.cached_body_done_indices = self.body_done_indices;
800        ctx.attempted_endpoints = self.attempted_endpoints;
801    }
802}
803
804// -----------------------------------------------------------------------------
805// Tests
806// -----------------------------------------------------------------------------
807
808#[cfg(test)]
809#[expect(clippy::allow_attributes, reason = "blanket test suppressions")]
810#[allow(
811    clippy::unwrap_used,
812    clippy::expect_used,
813    clippy::indexing_slicing,
814    clippy::field_reassign_with_default,
815    clippy::too_many_lines,
816    clippy::cast_possible_truncation,
817    clippy::significant_drop_tightening,
818    reason = "tests"
819)]
820mod tests {
821    use praxis_core::connectivity::ConnectionOptions;
822
823    use super::*;
824
825    /// Maximum number of upstream connection retries for the legacy default policy.
826    const MAX_RETRIES: usize = praxis_core::config::DEFAULT_MAX_RETRIES as usize;
827
828    /// Default Pingora retry body buffer limit (64 `KiB`).
829    const RETRY_BODY_LIMIT: u64 = praxis_core::config::DEFAULT_RETRY_BODY_LIMIT_BYTES;
830
831    #[test]
832    fn first_failure_idempotent_sets_retry() {
833        let mut ctx = PingoraRequestCtx::default();
834        ctx.request_is_idempotent = true;
835        let e = handle_connect_failure(&mut ctx, make_error());
836        assert!(e.retry(), "first failure should set retry flag");
837        assert_eq!(ctx.retries, 1);
838    }
839
840    #[test]
841    fn large_body_skips_retry() {
842        let mut ctx = PingoraRequestCtx::default();
843        ctx.request_is_idempotent = true;
844        ctx.request_body_bytes = RETRY_BODY_LIMIT + 1;
845        let e = handle_connect_failure(&mut ctx, make_error());
846        assert!(!e.retry(), "should not retry when body exceeds retry buffer limit");
847        assert_eq!(ctx.retries, 0, "retry counter should not increment");
848    }
849
850    #[test]
851    fn mutated_body_exceeding_limit_skips_retry() {
852        let mut ctx = PingoraRequestCtx::default();
853        ctx.request_is_idempotent = true;
854        ctx.request_body_bytes = 1024;
855        ctx.mutated_request_body_len = Some((RETRY_BODY_LIMIT + 1) as usize);
856        let e = handle_connect_failure(&mut ctx, make_error());
857        assert!(
858            !e.retry(),
859            "should not retry when mutated body exceeds retry buffer limit"
860        );
861        assert_eq!(ctx.retries, 0);
862    }
863
864    #[test]
865    fn body_at_limit_allows_retry() {
866        let mut ctx = PingoraRequestCtx::default();
867        ctx.request_is_idempotent = true;
868        ctx.request_body_bytes = RETRY_BODY_LIMIT;
869        let e = handle_connect_failure(&mut ctx, make_error());
870        assert!(e.retry(), "body exactly at limit should allow retry");
871        assert_eq!(ctx.retries, 1);
872    }
873
874    #[test]
875    fn zero_body_allows_retry() {
876        let mut ctx = PingoraRequestCtx::default();
877        ctx.request_is_idempotent = true;
878        ctx.request_body_bytes = 0;
879        let e = handle_connect_failure(&mut ctx, make_error());
880        assert!(e.retry(), "zero-length body should allow retry");
881        assert_eq!(ctx.retries, 1);
882    }
883
884    #[test]
885    fn max_retries_exhausted_does_not_retry() {
886        let mut ctx = PingoraRequestCtx::default();
887        ctx.request_is_idempotent = true;
888        ctx.retries = MAX_RETRIES as u32;
889        let e = handle_connect_failure(&mut ctx, make_error());
890        assert!(!e.retry(), "should not retry after MAX_RETRIES");
891        assert_eq!(ctx.retries as usize, MAX_RETRIES);
892    }
893
894    #[test]
895    fn counter_increments_across_calls() {
896        let mut ctx = PingoraRequestCtx::default();
897        ctx.request_is_idempotent = true;
898        for expected in 1..=MAX_RETRIES {
899            let _result = handle_connect_failure(&mut ctx, make_error());
900            assert_eq!(ctx.retries as usize, expected);
901        }
902        let e = handle_connect_failure(&mut ctx, make_error());
903        assert!(!e.retry(), "should not retry after reaching MAX_RETRIES");
904        assert_eq!(ctx.retries as usize, MAX_RETRIES);
905    }
906
907    #[test]
908    fn non_idempotent_request_never_retries() {
909        let mut ctx = PingoraRequestCtx::default();
910        ctx.request_is_idempotent = false;
911        let e = handle_connect_failure(&mut ctx, make_error());
912        assert!(!e.retry(), "non-idempotent request should never retry");
913        assert_eq!(ctx.retries, 0);
914    }
915
916    #[test]
917    fn connect_failure_clears_upstream_connect_start() {
918        let mut ctx = PingoraRequestCtx::default();
919        ctx.upstream_connect_start = Some(std::time::Instant::now());
920        let _e = handle_connect_failure(&mut ctx, make_error());
921        assert!(
922            ctx.upstream_connect_start.is_none(),
923            "failed connect should consume upstream_connect_start for duration recording"
924        );
925    }
926
927    #[test]
928    fn response_503_retries_when_status5xx_enabled() {
929        let mut ctx = PingoraRequestCtx::default();
930        ctx.request_is_idempotent = true;
931        ctx.retry_policy = Some(Arc::new(praxis_core::config::RetryPolicy {
932            configured: true,
933            retriable_conditions: vec![praxis_core::config::RetriableCondition::Status5xx],
934            ..praxis_core::config::RetryPolicy::legacy_default()
935        }));
936        let e = maybe_retry_response(&mut ctx, 503).expect("503 should be retriable");
937        assert!(e.retry(), "503 under Status5xx should set retry");
938        assert_eq!(ctx.retries, 1);
939        assert!(ctx.reselect_on_retry);
940        assert!(ctx.pending_backoff.is_some());
941    }
942
943    #[test]
944    fn response_404_does_not_retry() {
945        let mut ctx = PingoraRequestCtx::default();
946        ctx.request_is_idempotent = true;
947        ctx.retry_policy = Some(Arc::new(praxis_core::config::RetryPolicy {
948            retriable_conditions: vec![praxis_core::config::RetriableCondition::Status5xx],
949            ..praxis_core::config::RetryPolicy::legacy_default()
950        }));
951        assert!(
952            maybe_retry_response(&mut ctx, 404).is_none(),
953            "404 must never trigger status-based retry"
954        );
955        assert_eq!(ctx.retries, 0);
956    }
957
958    #[test]
959    fn response_502_does_not_retry_under_legacy_default() {
960        let mut ctx = PingoraRequestCtx::default();
961        ctx.request_is_idempotent = true;
962        // Legacy default is connect_failure only — no Status5xx.
963        assert!(
964            maybe_retry_response(&mut ctx, 502).is_none(),
965            "legacy default must forward 5xx without retry"
966        );
967    }
968
969    #[test]
970    fn max_retries_zero_disables_connect_retry() {
971        let mut ctx = PingoraRequestCtx::default();
972        ctx.request_is_idempotent = true;
973        ctx.retry_policy = Some(Arc::new(praxis_core::config::RetryPolicy {
974            max_retries: Some(0),
975            ..praxis_core::config::RetryPolicy::legacy_default()
976        }));
977        let e = handle_connect_failure(&mut ctx, make_error());
978        assert!(!e.retry(), "max_retries: 0 must disable retries");
979        assert_eq!(ctx.retries, 0);
980    }
981
982    #[test]
983    fn non_idempotent_clears_pingora_default_retry_flag() {
984        let mut ctx = PingoraRequestCtx::default();
985        ctx.request_is_idempotent = false;
986        // Simulate Pingora marking the error retriable by default.
987        let mut e = make_error();
988        e.set_retry(true);
989        let e = handle_connect_failure(&mut ctx, e);
990        assert!(!e.retry(), "policy denial must clear Pingora's default retry flag");
991        assert_eq!(ctx.retries, 0);
992    }
993
994    #[tokio::test]
995    async fn logging_cleanup_noop_when_response_phase_done() {
996        let registry = praxis_filter::FilterRegistry::with_builtins();
997        let pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
998        let mut ctx = PingoraRequestCtx::default();
999        ctx.response_phase_done = true;
1000        ctx.request_snapshot = Some(praxis_filter::Request {
1001            method: http::Method::GET,
1002            uri: "/".parse().unwrap(),
1003            headers: http::HeaderMap::new(),
1004        });
1005        logging_cleanup(&pipeline, &mut ctx).await;
1006    }
1007
1008    #[tokio::test]
1009    async fn logging_cleanup_noop_when_no_snapshot() {
1010        let registry = praxis_filter::FilterRegistry::with_builtins();
1011        let pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
1012        let mut ctx = PingoraRequestCtx::default();
1013        ctx.response_phase_done = false;
1014        ctx.request_snapshot = None;
1015        logging_cleanup(&pipeline, &mut ctx).await;
1016    }
1017
1018    #[tokio::test]
1019    async fn logging_cleanup_runs_response_pipeline_when_needed() {
1020        let registry = praxis_filter::FilterRegistry::with_builtins();
1021        let pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
1022        let mut ctx = PingoraRequestCtx::default();
1023        ctx.response_phase_done = false;
1024        ctx.cluster = Some(Arc::from("test-cluster"));
1025        ctx.request_snapshot = Some(praxis_filter::Request {
1026            method: http::Method::GET,
1027            uri: "/test".parse().unwrap(),
1028            headers: http::HeaderMap::new(),
1029        });
1030        logging_cleanup(&pipeline, &mut ctx).await;
1031        assert_eq!(
1032            ctx.cluster.as_deref(),
1033            Some("test-cluster"),
1034            "cluster must be restored so the fallback access record can attribute the failure"
1035        );
1036    }
1037
1038    #[tokio::test]
1039    async fn logging_cleanup_preserves_filter_metadata() {
1040        let registry = praxis_filter::FilterRegistry::with_builtins();
1041        let pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
1042        let mut ctx = PingoraRequestCtx::default();
1043        ctx.response_phase_done = false;
1044        ctx.filter_metadata
1045            .insert("json_rpc.method".to_owned(), "service/invoke".to_owned());
1046        ctx.request_snapshot = Some(praxis_filter::Request {
1047            method: http::Method::POST,
1048            uri: "/api".parse().unwrap(),
1049            headers: http::HeaderMap::new(),
1050        });
1051        logging_cleanup(&pipeline, &mut ctx).await;
1052        assert_eq!(
1053            ctx.filter_metadata.get("json_rpc.method").map(String::as_str),
1054            Some("service/invoke"),
1055            "filter_metadata should survive logging_cleanup"
1056        );
1057    }
1058
1059    #[tokio::test]
1060    async fn logging_cleanup_preserves_extensions() {
1061        let registry = praxis_filter::FilterRegistry::with_builtins();
1062        let pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
1063        let mut ctx = PingoraRequestCtx::default();
1064        ctx.response_phase_done = false;
1065        ctx.extensions.insert(42_u32);
1066        ctx.request_snapshot = Some(praxis_filter::Request {
1067            method: http::Method::POST,
1068            uri: "/test".parse().unwrap(),
1069            headers: http::HeaderMap::new(),
1070        });
1071        logging_cleanup(&pipeline, &mut ctx).await;
1072        assert_eq!(
1073            ctx.extensions.get::<u32>(),
1074            Some(&42),
1075            "extensions should survive logging_cleanup"
1076        );
1077    }
1078
1079    #[test]
1080    fn passive_health_error_is_failure() {
1081        let (pipeline, ctx) = make_passive_scenario(Some(3), Some(2));
1082        let error = make_error();
1083        record_passive_health(&pipeline, Some(&error), &ctx);
1084
1085        let registry = pipeline.health_registry().unwrap();
1086        let entry = registry.get("test-cluster").unwrap();
1087        assert!(
1088            entry.endpoints()[0].is_healthy(),
1089            "single failure should not yet mark unhealthy (threshold=3)"
1090        );
1091    }
1092
1093    #[test]
1094    fn passive_health_status_500_is_failure() {
1095        let (pipeline, mut ctx) = make_passive_scenario(Some(3), Some(2));
1096        ctx.upstream_response_status = Some(500);
1097        record_passive_health(&pipeline, None, &ctx);
1098
1099        let registry = pipeline.health_registry().unwrap();
1100        let entry = registry.get("test-cluster").unwrap();
1101        assert!(
1102            entry.endpoints()[0].is_healthy(),
1103            "single 500 should not yet mark unhealthy (threshold=3)"
1104        );
1105    }
1106
1107    #[test]
1108    fn passive_health_status_below_500_is_success() {
1109        let (pipeline, mut ctx) = make_passive_scenario(Some(2), Some(1));
1110        ctx.upstream_response_status = Some(499);
1111        record_passive_health(&pipeline, None, &ctx);
1112
1113        let registry = pipeline.health_registry().unwrap();
1114        let entry = registry.get("test-cluster").unwrap();
1115        assert!(entry.endpoints()[0].is_healthy(), "status 499 should count as success");
1116    }
1117
1118    #[test]
1119    fn passive_unhealthy_threshold_transition() {
1120        let (pipeline, ctx) = make_passive_scenario(Some(2), Some(1));
1121        let error = make_error();
1122        record_passive_health(&pipeline, Some(&error), &ctx);
1123        record_passive_health(&pipeline, Some(&error), &ctx);
1124
1125        let registry = pipeline.health_registry().unwrap();
1126        let entry = registry.get("test-cluster").unwrap();
1127        assert!(
1128            !entry.endpoints()[0].is_healthy(),
1129            "2 consecutive failures should mark unhealthy (threshold=2)"
1130        );
1131    }
1132
1133    #[test]
1134    fn passive_healthy_threshold_recovery() {
1135        let (pipeline, ctx) = make_passive_scenario(Some(1), Some(2));
1136        let error = make_error();
1137        record_passive_health(&pipeline, Some(&error), &ctx);
1138
1139        let registry = pipeline.health_registry().unwrap();
1140        let entry = registry.get("test-cluster").unwrap();
1141        assert!(
1142            !entry.endpoints()[0].is_healthy(),
1143            "should be unhealthy after 1 failure"
1144        );
1145
1146        let ctx_ok = make_passive_ctx("test-cluster", 0, Some(200));
1147        record_passive_health(&pipeline, None, &ctx_ok);
1148        assert!(
1149            !entry.endpoints()[0].is_healthy(),
1150            "one success should not recover (threshold=2)"
1151        );
1152
1153        record_passive_health(&pipeline, None, &ctx_ok);
1154        assert!(
1155            entry.endpoints()[0].is_healthy(),
1156            "2 consecutive successes should recover (threshold=2)"
1157        );
1158    }
1159
1160    #[test]
1161    fn passive_health_no_thresholds_is_noop() {
1162        let (pipeline, ctx) = make_passive_scenario(None, None);
1163        let error = make_error();
1164        record_passive_health(&pipeline, Some(&error), &ctx);
1165
1166        let registry = pipeline.health_registry().unwrap();
1167        let entry = registry.get("test-cluster").unwrap();
1168        assert!(
1169            entry.endpoints()[0].is_healthy(),
1170            "no passive thresholds means failures are no-op"
1171        );
1172    }
1173
1174    #[test]
1175    fn passive_health_endpoint_index_out_of_bounds() {
1176        let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
1177        ctx.selected_endpoint_index = Some(999);
1178        let error = make_error();
1179        record_passive_health(&pipeline, Some(&error), &ctx);
1180
1181        let registry = pipeline.health_registry().unwrap();
1182        let entry = registry.get("test-cluster").unwrap();
1183        assert!(entry.endpoints()[0].is_healthy(), "out-of-bounds index should be no-op");
1184    }
1185
1186    #[test]
1187    fn passive_health_missing_cluster_is_noop() {
1188        let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
1189        ctx.cluster = None;
1190        ctx.metrics_cluster = None;
1191        let error = make_error();
1192        record_passive_health(&pipeline, Some(&error), &ctx);
1193    }
1194
1195    #[test]
1196    fn passive_health_falls_back_to_metrics_cluster() {
1197        let (pipeline, mut ctx) = make_passive_scenario(Some(2), Some(1));
1198        ctx.cluster = None;
1199        ctx.metrics_cluster = Some(Arc::from("test-cluster"));
1200        let error = make_error();
1201        record_passive_health(&pipeline, Some(&error), &ctx);
1202        record_passive_health(&pipeline, Some(&error), &ctx);
1203
1204        let registry = pipeline.health_registry().unwrap();
1205        let entry = registry.get("test-cluster").unwrap();
1206        assert!(
1207            !entry.endpoints()[0].is_healthy(),
1208            "fallback to metrics_cluster should still record passive health"
1209        );
1210    }
1211
1212    #[test]
1213    fn passive_health_missing_endpoint_index_is_noop() {
1214        let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
1215        ctx.selected_endpoint_index = None;
1216        let error = make_error();
1217        record_passive_health(&pipeline, Some(&error), &ctx);
1218    }
1219
1220    #[test]
1221    fn passive_health_missing_registry_is_noop() {
1222        let registry = praxis_filter::FilterRegistry::with_builtins();
1223        let pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
1224        let mut ctx = PingoraRequestCtx::default();
1225        ctx.cluster = Some(Arc::from("test-cluster"));
1226        ctx.selected_endpoint_index = Some(0);
1227        let error = make_error();
1228        record_passive_health(&pipeline, Some(&error), &ctx);
1229    }
1230
1231    #[test]
1232    fn passive_health_unknown_cluster_is_noop() {
1233        let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
1234        ctx.cluster = Some(Arc::from("nonexistent"));
1235        let error = make_error();
1236        record_passive_health(&pipeline, Some(&error), &ctx);
1237    }
1238
1239    #[test]
1240    fn size_limit_none_body_returns_false() {
1241        let mut bytes = 0_u64;
1242        assert!(!check_body_size_limit(None, &mut bytes, 100));
1243        assert_eq!(bytes, 0, "accumulated bytes unchanged for None body");
1244    }
1245
1246    #[test]
1247    fn size_limit_within_limit() {
1248        let mut bytes = 0_u64;
1249        let body = Some(Bytes::from_static(b"hello"));
1250        assert!(!check_body_size_limit(body.as_ref(), &mut bytes, 10));
1251        assert_eq!(bytes, 5);
1252    }
1253
1254    #[test]
1255    fn size_limit_at_exact_limit() {
1256        let mut bytes = 0_u64;
1257        let body = Some(Bytes::from_static(b"exact"));
1258        assert!(!check_body_size_limit(body.as_ref(), &mut bytes, 5));
1259        assert_eq!(bytes, 5);
1260    }
1261
1262    #[test]
1263    fn size_limit_exceeds_limit() {
1264        let mut bytes = 0_u64;
1265        let body = Some(Bytes::from_static(b"toolong"));
1266        assert!(check_body_size_limit(body.as_ref(), &mut bytes, 3));
1267    }
1268
1269    #[test]
1270    fn size_limit_cumulative_overflow() {
1271        let mut bytes = 0_u64;
1272        let first = Some(Bytes::from_static(b"aaa"));
1273        assert!(!check_body_size_limit(first.as_ref(), &mut bytes, 5));
1274
1275        let second = Some(Bytes::from_static(b"bbb"));
1276        assert!(check_body_size_limit(second.as_ref(), &mut bytes, 5));
1277        assert_eq!(bytes, 6);
1278    }
1279
1280    #[test]
1281    fn stream_buffer_accumulates_chunks() {
1282        let mut body = Some(Bytes::from_static(b"hello "));
1283        let mut buf: Option<BodyBuffer> = None;
1284        assert!(!accumulate_stream_buffer(&mut body, &mut buf, false, Some(100)));
1285        assert!(buf.is_some());
1286
1287        body = Some(Bytes::from_static(b"world"));
1288        assert!(!accumulate_stream_buffer(&mut body, &mut buf, false, Some(100)));
1289
1290        let frozen = buf.take().unwrap().freeze();
1291        assert_eq!(frozen, Bytes::from_static(b"hello world"));
1292    }
1293
1294    #[test]
1295    fn stream_buffer_freezes_at_eos() {
1296        let mut body = Some(Bytes::from_static(b"data"));
1297        let mut buf: Option<BodyBuffer> = None;
1298        assert!(!accumulate_stream_buffer(&mut body, &mut buf, false, Some(100)));
1299
1300        body = Some(Bytes::from_static(b" end"));
1301        assert!(!accumulate_stream_buffer(&mut body, &mut buf, true, Some(100)));
1302        assert!(buf.is_none(), "buffer should be taken at EOS");
1303        assert_eq!(body.unwrap(), Bytes::from_static(b"data end"));
1304    }
1305
1306    #[test]
1307    fn stream_buffer_overflow() {
1308        let mut body = Some(Bytes::from_static(b"too long"));
1309        let mut buf: Option<BodyBuffer> = None;
1310        assert!(accumulate_stream_buffer(&mut body, &mut buf, false, Some(5)));
1311    }
1312
1313    #[test]
1314    fn stream_buffer_none_body() {
1315        let mut body: Option<Bytes> = None;
1316        let mut buf: Option<BodyBuffer> = None;
1317        assert!(!accumulate_stream_buffer(&mut body, &mut buf, false, Some(100)));
1318        assert!(buf.is_none());
1319    }
1320
1321    #[test]
1322    fn stream_buffer_uses_absolute_max_when_none() {
1323        let mut body = Some(Bytes::from_static(b"data"));
1324        let mut buf: Option<BodyBuffer> = None;
1325        assert!(!accumulate_stream_buffer(&mut body, &mut buf, false, None));
1326        assert!(buf.is_some(), "should create buffer with absolute max");
1327    }
1328
1329    #[test]
1330    fn suppress_clears_body_when_buffering() {
1331        let mut body = Some(Bytes::from_static(b"data"));
1332        suppress_stream_buffer_chunk(&mut body, true, false, false);
1333        assert!(body.is_none());
1334    }
1335
1336    #[test]
1337    fn suppress_noop_when_not_stream_buffer() {
1338        let mut body = Some(Bytes::from_static(b"data"));
1339        suppress_stream_buffer_chunk(&mut body, false, false, false);
1340        assert!(body.is_some());
1341    }
1342
1343    #[test]
1344    fn suppress_noop_when_released() {
1345        let mut body = Some(Bytes::from_static(b"data"));
1346        suppress_stream_buffer_chunk(&mut body, true, true, false);
1347        assert!(body.is_some());
1348    }
1349
1350    #[test]
1351    fn suppress_noop_at_eos() {
1352        let mut body = Some(Bytes::from_static(b"data"));
1353        suppress_stream_buffer_chunk(&mut body, true, false, true);
1354        assert!(body.is_some());
1355    }
1356
1357    #[test]
1358    fn release_sets_flag_and_flushes_buffer() {
1359        let mut body: Option<Bytes> = None;
1360        let mut released = false;
1361        let mut buf = Some(BodyBuffer::new(100));
1362        buf.as_mut().unwrap().push(Bytes::from_static(b"buffered")).unwrap();
1363
1364        release_stream_buffer(&mut body, true, &mut released, &mut buf, false);
1365        assert!(released);
1366        assert_eq!(body.unwrap(), Bytes::from_static(b"buffered"));
1367        assert!(buf.is_none());
1368    }
1369
1370    #[test]
1371    fn release_noop_when_already_released() {
1372        let mut body: Option<Bytes> = None;
1373        let mut released = true;
1374        let mut buf: Option<BodyBuffer> = None;
1375
1376        release_stream_buffer(&mut body, true, &mut released, &mut buf, false);
1377        assert!(body.is_none(), "body should be unchanged when already released");
1378    }
1379
1380    #[test]
1381    fn release_noop_when_not_stream_buffer() {
1382        let mut body: Option<Bytes> = None;
1383        let mut released = false;
1384        let mut buf: Option<BodyBuffer> = None;
1385
1386        release_stream_buffer(&mut body, false, &mut released, &mut buf, false);
1387        assert!(!released, "released flag should be unchanged for non-stream-buffer");
1388    }
1389
1390    #[test]
1391    fn release_at_eos_sets_flag_but_no_flush() {
1392        let mut body: Option<Bytes> = None;
1393        let mut released = false;
1394        let mut buf = Some(BodyBuffer::new(100));
1395        buf.as_mut().unwrap().push(Bytes::from_static(b"data")).unwrap();
1396
1397        release_stream_buffer(&mut body, true, &mut released, &mut buf, true);
1398        assert!(released);
1399        assert!(body.is_none(), "body should not be overwritten at EOS");
1400        assert!(buf.is_some(), "buffer should not be taken at EOS");
1401    }
1402
1403    #[test]
1404    fn write_back_transfers_fields() {
1405        let mut ctx = PingoraRequestCtx::default();
1406
1407        let mut extensions = RequestExtensions::new();
1408        extensions.insert(42_u32);
1409
1410        let state_val: Box<dyn std::any::Any + Send + Sync> = Box::new(99_i32);
1411        let filter_state = HashMap::from([(0_usize, state_val)]);
1412
1413        let output = BodyFilterOutput {
1414            cluster: Some(Arc::from("test-cluster")),
1415            upstream: Some(Upstream {
1416                address: Arc::from("10.0.0.1:80"),
1417                authority: None,
1418                connection: Arc::new(ConnectionOptions::default()),
1419                tls: None,
1420            }),
1421            extensions,
1422            attempted_endpoints: Vec::new(),
1423            filter_metadata: HashMap::from([("key".to_owned(), "val".to_owned())]),
1424            filter_state,
1425            executed_filter_indices: vec![true, false],
1426            body_done_indices: vec![false, true],
1427        };
1428        output.write_back(&mut ctx);
1429
1430        assert_eq!(ctx.cluster.as_deref(), Some("test-cluster"));
1431        assert!(ctx.upstream.is_some(), "upstream should transfer");
1432        assert_eq!(ctx.upstream.as_ref().unwrap().address.as_ref(), "10.0.0.1:80");
1433        assert_eq!(ctx.extensions.get::<u32>(), Some(&42));
1434        assert_eq!(ctx.filter_metadata.get("key").map(String::as_str), Some("val"));
1435        assert_eq!(ctx.filter_state.len(), 1, "filter_state should transfer");
1436        assert_eq!(
1437            ctx.filter_state.get(&0).and_then(|v| v.downcast_ref::<i32>()),
1438            Some(&99)
1439        );
1440        assert_eq!(ctx.cached_executed_filter_indices, vec![true, false]);
1441        assert_eq!(ctx.cached_body_done_indices, vec![false, true]);
1442    }
1443
1444    // -------------------------------------------------------------------------
1445    // Fallback Access Log
1446    // -------------------------------------------------------------------------
1447
1448    #[test]
1449    fn fallback_access_log_emits_for_incomplete_request() {
1450        let pipeline = access_log_pipeline();
1451        let mut ctx = make_fallback_ctx();
1452
1453        let events = capture_access_events(|| maybe_emit_fallback_access_log(&pipeline, 502, &mut ctx));
1454        assert_eq!(
1455            events.len(),
1456            1,
1457            "incomplete request must produce a fallback access record"
1458        );
1459    }
1460
1461    #[test]
1462    fn fallback_access_log_skips_completed_delivery() {
1463        let pipeline = access_log_pipeline();
1464        let mut ctx = make_fallback_ctx();
1465        ctx.response_delivery_complete = true;
1466
1467        let events = capture_access_events(|| maybe_emit_fallback_access_log(&pipeline, 200, &mut ctx));
1468        assert!(events.is_empty(), "completed delivery already logged via the filter");
1469    }
1470
1471    #[test]
1472    fn fallback_access_log_skips_upgraded_connections() {
1473        let pipeline = access_log_pipeline();
1474        let mut ctx = make_fallback_ctx();
1475        ctx.connection_upgraded = true;
1476
1477        let events = capture_access_events(|| maybe_emit_fallback_access_log(&pipeline, 101, &mut ctx));
1478        assert!(events.is_empty(), "upgraded connections have no body completion");
1479    }
1480
1481    #[test]
1482    fn fallback_access_log_skips_without_access_log_filter() {
1483        let registry = praxis_filter::FilterRegistry::with_builtins();
1484        let pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
1485        let mut ctx = make_fallback_ctx();
1486
1487        let events = capture_access_events(|| maybe_emit_fallback_access_log(&pipeline, 502, &mut ctx));
1488        assert!(events.is_empty(), "no access_log filter means no fallback record");
1489    }
1490
1491    #[test]
1492    fn fallback_access_log_honors_entry_conditions() {
1493        let registry = praxis_filter::FilterRegistry::with_builtins();
1494        let mut entries = vec![praxis_filter::FilterEntry {
1495            branch_chains: None,
1496            conditions: vec![serde_yaml::from_str("when:\n  path_prefix: /api\n").unwrap()],
1497            failure_mode: praxis_filter::FailureMode::default(),
1498            filter_type: "access_log".to_owned(),
1499            config: serde_yaml::Value::Null,
1500            name: None,
1501            response_conditions: vec![],
1502        }];
1503        let pipeline = FilterPipeline::build(&mut entries, &registry).unwrap();
1504
1505        let mut excluded = make_fallback_ctx();
1506        let events = capture_access_events(|| maybe_emit_fallback_access_log(&pipeline, 502, &mut excluded));
1507        assert!(
1508            events.is_empty(),
1509            "requests the operator scoped out must not gain fallback records"
1510        );
1511
1512        let mut included = make_fallback_ctx();
1513        if let Some(snapshot) = included.request_snapshot.as_mut() {
1514            snapshot.uri = "/api/users".parse().unwrap();
1515        }
1516        let events = capture_access_events(|| maybe_emit_fallback_access_log(&pipeline, 502, &mut included));
1517        assert_eq!(events.len(), 1, "in-scope incomplete requests still get a record");
1518    }
1519
1520    #[test]
1521    fn aborted_response_body_at_eos_is_not_marked_delivered() {
1522        // access_log declares read-only response body access, so the hook
1523        // reaches the SizeLimit check instead of early-returning.
1524        let pipeline = access_log_pipeline();
1525        let mut ctx = make_fallback_ctx();
1526        ctx.response_body_mode = BodyMode::SizeLimit { max_bytes: 4 };
1527        let mut body = Some(Bytes::from_static(b"exceeds the limit"));
1528
1529        let result = response_body_filter::execute(&pipeline, &mut body, true, &mut ctx);
1530        assert!(result.is_err(), "over-limit body must abort");
1531        assert!(
1532            !ctx.response_delivery_complete,
1533            "a response aborted at end-of-stream was not delivered; the fallback record must fire"
1534        );
1535    }
1536
1537    #[test]
1538    fn response_body_eos_marks_delivery_complete() {
1539        let registry = praxis_filter::FilterRegistry::with_builtins();
1540        let pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
1541        let mut ctx = PingoraRequestCtx::default();
1542        let mut body: Option<Bytes> = None;
1543
1544        let _timeout = response_body_filter::execute(&pipeline, &mut body, false, &mut ctx).unwrap();
1545        assert!(
1546            !ctx.response_delivery_complete,
1547            "mid-stream chunks must not mark delivery complete"
1548        );
1549
1550        let _timeout = response_body_filter::execute(&pipeline, &mut body, true, &mut ctx).unwrap();
1551        assert!(ctx.response_delivery_complete, "end-of-stream marks delivery complete");
1552    }
1553
1554    // -------------------------------------------------------------------------
1555    // Span Attribute Helpers
1556    // -------------------------------------------------------------------------
1557
1558    #[test]
1559    fn http_version_label_http_09() {
1560        assert_eq!(
1561            http_version_label(http::Version::HTTP_09),
1562            "0.9",
1563            "HTTP/0.9 should map to '0.9'"
1564        );
1565    }
1566
1567    #[test]
1568    fn http_version_label_http_10() {
1569        assert_eq!(
1570            http_version_label(http::Version::HTTP_10),
1571            "1.0",
1572            "HTTP/1.0 should map to '1.0'"
1573        );
1574    }
1575
1576    #[test]
1577    fn http_version_label_http_11() {
1578        assert_eq!(
1579            http_version_label(http::Version::HTTP_11),
1580            "1.1",
1581            "HTTP/1.1 should map to '1.1'"
1582        );
1583    }
1584
1585    #[test]
1586    fn http_version_label_http_2() {
1587        assert_eq!(
1588            http_version_label(http::Version::HTTP_2),
1589            "2",
1590            "HTTP/2 should map to '2'"
1591        );
1592    }
1593
1594    #[test]
1595    fn http_version_label_http_3() {
1596        assert_eq!(
1597            http_version_label(http::Version::HTTP_3),
1598            "3",
1599            "HTTP/3 should map to '3'"
1600        );
1601    }
1602
1603    #[test]
1604    fn record_response_span_attributes_noop_for_disabled_span() {
1605        let ctx = PingoraRequestCtx::default();
1606        assert!(ctx.request_span.is_disabled(), "default span should be disabled");
1607    }
1608
1609    /// Layer that captures every `Span::record` call as `(field, value)` pairs.
1610    #[derive(Clone, Default)]
1611    struct RecordCapture(Arc<std::sync::Mutex<Vec<(String, String)>>>);
1612
1613    impl<S> tracing_subscriber::Layer<S> for RecordCapture
1614    where
1615        S: tracing::Subscriber + for<'a> tracing_subscriber::registry::LookupSpan<'a>,
1616    {
1617        fn on_record(
1618            &self,
1619            _id: &tracing::span::Id,
1620            values: &tracing::span::Record<'_>,
1621            _ctx: tracing_subscriber::layer::Context<'_, S>,
1622        ) {
1623            struct Visitor<'a>(&'a mut Vec<(String, String)>);
1624            impl tracing::field::Visit for Visitor<'_> {
1625                fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
1626                    self.0.push((field.name().to_owned(), format!("{value:?}")));
1627                }
1628            }
1629            let mut captured = self.0.lock().expect("capture lock");
1630            values.record(&mut Visitor(&mut captured));
1631        }
1632    }
1633
1634    #[test]
1635    fn record_response_span_fields_records_status_upstream_and_cluster() {
1636        use tracing_subscriber::layer::SubscriberExt as _;
1637
1638        let capture = RecordCapture::default();
1639        let subscriber = tracing_subscriber::registry().with(capture.clone());
1640        let _guard = tracing::subscriber::set_default(subscriber);
1641
1642        let mut ctx = PingoraRequestCtx::default();
1643        ctx.metrics_cluster = Some(Arc::from("api-cluster"));
1644        ctx.upstream_for_retry = Some(Upstream {
1645            address: Arc::from("10.0.0.1:80"),
1646            authority: None,
1647            connection: Arc::new(ConnectionOptions::default()),
1648            tls: None,
1649        });
1650        ctx.request_span = tracing::info_span!(
1651            "test_span",
1652            "http.response.status_code" = tracing::field::Empty,
1653            "otel.status_code" = tracing::field::Empty,
1654            "upstream.address" = tracing::field::Empty,
1655            "upstream.cluster" = tracing::field::Empty,
1656        );
1657
1658        record_response_span_fields(Some(http::StatusCode::SERVICE_UNAVAILABLE), "GET", None, &ctx);
1659
1660        let captured = capture.0.lock().expect("capture lock");
1661        let get = |name: &str| {
1662            let value = captured.iter().find(|(f, _)| f == name).map(|(_, v)| v.clone());
1663            assert!(value.is_some(), "field {name} not recorded; got {captured:?}");
1664            value.unwrap_or_default()
1665        };
1666        assert_eq!(get("http.response.status_code"), "503");
1667        assert_eq!(get("otel.status_code"), "\"ERROR\"", "5xx must set otel error status");
1668        assert_eq!(get("upstream.address"), "\"10.0.0.1:80\"");
1669        assert_eq!(get("upstream.cluster"), "\"api-cluster\"");
1670    }
1671
1672    #[test]
1673    fn record_response_span_fields_success_has_no_error_status() {
1674        use tracing_subscriber::layer::SubscriberExt as _;
1675
1676        let capture = RecordCapture::default();
1677        let subscriber = tracing_subscriber::registry().with(capture.clone());
1678        let _guard = tracing::subscriber::set_default(subscriber);
1679
1680        let mut ctx = PingoraRequestCtx::default();
1681        ctx.request_span = tracing::info_span!(
1682            "test_span",
1683            "http.response.status_code" = tracing::field::Empty,
1684            "otel.status_code" = tracing::field::Empty,
1685        );
1686
1687        record_response_span_fields(Some(http::StatusCode::OK), "GET", None, &ctx);
1688
1689        let captured = capture.0.lock().expect("capture lock");
1690        assert!(
1691            captured
1692                .iter()
1693                .any(|(f, v)| f == "http.response.status_code" && v == "200"),
1694            "status should be recorded: {captured:?}"
1695        );
1696        assert!(
1697            !captured.iter().any(|(f, _)| f == "otel.status_code"),
1698            "2xx must not set otel error status: {captured:?}"
1699        );
1700    }
1701
1702    #[test]
1703    fn record_response_span_attributes_records_exchange_span_fields() {
1704        let mut ctx = PingoraRequestCtx::default();
1705        ctx.request_span = tracing::info_span!(
1706            "test_request",
1707            "http.response.status_code" = tracing::field::Empty,
1708            "server.address" = tracing::field::Empty,
1709            "upstream.cluster" = tracing::field::Empty,
1710        );
1711        ctx.upstream_exchange_span = tracing::info_span!(
1712            parent: &ctx.request_span,
1713            "upstream_exchange",
1714            "http.response.status_code" = tracing::field::Empty,
1715            "http.response.body.size" = tracing::field::Empty,
1716        );
1717        ctx.response_body_bytes = 4096;
1718
1719        // Verify recording on exchange span does not panic.
1720        ctx.upstream_exchange_span.record("http.response.status_code", 200_u16);
1721        ctx.upstream_exchange_span.record("http.response.body.size", 4096_u64);
1722    }
1723
1724    #[test]
1725    fn record_response_span_attributes_skips_exchange_when_disabled() {
1726        let mut ctx = PingoraRequestCtx::default();
1727        ctx.request_span = tracing::info_span!(
1728            "test_request",
1729            "http.response.status_code" = tracing::field::Empty,
1730            "server.address" = tracing::field::Empty,
1731            "upstream.cluster" = tracing::field::Empty,
1732        );
1733        // Exchange span remains disabled (default).
1734        assert!(
1735            ctx.upstream_exchange_span.is_disabled(),
1736            "exchange span should be disabled by default"
1737        );
1738        // Should not panic when exchange span is disabled.
1739    }
1740
1741    // -------------------------------------------------------------------------
1742    // Span Event Tests
1743    // -------------------------------------------------------------------------
1744
1745    #[test]
1746    fn retry_with_upstream_address_sets_retry_flag() {
1747        let mut ctx = PingoraRequestCtx::default();
1748        ctx.request_is_idempotent = true;
1749        ctx.upstream_for_retry = Some(Upstream {
1750            address: Arc::from("10.0.0.1:8080"),
1751            connection: Arc::new(ConnectionOptions::default()),
1752            tls: None,
1753            authority: None,
1754        });
1755        let e = handle_connect_failure(&mut ctx, make_error());
1756        assert!(e.retry(), "should retry with upstream address present");
1757        assert_eq!(ctx.retries, 1, "retry counter should increment to 1");
1758    }
1759
1760    #[test]
1761    fn retry_without_upstream_address_uses_fallback() {
1762        let mut ctx = PingoraRequestCtx::default();
1763        ctx.request_is_idempotent = true;
1764        ctx.upstream_for_retry = None;
1765        let e = handle_connect_failure(&mut ctx, make_error());
1766        assert!(
1767            e.retry(),
1768            "should retry even when upstream_for_retry is None (address defaults to unknown)"
1769        );
1770        assert_eq!(ctx.retries, 1, "retry counter should increment to 1");
1771    }
1772
1773    #[test]
1774    fn retry_exhausted_with_upstream_address_does_not_retry() {
1775        let mut ctx = PingoraRequestCtx::default();
1776        ctx.request_is_idempotent = true;
1777        ctx.retries = MAX_RETRIES as u32;
1778        ctx.upstream_for_retry = Some(Upstream {
1779            address: Arc::from("10.0.0.2:443"),
1780            connection: Arc::new(ConnectionOptions::default()),
1781            tls: None,
1782            authority: None,
1783        });
1784        let e = handle_connect_failure(&mut ctx, make_error());
1785        assert!(
1786            !e.retry(),
1787            "should not retry after MAX_RETRIES even with upstream address"
1788        );
1789    }
1790
1791    #[test]
1792    fn large_body_skip_with_upstream_address() {
1793        let mut ctx = PingoraRequestCtx::default();
1794        ctx.request_is_idempotent = true;
1795        ctx.request_body_bytes = RETRY_BODY_LIMIT + 1;
1796        ctx.upstream_for_retry = Some(Upstream {
1797            address: Arc::from("10.0.0.3:8080"),
1798            connection: Arc::new(ConnectionOptions::default()),
1799            tls: None,
1800            authority: None,
1801        });
1802        let e = handle_connect_failure(&mut ctx, make_error());
1803        assert!(!e.retry(), "should not retry large body even with upstream address");
1804        assert_eq!(ctx.retries, 0, "retry counter should not increment");
1805    }
1806
1807    // -------------------------------------------------------------------------
1808    // Test Utilities
1809    // -------------------------------------------------------------------------
1810
1811    /// Create a connect error for tests.
1812    fn make_error() -> Box<pingora_core::Error> {
1813        pingora_core::Error::explain(pingora_core::ErrorType::ConnectError, "test connect failure")
1814    }
1815
1816    /// Build a pipeline containing an `access_log` filter.
1817    fn access_log_pipeline() -> FilterPipeline {
1818        let registry = praxis_filter::FilterRegistry::with_builtins();
1819        let mut entries = vec![praxis_filter::FilterEntry {
1820            branch_chains: None,
1821            conditions: vec![],
1822            failure_mode: praxis_filter::FailureMode::default(),
1823            filter_type: "access_log".to_owned(),
1824            config: serde_yaml::Value::Null,
1825            name: None,
1826            response_conditions: vec![],
1827        }];
1828        FilterPipeline::build(&mut entries, &registry).unwrap()
1829    }
1830
1831    /// Build a context with a request snapshot for fallback logging tests.
1832    fn make_fallback_ctx() -> PingoraRequestCtx {
1833        let mut ctx = PingoraRequestCtx::default();
1834        ctx.request_snapshot = Some(praxis_filter::Request {
1835            method: http::Method::GET,
1836            uri: "/incomplete".parse().unwrap(),
1837            headers: http::HeaderMap::new(),
1838        });
1839        ctx
1840    }
1841
1842    /// Capture `access` info events emitted while running `f`.
1843    fn capture_access_events<F: FnOnce()>(f: F) -> Vec<String> {
1844        use tracing_subscriber::layer::SubscriberExt as _;
1845
1846        let messages = Arc::new(std::sync::Mutex::new(Vec::<String>::new()));
1847        let capture = AccessCapture(Arc::clone(&messages));
1848        let subscriber = tracing_subscriber::registry().with(capture);
1849        tracing::subscriber::with_default(subscriber, f);
1850        let mut guard = messages.lock().unwrap();
1851        std::mem::take(&mut *guard)
1852    }
1853
1854    /// Layer capturing `access` records for assertions.
1855    struct AccessCapture(Arc<std::sync::Mutex<Vec<String>>>);
1856
1857    impl<S: tracing::Subscriber> tracing_subscriber::Layer<S> for AccessCapture {
1858        fn on_event(&self, event: &tracing::Event<'_>, _ctx: tracing_subscriber::layer::Context<'_, S>) {
1859            let mut visitor = AccessMessageVisitor(String::new());
1860            event.record(&mut visitor);
1861            if visitor.0.contains("access") {
1862                self.0.lock().unwrap().push(visitor.0);
1863            }
1864        }
1865    }
1866
1867    /// Visitor extracting the `message` field from an event.
1868    struct AccessMessageVisitor(String);
1869
1870    impl tracing::field::Visit for AccessMessageVisitor {
1871        fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
1872            if field.name() == "message" {
1873                self.0 = format!("{value:?}");
1874            }
1875        }
1876    }
1877
1878    /// Build a [`PingoraRequestCtx`] for passive health testing.
1879    fn make_passive_ctx(cluster: &str, endpoint_idx: usize, status: Option<u16>) -> PingoraRequestCtx {
1880        let mut ctx = PingoraRequestCtx::default();
1881        ctx.cluster = Some(Arc::from(cluster));
1882        ctx.selected_endpoint_index = Some(endpoint_idx);
1883        ctx.upstream_response_status = status;
1884        ctx
1885    }
1886
1887    /// Build a pipeline with a health registry and a matching context
1888    /// for passive health testing.
1889    fn make_passive_scenario(
1890        passive_unhealthy: Option<u32>,
1891        passive_healthy: Option<u32>,
1892    ) -> (FilterPipeline, PingoraRequestCtx) {
1893        use praxis_core::health::{ClusterHealthEntry, EndpointHealth};
1894
1895        let entry = ClusterHealthEntry::new(
1896            vec![EndpointHealth::new()],
1897            vec![Arc::from("10.0.0.1:80")],
1898            passive_unhealthy,
1899            passive_healthy,
1900        );
1901        let mut map = HashMap::new();
1902        map.insert(Arc::from("test-cluster"), Arc::new(entry));
1903        let health_registry = Arc::new(map);
1904
1905        let registry = praxis_filter::FilterRegistry::with_builtins();
1906        let mut pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
1907        pipeline.set_health_registry(health_registry);
1908
1909        let ctx = make_passive_ctx("test-cluster", 0, None);
1910
1911        (pipeline, ctx)
1912    }
1913}