From 18b23114f7ae455c676c80c6504a6a381f5a9f83 Mon Sep 17 00:00:00 2001 From: Chang Su Date: Wed, 18 Mar 2026 15:23:25 -0700 Subject: [PATCH 1/2] feat(core): add execute_with_resilience as the primary resilience API Add execute_with_resilience() function that encapsulates retry loop, circuit breaker checks, and outcome recording using the worker's own resolved config. This is the single entry point routers will use instead of calling RetryExecutor directly. - Checks circuit breaker before first attempt (when CB enabled) - Skips retry loop when retries are disabled (single attempt) - Uses worker's is_retryable() for retry decisions - Records backoff metrics via Metrics::record_worker_retry_backoff - 5 new async tests covering success, CB open, retries disabled, retry on retryable status, and CB disabled bypass Signed-off-by: Chang Su --- model_gateway/src/core/mod.rs | 4 +- model_gateway/src/core/resilience.rs | 186 ++++++++++++++++++++++++++- 2 files changed, 188 insertions(+), 2 deletions(-) diff --git a/model_gateway/src/core/mod.rs b/model_gateway/src/core/mod.rs index 9e6f4262c1..dba5b48fa7 100644 --- a/model_gateway/src/core/mod.rs +++ b/model_gateway/src/core/mod.rs @@ -39,7 +39,9 @@ pub use openai_protocol::{ model_type::{Endpoint, ModelType}, worker::{ProviderType, WorkerGroupKey}, }; -pub use resilience::{resolve_resilience, ResolvedResilience, DEFAULT_RETRYABLE_STATUS_CODES}; +pub use resilience::{ + execute_with_resilience, resolve_resilience, ResolvedResilience, DEFAULT_RETRYABLE_STATUS_CODES, +}; pub use retry::{is_retryable_status, RetryExecutor}; pub use worker::{ AttachedBody, BasicWorker, ConnectionMode, RuntimeType, Worker, WorkerLoadGuard, WorkerType, diff --git a/model_gateway/src/core/resilience.rs b/model_gateway/src/core/resilience.rs index c310669454..9eb9c0cc47 100644 --- a/model_gateway/src/core/resilience.rs +++ b/model_gateway/src/core/resilience.rs @@ -5,9 +5,15 @@ use std::{collections::HashSet, time::Duration}; +use axum::{http::StatusCode, response::IntoResponse}; use openai_protocol::worker::ResilienceUpdate; +use tracing::debug; -use crate::{config::types::RetryConfig, core::circuit_breaker::CircuitBreakerConfig}; +use crate::{ + config::types::RetryConfig, + core::{circuit_breaker::CircuitBreakerConfig, retry::RetryExecutor, worker::Worker}, + observability::metrics::Metrics, +}; /// Default retryable HTTP status codes. pub const DEFAULT_RETRYABLE_STATUS_CODES: &[u16] = &[408, 429, 500, 502, 503, 504]; @@ -108,9 +114,187 @@ pub fn resolve_resilience( (resolved, cb_config) } +/// Execute a request with the worker's retry and circuit breaker config. +/// +/// This is the primary resilience API. Routers call this instead of +/// using `RetryExecutor` directly. +/// +/// Circuit breaker outcomes are recorded automatically. +/// Retry decisions use the worker's `is_retryable()` predicate. +pub async fn execute_with_resilience( + worker: &(dyn Worker + '_), + operation: F, +) -> axum::response::Response +where + F: FnMut(u32) -> Fut + Send, + Fut: std::future::Future + Send, +{ + let resilience = worker.resilience(); + let worker_url = worker.url(); + + // Check circuit breaker before first attempt + if resilience.circuit_breaker_enabled && !worker.circuit_breaker().can_execute() { + debug!( + worker_url = worker_url, + "Circuit breaker open, rejecting request" + ); + return ( + StatusCode::SERVICE_UNAVAILABLE, + format!("Circuit breaker open for worker {worker_url}"), + ) + .into_response(); + } + + if !resilience.retry_enabled { + // Single attempt, no retries + let response = execute_single(worker, operation).await; + if resilience.circuit_breaker_enabled { + let success = !worker.is_retryable(&response); + worker.circuit_breaker().record_outcome(success); + } + return response; + } + + // Retry loop — delegate to RetryExecutor with worker-aware hooks + RetryExecutor::execute_response_with_retry( + &resilience.retry, + operation, + |res, _attempt| worker.is_retryable(res), + |delay, attempt| { + Metrics::record_worker_retry_backoff(attempt, delay); + debug!( + worker_url = worker_url, + attempt = attempt, + delay_ms = delay.as_millis() as u64, + "Worker retry backoff" + ); + }, + || { + debug!(worker_url = worker_url, "Worker retries exhausted"); + }, + ) + .await +} + +/// Execute a single attempt of an operation (no retry). +async fn execute_single( + _worker: &(dyn Worker + '_), + mut operation: F, +) -> axum::response::Response +where + F: FnMut(u32) -> Fut + Send, + Fut: std::future::Future + Send, +{ + operation(0).await +} + #[cfg(test)] mod tests { + use std::sync::{ + atomic::{AtomicU32, Ordering}, + Arc, + }; + + use axum::{http::StatusCode, response::IntoResponse}; + use super::*; + use crate::core::worker_builder::BasicWorkerBuilder; + + #[tokio::test] + async fn test_execute_with_resilience_success() { + let worker = BasicWorkerBuilder::new("http://test:8080").build(); + let response = execute_with_resilience(&worker, |_attempt| async { + (StatusCode::OK, "ok").into_response() + }) + .await; + assert_eq!(response.status(), StatusCode::OK); + } + + #[tokio::test] + async fn test_execute_with_resilience_circuit_open() { + let worker = BasicWorkerBuilder::new("http://test:8080").build(); + worker.circuit_breaker().force_open(); + let response = execute_with_resilience(&worker, |_attempt| async { + (StatusCode::OK, "ok").into_response() + }) + .await; + assert_eq!(response.status(), StatusCode::SERVICE_UNAVAILABLE); + } + + #[tokio::test] + async fn test_execute_with_resilience_retries_disabled() { + let resolved = ResolvedResilience { + retry_enabled: false, + ..Default::default() + }; + let worker = BasicWorkerBuilder::new("http://test:8080") + .resilience(resolved) + .build(); + + let call_count = Arc::new(AtomicU32::new(0)); + let cc = call_count.clone(); + let response = execute_with_resilience(&worker, move |_attempt| { + cc.fetch_add(1, Ordering::Relaxed); + async { (StatusCode::SERVICE_UNAVAILABLE, "fail").into_response() } + }) + .await; + + assert_eq!(response.status(), StatusCode::SERVICE_UNAVAILABLE); + assert_eq!(call_count.load(Ordering::Relaxed), 1); // No retries + } + + #[tokio::test] + async fn test_execute_with_resilience_retries_on_retryable_status() { + let resolved = ResolvedResilience { + retry: RetryConfig { + max_retries: 3, + initial_backoff_ms: 1, + max_backoff_ms: 2, + backoff_multiplier: 1.0, + jitter_factor: 0.0, + }, + ..Default::default() + }; + let worker = BasicWorkerBuilder::new("http://test:8080") + .resilience(resolved) + .build(); + + let call_count = Arc::new(AtomicU32::new(0)); + let cc = call_count.clone(); + let response = execute_with_resilience(&worker, move |_attempt| { + let count = cc.fetch_add(1, Ordering::Relaxed); + async move { + if count < 2 { + (StatusCode::SERVICE_UNAVAILABLE, "fail").into_response() + } else { + (StatusCode::OK, "ok").into_response() + } + } + }) + .await; + + assert_eq!(response.status(), StatusCode::OK); + assert_eq!(call_count.load(Ordering::Relaxed), 3); + } + + #[tokio::test] + async fn test_execute_with_resilience_cb_disabled_ignores_open() { + let resolved = ResolvedResilience { + circuit_breaker_enabled: false, + ..Default::default() + }; + let worker = BasicWorkerBuilder::new("http://test:8080") + .resilience(resolved) + .build(); + worker.circuit_breaker().force_open(); + + // CB is disabled — request should go through even though CB is open + let response = execute_with_resilience(&worker, |_attempt| async { + (StatusCode::OK, "ok").into_response() + }) + .await; + assert_eq!(response.status(), StatusCode::OK); + } #[test] fn test_default_resilience() { From c320e017e0c5259dd93226c4cd2797299f71944a Mon Sep 17 00:00:00 2001 From: Chang Su Date: Wed, 18 Mar 2026 16:04:00 -0700 Subject: [PATCH 2/2] fix(core): record CB outcome in retry path, remove URL from 503 body MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Record circuit breaker outcome after retry loop completes, not just in the no-retry path — fixes bug where CB never learned from retried requests - Replace internal worker URL in 503 body with generic message to avoid leaking backend addresses to clients - Inline execute_single helper (was just operation(0).await) Signed-off-by: Chang Su --- model_gateway/src/core/resilience.rs | 28 ++++++++++++---------------- 1 file changed, 12 insertions(+), 16 deletions(-) diff --git a/model_gateway/src/core/resilience.rs b/model_gateway/src/core/resilience.rs index 9eb9c0cc47..32d7175a53 100644 --- a/model_gateway/src/core/resilience.rs +++ b/model_gateway/src/core/resilience.rs @@ -123,7 +123,7 @@ pub fn resolve_resilience( /// Retry decisions use the worker's `is_retryable()` predicate. pub async fn execute_with_resilience( worker: &(dyn Worker + '_), - operation: F, + mut operation: F, ) -> axum::response::Response where F: FnMut(u32) -> Fut + Send, @@ -140,14 +140,14 @@ where ); return ( StatusCode::SERVICE_UNAVAILABLE, - format!("Circuit breaker open for worker {worker_url}"), + "Upstream worker temporarily unavailable", ) .into_response(); } if !resilience.retry_enabled { // Single attempt, no retries - let response = execute_single(worker, operation).await; + let response = operation(0).await; if resilience.circuit_breaker_enabled { let success = !worker.is_retryable(&response); worker.circuit_breaker().record_outcome(success); @@ -156,7 +156,7 @@ where } // Retry loop — delegate to RetryExecutor with worker-aware hooks - RetryExecutor::execute_response_with_retry( + let response = RetryExecutor::execute_response_with_retry( &resilience.retry, operation, |res, _attempt| worker.is_retryable(res), @@ -173,19 +173,15 @@ where debug!(worker_url = worker_url, "Worker retries exhausted"); }, ) - .await -} + .await; -/// Execute a single attempt of an operation (no retry). -async fn execute_single( - _worker: &(dyn Worker + '_), - mut operation: F, -) -> axum::response::Response -where - F: FnMut(u32) -> Fut + Send, - Fut: std::future::Future + Send, -{ - operation(0).await + // Record final outcome for circuit breaker after all retries + if resilience.circuit_breaker_enabled { + let success = !worker.is_retryable(&response); + worker.circuit_breaker().record_outcome(success); + } + + response } #[cfg(test)]