diff --git a/sgl-model-gateway/src/core/steps/worker/local/remove_from_worker_registry.rs b/sgl-model-gateway/src/core/steps/worker/local/remove_from_worker_registry.rs index d9bd39012..4693687a9 100644 --- a/sgl-model-gateway/src/core/steps/worker/local/remove_from_worker_registry.rs +++ b/sgl-model-gateway/src/core/steps/worker/local/remove_from_worker_registry.rs @@ -7,6 +7,7 @@ use tracing::{debug, warn}; use crate::{ app_context::AppContext, + observability::metrics::RouterMetrics, workflow::{StepExecutor, StepResult, WorkflowContext, WorkflowError, WorkflowResult}, }; @@ -48,6 +49,9 @@ impl StepExecutor for RemoveFromWorkerRegistryStep { debug!("Removed {} worker(s) from registry", removed_count); } + // Update active workers metric + RouterMetrics::set_active_workers(app_context.worker_registry.len()); + Ok(StepResult::Success) } diff --git a/sgl-model-gateway/src/core/steps/worker/shared/register.rs b/sgl-model-gateway/src/core/steps/worker/shared/register.rs index b63962647..4aa01eeb8 100644 --- a/sgl-model-gateway/src/core/steps/worker/shared/register.rs +++ b/sgl-model-gateway/src/core/steps/worker/shared/register.rs @@ -8,6 +8,7 @@ use tracing::debug; use crate::{ app_context::AppContext, core::Worker, + observability::metrics::RouterMetrics, workflow::{StepExecutor, StepResult, WorkflowContext, WorkflowResult}, }; @@ -36,6 +37,9 @@ impl StepExecutor for RegisterWorkersStep { worker_ids.push(worker_id); } + // Update active workers metric + RouterMetrics::set_active_workers(app_context.worker_registry.len()); + context.set("worker_ids", worker_ids); Ok(StepResult::Success) } diff --git a/sgl-model-gateway/src/core/worker_registry.rs b/sgl-model-gateway/src/core/worker_registry.rs index f96b88230..4bbd38f36 100644 --- a/sgl-model-gateway/src/core/worker_registry.rs +++ b/sgl-model-gateway/src/core/worker_registry.rs @@ -224,6 +224,16 @@ impl WorkerRegistry { .unwrap_or_default() } + /// Get the number of workers in the registry + pub fn len(&self) -> usize { + self.workers.len() + } + + /// Check if the registry is empty + pub fn is_empty(&self) -> bool { + self.workers.is_empty() + } + /// Get all workers pub fn get_all(&self) -> Vec> { self.workers diff --git a/sgl-model-gateway/src/observability/metrics.rs b/sgl-model-gateway/src/observability/metrics.rs index 6426c612c..77e5bfa39 100644 --- a/sgl-model-gateway/src/observability/metrics.rs +++ b/sgl-model-gateway/src/observability/metrics.rs @@ -37,7 +37,7 @@ pub fn init_metrics() { "Total number of request errors by route and error type" ); describe_counter!( - "sgl_router_upstream_http_responses_total", + "sgl_router_attempt_http_responses_total", "Total number of upstream engine HTTP responses by status code" ); describe_counter!( @@ -281,6 +281,10 @@ impl RouterMetrics { .set(if healthy { 1.0 } else { 0.0 }); } + pub fn set_active_workers(count: usize) { + gauge!("sgl_router_active_workers").set(count as f64); + } + pub fn record_processed_request(worker_url: &str) { counter!("sgl_router_processed_requests_total", "worker" => worker_url.to_string() diff --git a/sgl-model-gateway/src/policies/cache_aware.rs b/sgl-model-gateway/src/policies/cache_aware.rs index af2648dee..1a1a8b598 100644 --- a/sgl-model-gateway/src/policies/cache_aware.rs +++ b/sgl-model-gateway/src/policies/cache_aware.rs @@ -134,6 +134,12 @@ impl CacheAwarePolicy { let model_id = tree_ref.key(); let tree = tree_ref.value(); tree.evict_tenant_by_size(max_tree_size); + + // Update tree size metrics per worker (tenant) + for entry in tree.tenant_char_count.iter() { + RouterMetrics::set_tree_size(entry.key(), *entry.value()); + } + debug!( "Cache eviction completed for model {}, max_size: {}", model_id, max_tree_size diff --git a/sgl-model-gateway/src/service_discovery.rs b/sgl-model-gateway/src/service_discovery.rs index aa8cd1d7a..7ffb71799 100644 --- a/sgl-model-gateway/src/service_discovery.rs +++ b/sgl-model-gateway/src/service_discovery.rs @@ -18,7 +18,10 @@ use rustls; use tokio::{task, time}; use tracing::{debug, error, info, warn}; -use crate::{app_context::AppContext, core::Job, protocols::worker_spec::WorkerConfigRequest}; +use crate::{ + app_context::AppContext, core::Job, observability::metrics::RouterMetrics, + protocols::worker_spec::WorkerConfigRequest, +}; #[derive(Debug, Clone)] pub struct ServiceDiscoveryConfig { @@ -404,6 +407,7 @@ async fn handle_pod_event( match job_queue.submit(job).await { Ok(_) => { debug!("Worker addition job submitted for: {}", worker_url); + RouterMetrics::record_discovery_update(1, 0); } Err(e) => { error!( @@ -462,6 +466,7 @@ async fn handle_pod_deletion( ); } else { debug!("Submitted worker removal job for {}", worker_url); + RouterMetrics::record_discovery_update(0, 1); } } else { error!(