Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
31 changes: 16 additions & 15 deletions bindings/golang/src/policy.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@ use std::{
os::raw::c_char,
ptr,
sync::{
atomic::{AtomicBool, AtomicUsize, Ordering},
atomic::{AtomicU8, AtomicUsize, Ordering},
Arc,
},
};
Expand All @@ -19,7 +19,7 @@ use async_trait::async_trait;
use llm_tokenizer::{create_tokenizer_from_file, traits::Tokenizer};
use openai_protocol::{
chat::ChatCompletionRequest,
worker::{HealthCheckConfig, WorkerSpec},
worker::{HealthCheckConfig, WorkerSpec, WorkerStatus},
};
use smg::{
policies::{
Expand Down Expand Up @@ -54,7 +54,7 @@ use super::{
pub struct GrpcWorker {
pub(crate) client: Arc<SglangSchedulerClient>,
pub(crate) endpoint: String,
pub(crate) healthy: AtomicBool,
pub(crate) status: AtomicU8,
pub(crate) load: AtomicUsize,
pub(crate) processed: AtomicUsize,
pub(crate) circuit_breaker: CircuitBreaker,
Expand All @@ -80,7 +80,7 @@ impl GrpcWorker {
client,
routing_key_load: WorkerRoutingKeyLoad::new(&endpoint),
endpoint,
healthy: AtomicBool::new(true),
status: AtomicU8::new(WorkerStatus::Ready as u8),
load: AtomicUsize::new(0),
processed: AtomicUsize::new(0),
circuit_breaker: CircuitBreaker::new(),
Expand All @@ -96,7 +96,10 @@ impl std::fmt::Debug for GrpcWorker {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("GrpcWorker")
.field("endpoint", &self.endpoint)
.field("healthy", &self.healthy.load(Ordering::Relaxed))
.field(
"status",
&WorkerStatus::from_u8(self.status.load(Ordering::Relaxed)),
)
.finish()
}
}
Expand All @@ -119,12 +122,12 @@ impl Worker for GrpcWorker {
&self.metadata.spec.connection_mode
}

fn is_healthy(&self) -> bool {
self.healthy.load(Ordering::Relaxed)
fn status(&self) -> WorkerStatus {
WorkerStatus::from_u8(self.status.load(Ordering::Relaxed))

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

The status() method currently uses Ordering::Relaxed for loading the worker status. While this might be acceptable for a simple health flag, using Ordering::Acquire is more consistent with the BasicWorker implementation in the core gateway and ensures that any state changes synchronized via the status update are visible to the thread performing the load. This is particularly important when the status is used by load balancing policies to make routing decisions.

Suggested change
WorkerStatus::from_u8(self.status.load(Ordering::Relaxed))
WorkerStatus::from_u8(self.status.load(Ordering::Acquire))
References
  1. To maintain consistency, changes to a feature on one execution path should align with its behavior on related paths or core implementations.

}

fn set_healthy(&self, healthy: bool) {
self.healthy.store(healthy, Ordering::Relaxed);
fn set_status(&self, status: WorkerStatus) {
self.status.store(status as u8, Ordering::Relaxed);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

The set_status() method currently uses Ordering::Relaxed for storing the worker status. To ensure proper synchronization with threads performing an Acquire load (as suggested for the status() method), Ordering::Release should be used. This ensures that all previous memory operations in the current thread are visible to other threads that subsequently load the status with Acquire ordering.

Suggested change
self.status.store(status as u8, Ordering::Relaxed);
self.status.store(status as u8, Ordering::Release);
References
  1. To maintain consistency, changes to a feature on one execution path should align with its behavior on related paths or core implementations.

}

async fn check_health_async(&self) -> WorkerResult<()> {
Expand Down Expand Up @@ -178,11 +181,11 @@ impl Worker for GrpcWorker {
}

async fn grpc_health_check(&self) -> WorkerResult<bool> {
Ok(self.healthy.load(Ordering::Relaxed))
Ok(self.is_healthy())
}

async fn http_health_check(&self) -> WorkerResult<bool> {
Ok(self.healthy.load(Ordering::Relaxed))
Ok(self.is_healthy())
}
Comment on lines 183 to 189

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🧹 Nitpick | 🔵 Trivial

Health check stubs return current health state — verify intent.

Both grpc_health_check and http_health_check now return self.is_healthy() instead of always returning true. This changes the semantics: previously, these would indicate "check succeeded", now they report current health state.

Given the comment "FFI workers don't do their own health checks", the intent seems correct — but returning the current health state from a "health check" method is a bit unusual. The caller expects this to perform a check, not reflect cached state.

Consider clarifying with a comment or renaming the semantics in a follow-up:

async fn grpc_health_check(&self) -> WorkerResult<bool> {
    // FFI workers don't perform active health checks; report current state
    Ok(self.is_healthy())
}
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@bindings/golang/src/policy.rs` around lines 183 - 189, The health-check
methods grpc_health_check and http_health_check now return the cached state from
is_healthy() which can be confusing since callers expect an active check; update
both functions to include a short clarifying comment like "FFI workers don't
perform active health checks; report current state" above the
Ok(self.is_healthy()) return (or alternatively rename the methods to reflect
"is_healthy" semantics), and keep the existing behavior if that was intended so
readers understand this is a deliberate design choice.

}

Expand Down Expand Up @@ -377,7 +380,7 @@ pub unsafe extern "C" fn sgl_multi_client_healthy_count(
(*handle)
.grpc_workers
.iter()
.filter(|w| w.healthy.load(Ordering::Relaxed))
.filter(|w| w.is_healthy())
.count()
}

Expand All @@ -398,9 +401,7 @@ pub unsafe extern "C" fn sgl_multi_client_set_worker_health(
if worker_index >= client.grpc_workers.len() {
return SglErrorCode::InvalidArgument;
}
client.grpc_workers[worker_index]
.healthy
.store(healthy, Ordering::Relaxed);
client.grpc_workers[worker_index].set_healthy(healthy);
SglErrorCode::Success
}

Expand Down
10 changes: 6 additions & 4 deletions model_gateway/src/worker/builder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -207,7 +207,7 @@ impl BasicWorkerBuilder {
/// Build the BasicWorker instance
pub fn build(mut self) -> BasicWorker {
use std::sync::{
atomic::{AtomicBool, AtomicUsize},
atomic::{AtomicU8, AtomicUsize},
Arc,
};

Expand Down Expand Up @@ -239,8 +239,10 @@ impl BasicWorkerBuilder {
None => OnceCell::new(),
});

let healthy = true;
Metrics::set_worker_health(&metadata.spec.url, healthy);
// Workers start Ready (routable). PR 6b will change this to Pending
// for health-checked workers.
let initial_status = openai_protocol::worker::WorkerStatus::Ready;
Metrics::set_worker_health(&metadata.spec.url, true);

let http_client = self.http_client.unwrap_or_else(|| {
reqwest::Client::builder()
Expand All @@ -258,7 +260,7 @@ impl BasicWorkerBuilder {
load_counter: Arc::new(AtomicUsize::new(0)),
worker_routing_key_load: Arc::new(WorkerRoutingKeyLoad::new(&metadata.spec.url)),
processed_counter: Arc::new(AtomicUsize::new(0)),
healthy: Arc::new(AtomicBool::new(healthy)),
status: Arc::new(AtomicU8::new(initial_status as u8)),
consecutive_failures: Arc::new(AtomicUsize::new(0)),
consecutive_successes: Arc::new(AtomicUsize::new(0)),
circuit_breaker: CircuitBreaker::with_config_and_label(
Expand Down
51 changes: 39 additions & 12 deletions model_gateway/src/worker/worker.rs
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
use std::{
fmt,
sync::{
atomic::{AtomicBool, AtomicUsize, Ordering},
atomic::{AtomicU8, AtomicUsize, Ordering},
Arc, LazyLock,
},
time::Duration,
Expand Down Expand Up @@ -141,11 +141,35 @@ pub trait Worker: Send + Sync + fmt::Debug {
self.metadata().spec.bootstrap_port
}

/// Check if the worker is currently healthy
fn is_healthy(&self) -> bool;
/// Get the worker's lifecycle status.
fn status(&self) -> WorkerStatus;

/// Set the worker's health status
fn set_healthy(&self, healthy: bool);
/// Set the worker's lifecycle status.
fn set_status(&self, status: WorkerStatus);

/// Check if the worker is currently healthy (status == Ready).
///
/// This is a routing predicate — returns true only for `Ready` workers.
/// A `Pending` worker is not "unhealthy", just unverified.
fn is_healthy(&self) -> bool {
self.status() == WorkerStatus::Ready
}

/// Set the worker's health status (compatibility shim).
///
/// Maps `true` → `Ready`, `false` → `NotReady`.
/// Prefer `set_status()` for explicit state transitions.
fn set_healthy(&self, healthy: bool) {
if healthy {
self.set_status(WorkerStatus::Ready);
} else {
// Only transition to NotReady if currently Ready.
// Don't transition Pending→NotReady (hasn't proven itself).
if self.status() == WorkerStatus::Ready {
self.set_status(WorkerStatus::NotReady);
}
}
}
Comment on lines +162 to +172

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🧹 Nitpick | 🔵 Trivial

Minor TOCTOU window in set_healthy(false) — acceptable for current usage.

The read-then-write pattern between self.status() and self.set_status() has a theoretical race window. However:

  • Health checks run sequentially per worker
  • The guard only prevents Pending → NotReady
  • Worst case: a concurrent set_healthy(true) wins, which is benign

If future usage requires stricter atomicity, consider a compare-and-swap pattern. For now, this is fine.

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@model_gateway/src/worker/worker.rs` around lines 162 - 172, The current
set_healthy(…) does a non-atomic read-then-write (status() then set_status())
which has a TOCTOU window; replace this with an atomic compare-and-swap on the
worker status so the transition to NotReady only occurs if the current status is
still Ready. Implement or use a method like compare_and_set_status(expected:
WorkerStatus, new: WorkerStatus) (or add an atomic field and CAS helper) and
call compare_and_set_status(WorkerStatus::Ready, WorkerStatus::NotReady) inside
set_healthy(false) instead of calling status() then set_status(); keep
WorkerStatus enum values (Ready, NotReady, Pending) for comparisons.


/// Perform an async health check on the worker
async fn check_health_async(&self) -> WorkerResult<()>;
Expand Down Expand Up @@ -500,7 +524,7 @@ pub struct BasicWorker {
pub load_counter: Arc<AtomicUsize>,
pub worker_routing_key_load: Arc<WorkerRoutingKeyLoad>,
pub processed_counter: Arc<AtomicUsize>,
pub healthy: Arc<AtomicBool>,
pub status: Arc<AtomicU8>,
pub consecutive_failures: Arc<AtomicUsize>,
pub consecutive_successes: Arc<AtomicUsize>,
pub circuit_breaker: CircuitBreaker,
Expand All @@ -521,7 +545,10 @@ impl fmt::Debug for BasicWorker {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("BasicWorker")
.field("metadata", &self.metadata)
.field("healthy", &self.healthy.load(Ordering::Relaxed))
.field(
"status",
&WorkerStatus::from_u8(self.status.load(Ordering::Relaxed)),
)
.field("circuit_breaker", &self.circuit_breaker)
.field("grpc_client", &"<OnceCell>")
.finish()
Expand Down Expand Up @@ -553,13 +580,13 @@ impl Worker for BasicWorker {
&self.metadata.spec.connection_mode
}

fn is_healthy(&self) -> bool {
self.healthy.load(Ordering::Acquire)
fn status(&self) -> WorkerStatus {
WorkerStatus::from_u8(self.status.load(Ordering::Acquire))
}

fn set_healthy(&self, healthy: bool) {
self.healthy.store(healthy, Ordering::Release);
Metrics::set_worker_health(self.url(), healthy);
fn set_status(&self, status: WorkerStatus) {
self.status.store(status as u8, Ordering::Release);
Metrics::set_worker_health(self.url(), status == WorkerStatus::Ready);
}

async fn check_health_async(&self) -> WorkerResult<()> {
Expand Down
Loading