mirror of
https://github.com/fawney19/Aether.git
synced 2026-09-02 01:10:23 +08:00
perf: 并行化 admin 聚合路由并完善前端缓存预取
- gateway: usage detail / provider summary / pool overview / users list 改为 tokio join 并行拉取依赖数据 - usage: interval timeline 支持自动刷新并按查询区间动态展示,取消服务端 120 分钟过滤并在 ScatterChart 统一封顶 - frontend: 新增管理端导航预取工具及 SidebarNav/MainLayout 触发,admin 读接口统一走 cachedRequest 的短期缓存 - dashboard: request detail 支持短 TTL 缓存并在 UsageRecordsTable mousedown 时预取 - data: migrate 测试在 wait_for_postgres 失败时清理子进程,避免遗留
This commit is contained in:
@@ -65,9 +65,6 @@ pub(super) async fn build_admin_usage_cache_affinity_interval_timeline_response(
|
||||
let mut usernames_by_user_id = BTreeMap::new();
|
||||
|
||||
for row in intervals {
|
||||
if row.interval_minutes > 120.0 {
|
||||
continue;
|
||||
}
|
||||
let mut point = json!({
|
||||
"x": unix_secs_to_rfc3339(row.created_at_unix_secs),
|
||||
"y": ((row.interval_minutes * 100.0).round()) / 100.0,
|
||||
|
||||
@@ -23,6 +23,7 @@ use axum::{
|
||||
};
|
||||
use serde_json::json;
|
||||
use std::collections::BTreeMap;
|
||||
use tokio::try_join;
|
||||
|
||||
pub(super) async fn maybe_build_local_admin_usage_detail_response(
|
||||
state: &AdminAppState<'_>,
|
||||
@@ -160,43 +161,50 @@ pub(super) async fn maybe_build_local_admin_usage_detail_response(
|
||||
));
|
||||
};
|
||||
|
||||
let users_by_id: BTreeMap<String, aether_data::repository::users::StoredUserSummary> =
|
||||
state
|
||||
.resolve_auth_user_summaries_by_ids(
|
||||
&item.user_id.clone().into_iter().collect::<Vec<_>>(),
|
||||
)
|
||||
.await?;
|
||||
let provider_key_names =
|
||||
admin_usage_provider_key_names(state, std::slice::from_ref(&item)).await?;
|
||||
let api_key_names =
|
||||
admin_usage_api_key_names(state, std::slice::from_ref(&item)).await?;
|
||||
let user_ids = item.user_id.clone().into_iter().collect::<Vec<_>>();
|
||||
let (users_by_id, provider_key_names, api_key_names): (
|
||||
BTreeMap<String, aether_data::repository::users::StoredUserSummary>,
|
||||
BTreeMap<String, String>,
|
||||
BTreeMap<String, String>,
|
||||
) = try_join!(
|
||||
state.resolve_auth_user_summaries_by_ids(&user_ids),
|
||||
admin_usage_provider_key_names(state, std::slice::from_ref(&item)),
|
||||
admin_usage_api_key_names(state, std::slice::from_ref(&item)),
|
||||
)?;
|
||||
let provider_key_name = admin_usage_provider_key_name(&item, &provider_key_names);
|
||||
|
||||
let request_body =
|
||||
admin_usage_resolve_request_capture_body_for_item(state, &item, None).await?;
|
||||
let mut detail_item = item.clone();
|
||||
let request_body = if include_bodies {
|
||||
let (request_body, provider_request_body, response_body, client_response_body) = try_join!(
|
||||
admin_usage_resolve_request_capture_body_for_item(state, &item, None),
|
||||
admin_usage_resolve_body_value(
|
||||
state,
|
||||
&item,
|
||||
item.provider_request_body.as_ref(),
|
||||
UsageBodyField::ProviderRequestBody,
|
||||
),
|
||||
admin_usage_resolve_body_value(
|
||||
state,
|
||||
&item,
|
||||
item.response_body.as_ref(),
|
||||
UsageBodyField::ResponseBody,
|
||||
),
|
||||
admin_usage_resolve_body_value(
|
||||
state,
|
||||
&item,
|
||||
item.client_response_body.as_ref(),
|
||||
UsageBodyField::ClientResponseBody,
|
||||
),
|
||||
)?;
|
||||
detail_item.provider_request_body = provider_request_body;
|
||||
detail_item.response_body = response_body;
|
||||
detail_item.client_response_body = client_response_body;
|
||||
request_body
|
||||
} else {
|
||||
None
|
||||
};
|
||||
if include_bodies {
|
||||
detail_item.provider_request_body = admin_usage_resolve_body_value(
|
||||
state,
|
||||
&item,
|
||||
item.provider_request_body.as_ref(),
|
||||
UsageBodyField::ProviderRequestBody,
|
||||
)
|
||||
.await?;
|
||||
detail_item.response_body = admin_usage_resolve_body_value(
|
||||
state,
|
||||
&item,
|
||||
item.response_body.as_ref(),
|
||||
UsageBodyField::ResponseBody,
|
||||
)
|
||||
.await?;
|
||||
detail_item.client_response_body = admin_usage_resolve_body_value(
|
||||
state,
|
||||
&item,
|
||||
item.client_response_body.as_ref(),
|
||||
UsageBodyField::ClientResponseBody,
|
||||
)
|
||||
.await?;
|
||||
// request_body 已通过 request capture 解析;其余 detached body 在上方并行加载。
|
||||
}
|
||||
let default_headers = admin_usage_curl_headers();
|
||||
let payload = build_admin_usage_detail_payload(
|
||||
|
||||
@@ -35,24 +35,31 @@ pub(super) async fn build_admin_pool_overview_response(
|
||||
.iter()
|
||||
.map(|(provider, _)| provider.id.clone())
|
||||
.collect::<Vec<_>>();
|
||||
let key_stats = if provider_ids.is_empty() {
|
||||
Vec::new()
|
||||
} else {
|
||||
state
|
||||
.list_provider_catalog_key_stats_by_provider_ids(&provider_ids)
|
||||
.await?
|
||||
};
|
||||
let redis_runner = state.redis_kv_runner();
|
||||
let (key_stats_result, cooldown_counts_by_provider) = tokio::join!(
|
||||
async {
|
||||
if provider_ids.is_empty() {
|
||||
Ok(Vec::new())
|
||||
} else {
|
||||
state
|
||||
.list_provider_catalog_key_stats_by_provider_ids(&provider_ids)
|
||||
.await
|
||||
}
|
||||
},
|
||||
async {
|
||||
match redis_runner.as_ref() {
|
||||
Some(runner) if !provider_ids.is_empty() => {
|
||||
read_admin_provider_pool_cooldown_counts(runner, &provider_ids).await
|
||||
}
|
||||
_ => BTreeMap::new(),
|
||||
}
|
||||
},
|
||||
);
|
||||
let key_stats = key_stats_result?;
|
||||
let key_stats_by_provider = key_stats
|
||||
.into_iter()
|
||||
.map(|item| (item.provider_id.clone(), item))
|
||||
.collect::<BTreeMap<_, _>>();
|
||||
let redis_runner = state.redis_kv_runner();
|
||||
let cooldown_counts_by_provider = match redis_runner.as_ref() {
|
||||
Some(runner) if !provider_ids.is_empty() => {
|
||||
read_admin_provider_pool_cooldown_counts(runner, &provider_ids).await
|
||||
}
|
||||
_ => BTreeMap::new(),
|
||||
};
|
||||
|
||||
let providers = pool_enabled_providers
|
||||
.into_iter()
|
||||
|
||||
@@ -3,6 +3,8 @@ use crate::handlers::admin::request::AdminAppState;
|
||||
use aether_data_contracts::repository::provider_catalog::{
|
||||
StoredProviderCatalogEndpoint, StoredProviderCatalogKey,
|
||||
};
|
||||
use aether_data_contracts::repository::quota::StoredProviderQuotaSnapshot;
|
||||
use futures_util::future::join_all;
|
||||
use serde_json::json;
|
||||
use std::collections::{BTreeMap, BTreeSet};
|
||||
use std::time::{SystemTime, UNIX_EPOCH};
|
||||
@@ -90,6 +92,9 @@ pub(crate) async fn build_admin_providers_summary_payload(
|
||||
let normalized_status = status.trim().to_ascii_lowercase();
|
||||
let normalized_api_format = api_format.trim();
|
||||
let normalized_model_id = model_id.trim();
|
||||
let requires_api_format_filter =
|
||||
normalized_api_format != "all" && !normalized_api_format.is_empty();
|
||||
let requires_model_filter = normalized_model_id != "all" && !normalized_model_id.is_empty();
|
||||
|
||||
let mut providers = state
|
||||
.list_provider_catalog_providers(false)
|
||||
@@ -100,7 +105,7 @@ pub(crate) async fn build_admin_providers_summary_payload(
|
||||
.iter()
|
||||
.map(|provider| provider.id.clone())
|
||||
.collect::<Vec<_>>();
|
||||
let all_endpoints = if all_provider_ids.is_empty() {
|
||||
let all_endpoints = if !requires_api_format_filter || all_provider_ids.is_empty() {
|
||||
Vec::new()
|
||||
} else {
|
||||
state
|
||||
@@ -109,7 +114,7 @@ pub(crate) async fn build_admin_providers_summary_payload(
|
||||
.ok()
|
||||
.unwrap_or_default()
|
||||
};
|
||||
let active_global_model_refs = if all_provider_ids.is_empty() {
|
||||
let active_global_model_refs = if !requires_model_filter || all_provider_ids.is_empty() {
|
||||
Vec::new()
|
||||
} else {
|
||||
state
|
||||
@@ -151,8 +156,7 @@ pub(crate) async fn build_admin_providers_summary_payload(
|
||||
_ => {}
|
||||
}
|
||||
|
||||
if normalized_api_format != "all"
|
||||
&& !normalized_api_format.is_empty()
|
||||
if requires_api_format_filter
|
||||
&& !api_formats_by_provider
|
||||
.get(&provider.id)
|
||||
.is_some_and(|items| items.contains(normalized_api_format))
|
||||
@@ -160,8 +164,7 @@ pub(crate) async fn build_admin_providers_summary_payload(
|
||||
return false;
|
||||
}
|
||||
|
||||
if normalized_model_id != "all"
|
||||
&& !normalized_model_id.is_empty()
|
||||
if requires_model_filter
|
||||
&& !active_global_model_ids_by_provider
|
||||
.get(&provider.id)
|
||||
.is_some_and(|items| items.contains(normalized_model_id))
|
||||
@@ -191,32 +194,21 @@ pub(crate) async fn build_admin_providers_summary_payload(
|
||||
.iter()
|
||||
.map(|provider| provider.id.clone())
|
||||
.collect::<Vec<_>>();
|
||||
let endpoints = if provider_ids.is_empty() {
|
||||
Vec::new()
|
||||
let (endpoints, keys, model_stats, page_active_global_model_refs) = if provider_ids.is_empty() {
|
||||
(Vec::new(), Vec::new(), Vec::new(), Vec::new())
|
||||
} else {
|
||||
state
|
||||
.list_provider_catalog_endpoints_by_provider_ids(&provider_ids)
|
||||
.await
|
||||
.ok()
|
||||
.unwrap_or_default()
|
||||
};
|
||||
let keys = if provider_ids.is_empty() {
|
||||
Vec::new()
|
||||
} else {
|
||||
state
|
||||
.list_provider_catalog_key_summaries_by_provider_ids(&provider_ids)
|
||||
.await
|
||||
.ok()
|
||||
.unwrap_or_default()
|
||||
};
|
||||
let model_stats = if provider_ids.is_empty() {
|
||||
Vec::new()
|
||||
} else {
|
||||
state
|
||||
.list_provider_model_stats(&provider_ids)
|
||||
.await
|
||||
.ok()
|
||||
.unwrap_or_default()
|
||||
let (endpoints_result, keys_result, model_stats_result, active_global_model_refs_result) = tokio::join!(
|
||||
state.list_provider_catalog_endpoints_by_provider_ids(&provider_ids),
|
||||
state.list_provider_catalog_key_summaries_by_provider_ids(&provider_ids),
|
||||
state.list_provider_model_stats(&provider_ids),
|
||||
state.list_active_global_model_ids_by_provider_ids(&provider_ids),
|
||||
);
|
||||
(
|
||||
endpoints_result.ok().unwrap_or_default(),
|
||||
keys_result.ok().unwrap_or_default(),
|
||||
model_stats_result.ok().unwrap_or_default(),
|
||||
active_global_model_refs_result.ok().unwrap_or_default(),
|
||||
)
|
||||
};
|
||||
let mut endpoints_by_provider = BTreeMap::<String, Vec<StoredProviderCatalogEndpoint>>::new();
|
||||
for endpoint in endpoints {
|
||||
@@ -236,17 +228,30 @@ pub(crate) async fn build_admin_providers_summary_payload(
|
||||
.into_iter()
|
||||
.map(|stats| (stats.provider_id.clone(), stats))
|
||||
.collect::<BTreeMap<_, _>>();
|
||||
let mut active_global_model_ids_by_provider = BTreeMap::<String, BTreeSet<String>>::new();
|
||||
for row in page_active_global_model_refs {
|
||||
active_global_model_ids_by_provider
|
||||
.entry(row.provider_id)
|
||||
.or_default()
|
||||
.insert(row.global_model_id);
|
||||
}
|
||||
let quota_snapshots_by_provider = join_all(provider_ids.iter().map(|provider_id| async {
|
||||
let quota_snapshot = state
|
||||
.read_provider_quota_snapshot(provider_id)
|
||||
.await
|
||||
.ok()
|
||||
.flatten();
|
||||
(provider_id.clone(), quota_snapshot)
|
||||
}))
|
||||
.await
|
||||
.into_iter()
|
||||
.collect::<BTreeMap<String, Option<StoredProviderQuotaSnapshot>>>();
|
||||
let now_unix_secs = SystemTime::now()
|
||||
.duration_since(UNIX_EPOCH)
|
||||
.unwrap_or_default()
|
||||
.as_secs();
|
||||
let mut items = Vec::with_capacity(providers.len());
|
||||
for provider in providers {
|
||||
let quota_snapshot = state
|
||||
.read_provider_quota_snapshot(&provider.id)
|
||||
.await
|
||||
.ok()
|
||||
.flatten();
|
||||
let active_global_model_ids = active_global_model_ids_by_provider
|
||||
.get(&provider.id)
|
||||
.cloned()
|
||||
@@ -263,7 +268,9 @@ pub(crate) async fn build_admin_providers_summary_payload(
|
||||
.get(&provider.id)
|
||||
.map(Vec::as_slice)
|
||||
.unwrap_or(&[]),
|
||||
quota_snapshot.as_ref(),
|
||||
quota_snapshots_by_provider
|
||||
.get(&provider.id)
|
||||
.and_then(Option::as_ref),
|
||||
model_stats_by_provider.get(&provider.id),
|
||||
active_global_model_ids,
|
||||
now_unix_secs,
|
||||
|
||||
@@ -42,21 +42,20 @@ pub(in super::super) async fn build_admin_list_users_response(
|
||||
.iter()
|
||||
.map(|row| row.id.clone())
|
||||
.collect::<Vec<_>>();
|
||||
let auth_by_user_id = state
|
||||
.list_user_auth_by_ids(&user_ids)
|
||||
.await?
|
||||
let (auth_rows_result, wallet_rows_result, usage_totals_result) = tokio::join!(
|
||||
state.list_user_auth_by_ids(&user_ids),
|
||||
state.list_wallet_snapshots_by_user_ids(&user_ids),
|
||||
state.summarize_usage_totals_by_user_ids(&user_ids),
|
||||
);
|
||||
let auth_by_user_id = auth_rows_result?
|
||||
.into_iter()
|
||||
.map(|user| (user.id.clone(), user))
|
||||
.collect::<BTreeMap<_, _>>();
|
||||
let wallet_by_user_id = state
|
||||
.list_wallet_snapshots_by_user_ids(&user_ids)
|
||||
.await?
|
||||
let wallet_by_user_id = wallet_rows_result?
|
||||
.into_iter()
|
||||
.filter_map(|wallet| wallet.user_id.clone().map(|user_id| (user_id, wallet)))
|
||||
.collect::<BTreeMap<_, _>>();
|
||||
let usage_totals_by_user_id = state
|
||||
.summarize_usage_totals_by_user_ids(&user_ids)
|
||||
.await?
|
||||
let usage_totals_by_user_id = usage_totals_result?
|
||||
.into_iter()
|
||||
.map(|item| (item.user_id.clone(), item))
|
||||
.collect::<BTreeMap<_, _>>();
|
||||
|
||||
@@ -909,15 +909,13 @@ pub(super) async fn handle_users_me_usage_interval_timeline_get(
|
||||
|
||||
let mut points = Vec::new();
|
||||
for row in intervals {
|
||||
if row.interval_minutes <= 120.0 {
|
||||
points.push(json!({
|
||||
"x": unix_secs_to_rfc3339(row.created_at_unix_secs),
|
||||
"y": round_to(row.interval_minutes, 2),
|
||||
"model": row.model,
|
||||
}));
|
||||
if points.len() >= limit {
|
||||
break;
|
||||
}
|
||||
points.push(json!({
|
||||
"x": unix_secs_to_rfc3339(row.created_at_unix_secs),
|
||||
"y": round_to(row.interval_minutes, 2),
|
||||
"model": row.model,
|
||||
}));
|
||||
if points.len() >= limit {
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -5245,7 +5245,7 @@ async fn gateway_handles_users_me_usage_interval_timeline_and_heatmap_locally_wi
|
||||
"gpt-4.1",
|
||||
"OpenAI",
|
||||
"completed",
|
||||
now - chrono::Duration::days(1),
|
||||
now - chrono::Duration::days(1) - chrono::Duration::minutes(1),
|
||||
),
|
||||
]));
|
||||
let (gateway_url, upstream_hits, gateway_handle, upstream_handle) =
|
||||
|
||||
Reference in New Issue
Block a user