mirror of
https://github.com/fawney19/Aether.git
synced 2026-10-04 08:27:46 +08:00
Merge remote-tracking branch 'origin/pr/587'
This commit is contained in:
@@ -1,7 +1,7 @@
|
||||
use super::extractors::{admin_health_key_id, admin_recover_key_id};
|
||||
use crate::handlers::admin::request::{AdminAppState, AdminRequestContext};
|
||||
use crate::handlers::admin::shared::query_param_value;
|
||||
use crate::handlers::public::ApiFormatHealthMonitorOptions;
|
||||
use crate::handlers::public::{ApiFormatHealthMonitorOptions, ModelHealthMonitorOptions};
|
||||
use crate::GatewayError;
|
||||
use axum::{
|
||||
body::Body,
|
||||
@@ -157,6 +157,80 @@ pub(super) async fn maybe_build_local_admin_endpoints_health_response(
|
||||
return Ok(Some(Json(payload).into_response()));
|
||||
}
|
||||
|
||||
if decision.route_family.as_deref() == Some("endpoints_health")
|
||||
&& decision.route_kind.as_deref() == Some("health_models")
|
||||
&& request_context.path() == "/api/admin/endpoints/health/models"
|
||||
{
|
||||
if !state.has_usage_data_reader() {
|
||||
return Ok(Some(build_admin_endpoint_health_data_unavailable_response()));
|
||||
}
|
||||
let lookback_hours = query_param_value(request_context.query_string(), "lookback_hours")
|
||||
.and_then(|value| value.parse::<u64>().ok())
|
||||
.filter(|value| (1..=72).contains(value))
|
||||
.unwrap_or(6);
|
||||
let model_limit = query_param_value(request_context.query_string(), "model_limit")
|
||||
.and_then(|value| value.parse::<usize>().ok())
|
||||
.filter(|value| (1..=50).contains(value))
|
||||
.unwrap_or(12);
|
||||
let per_model_limit = query_param_value(request_context.query_string(), "per_model_limit")
|
||||
.and_then(|value| value.parse::<usize>().ok())
|
||||
.filter(|value| (10..=200).contains(value))
|
||||
.unwrap_or(60);
|
||||
let Some(payload) = state
|
||||
.build_model_health_monitor_payload(
|
||||
lookback_hours,
|
||||
model_limit,
|
||||
per_model_limit,
|
||||
ModelHealthMonitorOptions {
|
||||
include_provider_count: true,
|
||||
},
|
||||
)
|
||||
.await
|
||||
else {
|
||||
return Ok(Some(build_admin_endpoint_health_data_unavailable_response()));
|
||||
};
|
||||
return Ok(Some(Json(payload).into_response()));
|
||||
}
|
||||
|
||||
if decision.route_family.as_deref() == Some("endpoints_health")
|
||||
&& decision.route_kind.as_deref() == Some("health_providers")
|
||||
&& request_context.path() == "/api/admin/endpoints/health/providers"
|
||||
{
|
||||
if !state.has_provider_catalog_data_reader() || !state.has_usage_data_reader() {
|
||||
return Ok(Some(build_admin_endpoint_health_data_unavailable_response()));
|
||||
}
|
||||
let lookback_hours = query_param_value(request_context.query_string(), "lookback_hours")
|
||||
.and_then(|value| value.parse::<u64>().ok())
|
||||
.filter(|value| (1..=72).contains(value))
|
||||
.unwrap_or(6);
|
||||
let provider_limit = query_param_value(request_context.query_string(), "provider_limit")
|
||||
.and_then(|value| value.parse::<usize>().ok())
|
||||
.filter(|value| (1..=100).contains(value))
|
||||
.unwrap_or(50);
|
||||
let per_provider_model_limit =
|
||||
query_param_value(request_context.query_string(), "per_provider_model_limit")
|
||||
.and_then(|value| value.parse::<usize>().ok())
|
||||
.filter(|value| (1..=50).contains(value))
|
||||
.unwrap_or(12);
|
||||
let per_model_event_limit =
|
||||
query_param_value(request_context.query_string(), "per_model_limit")
|
||||
.and_then(|value| value.parse::<usize>().ok())
|
||||
.filter(|value| (10..=200).contains(value))
|
||||
.unwrap_or(100);
|
||||
let Some(payload) = state
|
||||
.build_provider_health_monitor_payload(
|
||||
lookback_hours,
|
||||
provider_limit,
|
||||
per_provider_model_limit,
|
||||
per_model_event_limit,
|
||||
)
|
||||
.await
|
||||
else {
|
||||
return Ok(Some(build_admin_endpoint_health_data_unavailable_response()));
|
||||
};
|
||||
return Ok(Some(Json(payload).into_response()));
|
||||
}
|
||||
|
||||
if decision.route_family.as_deref() == Some("endpoints_health")
|
||||
&& decision.route_kind.as_deref() == Some("health_status")
|
||||
&& request_context.path() == "/api/admin/endpoints/health/status"
|
||||
|
||||
@@ -242,6 +242,40 @@ impl<'a> AdminAppState<'a> {
|
||||
.await
|
||||
}
|
||||
|
||||
pub(crate) async fn build_model_health_monitor_payload(
|
||||
&self,
|
||||
lookback_hours: u64,
|
||||
model_limit: usize,
|
||||
per_model_limit: usize,
|
||||
options: crate::handlers::public::ModelHealthMonitorOptions,
|
||||
) -> Option<serde_json::Value> {
|
||||
crate::handlers::public::build_model_health_monitor_payload(
|
||||
self.app,
|
||||
lookback_hours,
|
||||
model_limit,
|
||||
per_model_limit,
|
||||
options,
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
pub(crate) async fn build_provider_health_monitor_payload(
|
||||
&self,
|
||||
lookback_hours: u64,
|
||||
provider_limit: usize,
|
||||
per_provider_model_limit: usize,
|
||||
per_model_event_limit: usize,
|
||||
) -> Option<serde_json::Value> {
|
||||
crate::handlers::public::build_provider_health_monitor_payload(
|
||||
self.app,
|
||||
lookback_hours,
|
||||
provider_limit,
|
||||
per_provider_model_limit,
|
||||
per_model_event_limit,
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
pub(crate) async fn execute_execution_runtime_sync_plan(
|
||||
&self,
|
||||
trace_id: Option<&str>,
|
||||
|
||||
@@ -12,7 +12,13 @@ use aether_data_contracts::repository::candidates::{
|
||||
use aether_data_contracts::repository::global_models::{
|
||||
PublicCatalogModelListQuery, PublicCatalogModelSearchQuery, StoredPublicCatalogModel,
|
||||
};
|
||||
use aether_data_contracts::repository::provider_catalog::StoredProviderCatalogKey;
|
||||
use aether_data_contracts::repository::provider_catalog::{
|
||||
StoredProviderCatalogKey, StoredProviderCatalogProvider,
|
||||
};
|
||||
use aether_data_contracts::repository::usage::{
|
||||
StoredRequestUsageAudit, StoredUsageBreakdownSummaryRow, UsageAuditListQuery,
|
||||
UsageBreakdownGroupBy, UsageBreakdownSummaryQuery,
|
||||
};
|
||||
use serde_json::json;
|
||||
use std::collections::{BTreeMap, BTreeSet};
|
||||
use std::time::{SystemTime, UNIX_EPOCH};
|
||||
@@ -91,6 +97,13 @@ pub(crate) struct ApiFormatHealthMonitorOptions {
|
||||
pub(crate) include_key_count: bool,
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy)]
|
||||
pub(crate) struct ModelHealthMonitorOptions {
|
||||
pub(crate) include_provider_count: bool,
|
||||
}
|
||||
|
||||
const MODEL_HEALTH_TIMELINE_SEGMENTS: u32 = 60;
|
||||
|
||||
pub(crate) fn provider_key_api_formats(key: &StoredProviderCatalogKey) -> Vec<String> {
|
||||
provider_key_configured_api_formats(key)
|
||||
}
|
||||
@@ -540,6 +553,507 @@ pub(crate) async fn build_api_format_health_monitor_payload(
|
||||
}))
|
||||
}
|
||||
|
||||
pub(crate) async fn build_model_health_monitor_payload(
|
||||
state: &AppState,
|
||||
lookback_hours: u64,
|
||||
model_limit: usize,
|
||||
per_model_limit: usize,
|
||||
options: ModelHealthMonitorOptions,
|
||||
) -> Option<serde_json::Value> {
|
||||
if !state.has_usage_data_reader() {
|
||||
return None;
|
||||
}
|
||||
|
||||
let now_unix_secs = SystemTime::now()
|
||||
.duration_since(UNIX_EPOCH)
|
||||
.map(|duration| duration.as_secs())
|
||||
.unwrap_or_default();
|
||||
let since_unix_secs = now_unix_secs.saturating_sub(lookback_hours * 3600);
|
||||
|
||||
let breakdown = state
|
||||
.summarize_usage_breakdown(&UsageBreakdownSummaryQuery {
|
||||
created_from_unix_secs: since_unix_secs,
|
||||
created_until_unix_secs: now_unix_secs,
|
||||
user_id: None,
|
||||
provider_name: None,
|
||||
group_by: UsageBreakdownGroupBy::Model,
|
||||
})
|
||||
.await
|
||||
.ok()
|
||||
.unwrap_or_default();
|
||||
|
||||
let selected_models = breakdown
|
||||
.into_iter()
|
||||
.filter(|row| !row.group_key.trim().is_empty())
|
||||
.take(model_limit)
|
||||
.collect::<Vec<_>>();
|
||||
|
||||
let mut models = Vec::with_capacity(selected_models.len());
|
||||
for row in selected_models {
|
||||
let events = state
|
||||
.list_usage_audits(&UsageAuditListQuery {
|
||||
created_from_unix_secs: Some(since_unix_secs),
|
||||
created_until_unix_secs: Some(now_unix_secs),
|
||||
model: Some(row.group_key.clone()),
|
||||
limit: Some(per_model_limit),
|
||||
newest_first: true,
|
||||
..UsageAuditListQuery::default()
|
||||
})
|
||||
.await
|
||||
.ok()
|
||||
.unwrap_or_default();
|
||||
|
||||
let (timeline, time_range_start, time_range_end) = build_model_health_timeline(
|
||||
&events,
|
||||
since_unix_secs,
|
||||
now_unix_secs,
|
||||
MODEL_HEALTH_TIMELINE_SEGMENTS,
|
||||
);
|
||||
let provider_count = model_health_provider_count(&events);
|
||||
let first_byte_average = model_health_average_first_byte_ms(&events);
|
||||
let last_event_at = events
|
||||
.first()
|
||||
.and_then(|item| unix_secs_to_rfc3339(item.created_at_unix_ms));
|
||||
let event_payload = events
|
||||
.iter()
|
||||
.rev()
|
||||
.map(model_health_event_payload)
|
||||
.collect::<Vec<_>>();
|
||||
|
||||
let total_attempts = row.request_count;
|
||||
let success_count = row.success_count.min(total_attempts);
|
||||
let failed_count = total_attempts.saturating_sub(success_count);
|
||||
let success_rate = if total_attempts > 0 {
|
||||
success_count as f64 / total_attempts as f64
|
||||
} else {
|
||||
1.0
|
||||
};
|
||||
let avg_latency_ms = model_health_average_latency_ms(&row);
|
||||
|
||||
let model_name = row.group_key.clone();
|
||||
let mut model_payload = json!({
|
||||
"model": model_name,
|
||||
"display_name": model_health_display_name(&row.group_key),
|
||||
"total_attempts": total_attempts,
|
||||
"success_count": success_count,
|
||||
"failed_count": failed_count,
|
||||
"success_rate": success_rate,
|
||||
"avg_latency_ms": avg_latency_ms,
|
||||
"avg_first_byte_ms": first_byte_average,
|
||||
"last_event_at": last_event_at,
|
||||
"events": event_payload,
|
||||
"timeline": timeline,
|
||||
"time_range_start": unix_secs_to_rfc3339(time_range_start),
|
||||
"time_range_end": unix_secs_to_rfc3339(time_range_end),
|
||||
});
|
||||
if options.include_provider_count {
|
||||
model_payload["provider_count"] = json!(provider_count);
|
||||
}
|
||||
models.push(model_payload);
|
||||
}
|
||||
|
||||
Some(json!({
|
||||
"generated_at": unix_secs_to_rfc3339(now_unix_secs),
|
||||
"models": models,
|
||||
}))
|
||||
}
|
||||
|
||||
pub(crate) async fn build_provider_health_monitor_payload(
|
||||
state: &AppState,
|
||||
lookback_hours: u64,
|
||||
provider_limit: usize,
|
||||
per_provider_model_limit: usize,
|
||||
per_model_event_limit: usize,
|
||||
) -> Option<serde_json::Value> {
|
||||
if !state.has_provider_catalog_data_reader() || !state.has_usage_data_reader() {
|
||||
return None;
|
||||
}
|
||||
|
||||
let now_unix_secs = SystemTime::now()
|
||||
.duration_since(UNIX_EPOCH)
|
||||
.map(|duration| duration.as_secs())
|
||||
.unwrap_or_default();
|
||||
let since_unix_secs = now_unix_secs.saturating_sub(lookback_hours * 3600);
|
||||
|
||||
let providers = state
|
||||
.list_provider_catalog_providers(true)
|
||||
.await
|
||||
.ok()
|
||||
.unwrap_or_default()
|
||||
.into_iter()
|
||||
.filter(|provider| provider.is_active)
|
||||
.take(provider_limit)
|
||||
.collect::<Vec<_>>();
|
||||
|
||||
let provider_breakdown = state
|
||||
.summarize_usage_breakdown(&UsageBreakdownSummaryQuery {
|
||||
created_from_unix_secs: since_unix_secs,
|
||||
created_until_unix_secs: now_unix_secs,
|
||||
user_id: None,
|
||||
provider_name: None,
|
||||
group_by: UsageBreakdownGroupBy::Provider,
|
||||
})
|
||||
.await
|
||||
.ok()
|
||||
.unwrap_or_default()
|
||||
.into_iter()
|
||||
.map(|row| (row.group_key.clone(), row))
|
||||
.collect::<BTreeMap<_, _>>();
|
||||
|
||||
let mut payload = Vec::with_capacity(providers.len());
|
||||
for provider in providers {
|
||||
let provider_stats = provider_breakdown.get(&provider.name);
|
||||
payload.push(
|
||||
build_provider_health_payload(
|
||||
state,
|
||||
provider,
|
||||
provider_stats,
|
||||
since_unix_secs,
|
||||
now_unix_secs,
|
||||
per_provider_model_limit,
|
||||
per_model_event_limit,
|
||||
)
|
||||
.await,
|
||||
);
|
||||
}
|
||||
|
||||
payload.sort_by(|left, right| {
|
||||
let left_rank = provider_health_sort_rank(left);
|
||||
let right_rank = provider_health_sort_rank(right);
|
||||
left_rank.cmp(&right_rank).then_with(|| {
|
||||
left.get("provider_name")
|
||||
.and_then(serde_json::Value::as_str)
|
||||
.unwrap_or_default()
|
||||
.cmp(
|
||||
right
|
||||
.get("provider_name")
|
||||
.and_then(serde_json::Value::as_str)
|
||||
.unwrap_or_default(),
|
||||
)
|
||||
})
|
||||
});
|
||||
|
||||
Some(json!({
|
||||
"generated_at": unix_secs_to_rfc3339(now_unix_secs),
|
||||
"providers": payload,
|
||||
}))
|
||||
}
|
||||
|
||||
async fn build_provider_health_payload(
|
||||
state: &AppState,
|
||||
provider: StoredProviderCatalogProvider,
|
||||
provider_stats: Option<&StoredUsageBreakdownSummaryRow>,
|
||||
since_unix_secs: u64,
|
||||
now_unix_secs: u64,
|
||||
per_provider_model_limit: usize,
|
||||
per_model_event_limit: usize,
|
||||
) -> serde_json::Value {
|
||||
let model_breakdown = state
|
||||
.summarize_usage_breakdown(&UsageBreakdownSummaryQuery {
|
||||
created_from_unix_secs: since_unix_secs,
|
||||
created_until_unix_secs: now_unix_secs,
|
||||
user_id: None,
|
||||
provider_name: Some(provider.name.clone()),
|
||||
group_by: UsageBreakdownGroupBy::Model,
|
||||
})
|
||||
.await
|
||||
.ok()
|
||||
.unwrap_or_default();
|
||||
|
||||
let provider_event_limit = per_model_event_limit
|
||||
.saturating_mul(per_provider_model_limit.max(1))
|
||||
.max(per_model_event_limit);
|
||||
let provider_events = state
|
||||
.list_usage_audits(&UsageAuditListQuery {
|
||||
created_from_unix_secs: Some(since_unix_secs),
|
||||
created_until_unix_secs: Some(now_unix_secs),
|
||||
provider_name: Some(provider.name.clone()),
|
||||
limit: Some(provider_event_limit),
|
||||
newest_first: true,
|
||||
..UsageAuditListQuery::default()
|
||||
})
|
||||
.await
|
||||
.ok()
|
||||
.unwrap_or_default();
|
||||
let mut models = Vec::new();
|
||||
for row in model_breakdown
|
||||
.iter()
|
||||
.filter(|row| !row.group_key.trim().is_empty())
|
||||
.take(per_provider_model_limit)
|
||||
{
|
||||
let events = state
|
||||
.list_usage_audits(&UsageAuditListQuery {
|
||||
created_from_unix_secs: Some(since_unix_secs),
|
||||
created_until_unix_secs: Some(now_unix_secs),
|
||||
provider_name: Some(provider.name.clone()),
|
||||
model: Some(row.group_key.clone()),
|
||||
limit: Some(per_model_event_limit),
|
||||
newest_first: true,
|
||||
..UsageAuditListQuery::default()
|
||||
})
|
||||
.await
|
||||
.ok()
|
||||
.unwrap_or_default();
|
||||
models.push(model_health_payload_from_row(
|
||||
row,
|
||||
&events,
|
||||
since_unix_secs,
|
||||
now_unix_secs,
|
||||
None,
|
||||
));
|
||||
}
|
||||
|
||||
let (total_attempts, success_count, failed_count, success_rate, avg_latency_ms) =
|
||||
if let Some(row) = provider_stats {
|
||||
let total_attempts = row.request_count;
|
||||
let success_count = row.success_count.min(total_attempts);
|
||||
let failed_count = total_attempts.saturating_sub(success_count);
|
||||
let success_rate = if total_attempts > 0 {
|
||||
success_count as f64 / total_attempts as f64
|
||||
} else {
|
||||
1.0
|
||||
};
|
||||
(
|
||||
total_attempts,
|
||||
success_count,
|
||||
failed_count,
|
||||
success_rate,
|
||||
model_health_average_latency_ms(row),
|
||||
)
|
||||
} else {
|
||||
(0, 0, 0, 1.0, None)
|
||||
};
|
||||
|
||||
let (timeline, time_range_start, time_range_end) = build_model_health_timeline(
|
||||
&provider_events,
|
||||
since_unix_secs,
|
||||
now_unix_secs,
|
||||
MODEL_HEALTH_TIMELINE_SEGMENTS,
|
||||
);
|
||||
let last_event_at = provider_events
|
||||
.iter()
|
||||
.max_by_key(|event| event.created_at_unix_ms)
|
||||
.and_then(|event| unix_secs_to_rfc3339(event.created_at_unix_ms));
|
||||
|
||||
json!({
|
||||
"provider_id": provider.id,
|
||||
"provider_name": provider.name,
|
||||
"provider_type": provider.provider_type,
|
||||
"is_active": provider.is_active,
|
||||
"total_attempts": total_attempts,
|
||||
"success_count": success_count,
|
||||
"failed_count": failed_count,
|
||||
"success_rate": success_rate,
|
||||
"avg_latency_ms": avg_latency_ms,
|
||||
"avg_first_byte_ms": model_health_average_first_byte_ms(&provider_events),
|
||||
"model_count": model_breakdown.len(),
|
||||
"last_event_at": last_event_at,
|
||||
"timeline": timeline,
|
||||
"time_range_start": unix_secs_to_rfc3339(time_range_start),
|
||||
"time_range_end": unix_secs_to_rfc3339(time_range_end),
|
||||
"models": models,
|
||||
})
|
||||
}
|
||||
|
||||
fn provider_health_sort_rank(provider: &serde_json::Value) -> u8 {
|
||||
let total_attempts = provider
|
||||
.get("total_attempts")
|
||||
.and_then(serde_json::Value::as_u64)
|
||||
.unwrap_or(0);
|
||||
if total_attempts == 0 {
|
||||
return 3;
|
||||
}
|
||||
let success_rate = provider
|
||||
.get("success_rate")
|
||||
.and_then(serde_json::Value::as_f64)
|
||||
.unwrap_or(1.0);
|
||||
if success_rate < 0.8 {
|
||||
0
|
||||
} else if success_rate < 0.95 {
|
||||
1
|
||||
} else {
|
||||
2
|
||||
}
|
||||
}
|
||||
|
||||
fn model_health_payload_from_row(
|
||||
row: &StoredUsageBreakdownSummaryRow,
|
||||
events: &[StoredRequestUsageAudit],
|
||||
since_unix_secs: u64,
|
||||
now_unix_secs: u64,
|
||||
provider_count: Option<usize>,
|
||||
) -> serde_json::Value {
|
||||
let (timeline, time_range_start, time_range_end) = build_model_health_timeline(
|
||||
events,
|
||||
since_unix_secs,
|
||||
now_unix_secs,
|
||||
MODEL_HEALTH_TIMELINE_SEGMENTS,
|
||||
);
|
||||
let last_event_at = events
|
||||
.first()
|
||||
.and_then(|item| unix_secs_to_rfc3339(item.created_at_unix_ms));
|
||||
let event_payload = events
|
||||
.iter()
|
||||
.rev()
|
||||
.map(model_health_event_payload)
|
||||
.collect::<Vec<_>>();
|
||||
|
||||
let total_attempts = row.request_count;
|
||||
let success_count = row.success_count.min(total_attempts);
|
||||
let failed_count = total_attempts.saturating_sub(success_count);
|
||||
let success_rate = if total_attempts > 0 {
|
||||
success_count as f64 / total_attempts as f64
|
||||
} else {
|
||||
1.0
|
||||
};
|
||||
|
||||
let mut model_payload = json!({
|
||||
"model": row.group_key.clone(),
|
||||
"display_name": model_health_display_name(&row.group_key),
|
||||
"total_attempts": total_attempts,
|
||||
"success_count": success_count,
|
||||
"failed_count": failed_count,
|
||||
"success_rate": success_rate,
|
||||
"avg_latency_ms": model_health_average_latency_ms(row),
|
||||
"avg_first_byte_ms": model_health_average_first_byte_ms(events),
|
||||
"last_event_at": last_event_at,
|
||||
"events": event_payload,
|
||||
"timeline": timeline,
|
||||
"time_range_start": unix_secs_to_rfc3339(time_range_start),
|
||||
"time_range_end": unix_secs_to_rfc3339(time_range_end),
|
||||
});
|
||||
if let Some(provider_count) = provider_count {
|
||||
model_payload["provider_count"] = json!(provider_count);
|
||||
}
|
||||
model_payload
|
||||
}
|
||||
|
||||
fn model_health_average_latency_ms(row: &StoredUsageBreakdownSummaryRow) -> Option<f64> {
|
||||
if row.overall_response_time_samples > 0 {
|
||||
return Some(row.overall_response_time_sum_ms / row.overall_response_time_samples as f64);
|
||||
}
|
||||
if row.response_time_samples > 0 {
|
||||
return Some(row.response_time_sum_ms / row.response_time_samples as f64);
|
||||
}
|
||||
None
|
||||
}
|
||||
|
||||
fn model_health_average_first_byte_ms(events: &[StoredRequestUsageAudit]) -> Option<f64> {
|
||||
let mut sum = 0u64;
|
||||
let mut count = 0u64;
|
||||
for event in events {
|
||||
if let Some(first_byte_time_ms) = event.first_byte_time_ms {
|
||||
sum = sum.saturating_add(first_byte_time_ms);
|
||||
count = count.saturating_add(1);
|
||||
}
|
||||
}
|
||||
if count == 0 {
|
||||
None
|
||||
} else {
|
||||
Some(sum as f64 / count as f64)
|
||||
}
|
||||
}
|
||||
|
||||
fn model_health_provider_count(events: &[StoredRequestUsageAudit]) -> usize {
|
||||
let mut providers = BTreeSet::new();
|
||||
for event in events {
|
||||
if let Some(provider_id) = event.provider_id.as_deref() {
|
||||
providers.insert(provider_id.to_string());
|
||||
} else if !event.provider_name.trim().is_empty() {
|
||||
providers.insert(event.provider_name.clone());
|
||||
}
|
||||
}
|
||||
providers.len()
|
||||
}
|
||||
|
||||
fn model_health_event_payload(event: &StoredRequestUsageAudit) -> serde_json::Value {
|
||||
json!({
|
||||
"timestamp": unix_secs_to_rfc3339(event.created_at_unix_ms),
|
||||
"status": model_health_event_status(event),
|
||||
"status_code": event.status_code,
|
||||
"latency_ms": event.response_time_ms,
|
||||
"first_byte_time_ms": event.first_byte_time_ms,
|
||||
"error_type": event.error_category,
|
||||
})
|
||||
}
|
||||
|
||||
fn model_health_event_status(event: &StoredRequestUsageAudit) -> &'static str {
|
||||
if model_health_event_success(event) {
|
||||
"success"
|
||||
} else {
|
||||
"failed"
|
||||
}
|
||||
}
|
||||
|
||||
fn model_health_event_success(event: &StoredRequestUsageAudit) -> bool {
|
||||
!event.status.eq_ignore_ascii_case("failed")
|
||||
&& event.status_code.is_none_or(|status| status < 400)
|
||||
&& event
|
||||
.error_message
|
||||
.as_deref()
|
||||
.map(str::trim)
|
||||
.unwrap_or_default()
|
||||
.is_empty()
|
||||
}
|
||||
|
||||
fn build_model_health_timeline(
|
||||
events: &[StoredRequestUsageAudit],
|
||||
since_unix_secs: u64,
|
||||
until_unix_secs: u64,
|
||||
segments: u32,
|
||||
) -> (Vec<&'static str>, u64, u64) {
|
||||
#[derive(Default)]
|
||||
struct Bucket {
|
||||
success_count: u64,
|
||||
failed_count: u64,
|
||||
}
|
||||
|
||||
let safe_range = until_unix_secs.saturating_sub(since_unix_secs).max(1);
|
||||
let mut buckets = (0..segments).map(|_| Bucket::default()).collect::<Vec<_>>();
|
||||
|
||||
for event in events {
|
||||
let timestamp = event.created_at_unix_ms;
|
||||
if timestamp < since_unix_secs || timestamp > until_unix_secs {
|
||||
continue;
|
||||
}
|
||||
let offset = timestamp.saturating_sub(since_unix_secs);
|
||||
let mut segment_idx = ((offset as u128 * segments as u128) / safe_range as u128) as usize;
|
||||
if segment_idx >= segments as usize {
|
||||
segment_idx = segments.saturating_sub(1) as usize;
|
||||
}
|
||||
let bucket = &mut buckets[segment_idx];
|
||||
if model_health_event_success(event) {
|
||||
bucket.success_count = bucket.success_count.saturating_add(1);
|
||||
} else {
|
||||
bucket.failed_count = bucket.failed_count.saturating_add(1);
|
||||
}
|
||||
}
|
||||
|
||||
let timeline = buckets
|
||||
.into_iter()
|
||||
.map(|bucket| {
|
||||
let total = bucket.success_count.saturating_add(bucket.failed_count);
|
||||
if total == 0 {
|
||||
return "unknown";
|
||||
}
|
||||
let success_rate = bucket.success_count as f64 / total as f64;
|
||||
if success_rate >= 0.95 {
|
||||
"healthy"
|
||||
} else if success_rate >= 0.7 {
|
||||
"warning"
|
||||
} else {
|
||||
"unhealthy"
|
||||
}
|
||||
})
|
||||
.collect::<Vec<_>>();
|
||||
|
||||
(timeline, since_unix_secs, until_unix_secs)
|
||||
}
|
||||
|
||||
fn model_health_display_name(model: &str) -> String {
|
||||
model.trim().to_string()
|
||||
}
|
||||
|
||||
pub(crate) fn build_public_health_timeline(
|
||||
buckets_by_segment: &BTreeMap<u32, PublicHealthTimelineBucket>,
|
||||
segments: u32,
|
||||
|
||||
@@ -8,10 +8,12 @@ pub(crate) use self::ai_public::{
|
||||
};
|
||||
pub(crate) use self::catalog_helpers::{
|
||||
admin_requested_force_stream, api_format_display_name, build_api_format_health_monitor_payload,
|
||||
build_model_health_monitor_payload, build_provider_health_monitor_payload,
|
||||
build_public_catalog_models_payload, build_public_catalog_search_models_payload,
|
||||
build_public_health_timeline, build_public_providers_payload, normalize_admin_base_url,
|
||||
provider_key_api_formats, request_candidate_event_unix_ms, request_candidate_status_label,
|
||||
sanitize_public_model_config_for_user, ApiFormatHealthMonitorOptions,
|
||||
ModelHealthMonitorOptions,
|
||||
};
|
||||
pub(crate) use self::system_modules_helpers::{
|
||||
build_admin_keys_grouped_by_format_payload, build_public_auth_modules_status_payload,
|
||||
|
||||
@@ -1,9 +1,10 @@
|
||||
use super::{
|
||||
build_api_format_health_monitor_payload, build_public_auth_modules_status_payload,
|
||||
build_public_catalog_models_payload, build_public_catalog_search_models_payload,
|
||||
build_public_providers_payload, capability_detail_by_name, ldap_module_config_is_valid,
|
||||
sanitize_public_model_config_for_user, serialize_public_capability, supported_capability_names,
|
||||
ApiFormatHealthMonitorOptions, PUBLIC_CAPABILITY_DEFINITIONS,
|
||||
build_api_format_health_monitor_payload, build_model_health_monitor_payload,
|
||||
build_public_auth_modules_status_payload, build_public_catalog_models_payload,
|
||||
build_public_catalog_search_models_payload, build_public_providers_payload,
|
||||
capability_detail_by_name, ldap_module_config_is_valid, sanitize_public_model_config_for_user,
|
||||
serialize_public_capability, supported_capability_names, ApiFormatHealthMonitorOptions,
|
||||
ModelHealthMonitorOptions, PUBLIC_CAPABILITY_DEFINITIONS,
|
||||
};
|
||||
use crate::control::GatewayPublicRequestContext;
|
||||
use crate::handlers::shared::{
|
||||
@@ -426,6 +427,43 @@ pub(crate) async fn maybe_build_local_public_support_response(
|
||||
.await?;
|
||||
return Some(Json(payload).into_response());
|
||||
}
|
||||
|
||||
if decision.route_kind.as_deref() == Some("health_models")
|
||||
&& request_context.request_path == "/api/public/health/models"
|
||||
{
|
||||
let lookback_hours = query_param_value(
|
||||
request_context.request_query_string.as_deref(),
|
||||
"lookback_hours",
|
||||
)
|
||||
.and_then(|value| value.parse::<u64>().ok())
|
||||
.filter(|value| (1..=168).contains(value))
|
||||
.unwrap_or(6);
|
||||
let model_limit = query_param_value(
|
||||
request_context.request_query_string.as_deref(),
|
||||
"model_limit",
|
||||
)
|
||||
.and_then(|value| value.parse::<usize>().ok())
|
||||
.filter(|value| (1..=50).contains(value))
|
||||
.unwrap_or(12);
|
||||
let per_model_limit = query_param_value(
|
||||
request_context.request_query_string.as_deref(),
|
||||
"per_model_limit",
|
||||
)
|
||||
.and_then(|value| value.parse::<usize>().ok())
|
||||
.filter(|value| (10..=500).contains(value))
|
||||
.unwrap_or(100);
|
||||
let payload = build_model_health_monitor_payload(
|
||||
state,
|
||||
lookback_hours,
|
||||
model_limit,
|
||||
per_model_limit,
|
||||
ModelHealthMonitorOptions {
|
||||
include_provider_count: false,
|
||||
},
|
||||
)
|
||||
.await?;
|
||||
return Some(Json(payload).into_response());
|
||||
}
|
||||
}
|
||||
|
||||
if decision.route_family.as_deref() == Some("capabilities") {
|
||||
|
||||
@@ -847,6 +847,7 @@ pub(super) async fn handle_users_me_usage_get(
|
||||
created_from_unix_secs,
|
||||
created_until_unix_secs,
|
||||
user_id: Some(auth.user.id.clone()),
|
||||
provider_name: None,
|
||||
group_by: UsageBreakdownGroupBy::Model,
|
||||
})
|
||||
.await
|
||||
@@ -866,6 +867,7 @@ pub(super) async fn handle_users_me_usage_get(
|
||||
created_from_unix_secs,
|
||||
created_until_unix_secs,
|
||||
user_id: Some(auth.user.id.clone()),
|
||||
provider_name: None,
|
||||
group_by: UsageBreakdownGroupBy::Provider,
|
||||
})
|
||||
.await
|
||||
@@ -885,6 +887,7 @@ pub(super) async fn handle_users_me_usage_get(
|
||||
created_from_unix_secs,
|
||||
created_until_unix_secs,
|
||||
user_id: Some(auth.user.id.clone()),
|
||||
provider_name: None,
|
||||
group_by: UsageBreakdownGroupBy::ApiFormat,
|
||||
})
|
||||
.await
|
||||
|
||||
Reference in New Issue
Block a user