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