From 4449c170113ee3661c3b8daf91266aa21df69164 Mon Sep 17 00:00:00 2001 From: fzyzcjy <5236035+fzyzcjy@users.noreply.github.com> Date: Mon, 15 Dec 2025 07:43:04 +0800 Subject: [PATCH] [model-gateway] fix circuit breaker metrics (#15099) Co-authored-by: Simo Lin --- sgl-model-gateway/src/core/circuit_breaker.rs | 44 ++++++++++++++----- sgl-model-gateway/src/core/worker.rs | 4 -- sgl-model-gateway/src/core/worker_builder.rs | 5 ++- 3 files changed, 38 insertions(+), 15 deletions(-) diff --git a/sgl-model-gateway/src/core/circuit_breaker.rs b/sgl-model-gateway/src/core/circuit_breaker.rs index ec00e1467..3138c057c 100644 --- a/sgl-model-gateway/src/core/circuit_breaker.rs +++ b/sgl-model-gateway/src/core/circuit_breaker.rs @@ -8,6 +8,8 @@ use std::{ use tracing::info; +use crate::observability::metrics::RouterMetrics; + /// Circuit breaker configuration #[derive(Debug, Clone)] pub struct CircuitBreakerConfig { @@ -82,16 +84,17 @@ pub struct CircuitBreaker { last_failure_time: Arc>>, last_state_change: Arc>, config: CircuitBreakerConfig, + metric_label: String, } impl CircuitBreaker { /// Create a new circuit breaker with default configuration pub fn new() -> Self { - Self::with_config(CircuitBreakerConfig::default()) + Self::with_config_and_label(CircuitBreakerConfig::default(), String::new()) } - /// Create a new circuit breaker with custom configuration - pub fn with_config(config: CircuitBreakerConfig) -> Self { + /// Create a new circuit breaker with custom configuration and metric label + pub fn with_config_and_label(config: CircuitBreakerConfig, metric_label: String) -> Self { Self { state: Arc::new(RwLock::new(CircuitState::Closed)), consecutive_failures: Arc::new(AtomicU32::new(0)), @@ -101,9 +104,15 @@ impl CircuitBreaker { last_failure_time: Arc::new(RwLock::new(None)), last_state_change: Arc::new(RwLock::new(Instant::now())), config, + metric_label, } } + /// Get the metric label + pub fn metric_label(&self) -> &str { + &self.metric_label + } + /// Check if a request can be executed pub fn can_execute(&self) -> bool { let state = self.state(); @@ -140,6 +149,10 @@ impl CircuitBreaker { } else { self.record_failure(); } + + let outcome_str = if success { "success" } else { "failure" }; + RouterMetrics::record_cb_outcome(&self.metric_label, outcome_str); + self.publish_gauge_metrics(); } /// Record a successful request @@ -217,6 +230,9 @@ impl CircuitBreaker { let from = old_state.as_str(); let to = new_state.as_str(); info!("Circuit breaker state transition: {} -> {}", from, to); + RouterMetrics::record_cb_state_transition(&self.metric_label, from, to); + RouterMetrics::set_cb_state(&self.metric_label, new_state.to_int()); + self.publish_gauge_metrics(); } } @@ -274,6 +290,7 @@ impl CircuitBreaker { self.transition_to(CircuitState::Closed); self.consecutive_failures.store(0, Ordering::Release); self.consecutive_successes.store(0, Ordering::Release); + self.publish_gauge_metrics(); } /// Force the circuit to open (for manual intervention) @@ -293,6 +310,12 @@ impl CircuitBreaker { time_since_last_state_change: self.time_since_last_state_change(), } } + + // TODO maybe publish whenever the variable is changed + fn publish_gauge_metrics(&self) { + RouterMetrics::set_cb_consecutive_failures(&self.metric_label, self.failure_count()); + RouterMetrics::set_cb_consecutive_successes(&self.metric_label, self.success_count()); + } } impl Clone for CircuitBreaker { @@ -306,6 +329,7 @@ impl Clone for CircuitBreaker { last_failure_time: Arc::clone(&self.last_failure_time), last_state_change: Arc::clone(&self.last_state_change), config: self.config.clone(), + metric_label: self.metric_label.clone(), } } } @@ -349,7 +373,7 @@ mod tests { failure_threshold: 3, ..Default::default() }; - let cb = CircuitBreaker::with_config(config); + let cb = CircuitBreaker::with_config_and_label(config, String::new()); assert_eq!(cb.state(), CircuitState::Closed); cb.record_failure(); @@ -370,7 +394,7 @@ mod tests { timeout_duration: Duration::from_millis(100), ..Default::default() }; - let cb = CircuitBreaker::with_config(config); + let cb = CircuitBreaker::with_config_and_label(config, String::new()); cb.record_failure(); assert_eq!(cb.state(), CircuitState::Open); @@ -389,7 +413,7 @@ mod tests { timeout_duration: Duration::from_millis(50), ..Default::default() }; - let cb = CircuitBreaker::with_config(config); + let cb = CircuitBreaker::with_config_and_label(config, String::new()); cb.record_failure(); assert_eq!(cb.state(), CircuitState::Open); @@ -412,7 +436,7 @@ mod tests { timeout_duration: Duration::from_millis(50), ..Default::default() }; - let cb = CircuitBreaker::with_config(config); + let cb = CircuitBreaker::with_config_and_label(config, String::new()); cb.record_failure(); assert_eq!(cb.state(), CircuitState::Open); @@ -432,7 +456,7 @@ mod tests { failure_threshold: 3, ..Default::default() }; - let cb = CircuitBreaker::with_config(config); + let cb = CircuitBreaker::with_config_and_label(config, String::new()); cb.record_failure(); cb.record_failure(); @@ -453,7 +477,7 @@ mod tests { failure_threshold: 1, ..Default::default() }; - let cb = CircuitBreaker::with_config(config); + let cb = CircuitBreaker::with_config_and_label(config, String::new()); cb.record_failure(); assert_eq!(cb.state(), CircuitState::Open); @@ -480,7 +504,7 @@ mod tests { failure_threshold: 2, ..Default::default() }; - let cb = CircuitBreaker::with_config(config); + let cb = CircuitBreaker::with_config_and_label(config, String::new()); cb.record_success(); cb.record_failure(); diff --git a/sgl-model-gateway/src/core/worker.rs b/sgl-model-gateway/src/core/worker.rs index 28e596a9e..2e1eb0e1c 100644 --- a/sgl-model-gateway/src/core/worker.rs +++ b/sgl-model-gateway/src/core/worker.rs @@ -129,10 +129,6 @@ pub trait Worker: Send + Sync + fmt::Debug { /// Record the outcome of a request to this worker fn record_outcome(&self, success: bool) { - let outcome_str = if success { "success" } else { "failure" }; - RouterMetrics::record_cb_outcome(self.url(), outcome_str); - - let before = self.circuit_breaker().state(); self.circuit_breaker().record_outcome(success); let after = self.circuit_breaker().state(); diff --git a/sgl-model-gateway/src/core/worker_builder.rs b/sgl-model-gateway/src/core/worker_builder.rs index 63ff5b099..5784050c5 100644 --- a/sgl-model-gateway/src/core/worker_builder.rs +++ b/sgl-model-gateway/src/core/worker_builder.rs @@ -185,7 +185,10 @@ impl BasicWorkerBuilder { healthy: Arc::new(AtomicBool::new(true)), consecutive_failures: Arc::new(AtomicUsize::new(0)), consecutive_successes: Arc::new(AtomicUsize::new(0)), - circuit_breaker: CircuitBreaker::with_config(self.circuit_breaker_config), + circuit_breaker: CircuitBreaker::with_config_and_label( + self.circuit_breaker_config, + self.url.clone(), + ), grpc_client, models_override: Arc::new(StdRwLock::new(None)), }