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 watcher tasks run for the process
74/// lifetime; the caller keeps this `Vec` to stop them early via
75/// `send(true)` (dropping the senders does not stop them).
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 method = request_method_label(session, ctx);
463
464    let cluster = ctx.metrics_cluster_shared.clone().unwrap_or_else(metrics::cluster_none);
465
466    let route = ctx.metrics_route.clone().unwrap_or_else(metrics::route_unknown);
467
468    let labels = metrics::RequestMetricLabels {
469        cluster: cluster.clone(),
470        method,
471        route,
472        status_class,
473    };
474
475    emit_upstream_request_metric(ctx, &cluster);
476
477    if let Some(error_type) = ctx.error_type {
478        metrics::record_error(error_type);
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/// Resolve the bounded `method` label for a request.
493///
494/// Falls back to the request snapshot when the session header has already
495/// been consumed, and to `"UNKNOWN"` when neither carries a method.
496fn request_method_label(session: &Session, ctx: &PingoraRequestCtx) -> &'static str {
497    let request_method = session.req_header().method.as_str();
498    let raw_method = if request_method.is_empty() {
499        ctx.request_snapshot.as_ref().map_or("UNKNOWN", |r| r.method.as_str())
500    } else {
501        request_method
502    };
503    metrics::method_label(raw_method)
504}
505
506/// Count a request that reached an upstream endpoint.
507///
508/// Only requests that actually reached an upstream carry an upstream
509/// response status; filter rejections and connect failures never do, and
510/// must not inflate the upstream denominator.
511fn emit_upstream_request_metric(ctx: &PingoraRequestCtx, cluster: &::metrics::SharedString) {
512    if let Some(upstream_status) = ctx.upstream_response_status
513        && let Some(upstream) = ctx.upstream_for_retry.as_ref()
514    {
515        metrics::record_upstream_request(
516            cluster.clone(),
517            ::metrics::SharedString::from(Arc::clone(&upstream.address)),
518            metrics::status_class(upstream_status),
519        );
520    }
521}
522
523/// Record a passive health observation for the selected upstream endpoint.
524///
525/// Called from the `logging` hook on every completed request. Determines
526/// success/failure from the error argument and the stashed upstream
527/// response status code.
528///
529/// No-op when no upstream was selected, no health registry is available,
530/// or passive checking is not configured for the cluster.
531fn record_passive_health(pipeline: &FilterPipeline, error: Option<&pingora_core::Error>, ctx: &PingoraRequestCtx) {
532    let cluster_name = ctx.cluster.as_ref().or(ctx.metrics_cluster.as_ref());
533    let Some(cluster_name) = cluster_name else {
534        return;
535    };
536    let Some(idx) = ctx.selected_endpoint_index else {
537        return;
538    };
539    let Some(registry) = pipeline.health_registry() else {
540        return;
541    };
542    let Some(health) = registry.get(cluster_name) else {
543        return;
544    };
545
546    // A request that never contacted the upstream carries no signal about the
547    // endpoint, so skip it. upstream_contacted is set once a peer is resolved
548    // and stays set across retries (unlike upstream_for_retry, which a retry
549    // clears to force reselection), so it distinguishes a genuine connect or
550    // read failure from a filter reject or a proxy-generated terminal response
551    // after endpoint selection, which would otherwise record a spurious
552    // observation against the untouched endpoint and skew a real failure
553    // streak.
554    if !ctx.upstream_contacted {
555        return;
556    }
557
558    // Classify the observation by the error's origin. A client-sourced
559    // (Downstream) error carries no signal about the endpoint, so when it
560    // arrives without an upstream response we skip the observation
561    // entirely: recording a failure would eject a healthy upstream, and
562    // recording a success would clear a real failure streak and mask a
563    // failing one. Upstream/Internal/Unset errors and 5xx responses count
564    // as failures (Internal/Unset are kept because a real endpoint failure
565    // is not always tagged Upstream, and missing one is worse here than an
566    // occasional false positive).
567    let is_downstream_error = error.is_some_and(|e| matches!(e.esource(), pingora_core::ErrorSource::Downstream));
568    if is_downstream_error && ctx.upstream_response_status.is_none() {
569        return;
570    }
571    let is_failure =
572        ctx.upstream_response_status.is_some_and(|s| s >= 500) || (error.is_some() && !is_downstream_error);
573    apply_passive_threshold(health, idx, cluster_name, is_failure);
574}
575
576/// Apply passive health threshold for a single endpoint observation.
577fn apply_passive_threshold(
578    health: &praxis_core::health::ClusterHealthEntry,
579    idx: usize,
580    cluster_name: &Arc<str>,
581    is_failure: bool,
582) {
583    if is_failure {
584        if let Some(threshold) = health.passive_unhealthy_threshold()
585            && health
586                .endpoints()
587                .get(idx)
588                .is_some_and(|ep| ep.record_failure(threshold))
589        {
590            tracing::warn!(
591                cluster = %cluster_name,
592                endpoint_index = idx,
593                threshold,
594                "passive health: endpoint marked unhealthy"
595            );
596            emit_passive_health_transition(health, cluster_name, metrics::HEALTH_RESULT_UNHEALTHY);
597        }
598    } else if let Some(threshold) = health.passive_healthy_threshold()
599        && health
600            .endpoints()
601            .get(idx)
602            .is_some_and(|ep| ep.record_success(threshold))
603    {
604        tracing::info!(
605            cluster = %cluster_name,
606            endpoint_index = idx,
607            threshold,
608            "passive health: endpoint recovered"
609        );
610        emit_passive_health_transition(health, cluster_name, metrics::HEALTH_RESULT_HEALTHY);
611    }
612}
613
614/// Refresh health gauges and increment the transition counter after a passive flip.
615fn emit_passive_health_transition(
616    health: &praxis_core::health::ClusterHealthEntry,
617    cluster_name: &Arc<str>,
618    result: &'static str,
619) {
620    let (healthy, total) = metrics::count_healthy_endpoints(health);
621    metrics::record_health_transition(
622        ::metrics::SharedString::from(Arc::clone(cluster_name)),
623        result,
624        healthy,
625        total,
626    );
627}
628
629/// Map an [`http::Version`] to the [OTel `network.protocol.version`] value.
630///
631/// [OTel `network.protocol.version`]: https://opentelemetry.io/docs/specs/semconv/attributes-registry/network/
632pub(super) fn http_version_label(version: http::Version) -> &'static str {
633    match version {
634        http::Version::HTTP_09 => "0.9",
635        http::Version::HTTP_10 => "1.0",
636        http::Version::HTTP_11 => "1.1",
637        http::Version::HTTP_2 => "2",
638        http::Version::HTTP_3 => "3",
639        _ => "unknown",
640    }
641}
642
643/// Record response-phase span attributes that are only available after
644/// the upstream exchange.
645///
646/// Called from the `logging` hook to fill in `http.response.status_code`,
647/// `otel.status_code` and `error.type` (5xx only), `http.route` and the
648/// `otel.name` upgrade to `{method} {route}` (when a route matched),
649/// `upstream.address`, and `upstream.cluster` on the root request span,
650/// and response attributes on the upstream exchange span.
651fn record_response_span_attributes(session: &Session, ctx: &PingoraRequestCtx) {
652    if ctx.request_span.is_disabled() {
653        return;
654    }
655    let response = session.response_written();
656    let status = response.map(|resp| resp.status);
657    let method = session.req_header().method.as_str();
658    record_response_span_fields(status, method, response, ctx);
659}
660
661/// Record the response-phase fields once the status and method have been extracted.
662///
663/// Split from [`record_response_span_attributes`] so the recording logic is
664/// unit-testable without constructing a live Pingora session.
665fn record_response_span_fields(
666    status: Option<http::StatusCode>,
667    method: &str,
668    response: Option<&pingora_http::ResponseHeader>,
669    ctx: &PingoraRequestCtx,
670) {
671    if let Some(status) = status {
672        let code = status.as_u16();
673        if code > 0 {
674            ctx.request_span.record("http.response.status_code", code);
675        }
676        if status.is_server_error() {
677            ctx.request_span.record("otel.status_code", "ERROR");
678            // OTel semconv: error.type for an HTTP status is the numeric code
679            // as a string, not StatusCode's "{code} {reason}" Display form.
680            ctx.request_span.record("error.type", code.to_string().as_str());
681        }
682    }
683
684    if let Some(route) = &ctx.metrics_route {
685        ctx.request_span.record("http.route", route.as_ref());
686        ctx.request_span
687            .record("otel.name", format!("{method} {route}").as_str());
688    }
689
690    if let Some(upstream) = &ctx.upstream_for_retry {
691        ctx.request_span.record("upstream.address", upstream.address.as_ref());
692    }
693
694    if let Some(cluster) = &ctx.metrics_cluster {
695        ctx.request_span.record("upstream.cluster", cluster.as_ref());
696    }
697
698    record_upstream_exchange_span(ctx, response);
699}
700
701/// Record the upstream-exchange child span's response fields.
702fn record_upstream_exchange_span(ctx: &PingoraRequestCtx, response: Option<&pingora_http::ResponseHeader>) {
703    if ctx.upstream_exchange_span.is_disabled() {
704        return;
705    }
706    // Prefer the upstream's own status (captured before any response-phase
707    // rewrite); fall back to the written response when it was not captured.
708    if let Some(status) = ctx
709        .upstream_response_status
710        .or_else(|| response.map(|resp| resp.status.as_u16()))
711    {
712        ctx.upstream_exchange_span.record("http.response.status_code", status);
713    }
714    ctx.upstream_exchange_span
715        .record("http.response.body.size", ctx.response_body_bytes);
716}
717
718/// Build [`HttpServerOptions`] with h2c enabled.
719///
720/// [`HttpServerOptions`]: pingora_core::apps::HttpServerOptions
721fn h2c_server_options() -> HttpServerOptions {
722    let mut opts = HttpServerOptions::default();
723    opts.h2c = true;
724    opts
725}
726
727/// Build [`H2Options`] with limits to mitigate HPACK amplification attacks
728/// (CWE-409).
729///
730/// Without explicit limits the `h2` crate defaults allow unbounded header
731/// list sizes and concurrent streams, enabling a small compressed request
732/// to allocate hundreds of megabytes on the server.
733///
734/// [`H2Options`]: pingora_core::protocols::http::v2::server::H2Options
735fn h2_server_options() -> H2Options {
736    let mut opts = H2Options::new();
737    opts.max_header_list_size(65_536); // 64 KiB
738    opts.max_concurrent_streams(128);
739    opts
740}
741
742/// Accumulate `chunk.len()` into `accumulated_bytes` and return `true` when
743/// the total exceeds `max_bytes`. Returns `false` when the body is `None`.
744fn check_body_size_limit(body: Option<&Bytes>, accumulated_bytes: &mut u64, max_bytes: usize) -> bool {
745    if let Some(chunk) = body {
746        let chunk_len = chunk.len() as u64;
747        *accumulated_bytes += chunk_len;
748
749        let limit = max_bytes as u64;
750        return *accumulated_bytes > limit;
751    }
752    false
753}
754
755/// Push `chunk` into the stream buffer, creating it if absent. At end-of-stream
756/// the buffer is frozen into `body`. Returns `true` when the push overflows.
757fn accumulate_stream_buffer(
758    body: &mut Option<Bytes>,
759    body_buffer: &mut Option<BodyBuffer>,
760    end_of_stream: bool,
761    max_bytes: Option<usize>,
762) -> bool {
763    if let Some(chunk) = &*body {
764        let limit = max_bytes.unwrap_or(ABSOLUTE_MAX_BODY_BYTES);
765        let buf = body_buffer.get_or_insert_with(|| BodyBuffer::new(limit));
766
767        if buf.push(chunk.clone()).is_err() {
768            return true;
769        }
770    }
771
772    if end_of_stream {
773        tracing::trace!("stream buffer: freezing accumulated body before pipeline at EOS");
774        *body = body_buffer.take().map(BodyBuffer::freeze);
775    } else {
776        tracing::trace!("stream buffer: filters see the original chunk");
777    }
778    false
779}
780
781/// Suppress the body chunk while the stream buffer is still accumulating
782/// (i.e. `Continue`/`BodyDone` before release).
783#[expect(
784    clippy::fn_params_excessive_bools,
785    reason = "mirrors the caller's existing condition flags"
786)]
787fn suppress_stream_buffer_chunk(body: &mut Option<Bytes>, is_stream_buffer: bool, released: bool, end_of_stream: bool) {
788    if is_stream_buffer && !released && !end_of_stream {
789        *body = None;
790    }
791}
792
793/// Release the accumulated stream buffer on `FilterAction::Release`.
794fn release_stream_buffer(
795    body: &mut Option<Bytes>,
796    is_stream_buffer: bool,
797    released: &mut bool,
798    body_buffer: &mut Option<BodyBuffer>,
799    end_of_stream: bool,
800) {
801    if is_stream_buffer && !*released {
802        *released = true;
803        if !end_of_stream {
804            *body = body_buffer.take().map(BodyBuffer::freeze);
805        }
806    }
807}
808
809/// Shared fields extracted from an `HttpFilterContext` after body filter
810/// execution. Written back to `PingoraRequestCtx` via [`write_back`].
811///
812/// [`write_back`]: BodyFilterOutput::write_back
813struct BodyFilterOutput {
814    /// Cluster selected by the filter pipeline.
815    cluster: Option<Arc<str>>,
816    /// Upstream endpoint selected by the load balancer.
817    upstream: Option<Upstream>,
818    /// Type-safe request-scoped extension container.
819    extensions: RequestExtensions,
820    /// Durable per-request metadata that persists across phases.
821    filter_metadata: HashMap<String, String>,
822    /// Typed per-filter state keyed by stable filter invocation ID.
823    filter_state: HashMap<usize, Box<dyn std::any::Any + Send + Sync>>,
824    /// Per-filter execution tracking indices.
825    executed_filter_indices: Vec<bool>,
826    /// Per-filter body-done tracking indices.
827    body_done_indices: Vec<bool>,
828    /// Endpoints already attempted for this request (retry exclusion set).
829    attempted_endpoints: Vec<Arc<str>>,
830}
831
832impl BodyFilterOutput {
833    /// Move the shared fields out of the filter context, replacing each
834    /// with its `Default` value (zero-allocation no-ops for the types involved).
835    fn take_from(fctx: &mut HttpFilterContext<'_>) -> Self {
836        Self {
837            cluster: fctx.cluster.take(),
838            upstream: fctx.upstream.take(),
839            extensions: std::mem::take(&mut fctx.extensions),
840            filter_metadata: std::mem::take(&mut fctx.filter_metadata),
841            filter_state: std::mem::take(&mut fctx.filter_state),
842            executed_filter_indices: std::mem::take(&mut fctx.executed_filter_indices),
843            body_done_indices: std::mem::take(&mut fctx.body_done_indices),
844            attempted_endpoints: std::mem::take(&mut fctx.attempted_endpoints),
845        }
846    }
847
848    /// Write the shared fields back to the protocol context.
849    fn write_back(self, ctx: &mut PingoraRequestCtx) {
850        ctx.cluster = self.cluster;
851        ctx.upstream = self.upstream;
852        ctx.extensions = self.extensions;
853        ctx.filter_metadata = self.filter_metadata;
854        ctx.filter_state = self.filter_state;
855        ctx.cached_executed_filter_indices = self.executed_filter_indices;
856        ctx.cached_body_done_indices = self.body_done_indices;
857        ctx.attempted_endpoints = self.attempted_endpoints;
858    }
859}
860
861// -----------------------------------------------------------------------------
862// Tests
863// -----------------------------------------------------------------------------
864
865#[cfg(test)]
866#[expect(clippy::allow_attributes, reason = "blanket test suppressions")]
867#[allow(
868    clippy::unwrap_used,
869    clippy::expect_used,
870    clippy::indexing_slicing,
871    clippy::field_reassign_with_default,
872    clippy::too_many_lines,
873    clippy::cast_possible_truncation,
874    clippy::significant_drop_tightening,
875    reason = "tests"
876)]
877mod tests {
878    use praxis_core::connectivity::ConnectionOptions;
879
880    use super::*;
881
882    /// Maximum number of upstream connection retries for the legacy default policy.
883    const MAX_RETRIES: usize = praxis_core::config::DEFAULT_MAX_RETRIES as usize;
884
885    /// Default Pingora retry body buffer limit (64 `KiB`).
886    const RETRY_BODY_LIMIT: u64 = praxis_core::config::DEFAULT_RETRY_BODY_LIMIT_BYTES;
887
888    #[test]
889    fn first_failure_idempotent_sets_retry() {
890        let mut ctx = PingoraRequestCtx::default();
891        ctx.request_is_idempotent = true;
892        let e = handle_connect_failure(&mut ctx, make_error());
893        assert!(e.retry(), "first failure should set retry flag");
894        assert_eq!(ctx.retries, 1);
895    }
896
897    #[test]
898    fn large_body_skips_retry() {
899        let mut ctx = PingoraRequestCtx::default();
900        ctx.request_is_idempotent = true;
901        ctx.request_body_bytes = RETRY_BODY_LIMIT + 1;
902        let e = handle_connect_failure(&mut ctx, make_error());
903        assert!(!e.retry(), "should not retry when body exceeds retry buffer limit");
904        assert_eq!(ctx.retries, 0, "retry counter should not increment");
905    }
906
907    #[test]
908    fn mutated_body_exceeding_limit_skips_retry() {
909        let mut ctx = PingoraRequestCtx::default();
910        ctx.request_is_idempotent = true;
911        ctx.request_body_bytes = 1024;
912        ctx.mutated_request_body_len = Some((RETRY_BODY_LIMIT + 1) as usize);
913        let e = handle_connect_failure(&mut ctx, make_error());
914        assert!(
915            !e.retry(),
916            "should not retry when mutated body exceeds retry buffer limit"
917        );
918        assert_eq!(ctx.retries, 0);
919    }
920
921    #[test]
922    fn body_at_limit_allows_retry() {
923        let mut ctx = PingoraRequestCtx::default();
924        ctx.request_is_idempotent = true;
925        ctx.request_body_bytes = RETRY_BODY_LIMIT;
926        let e = handle_connect_failure(&mut ctx, make_error());
927        assert!(e.retry(), "body exactly at limit should allow retry");
928        assert_eq!(ctx.retries, 1);
929    }
930
931    #[test]
932    fn zero_body_allows_retry() {
933        let mut ctx = PingoraRequestCtx::default();
934        ctx.request_is_idempotent = true;
935        ctx.request_body_bytes = 0;
936        let e = handle_connect_failure(&mut ctx, make_error());
937        assert!(e.retry(), "zero-length body should allow retry");
938        assert_eq!(ctx.retries, 1);
939    }
940
941    #[test]
942    fn max_retries_exhausted_does_not_retry() {
943        let mut ctx = PingoraRequestCtx::default();
944        ctx.request_is_idempotent = true;
945        ctx.retries = MAX_RETRIES as u32;
946        let e = handle_connect_failure(&mut ctx, make_error());
947        assert!(!e.retry(), "should not retry after MAX_RETRIES");
948        assert_eq!(ctx.retries as usize, MAX_RETRIES);
949    }
950
951    #[test]
952    fn counter_increments_across_calls() {
953        let mut ctx = PingoraRequestCtx::default();
954        ctx.request_is_idempotent = true;
955        for expected in 1..=MAX_RETRIES {
956            let _result = handle_connect_failure(&mut ctx, make_error());
957            assert_eq!(ctx.retries as usize, expected);
958        }
959        let e = handle_connect_failure(&mut ctx, make_error());
960        assert!(!e.retry(), "should not retry after reaching MAX_RETRIES");
961        assert_eq!(ctx.retries as usize, MAX_RETRIES);
962    }
963
964    #[test]
965    fn non_idempotent_request_never_retries() {
966        let mut ctx = PingoraRequestCtx::default();
967        ctx.request_is_idempotent = false;
968        let e = handle_connect_failure(&mut ctx, make_error());
969        assert!(!e.retry(), "non-idempotent request should never retry");
970        assert_eq!(ctx.retries, 0);
971    }
972
973    #[test]
974    fn connect_failure_clears_upstream_connect_start() {
975        let mut ctx = PingoraRequestCtx::default();
976        ctx.upstream_connect_start = Some(std::time::Instant::now());
977        let _e = handle_connect_failure(&mut ctx, make_error());
978        assert!(
979            ctx.upstream_connect_start.is_none(),
980            "failed connect should consume upstream_connect_start for duration recording"
981        );
982    }
983
984    #[test]
985    fn response_503_retries_when_status5xx_enabled() {
986        let mut ctx = PingoraRequestCtx::default();
987        ctx.request_is_idempotent = true;
988        ctx.retry_policy = Some(Arc::new(praxis_core::config::RetryPolicy {
989            configured: true,
990            retriable_conditions: vec![praxis_core::config::RetriableCondition::Status5xx],
991            ..praxis_core::config::RetryPolicy::legacy_default()
992        }));
993        let e = maybe_retry_response(&mut ctx, 503).expect("503 should be retriable");
994        assert!(e.retry(), "503 under Status5xx should set retry");
995        assert_eq!(ctx.retries, 1);
996        assert!(ctx.reselect_on_retry);
997        assert!(ctx.pending_backoff.is_some());
998    }
999
1000    #[test]
1001    fn response_404_does_not_retry() {
1002        let mut ctx = PingoraRequestCtx::default();
1003        ctx.request_is_idempotent = true;
1004        ctx.retry_policy = Some(Arc::new(praxis_core::config::RetryPolicy {
1005            retriable_conditions: vec![praxis_core::config::RetriableCondition::Status5xx],
1006            ..praxis_core::config::RetryPolicy::legacy_default()
1007        }));
1008        assert!(
1009            maybe_retry_response(&mut ctx, 404).is_none(),
1010            "404 must never trigger status-based retry"
1011        );
1012        assert_eq!(ctx.retries, 0);
1013    }
1014
1015    #[test]
1016    fn response_502_does_not_retry_under_legacy_default() {
1017        let mut ctx = PingoraRequestCtx::default();
1018        ctx.request_is_idempotent = true;
1019        // Legacy default is connect_failure only — no Status5xx.
1020        assert!(
1021            maybe_retry_response(&mut ctx, 502).is_none(),
1022            "legacy default must forward 5xx without retry"
1023        );
1024    }
1025
1026    #[test]
1027    fn max_retries_zero_disables_connect_retry() {
1028        let mut ctx = PingoraRequestCtx::default();
1029        ctx.request_is_idempotent = true;
1030        ctx.retry_policy = Some(Arc::new(praxis_core::config::RetryPolicy {
1031            max_retries: Some(0),
1032            ..praxis_core::config::RetryPolicy::legacy_default()
1033        }));
1034        let e = handle_connect_failure(&mut ctx, make_error());
1035        assert!(!e.retry(), "max_retries: 0 must disable retries");
1036        assert_eq!(ctx.retries, 0);
1037    }
1038
1039    #[test]
1040    fn non_idempotent_clears_pingora_default_retry_flag() {
1041        let mut ctx = PingoraRequestCtx::default();
1042        ctx.request_is_idempotent = false;
1043        // Simulate Pingora marking the error retriable by default.
1044        let mut e = make_error();
1045        e.set_retry(true);
1046        let e = handle_connect_failure(&mut ctx, e);
1047        assert!(!e.retry(), "policy denial must clear Pingora's default retry flag");
1048        assert_eq!(ctx.retries, 0);
1049    }
1050
1051    #[tokio::test]
1052    async fn logging_cleanup_noop_when_response_phase_done() {
1053        let registry = praxis_filter::FilterRegistry::with_builtins();
1054        let pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
1055        let mut ctx = PingoraRequestCtx::default();
1056        ctx.response_phase_done = true;
1057        ctx.request_snapshot = Some(praxis_filter::Request {
1058            method: http::Method::GET,
1059            uri: "/".parse().unwrap(),
1060            headers: http::HeaderMap::new(),
1061        });
1062        logging_cleanup(&pipeline, &mut ctx).await;
1063    }
1064
1065    #[tokio::test]
1066    async fn logging_cleanup_noop_when_no_snapshot() {
1067        let registry = praxis_filter::FilterRegistry::with_builtins();
1068        let pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
1069        let mut ctx = PingoraRequestCtx::default();
1070        ctx.response_phase_done = false;
1071        ctx.request_snapshot = None;
1072        logging_cleanup(&pipeline, &mut ctx).await;
1073    }
1074
1075    #[tokio::test]
1076    async fn logging_cleanup_runs_response_pipeline_when_needed() {
1077        let registry = praxis_filter::FilterRegistry::with_builtins();
1078        let pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
1079        let mut ctx = PingoraRequestCtx::default();
1080        ctx.response_phase_done = false;
1081        ctx.cluster = Some(Arc::from("test-cluster"));
1082        ctx.request_snapshot = Some(praxis_filter::Request {
1083            method: http::Method::GET,
1084            uri: "/test".parse().unwrap(),
1085            headers: http::HeaderMap::new(),
1086        });
1087        logging_cleanup(&pipeline, &mut ctx).await;
1088        assert_eq!(
1089            ctx.cluster.as_deref(),
1090            Some("test-cluster"),
1091            "cluster must be restored so the fallback access record can attribute the failure"
1092        );
1093    }
1094
1095    #[tokio::test]
1096    async fn logging_cleanup_preserves_filter_metadata() {
1097        let registry = praxis_filter::FilterRegistry::with_builtins();
1098        let pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
1099        let mut ctx = PingoraRequestCtx::default();
1100        ctx.response_phase_done = false;
1101        ctx.filter_metadata
1102            .insert("json_rpc.method".to_owned(), "service/invoke".to_owned());
1103        ctx.request_snapshot = Some(praxis_filter::Request {
1104            method: http::Method::POST,
1105            uri: "/api".parse().unwrap(),
1106            headers: http::HeaderMap::new(),
1107        });
1108        logging_cleanup(&pipeline, &mut ctx).await;
1109        assert_eq!(
1110            ctx.filter_metadata.get("json_rpc.method").map(String::as_str),
1111            Some("service/invoke"),
1112            "filter_metadata should survive logging_cleanup"
1113        );
1114    }
1115
1116    #[tokio::test]
1117    async fn logging_cleanup_preserves_extensions() {
1118        let registry = praxis_filter::FilterRegistry::with_builtins();
1119        let pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
1120        let mut ctx = PingoraRequestCtx::default();
1121        ctx.response_phase_done = false;
1122        ctx.extensions.insert(42_u32);
1123        ctx.request_snapshot = Some(praxis_filter::Request {
1124            method: http::Method::POST,
1125            uri: "/test".parse().unwrap(),
1126            headers: http::HeaderMap::new(),
1127        });
1128        logging_cleanup(&pipeline, &mut ctx).await;
1129        assert_eq!(
1130            ctx.extensions.get::<u32>(),
1131            Some(&42),
1132            "extensions should survive logging_cleanup"
1133        );
1134    }
1135
1136    #[test]
1137    fn passive_health_error_is_failure() {
1138        let (pipeline, ctx) = make_passive_scenario(Some(3), Some(2));
1139        let error = make_error();
1140        record_passive_health(&pipeline, Some(&error), &ctx);
1141
1142        let registry = pipeline.health_registry().unwrap();
1143        let entry = registry.get("test-cluster").unwrap();
1144        assert!(
1145            entry.endpoints()[0].is_healthy(),
1146            "single failure should not yet mark unhealthy (threshold=3)"
1147        );
1148    }
1149
1150    #[test]
1151    fn passive_health_downstream_error_is_not_failure() {
1152        // A client-sourced (Downstream) error must not be charged
1153        // against a healthy endpoint even at an unhealthy-threshold of 1.
1154        let (pipeline, ctx) = make_passive_scenario(Some(1), Some(1));
1155        let error = make_error().into_down();
1156        record_passive_health(&pipeline, Some(&error), &ctx);
1157
1158        let registry = pipeline.health_registry().unwrap();
1159        let entry = registry.get("test-cluster").unwrap();
1160        assert!(
1161            entry.endpoints()[0].is_healthy(),
1162            "a downstream/client error must not mark the endpoint unhealthy"
1163        );
1164    }
1165
1166    #[test]
1167    fn passive_health_downstream_error_with_5xx_still_counts_as_failure() {
1168        // The downstream-error skip only applies when NO upstream status was
1169        // seen. A client (Downstream) error alongside a genuine upstream 5xx
1170        // must still be charged as a failure; this guards the
1171        // `&& ctx.upstream_response_status.is_none()` conjunct against removal.
1172        let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
1173        ctx.upstream_response_status = Some(503);
1174        let error = make_error().into_down();
1175        record_passive_health(&pipeline, Some(&error), &ctx);
1176
1177        let registry = pipeline.health_registry().unwrap();
1178        let entry = registry.get("test-cluster").unwrap();
1179        assert!(
1180            !entry.endpoints()[0].is_healthy(),
1181            "a downstream error with an upstream 503 must still mark the endpoint unhealthy"
1182        );
1183    }
1184
1185    #[test]
1186    fn passive_health_downstream_error_does_not_reset_failure_streak() {
1187        // A client disconnect between two genuine upstream failures must
1188        // not clear the endpoint's failure streak (which recording it as a
1189        // success would): the endpoint must still be ejected on the second
1190        // real failure.
1191        let (pipeline, ctx) = make_passive_scenario(Some(2), Some(1));
1192        let mut upstream_err = make_error();
1193        upstream_err.as_up();
1194        let downstream_err = make_error().into_down();
1195
1196        record_passive_health(&pipeline, Some(&upstream_err), &ctx);
1197        record_passive_health(&pipeline, Some(&downstream_err), &ctx);
1198        record_passive_health(&pipeline, Some(&upstream_err), &ctx);
1199
1200        let registry = pipeline.health_registry().unwrap();
1201        let entry = registry.get("test-cluster").unwrap();
1202        assert!(
1203            !entry.endpoints()[0].is_healthy(),
1204            "two upstream failures must eject the endpoint even with an interleaved client error"
1205        );
1206    }
1207
1208    #[test]
1209    fn passive_health_upstream_error_is_failure() {
1210        // An upstream-sourced error still counts as an endpoint failure.
1211        let (pipeline, ctx) = make_passive_scenario(Some(1), Some(1));
1212        let mut error = make_error();
1213        error.as_up();
1214        record_passive_health(&pipeline, Some(&error), &ctx);
1215
1216        let registry = pipeline.health_registry().unwrap();
1217        let entry = registry.get("test-cluster").unwrap();
1218        assert!(
1219            !entry.endpoints()[0].is_healthy(),
1220            "an upstream error at unhealthy-threshold 1 must mark the endpoint unhealthy"
1221        );
1222    }
1223
1224    #[test]
1225    fn passive_health_skips_observations_without_upstream_contact() {
1226        // A request that never contacted the upstream (a filter reject or a
1227        // proxy-generated terminal response after endpoint selection) must not
1228        // record a passive observation; recording a success would reset a real
1229        // failure streak. upstream_contacted is the "upstream contacted" signal.
1230        let (pipeline, mut ctx) = make_passive_scenario(Some(2), Some(1));
1231        let mut upstream_err = make_error();
1232        upstream_err.as_up();
1233
1234        // Contacted: a genuine upstream failure.
1235        record_passive_health(&pipeline, Some(&upstream_err), &ctx);
1236
1237        // Filter reject after selection: never contacted -> skipped.
1238        ctx.upstream_contacted = false;
1239        record_passive_health(&pipeline, None, &ctx);
1240
1241        // Proxy-generated terminal response: a status is set but the upstream
1242        // was never contacted -> also skipped.
1243        ctx.upstream_response_status = Some(200);
1244        record_passive_health(&pipeline, None, &ctx);
1245        ctx.upstream_response_status = None;
1246
1247        // Contacted again: the second genuine failure ejects the endpoint.
1248        ctx.upstream_contacted = true;
1249        record_passive_health(&pipeline, Some(&upstream_err), &ctx);
1250
1251        let registry = pipeline.health_registry().unwrap();
1252        let entry = registry.get("test-cluster").unwrap();
1253        assert!(
1254            !entry.endpoints()[0].is_healthy(),
1255            "observations without upstream contact must not reset the failure streak"
1256        );
1257    }
1258
1259    #[test]
1260    fn passive_health_records_connect_failure_after_reselect_clears_upstream() {
1261        // A retry decision clears upstream_for_retry to force reselection; if
1262        // no alternate endpoint exists the request ends with a connect error
1263        // and upstream_for_retry None. The sticky upstream_contacted signal
1264        // must still let that connect failure count toward ejection.
1265        let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
1266        ctx.upstream_for_retry = None;
1267        ctx.upstream_contacted = true;
1268        let mut error = make_error();
1269        error.as_up();
1270        record_passive_health(&pipeline, Some(&error), &ctx);
1271
1272        let registry = pipeline.health_registry().unwrap();
1273        let entry = registry.get("test-cluster").unwrap();
1274        assert!(
1275            !entry.endpoints()[0].is_healthy(),
1276            "a connect failure after reselection cleared upstream_for_retry must still count"
1277        );
1278    }
1279
1280    #[test]
1281    fn passive_health_status_500_is_failure() {
1282        let (pipeline, mut ctx) = make_passive_scenario(Some(3), Some(2));
1283        ctx.upstream_response_status = Some(500);
1284        record_passive_health(&pipeline, None, &ctx);
1285
1286        let registry = pipeline.health_registry().unwrap();
1287        let entry = registry.get("test-cluster").unwrap();
1288        assert!(
1289            entry.endpoints()[0].is_healthy(),
1290            "single 500 should not yet mark unhealthy (threshold=3)"
1291        );
1292    }
1293
1294    #[test]
1295    fn passive_health_status_below_500_is_success() {
1296        let (pipeline, mut ctx) = make_passive_scenario(Some(2), Some(1));
1297        ctx.upstream_response_status = Some(499);
1298        record_passive_health(&pipeline, None, &ctx);
1299
1300        let registry = pipeline.health_registry().unwrap();
1301        let entry = registry.get("test-cluster").unwrap();
1302        assert!(entry.endpoints()[0].is_healthy(), "status 499 should count as success");
1303    }
1304
1305    #[test]
1306    fn passive_unhealthy_threshold_transition() {
1307        let (pipeline, ctx) = make_passive_scenario(Some(2), Some(1));
1308        let error = make_error();
1309        record_passive_health(&pipeline, Some(&error), &ctx);
1310        record_passive_health(&pipeline, Some(&error), &ctx);
1311
1312        let registry = pipeline.health_registry().unwrap();
1313        let entry = registry.get("test-cluster").unwrap();
1314        assert!(
1315            !entry.endpoints()[0].is_healthy(),
1316            "2 consecutive failures should mark unhealthy (threshold=2)"
1317        );
1318    }
1319
1320    #[test]
1321    fn passive_healthy_threshold_recovery() {
1322        let (pipeline, ctx) = make_passive_scenario(Some(1), Some(2));
1323        let error = make_error();
1324        record_passive_health(&pipeline, Some(&error), &ctx);
1325
1326        let registry = pipeline.health_registry().unwrap();
1327        let entry = registry.get("test-cluster").unwrap();
1328        assert!(
1329            !entry.endpoints()[0].is_healthy(),
1330            "should be unhealthy after 1 failure"
1331        );
1332
1333        let ctx_ok = make_passive_ctx("test-cluster", 0, Some(200));
1334        record_passive_health(&pipeline, None, &ctx_ok);
1335        assert!(
1336            !entry.endpoints()[0].is_healthy(),
1337            "one success should not recover (threshold=2)"
1338        );
1339
1340        record_passive_health(&pipeline, None, &ctx_ok);
1341        assert!(
1342            entry.endpoints()[0].is_healthy(),
1343            "2 consecutive successes should recover (threshold=2)"
1344        );
1345    }
1346
1347    #[test]
1348    fn passive_health_no_thresholds_is_noop() {
1349        let (pipeline, ctx) = make_passive_scenario(None, None);
1350        let error = make_error();
1351        record_passive_health(&pipeline, Some(&error), &ctx);
1352
1353        let registry = pipeline.health_registry().unwrap();
1354        let entry = registry.get("test-cluster").unwrap();
1355        assert!(
1356            entry.endpoints()[0].is_healthy(),
1357            "no passive thresholds means failures are no-op"
1358        );
1359    }
1360
1361    #[test]
1362    fn passive_health_endpoint_index_out_of_bounds() {
1363        let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
1364        ctx.selected_endpoint_index = Some(999);
1365        let error = make_error();
1366        record_passive_health(&pipeline, Some(&error), &ctx);
1367
1368        let registry = pipeline.health_registry().unwrap();
1369        let entry = registry.get("test-cluster").unwrap();
1370        assert!(entry.endpoints()[0].is_healthy(), "out-of-bounds index should be no-op");
1371    }
1372
1373    #[test]
1374    fn passive_health_missing_cluster_is_noop() {
1375        let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
1376        ctx.cluster = None;
1377        ctx.metrics_cluster = None;
1378        let error = make_error();
1379        record_passive_health(&pipeline, Some(&error), &ctx);
1380    }
1381
1382    #[test]
1383    fn passive_health_falls_back_to_metrics_cluster() {
1384        let (pipeline, mut ctx) = make_passive_scenario(Some(2), Some(1));
1385        ctx.cluster = None;
1386        ctx.metrics_cluster = Some(Arc::from("test-cluster"));
1387        let error = make_error();
1388        record_passive_health(&pipeline, Some(&error), &ctx);
1389        record_passive_health(&pipeline, Some(&error), &ctx);
1390
1391        let registry = pipeline.health_registry().unwrap();
1392        let entry = registry.get("test-cluster").unwrap();
1393        assert!(
1394            !entry.endpoints()[0].is_healthy(),
1395            "fallback to metrics_cluster should still record passive health"
1396        );
1397    }
1398
1399    #[test]
1400    fn passive_health_missing_endpoint_index_is_noop() {
1401        let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
1402        ctx.selected_endpoint_index = None;
1403        let error = make_error();
1404        record_passive_health(&pipeline, Some(&error), &ctx);
1405    }
1406
1407    #[test]
1408    fn passive_health_missing_registry_is_noop() {
1409        let registry = praxis_filter::FilterRegistry::with_builtins();
1410        let pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
1411        let mut ctx = PingoraRequestCtx::default();
1412        ctx.cluster = Some(Arc::from("test-cluster"));
1413        ctx.selected_endpoint_index = Some(0);
1414        let error = make_error();
1415        record_passive_health(&pipeline, Some(&error), &ctx);
1416    }
1417
1418    #[test]
1419    fn passive_health_unknown_cluster_is_noop() {
1420        let (pipeline, mut ctx) = make_passive_scenario(Some(1), Some(1));
1421        ctx.cluster = Some(Arc::from("nonexistent"));
1422        let error = make_error();
1423        record_passive_health(&pipeline, Some(&error), &ctx);
1424    }
1425
1426    #[test]
1427    fn size_limit_none_body_returns_false() {
1428        let mut bytes = 0_u64;
1429        assert!(!check_body_size_limit(None, &mut bytes, 100));
1430        assert_eq!(bytes, 0, "accumulated bytes unchanged for None body");
1431    }
1432
1433    #[test]
1434    fn size_limit_within_limit() {
1435        let mut bytes = 0_u64;
1436        let body = Some(Bytes::from_static(b"hello"));
1437        assert!(!check_body_size_limit(body.as_ref(), &mut bytes, 10));
1438        assert_eq!(bytes, 5);
1439    }
1440
1441    #[test]
1442    fn size_limit_at_exact_limit() {
1443        let mut bytes = 0_u64;
1444        let body = Some(Bytes::from_static(b"exact"));
1445        assert!(!check_body_size_limit(body.as_ref(), &mut bytes, 5));
1446        assert_eq!(bytes, 5);
1447    }
1448
1449    #[test]
1450    fn size_limit_exceeds_limit() {
1451        let mut bytes = 0_u64;
1452        let body = Some(Bytes::from_static(b"toolong"));
1453        assert!(check_body_size_limit(body.as_ref(), &mut bytes, 3));
1454    }
1455
1456    #[test]
1457    fn size_limit_cumulative_overflow() {
1458        let mut bytes = 0_u64;
1459        let first = Some(Bytes::from_static(b"aaa"));
1460        assert!(!check_body_size_limit(first.as_ref(), &mut bytes, 5));
1461
1462        let second = Some(Bytes::from_static(b"bbb"));
1463        assert!(check_body_size_limit(second.as_ref(), &mut bytes, 5));
1464        assert_eq!(bytes, 6);
1465    }
1466
1467    #[test]
1468    fn stream_buffer_accumulates_chunks() {
1469        let mut body = Some(Bytes::from_static(b"hello "));
1470        let mut buf: Option<BodyBuffer> = None;
1471        assert!(!accumulate_stream_buffer(&mut body, &mut buf, false, Some(100)));
1472        assert!(buf.is_some());
1473
1474        body = Some(Bytes::from_static(b"world"));
1475        assert!(!accumulate_stream_buffer(&mut body, &mut buf, false, Some(100)));
1476
1477        let frozen = buf.take().unwrap().freeze();
1478        assert_eq!(frozen, Bytes::from_static(b"hello world"));
1479    }
1480
1481    #[test]
1482    fn stream_buffer_freezes_at_eos() {
1483        let mut body = Some(Bytes::from_static(b"data"));
1484        let mut buf: Option<BodyBuffer> = None;
1485        assert!(!accumulate_stream_buffer(&mut body, &mut buf, false, Some(100)));
1486
1487        body = Some(Bytes::from_static(b" end"));
1488        assert!(!accumulate_stream_buffer(&mut body, &mut buf, true, Some(100)));
1489        assert!(buf.is_none(), "buffer should be taken at EOS");
1490        assert_eq!(body.unwrap(), Bytes::from_static(b"data end"));
1491    }
1492
1493    #[test]
1494    fn stream_buffer_overflow() {
1495        let mut body = Some(Bytes::from_static(b"too long"));
1496        let mut buf: Option<BodyBuffer> = None;
1497        assert!(accumulate_stream_buffer(&mut body, &mut buf, false, Some(5)));
1498    }
1499
1500    #[test]
1501    fn stream_buffer_none_body() {
1502        let mut body: Option<Bytes> = None;
1503        let mut buf: Option<BodyBuffer> = None;
1504        assert!(!accumulate_stream_buffer(&mut body, &mut buf, false, Some(100)));
1505        assert!(buf.is_none());
1506    }
1507
1508    #[test]
1509    fn stream_buffer_uses_absolute_max_when_none() {
1510        let mut body = Some(Bytes::from_static(b"data"));
1511        let mut buf: Option<BodyBuffer> = None;
1512        assert!(!accumulate_stream_buffer(&mut body, &mut buf, false, None));
1513        assert!(buf.is_some(), "should create buffer with absolute max");
1514    }
1515
1516    #[test]
1517    fn suppress_clears_body_when_buffering() {
1518        let mut body = Some(Bytes::from_static(b"data"));
1519        suppress_stream_buffer_chunk(&mut body, true, false, false);
1520        assert!(body.is_none());
1521    }
1522
1523    #[test]
1524    fn suppress_noop_when_not_stream_buffer() {
1525        let mut body = Some(Bytes::from_static(b"data"));
1526        suppress_stream_buffer_chunk(&mut body, false, false, false);
1527        assert!(body.is_some());
1528    }
1529
1530    #[test]
1531    fn suppress_noop_when_released() {
1532        let mut body = Some(Bytes::from_static(b"data"));
1533        suppress_stream_buffer_chunk(&mut body, true, true, false);
1534        assert!(body.is_some());
1535    }
1536
1537    #[test]
1538    fn suppress_noop_at_eos() {
1539        let mut body = Some(Bytes::from_static(b"data"));
1540        suppress_stream_buffer_chunk(&mut body, true, false, true);
1541        assert!(body.is_some());
1542    }
1543
1544    #[test]
1545    fn release_sets_flag_and_flushes_buffer() {
1546        let mut body: Option<Bytes> = None;
1547        let mut released = false;
1548        let mut buf = Some(BodyBuffer::new(100));
1549        buf.as_mut().unwrap().push(Bytes::from_static(b"buffered")).unwrap();
1550
1551        release_stream_buffer(&mut body, true, &mut released, &mut buf, false);
1552        assert!(released);
1553        assert_eq!(body.unwrap(), Bytes::from_static(b"buffered"));
1554        assert!(buf.is_none());
1555    }
1556
1557    #[test]
1558    fn release_noop_when_already_released() {
1559        let mut body: Option<Bytes> = None;
1560        let mut released = true;
1561        let mut buf: Option<BodyBuffer> = None;
1562
1563        release_stream_buffer(&mut body, true, &mut released, &mut buf, false);
1564        assert!(body.is_none(), "body should be unchanged when already released");
1565    }
1566
1567    #[test]
1568    fn release_noop_when_not_stream_buffer() {
1569        let mut body: Option<Bytes> = None;
1570        let mut released = false;
1571        let mut buf: Option<BodyBuffer> = None;
1572
1573        release_stream_buffer(&mut body, false, &mut released, &mut buf, false);
1574        assert!(!released, "released flag should be unchanged for non-stream-buffer");
1575    }
1576
1577    #[test]
1578    fn release_at_eos_sets_flag_but_no_flush() {
1579        let mut body: Option<Bytes> = None;
1580        let mut released = false;
1581        let mut buf = Some(BodyBuffer::new(100));
1582        buf.as_mut().unwrap().push(Bytes::from_static(b"data")).unwrap();
1583
1584        release_stream_buffer(&mut body, true, &mut released, &mut buf, true);
1585        assert!(released);
1586        assert!(body.is_none(), "body should not be overwritten at EOS");
1587        assert!(buf.is_some(), "buffer should not be taken at EOS");
1588    }
1589
1590    #[test]
1591    fn write_back_transfers_fields() {
1592        let mut ctx = PingoraRequestCtx::default();
1593
1594        let mut extensions = RequestExtensions::new();
1595        extensions.insert(42_u32);
1596
1597        let state_val: Box<dyn std::any::Any + Send + Sync> = Box::new(99_i32);
1598        let filter_state = HashMap::from([(0_usize, state_val)]);
1599
1600        let output = BodyFilterOutput {
1601            cluster: Some(Arc::from("test-cluster")),
1602            upstream: Some(Upstream {
1603                address: Arc::from("10.0.0.1:80"),
1604                authority: None,
1605                connection: Arc::new(ConnectionOptions::default()),
1606                tls: None,
1607            }),
1608            extensions,
1609            attempted_endpoints: Vec::new(),
1610            filter_metadata: HashMap::from([("key".to_owned(), "val".to_owned())]),
1611            filter_state,
1612            executed_filter_indices: vec![true, false],
1613            body_done_indices: vec![false, true],
1614        };
1615        output.write_back(&mut ctx);
1616
1617        assert_eq!(ctx.cluster.as_deref(), Some("test-cluster"));
1618        assert!(ctx.upstream.is_some(), "upstream should transfer");
1619        assert_eq!(ctx.upstream.as_ref().unwrap().address.as_ref(), "10.0.0.1:80");
1620        assert_eq!(ctx.extensions.get::<u32>(), Some(&42));
1621        assert_eq!(ctx.filter_metadata.get("key").map(String::as_str), Some("val"));
1622        assert_eq!(ctx.filter_state.len(), 1, "filter_state should transfer");
1623        assert_eq!(
1624            ctx.filter_state.get(&0).and_then(|v| v.downcast_ref::<i32>()),
1625            Some(&99)
1626        );
1627        assert_eq!(ctx.cached_executed_filter_indices, vec![true, false]);
1628        assert_eq!(ctx.cached_body_done_indices, vec![false, true]);
1629    }
1630
1631    // -------------------------------------------------------------------------
1632    // Fallback Access Log
1633    // -------------------------------------------------------------------------
1634
1635    #[test]
1636    fn fallback_access_log_emits_for_incomplete_request() {
1637        let pipeline = access_log_pipeline();
1638        let mut ctx = make_fallback_ctx();
1639
1640        let events = capture_access_events(|| maybe_emit_fallback_access_log(&pipeline, 502, &mut ctx));
1641        assert_eq!(
1642            events.len(),
1643            1,
1644            "incomplete request must produce a fallback access record"
1645        );
1646    }
1647
1648    #[test]
1649    fn fallback_access_log_skips_completed_delivery() {
1650        let pipeline = access_log_pipeline();
1651        let mut ctx = make_fallback_ctx();
1652        ctx.response_delivery_complete = true;
1653
1654        let events = capture_access_events(|| maybe_emit_fallback_access_log(&pipeline, 200, &mut ctx));
1655        assert!(events.is_empty(), "completed delivery already logged via the filter");
1656    }
1657
1658    #[test]
1659    fn fallback_access_log_skips_upgraded_connections() {
1660        let pipeline = access_log_pipeline();
1661        let mut ctx = make_fallback_ctx();
1662        ctx.connection_upgraded = true;
1663
1664        let events = capture_access_events(|| maybe_emit_fallback_access_log(&pipeline, 101, &mut ctx));
1665        assert!(events.is_empty(), "upgraded connections have no body completion");
1666    }
1667
1668    #[test]
1669    fn fallback_access_log_skips_without_access_log_filter() {
1670        let registry = praxis_filter::FilterRegistry::with_builtins();
1671        let pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
1672        let mut ctx = make_fallback_ctx();
1673
1674        let events = capture_access_events(|| maybe_emit_fallback_access_log(&pipeline, 502, &mut ctx));
1675        assert!(events.is_empty(), "no access_log filter means no fallback record");
1676    }
1677
1678    #[test]
1679    fn fallback_access_log_honors_entry_conditions() {
1680        let registry = praxis_filter::FilterRegistry::with_builtins();
1681        let mut entries = vec![praxis_filter::FilterEntry {
1682            branch_chains: None,
1683            conditions: vec![serde_yaml::from_str("when:\n  path_prefix: /api\n").unwrap()],
1684            failure_mode: praxis_filter::FailureMode::default(),
1685            filter_type: "access_log".to_owned(),
1686            config: serde_yaml::Value::Null,
1687            name: None,
1688            response_conditions: vec![],
1689        }];
1690        let pipeline = FilterPipeline::build(&mut entries, &registry).unwrap();
1691
1692        let mut excluded = make_fallback_ctx();
1693        let events = capture_access_events(|| maybe_emit_fallback_access_log(&pipeline, 502, &mut excluded));
1694        assert!(
1695            events.is_empty(),
1696            "requests the operator scoped out must not gain fallback records"
1697        );
1698
1699        let mut included = make_fallback_ctx();
1700        if let Some(snapshot) = included.request_snapshot.as_mut() {
1701            snapshot.uri = "/api/users".parse().unwrap();
1702        }
1703        let events = capture_access_events(|| maybe_emit_fallback_access_log(&pipeline, 502, &mut included));
1704        assert_eq!(events.len(), 1, "in-scope incomplete requests still get a record");
1705    }
1706
1707    #[test]
1708    fn aborted_response_body_at_eos_is_not_marked_delivered() {
1709        // access_log declares read-only response body access, so the hook
1710        // reaches the SizeLimit check instead of early-returning.
1711        let pipeline = access_log_pipeline();
1712        let mut ctx = make_fallback_ctx();
1713        ctx.response_body_mode = BodyMode::SizeLimit { max_bytes: 4 };
1714        let mut body = Some(Bytes::from_static(b"exceeds the limit"));
1715
1716        let result = response_body_filter::execute(&pipeline, &mut body, true, &mut ctx);
1717        assert!(result.is_err(), "over-limit body must abort");
1718        assert!(
1719            !ctx.response_delivery_complete,
1720            "a response aborted at end-of-stream was not delivered; the fallback record must fire"
1721        );
1722    }
1723
1724    #[test]
1725    fn response_body_eos_marks_delivery_complete() {
1726        let registry = praxis_filter::FilterRegistry::with_builtins();
1727        let pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
1728        let mut ctx = PingoraRequestCtx::default();
1729        let mut body: Option<Bytes> = None;
1730
1731        let _timeout = response_body_filter::execute(&pipeline, &mut body, false, &mut ctx).unwrap();
1732        assert!(
1733            !ctx.response_delivery_complete,
1734            "mid-stream chunks must not mark delivery complete"
1735        );
1736
1737        let _timeout = response_body_filter::execute(&pipeline, &mut body, true, &mut ctx).unwrap();
1738        assert!(ctx.response_delivery_complete, "end-of-stream marks delivery complete");
1739    }
1740
1741    // -------------------------------------------------------------------------
1742    // Span Attribute Helpers
1743    // -------------------------------------------------------------------------
1744
1745    #[test]
1746    fn http_version_label_http_09() {
1747        assert_eq!(
1748            http_version_label(http::Version::HTTP_09),
1749            "0.9",
1750            "HTTP/0.9 should map to '0.9'"
1751        );
1752    }
1753
1754    #[test]
1755    fn http_version_label_http_10() {
1756        assert_eq!(
1757            http_version_label(http::Version::HTTP_10),
1758            "1.0",
1759            "HTTP/1.0 should map to '1.0'"
1760        );
1761    }
1762
1763    #[test]
1764    fn http_version_label_http_11() {
1765        assert_eq!(
1766            http_version_label(http::Version::HTTP_11),
1767            "1.1",
1768            "HTTP/1.1 should map to '1.1'"
1769        );
1770    }
1771
1772    #[test]
1773    fn http_version_label_http_2() {
1774        assert_eq!(
1775            http_version_label(http::Version::HTTP_2),
1776            "2",
1777            "HTTP/2 should map to '2'"
1778        );
1779    }
1780
1781    #[test]
1782    fn http_version_label_http_3() {
1783        assert_eq!(
1784            http_version_label(http::Version::HTTP_3),
1785            "3",
1786            "HTTP/3 should map to '3'"
1787        );
1788    }
1789
1790    #[test]
1791    fn record_response_span_attributes_noop_for_disabled_span() {
1792        let ctx = PingoraRequestCtx::default();
1793        assert!(ctx.request_span.is_disabled(), "default span should be disabled");
1794    }
1795
1796    /// Layer that captures every `Span::record` call as `(field, value)` pairs.
1797    #[derive(Clone, Default)]
1798    struct RecordCapture(Arc<std::sync::Mutex<Vec<(String, String)>>>);
1799
1800    impl<S> tracing_subscriber::Layer<S> for RecordCapture
1801    where
1802        S: tracing::Subscriber + for<'a> tracing_subscriber::registry::LookupSpan<'a>,
1803    {
1804        fn on_record(
1805            &self,
1806            _id: &tracing::span::Id,
1807            values: &tracing::span::Record<'_>,
1808            _ctx: tracing_subscriber::layer::Context<'_, S>,
1809        ) {
1810            struct Visitor<'a>(&'a mut Vec<(String, String)>);
1811            impl tracing::field::Visit for Visitor<'_> {
1812                fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
1813                    self.0.push((field.name().to_owned(), format!("{value:?}")));
1814                }
1815            }
1816            let mut captured = self.0.lock().expect("capture lock");
1817            values.record(&mut Visitor(&mut captured));
1818        }
1819    }
1820
1821    #[test]
1822    fn record_response_span_fields_records_status_upstream_and_cluster() {
1823        use tracing_subscriber::layer::SubscriberExt as _;
1824
1825        let capture = RecordCapture::default();
1826        let subscriber = tracing_subscriber::registry().with(capture.clone());
1827        let _guard = tracing::subscriber::set_default(subscriber);
1828
1829        let mut ctx = PingoraRequestCtx::default();
1830        ctx.metrics_cluster = Some(Arc::from("api-cluster"));
1831        ctx.upstream_for_retry = Some(Upstream {
1832            address: Arc::from("10.0.0.1:80"),
1833            authority: None,
1834            connection: Arc::new(ConnectionOptions::default()),
1835            tls: None,
1836        });
1837        ctx.request_span = tracing::info_span!(
1838            "test_span",
1839            "http.response.status_code" = tracing::field::Empty,
1840            "otel.status_code" = tracing::field::Empty,
1841            "upstream.address" = tracing::field::Empty,
1842            "upstream.cluster" = tracing::field::Empty,
1843        );
1844
1845        record_response_span_fields(Some(http::StatusCode::SERVICE_UNAVAILABLE), "GET", None, &ctx);
1846
1847        let captured = capture.0.lock().expect("capture lock");
1848        let get = |name: &str| {
1849            let value = captured.iter().find(|(f, _)| f == name).map(|(_, v)| v.clone());
1850            assert!(value.is_some(), "field {name} not recorded; got {captured:?}");
1851            value.unwrap_or_default()
1852        };
1853        assert_eq!(get("http.response.status_code"), "503");
1854        assert_eq!(get("otel.status_code"), "\"ERROR\"", "5xx must set otel error status");
1855        assert_eq!(get("upstream.address"), "\"10.0.0.1:80\"");
1856        assert_eq!(get("upstream.cluster"), "\"api-cluster\"");
1857    }
1858
1859    #[test]
1860    fn record_response_span_fields_success_has_no_error_status() {
1861        use tracing_subscriber::layer::SubscriberExt as _;
1862
1863        let capture = RecordCapture::default();
1864        let subscriber = tracing_subscriber::registry().with(capture.clone());
1865        let _guard = tracing::subscriber::set_default(subscriber);
1866
1867        let mut ctx = PingoraRequestCtx::default();
1868        ctx.request_span = tracing::info_span!(
1869            "test_span",
1870            "http.response.status_code" = tracing::field::Empty,
1871            "otel.status_code" = tracing::field::Empty,
1872        );
1873
1874        record_response_span_fields(Some(http::StatusCode::OK), "GET", None, &ctx);
1875
1876        let captured = capture.0.lock().expect("capture lock");
1877        assert!(
1878            captured
1879                .iter()
1880                .any(|(f, v)| f == "http.response.status_code" && v == "200"),
1881            "status should be recorded: {captured:?}"
1882        );
1883        assert!(
1884            !captured.iter().any(|(f, _)| f == "otel.status_code"),
1885            "2xx must not set otel error status: {captured:?}"
1886        );
1887    }
1888
1889    #[test]
1890    fn record_response_span_attributes_records_exchange_span_fields() {
1891        let mut ctx = PingoraRequestCtx::default();
1892        ctx.request_span = tracing::info_span!(
1893            "test_request",
1894            "http.response.status_code" = tracing::field::Empty,
1895            "server.address" = tracing::field::Empty,
1896            "upstream.cluster" = tracing::field::Empty,
1897        );
1898        ctx.upstream_exchange_span = tracing::info_span!(
1899            parent: &ctx.request_span,
1900            "upstream_exchange",
1901            "http.response.status_code" = tracing::field::Empty,
1902            "http.response.body.size" = tracing::field::Empty,
1903        );
1904        ctx.response_body_bytes = 4096;
1905
1906        // Verify recording on exchange span does not panic.
1907        ctx.upstream_exchange_span.record("http.response.status_code", 200_u16);
1908        ctx.upstream_exchange_span.record("http.response.body.size", 4096_u64);
1909    }
1910
1911    #[test]
1912    fn record_response_span_attributes_skips_exchange_when_disabled() {
1913        let mut ctx = PingoraRequestCtx::default();
1914        ctx.request_span = tracing::info_span!(
1915            "test_request",
1916            "http.response.status_code" = tracing::field::Empty,
1917            "server.address" = tracing::field::Empty,
1918            "upstream.cluster" = tracing::field::Empty,
1919        );
1920        // Exchange span remains disabled (default).
1921        assert!(
1922            ctx.upstream_exchange_span.is_disabled(),
1923            "exchange span should be disabled by default"
1924        );
1925        // Should not panic when exchange span is disabled.
1926    }
1927
1928    // -------------------------------------------------------------------------
1929    // Span Event Tests
1930    // -------------------------------------------------------------------------
1931
1932    #[test]
1933    fn retry_with_upstream_address_sets_retry_flag() {
1934        let mut ctx = PingoraRequestCtx::default();
1935        ctx.request_is_idempotent = true;
1936        ctx.upstream_for_retry = Some(Upstream {
1937            address: Arc::from("10.0.0.1:8080"),
1938            connection: Arc::new(ConnectionOptions::default()),
1939            tls: None,
1940            authority: None,
1941        });
1942        let e = handle_connect_failure(&mut ctx, make_error());
1943        assert!(e.retry(), "should retry with upstream address present");
1944        assert_eq!(ctx.retries, 1, "retry counter should increment to 1");
1945    }
1946
1947    #[test]
1948    fn retry_without_upstream_address_uses_fallback() {
1949        let mut ctx = PingoraRequestCtx::default();
1950        ctx.request_is_idempotent = true;
1951        ctx.upstream_for_retry = None;
1952        let e = handle_connect_failure(&mut ctx, make_error());
1953        assert!(
1954            e.retry(),
1955            "should retry even when upstream_for_retry is None (address defaults to unknown)"
1956        );
1957        assert_eq!(ctx.retries, 1, "retry counter should increment to 1");
1958    }
1959
1960    #[test]
1961    fn retry_exhausted_with_upstream_address_does_not_retry() {
1962        let mut ctx = PingoraRequestCtx::default();
1963        ctx.request_is_idempotent = true;
1964        ctx.retries = MAX_RETRIES as u32;
1965        ctx.upstream_for_retry = Some(Upstream {
1966            address: Arc::from("10.0.0.2:443"),
1967            connection: Arc::new(ConnectionOptions::default()),
1968            tls: None,
1969            authority: None,
1970        });
1971        let e = handle_connect_failure(&mut ctx, make_error());
1972        assert!(
1973            !e.retry(),
1974            "should not retry after MAX_RETRIES even with upstream address"
1975        );
1976    }
1977
1978    #[test]
1979    fn large_body_skip_with_upstream_address() {
1980        let mut ctx = PingoraRequestCtx::default();
1981        ctx.request_is_idempotent = true;
1982        ctx.request_body_bytes = RETRY_BODY_LIMIT + 1;
1983        ctx.upstream_for_retry = Some(Upstream {
1984            address: Arc::from("10.0.0.3:8080"),
1985            connection: Arc::new(ConnectionOptions::default()),
1986            tls: None,
1987            authority: None,
1988        });
1989        let e = handle_connect_failure(&mut ctx, make_error());
1990        assert!(!e.retry(), "should not retry large body even with upstream address");
1991        assert_eq!(ctx.retries, 0, "retry counter should not increment");
1992    }
1993
1994    // -------------------------------------------------------------------------
1995    // Test Utilities
1996    // -------------------------------------------------------------------------
1997
1998    /// Create a connect error for tests.
1999    fn make_error() -> Box<pingora_core::Error> {
2000        pingora_core::Error::explain(pingora_core::ErrorType::ConnectError, "test connect failure")
2001    }
2002
2003    /// Build a pipeline containing an `access_log` filter.
2004    fn access_log_pipeline() -> FilterPipeline {
2005        let registry = praxis_filter::FilterRegistry::with_builtins();
2006        let mut entries = vec![praxis_filter::FilterEntry {
2007            branch_chains: None,
2008            conditions: vec![],
2009            failure_mode: praxis_filter::FailureMode::default(),
2010            filter_type: "access_log".to_owned(),
2011            config: serde_yaml::Value::Null,
2012            name: None,
2013            response_conditions: vec![],
2014        }];
2015        FilterPipeline::build(&mut entries, &registry).unwrap()
2016    }
2017
2018    /// Build a context with a request snapshot for fallback logging tests.
2019    fn make_fallback_ctx() -> PingoraRequestCtx {
2020        let mut ctx = PingoraRequestCtx::default();
2021        ctx.request_snapshot = Some(praxis_filter::Request {
2022            method: http::Method::GET,
2023            uri: "/incomplete".parse().unwrap(),
2024            headers: http::HeaderMap::new(),
2025        });
2026        ctx
2027    }
2028
2029    /// Capture `access` info events emitted while running `f`.
2030    fn capture_access_events<F: FnOnce()>(f: F) -> Vec<String> {
2031        use tracing_subscriber::layer::SubscriberExt as _;
2032
2033        let messages = Arc::new(std::sync::Mutex::new(Vec::<String>::new()));
2034        let capture = AccessCapture(Arc::clone(&messages));
2035        let subscriber = tracing_subscriber::registry().with(capture);
2036        tracing::subscriber::with_default(subscriber, f);
2037        let mut guard = messages.lock().unwrap();
2038        std::mem::take(&mut *guard)
2039    }
2040
2041    /// Layer capturing `access` records for assertions.
2042    struct AccessCapture(Arc<std::sync::Mutex<Vec<String>>>);
2043
2044    impl<S: tracing::Subscriber> tracing_subscriber::Layer<S> for AccessCapture {
2045        fn on_event(&self, event: &tracing::Event<'_>, _ctx: tracing_subscriber::layer::Context<'_, S>) {
2046            let mut visitor = AccessMessageVisitor(String::new());
2047            event.record(&mut visitor);
2048            if visitor.0.contains("access") {
2049                self.0.lock().unwrap().push(visitor.0);
2050            }
2051        }
2052    }
2053
2054    /// Visitor extracting the `message` field from an event.
2055    struct AccessMessageVisitor(String);
2056
2057    impl tracing::field::Visit for AccessMessageVisitor {
2058        fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
2059            if field.name() == "message" {
2060                self.0 = format!("{value:?}");
2061            }
2062        }
2063    }
2064
2065    /// Build a [`PingoraRequestCtx`] for passive health testing.
2066    fn make_passive_ctx(cluster: &str, endpoint_idx: usize, status: Option<u16>) -> PingoraRequestCtx {
2067        let mut ctx = PingoraRequestCtx::default();
2068        ctx.cluster = Some(Arc::from(cluster));
2069        ctx.selected_endpoint_index = Some(endpoint_idx);
2070        ctx.upstream_response_status = status;
2071        // A passive-health observation implies the upstream was contacted.
2072        ctx.upstream_contacted = true;
2073        ctx
2074    }
2075
2076    /// Build a pipeline with a health registry and a matching context
2077    /// for passive health testing.
2078    fn make_passive_scenario(
2079        passive_unhealthy: Option<u32>,
2080        passive_healthy: Option<u32>,
2081    ) -> (FilterPipeline, PingoraRequestCtx) {
2082        use praxis_core::health::{ClusterHealthEntry, EndpointHealth};
2083
2084        let entry = ClusterHealthEntry::new(
2085            vec![EndpointHealth::new()],
2086            vec![Arc::from("10.0.0.1:80")],
2087            passive_unhealthy,
2088            passive_healthy,
2089        );
2090        let mut map = HashMap::new();
2091        map.insert(Arc::from("test-cluster"), Arc::new(entry));
2092        let health_registry = Arc::new(map);
2093
2094        let registry = praxis_filter::FilterRegistry::with_builtins();
2095        let mut pipeline = FilterPipeline::build(&mut [], &registry).unwrap();
2096        pipeline.set_health_registry(health_registry);
2097
2098        let ctx = make_passive_ctx("test-cluster", 0, None);
2099
2100        (pipeline, ctx)
2101    }
2102}