[model-gateway] Fix metric emission gaps and name mismatch (#15093)
This commit is contained in:
@@ -7,6 +7,7 @@ use tracing::{debug, warn};
|
|||||||
|
|
||||||
use crate::{
|
use crate::{
|
||||||
app_context::AppContext,
|
app_context::AppContext,
|
||||||
|
observability::metrics::RouterMetrics,
|
||||||
workflow::{StepExecutor, StepResult, WorkflowContext, WorkflowError, WorkflowResult},
|
workflow::{StepExecutor, StepResult, WorkflowContext, WorkflowError, WorkflowResult},
|
||||||
};
|
};
|
||||||
|
|
||||||
@@ -48,6 +49,9 @@ impl StepExecutor for RemoveFromWorkerRegistryStep {
|
|||||||
debug!("Removed {} worker(s) from registry", removed_count);
|
debug!("Removed {} worker(s) from registry", removed_count);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Update active workers metric
|
||||||
|
RouterMetrics::set_active_workers(app_context.worker_registry.len());
|
||||||
|
|
||||||
Ok(StepResult::Success)
|
Ok(StepResult::Success)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -8,6 +8,7 @@ use tracing::debug;
|
|||||||
use crate::{
|
use crate::{
|
||||||
app_context::AppContext,
|
app_context::AppContext,
|
||||||
core::Worker,
|
core::Worker,
|
||||||
|
observability::metrics::RouterMetrics,
|
||||||
workflow::{StepExecutor, StepResult, WorkflowContext, WorkflowResult},
|
workflow::{StepExecutor, StepResult, WorkflowContext, WorkflowResult},
|
||||||
};
|
};
|
||||||
|
|
||||||
@@ -36,6 +37,9 @@ impl StepExecutor for RegisterWorkersStep {
|
|||||||
worker_ids.push(worker_id);
|
worker_ids.push(worker_id);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Update active workers metric
|
||||||
|
RouterMetrics::set_active_workers(app_context.worker_registry.len());
|
||||||
|
|
||||||
context.set("worker_ids", worker_ids);
|
context.set("worker_ids", worker_ids);
|
||||||
Ok(StepResult::Success)
|
Ok(StepResult::Success)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -224,6 +224,16 @@ impl WorkerRegistry {
|
|||||||
.unwrap_or_default()
|
.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
|
/// Get all workers
|
||||||
pub fn get_all(&self) -> Vec<Arc<dyn Worker>> {
|
pub fn get_all(&self) -> Vec<Arc<dyn Worker>> {
|
||||||
self.workers
|
self.workers
|
||||||
|
|||||||
@@ -37,7 +37,7 @@ pub fn init_metrics() {
|
|||||||
"Total number of request errors by route and error type"
|
"Total number of request errors by route and error type"
|
||||||
);
|
);
|
||||||
describe_counter!(
|
describe_counter!(
|
||||||
"sgl_router_upstream_http_responses_total",
|
"sgl_router_attempt_http_responses_total",
|
||||||
"Total number of upstream engine HTTP responses by status code"
|
"Total number of upstream engine HTTP responses by status code"
|
||||||
);
|
);
|
||||||
describe_counter!(
|
describe_counter!(
|
||||||
@@ -281,6 +281,10 @@ impl RouterMetrics {
|
|||||||
.set(if healthy { 1.0 } else { 0.0 });
|
.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) {
|
pub fn record_processed_request(worker_url: &str) {
|
||||||
counter!("sgl_router_processed_requests_total",
|
counter!("sgl_router_processed_requests_total",
|
||||||
"worker" => worker_url.to_string()
|
"worker" => worker_url.to_string()
|
||||||
|
|||||||
@@ -134,6 +134,12 @@ impl CacheAwarePolicy {
|
|||||||
let model_id = tree_ref.key();
|
let model_id = tree_ref.key();
|
||||||
let tree = tree_ref.value();
|
let tree = tree_ref.value();
|
||||||
tree.evict_tenant_by_size(max_tree_size);
|
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!(
|
debug!(
|
||||||
"Cache eviction completed for model {}, max_size: {}",
|
"Cache eviction completed for model {}, max_size: {}",
|
||||||
model_id, max_tree_size
|
model_id, max_tree_size
|
||||||
|
|||||||
@@ -18,7 +18,10 @@ use rustls;
|
|||||||
use tokio::{task, time};
|
use tokio::{task, time};
|
||||||
use tracing::{debug, error, info, warn};
|
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)]
|
#[derive(Debug, Clone)]
|
||||||
pub struct ServiceDiscoveryConfig {
|
pub struct ServiceDiscoveryConfig {
|
||||||
@@ -404,6 +407,7 @@ async fn handle_pod_event(
|
|||||||
match job_queue.submit(job).await {
|
match job_queue.submit(job).await {
|
||||||
Ok(_) => {
|
Ok(_) => {
|
||||||
debug!("Worker addition job submitted for: {}", worker_url);
|
debug!("Worker addition job submitted for: {}", worker_url);
|
||||||
|
RouterMetrics::record_discovery_update(1, 0);
|
||||||
}
|
}
|
||||||
Err(e) => {
|
Err(e) => {
|
||||||
error!(
|
error!(
|
||||||
@@ -462,6 +466,7 @@ async fn handle_pod_deletion(
|
|||||||
);
|
);
|
||||||
} else {
|
} else {
|
||||||
debug!("Submitted worker removal job for {}", worker_url);
|
debug!("Submitted worker removal job for {}", worker_url);
|
||||||
|
RouterMetrics::record_discovery_update(0, 1);
|
||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
error!(
|
error!(
|
||||||
|
|||||||
Reference in New Issue
Block a user