use std::sync::Arc;
use praxis_core::config::RetryPolicy;
use tracing::{debug, warn};
use super::{metrics, retry};
use crate::http::pingora::context::PingoraRequestCtx;
pub(super) fn legacy_default_policy() -> Arc<RetryPolicy> {
static LEGACY_DEFAULT: std::sync::LazyLock<Arc<RetryPolicy>> =
std::sync::LazyLock::new(|| Arc::new(RetryPolicy::legacy_default()));
Arc::clone(&LEGACY_DEFAULT)
}
#[expect(clippy::too_many_lines, reason = "sequential guard checks")]
pub(super) fn handle_connect_failure(
ctx: &mut PingoraRequestCtx,
e: Box<pingora_core::Error>,
) -> Box<pingora_core::Error> {
let cluster = ctx.metrics_cluster_shared.clone().unwrap_or_else(metrics::cluster_none);
if let Some(start) = ctx.upstream_connect_start.take() {
metrics::record_upstream_connect_duration(cluster.clone(), start.elapsed().as_secs_f64());
}
metrics::record_upstream_connect_failure(cluster.clone());
let policy = ctx.retry_policy.clone().unwrap_or_else(legacy_default_policy);
let outcome = retry::classify_error(&e);
let decision = retry::should_retry(ctx, &policy, outcome, ctx.cluster_retry_state.as_deref());
match decision {
retry::RetryDecision::Retry { backoff } => {
ctx.retries += 1;
ctx.pending_backoff = Some(backoff);
ctx.reselect_on_retry = policy.configured;
if let Some(upstream) = ctx.upstream_for_retry.as_ref() {
let addr = Arc::clone(&upstream.address);
if !ctx.attempted_endpoints.iter().any(|e| e.as_ref() == addr.as_ref()) {
ctx.attempted_endpoints.push(addr);
}
}
if policy.configured {
if let Some(upstream) = ctx.upstream_for_retry.as_ref()
&& let Some(reselector) = ctx.endpoint_reselector.as_ref()
{
reselector.release(&upstream.address);
}
ctx.upstream_for_retry = None;
}
let upstream_address = ctx
.upstream_for_retry
.as_ref()
.map_or("unknown", |u| u.address.as_ref());
debug!(
retries = ctx.retries,
max = policy.effective_max_retries(),
?backoff,
upstream_address,
"retrying after connect failure"
);
let mut e = e;
e.set_retry(true);
e
},
retry::RetryDecision::DoNotRetry => {
if ctx.retries > 0 {
warn!(
retries = ctx.retries,
max = policy.effective_max_retries(),
upstream_address = ctx
.upstream_for_retry
.as_ref()
.map_or("unknown", |u| u.address.as_ref()),
"retry limit exhausted"
);
}
record_retry_exhausted_if_attempted(ctx, cluster);
let mut e = e;
e.set_retry(false);
e
},
}
}
#[expect(clippy::too_many_lines, reason = "sequential guard checks")]
pub(super) fn maybe_retry_response(ctx: &mut PingoraRequestCtx, status: u16) -> Option<Box<pingora_core::Error>> {
let policy = ctx.retry_policy.clone().unwrap_or_else(legacy_default_policy);
let outcome = retry::RetryOutcome::StatusCode(status);
let decision = retry::should_retry(ctx, &policy, outcome, ctx.cluster_retry_state.as_deref());
match decision {
retry::RetryDecision::Retry { backoff } => {
ctx.retries += 1;
ctx.pending_backoff = Some(backoff);
ctx.reselect_on_retry = policy.configured;
if let Some(upstream) = ctx.upstream_for_retry.as_ref() {
let addr = Arc::clone(&upstream.address);
if !ctx.attempted_endpoints.iter().any(|e| e.as_ref() == addr.as_ref()) {
ctx.attempted_endpoints.push(addr);
}
}
if policy.configured {
if let Some(upstream) = ctx.upstream_for_retry.as_ref()
&& let Some(reselector) = ctx.endpoint_reselector.as_ref()
{
reselector.release(&upstream.address);
}
ctx.upstream_for_retry = None;
}
debug!(
status,
retries = ctx.retries,
max = policy.effective_max_retries(),
?backoff,
"retrying after retriable response status"
);
let mut e =
pingora_core::Error::explain(pingora_core::ErrorType::HTTPStatus(status), "retriable upstream status");
e.set_retry(true);
Some(e)
},
retry::RetryDecision::DoNotRetry => None,
}
}
pub(super) fn release_retry_state(ctx: &mut PingoraRequestCtx) {
if !ctx.cluster_retry_state_released
&& let Some(state) = ctx.cluster_retry_state.take()
{
state.leave();
ctx.cluster_retry_state_released = true;
}
}
fn record_retry_exhausted_if_attempted(ctx: &PingoraRequestCtx, cluster: ::metrics::SharedString) {
if ctx.retries > 0 {
metrics::record_upstream_retry(cluster, metrics::RETRY_RESULT_EXHAUSTED);
}
}