From 4e0b56c8119a2673d0bf3bd748b8d9a895af79cb Mon Sep 17 00:00:00 2001 From: Kangyan-Zhou Date: Fri, 18 Sep 2026 09:46:29 -0700 Subject: [PATCH] [Router] Treat an upstream 503/429 as backpressure, not a breaker fault (2/3) (#39464) Co-authored-by: Kangyan Zhou Co-authored-by: Claude Opus 5 (1M context) --- .../sgl-router/src/health/circuit_breaker.rs | 92 +++++ experimental/sgl-router/src/proxy/mod.rs | 334 ++++++++++++++++-- 2 files changed, 400 insertions(+), 26 deletions(-) diff --git a/experimental/sgl-router/src/health/circuit_breaker.rs b/experimental/sgl-router/src/health/circuit_breaker.rs index d216d3ebb..26d8c73f0 100644 --- a/experimental/sgl-router/src/health/circuit_breaker.rs +++ b/experimental/sgl-router/src/health/circuit_breaker.rs @@ -172,6 +172,35 @@ impl CircuitBreaker { } } } + + /// Record a backpressure response (HTTP 503 / 429): the worker answered, so + /// it is responsive — busy, not faulty. + /// + /// - **Closed:** no-op. A busy worker must not open the breaker, and — + /// unlike [`record_success`](Self::record_success) — backpressure must + /// NOT reset an in-progress failure streak, so a worker interleaving real + /// 5xx faults with 503s still trips. + /// - **HalfOpen:** close. Any response observed here proves the worker is + /// answering, which is what the probe exists to find out. Leaving HalfOpen + /// unresolved would wedge the breaker permanently — the probe slot is + /// released only by a success or failure, and backpressure is neither — + /// shutting a recovered-but-busy worker out forever (a worse false-shed + /// than the one ignoring 503 removes). The responder is not necessarily + /// the probe: [`allow`](Self::allow) gates admission, not completion, so a + /// request admitted while Closed can land here. [`record_success`] has the + /// same property. + /// - **Open:** no-op, and reachable — `allow` gates admission, not + /// completion, so a request admitted while Closed can return after + /// concurrent failures have opened the breaker. A late backpressure answer + /// must not reset a breaker that has already tripped, exactly as + /// [`record_failure`](Self::record_failure) ignores failures while Open. + pub fn record_backpressure(&self) { + let mut g = self.inner.lock().unwrap(); + if matches!(g.state, State::HalfOpen { .. }) { + g.consecutive_failures = 0; + g.state = State::Closed; + } + } } impl Default for CircuitBreaker { @@ -245,4 +274,67 @@ mod tests { assert!(s.admit, "open past cooldown admits a probe"); assert_eq!(s.state_code, 1, "...but is still reported as open"); } + + #[tokio::test(start_paused = true)] + async fn backpressure_resolves_half_open_probe() { + // Regression guard: a backpressure (503/429) answer to a half-open + // probe must RESOLVE the probe, not wedge the breaker. The probe slot + // is otherwise released only by success/failure; without + // record_backpressure handling HalfOpen, a recovered-but-busy worker + // would be shut out forever. + let b = cb(1, 10); + b.record_failure(); // Open + assert_eq!(b.snapshot().state_code, 1); + tokio::time::advance(Duration::from_secs(11)).await; + assert!( + b.allow(), + "cooldown elapsed → claims the probe slot (HalfOpen)" + ); + assert_eq!(b.snapshot().state_code, 2); + + b.record_backpressure(); + assert_eq!( + b.snapshot().state_code, + 0, + "a 503 probe answer must close the breaker, not leave it wedged half-open", + ); + assert!( + b.would_allow(), + "worker must admit again after the probe resolves" + ); + } + + #[test] + fn backpressure_in_closed_state_preserves_failure_streak() { + // Unlike record_success, record_backpressure must NOT reset an + // in-progress streak: 2 faults + a 503 + 1 fault still hits threshold 3. + let b = cb(3, 30); + b.record_failure(); + b.record_failure(); + b.record_backpressure(); + assert_eq!( + b.snapshot().state_code, + 0, + "2 faults < threshold 3, still closed" + ); + b.record_failure(); + assert_eq!( + b.snapshot().state_code, + 1, + "the 503 must not have reset the streak; the 3rd fault opens the breaker", + ); + } + + #[test] + fn backpressure_alone_never_opens_a_closed_breaker() { + let b = cb(3, 30); + for _ in 0..10 { + b.record_backpressure(); + } + assert_eq!( + b.snapshot().state_code, + 0, + "backpressure alone must never open the breaker, regardless of volume", + ); + } } diff --git a/experimental/sgl-router/src/proxy/mod.rs b/experimental/sgl-router/src/proxy/mod.rs index 405018598..69542844d 100644 --- a/experimental/sgl-router/src/proxy/mod.rs +++ b/experimental/sgl-router/src/proxy/mod.rs @@ -31,6 +31,56 @@ fn parse_worker_url(worker_url: &str, breaker: &CircuitBreaker) -> Result BreakerOutcome { + use reqwest::StatusCode; + match status { + StatusCode::SERVICE_UNAVAILABLE | StatusCode::TOO_MANY_REQUESTS => BreakerOutcome::Neutral, + s if s.is_server_error() => BreakerOutcome::Failure, + _ => BreakerOutcome::Success, + } +} + #[derive(Debug)] pub struct Proxy { /// The negotiating client: HTTP/1.1 in cleartext, and ALPN `h2, http/1.1` @@ -120,9 +170,10 @@ impl Proxy { } } - /// Breaker-gated JSON POST: checks `breaker.allow()` first, records - /// success/failure based on response status, and returns - /// `ApiError::BreakerOpen` immediately when the breaker is Open. + /// Breaker-gated JSON POST: checks `breaker.allow()` first, classifies the + /// response status through [`breaker_outcome`] (success / failure / + /// backpressure), and returns `ApiError::BreakerOpen` immediately when the + /// breaker is Open. /// /// `worker_url` is the discovery-emitted worker URL string. It's parsed /// to [`reqwest::Url`] internally so we can use [`Url::join`] for clean @@ -185,10 +236,14 @@ impl Proxy { return Err(ApiError::UpstreamStatus { status }); } }; - if status.is_server_error() { - breaker.record_failure(); - } else { - breaker.record_success(); + match breaker_outcome(status) { + BreakerOutcome::Failure => breaker.record_failure(), + BreakerOutcome::Success => breaker.record_success(), + // Backpressure (503/429): the engine is healthy but busy. This never + // opens the breaker and (in Closed) leaves the failure streak + // intact, but it DOES resolve a half-open probe so a recovered + // worker that answers a probe with 503 isn't wedged shut. + BreakerOutcome::Neutral => breaker.record_backpressure(), } let mut out = Response::new(Body::from(bytes)); *out.status_mut() = status; @@ -199,8 +254,9 @@ impl Proxy { Ok(out) } - /// Breaker-gated streaming POST: checks `breaker.allow()` first, records - /// success/failure, and returns `ApiError::BreakerOpen` when Open. + /// Breaker-gated streaming POST: checks `breaker.allow()` first, classifies + /// the response status through [`breaker_outcome`], and returns + /// `ApiError::BreakerOpen` when Open. /// /// `stream_guards` — when `Some`, the value is threaded into the SSE /// pump task and held for the entire body lifetime (headers → last byte @@ -263,30 +319,42 @@ impl Proxy { }; // Breaker recording is deferred to the pump's completion hook so // an upstream that returns 2xx headers and then drops mid-stream - // is recorded as a failure. For 5xx headers we record_failure + // is recorded as a failure. For a genuine 5xx fault we record_failure // up front and skip the pump hook (the body we surface is the - // error response — its stream completing is not a worker win). + // error response — its stream completing is not a worker win). For a + // backpressure status (503/429) we record_backpressure up front and + // skip the hook: a busy-but-healthy engine's queue-full responses can't + // open the breaker, but a half-open probe answered with 503 is still + // resolved rather than wedged (see `breaker_outcome` / + // `record_backpressure`). let caller_end_hook = if status.is_success() { on_stream_end } else { None }; let on_complete: Option> = - if status.is_server_error() { - breaker.record_failure(); - None - } else { - let breaker_for_hook = Arc::clone(breaker); - Some(Box::new(move |end| { - if end.transport_ok { - breaker_for_hook.record_success(); - } else { - breaker_for_hook.record_failure(); - } - if let Some(hook) = caller_end_hook { - hook(end); - } - })) + match breaker_outcome(status) { + BreakerOutcome::Failure => { + breaker.record_failure(); + None + } + BreakerOutcome::Neutral => { + breaker.record_backpressure(); + None + } + BreakerOutcome::Success => { + let breaker_for_hook = Arc::clone(breaker); + Some(Box::new(move |end| { + if end.transport_ok { + breaker_for_hook.record_success(); + } else { + breaker_for_hook.record_failure(); + } + if let Some(hook) = caller_end_hook { + hook(end); + } + })) + } }; // Only record TTFT for successful streams; error-body chunks are not // generated tokens. @@ -315,7 +383,14 @@ impl Proxy { #[cfg(test)] mod tests { use super::*; + use crate::health::circuit_breaker::CircuitBreakerConfig; + use axum::routing::post; + use axum::Router; + use reqwest::StatusCode; + use std::num::NonZeroU32; use std::time::Duration; + use tokio::net::TcpListener; + use tokio::sync::oneshot; #[tokio::test] async fn new_returns_result_not_panic() { @@ -342,4 +417,211 @@ mod tests { p.admin_client() )); } + + #[test] + fn breaker_outcome_treats_backpressure_as_neutral() { + // Backpressure: healthy but busy — must not touch the breaker. + assert_eq!( + breaker_outcome(StatusCode::SERVICE_UNAVAILABLE), + BreakerOutcome::Neutral, + ); + assert_eq!( + breaker_outcome(StatusCode::TOO_MANY_REQUESTS), + BreakerOutcome::Neutral, + ); + // Genuine faults: still failures. + for s in [ + StatusCode::INTERNAL_SERVER_ERROR, + StatusCode::BAD_GATEWAY, + StatusCode::GATEWAY_TIMEOUT, + ] { + assert_eq!(breaker_outcome(s), BreakerOutcome::Failure, "{s}"); + } + // Non-5xx (incl. 4xx client errors): treated as success. + for s in [ + StatusCode::OK, + StatusCode::BAD_REQUEST, + StatusCode::NOT_FOUND, + ] { + assert_eq!(breaker_outcome(s), BreakerOutcome::Success, "{s}"); + } + } + + /// A fake upstream that answers every POST with a fixed status + tiny body. + async fn spawn_status_worker(status: u16) -> (String, oneshot::Sender<()>) { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let port = listener.local_addr().unwrap().port(); + let code = StatusCode::from_u16(status).unwrap(); + let app = Router::new().route( + "/v1/chat/completions", + post(move || async move { (code, "{\"error\":\"x\"}") }), + ); + let (tx, rx) = oneshot::channel::<()>(); + tokio::spawn(async move { + let _ = axum::serve(listener, app) + .with_graceful_shutdown(async move { + let _ = rx.await; + }) + .await; + }); + (format!("http://127.0.0.1:{port}"), tx) + } + + /// A saturated engine's own queue-full 503s must not trip the router's + /// circuit breaker. Dispatch far past any plausible failure threshold and + /// assert the breaker stays Closed and admitting. + #[tokio::test] + async fn engine_503_does_not_trip_breaker() { + let (url, _shutdown) = spawn_status_worker(503).await; + let proxy = Proxy::new(Duration::from_secs(5)).unwrap(); + let breaker = CircuitBreaker::new(); + let headers = HeaderMap::new(); + + for i in 0..50 { + let resp = proxy + .forward_json_to( + &url, + WireProtocol::Http1, + &breaker, + "/v1/chat/completions", + &headers, + Bytes::from_static(b"{}"), + ) + .await + .expect("dispatch should reach the worker (breaker must stay closed)"); + assert_eq!( + resp.status(), + StatusCode::SERVICE_UNAVAILABLE, + "iter {i}: client must still see the engine's 503", + ); + assert_eq!( + breaker.snapshot().state_code, + 0, + "iter {i}: 503 backpressure must leave the breaker Closed", + ); + } + assert!( + breaker.would_allow(), + "breaker must keep admitting after a burst of engine 503s", + ); + } + + /// Contrast guard, so the backpressure carve-out cannot disable fault + /// detection: a genuine 5xx fault (500) MUST still open the breaker. Loops + /// on `would_allow()` rather than a fixed count so the test stays correct if + /// the default `CircuitBreakerConfig` threshold changes. + #[tokio::test] + async fn engine_500_still_trips_breaker() { + let (url, _shutdown) = spawn_status_worker(500).await; + let proxy = Proxy::new(Duration::from_secs(5)).unwrap(); + let breaker = CircuitBreaker::new(); + let headers = HeaderMap::new(); + + for _ in 0..50 { + if !breaker.would_allow() { + break; + } + let _ = proxy + .forward_json_to( + &url, + WireProtocol::Http1, + &breaker, + "/v1/chat/completions", + &headers, + Bytes::from_static(b"{}"), + ) + .await; + } + assert_eq!( + breaker.snapshot().state_code, + 1, + "a run of 500s must open the breaker (fault detection still works)", + ); + } + + /// End-to-end wedge guard: a breaker that opened on real faults, then has + /// its half-open probe answered with a 503, must RECOVER — not stay shut + /// out forever. Exercises the `Neutral => record_backpressure` wiring in + /// `forward_json_to` through the half-open path. + #[tokio::test] + async fn engine_503_recovers_a_half_open_breaker() { + let (url, _shutdown) = spawn_status_worker(503).await; + let proxy = Proxy::new(Duration::from_secs(5)).unwrap(); + // threshold=1 so one prior fault opens it; a short cooldown so the probe + // is admitted quickly. The wait below is an order of magnitude longer + // than the cooldown rather than a thin margin, since this test needs a + // real socket and so cannot pause the clock. + let breaker = CircuitBreaker::with_config(CircuitBreakerConfig { + threshold: NonZeroU32::new(1).unwrap(), + cool_down: Duration::from_millis(20), + }); + let headers = HeaderMap::new(); + + // Simulate a prior genuine fault (e.g. a 500 / timeout) that tripped it. + breaker.record_failure(); + assert_eq!(breaker.snapshot().state_code, 1, "breaker should be Open"); + + // Let the cooldown elapse so the next dispatch claims the half-open probe. + tokio::time::sleep(Duration::from_millis(200)).await; + + let resp = proxy + .forward_json_to( + &url, + WireProtocol::Http1, + &breaker, + "/v1/chat/completions", + &headers, + Bytes::from_static(b"{}"), + ) + .await + .expect("the half-open probe must be admitted and reach the worker"); + assert_eq!(resp.status(), StatusCode::SERVICE_UNAVAILABLE); + assert_eq!( + breaker.snapshot().state_code, + 0, + "a 503 answer to the probe must close the breaker, not wedge it half-open", + ); + assert!( + breaker.would_allow(), + "worker must admit traffic again after recovering from the probe", + ); + } + + /// Streaming path parity: the engine's 503 on the streaming arm must also + /// leave the breaker untouched (no up-front failure, no completion hook). + #[tokio::test] + async fn engine_503_does_not_trip_breaker_streaming() { + use http_body_util::BodyExt; + + let (url, _shutdown) = spawn_status_worker(503).await; + let proxy = Proxy::new(Duration::from_secs(5)).unwrap(); + let breaker = Arc::new(CircuitBreaker::new()); + let headers = HeaderMap::new(); + + for i in 0..6 { + let resp = proxy + .forward_streaming_to( + &url, + WireProtocol::Http1, + &breaker, + "/v1/chat/completions", + &headers, + Bytes::from_static(b"{}"), + None, + None, + None, + ) + .await + .expect("streaming dispatch should reach the worker"); + assert_eq!(resp.status(), StatusCode::SERVICE_UNAVAILABLE, "iter {i}"); + // Drain the body so the pump task runs to completion (would fire any + // completion hook). For a 503 there is none, but draining proves it. + let _ = resp.into_body().collect().await; + assert_eq!( + breaker.snapshot().state_code, + 0, + "iter {i}: streaming 503 must leave the breaker Closed", + ); + } + } }