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