fix: speed up admin pool loading

This commit is contained in:
fawney19
2026-05-19 12:05:36 +08:00
parent c6a408d4e8
commit 6a1da6a5ff
18 changed files with 149 additions and 60 deletions

View File

@@ -13,6 +13,7 @@ use crate::provider_pool_demand::{
provider_pool_burst_pending, read_provider_pool_demand_snapshot,
};
use aether_runtime_state::{DataLayerError, RuntimeState};
use futures_util::future::join_all;
use std::collections::{BTreeMap, BTreeSet};
use std::time::{SystemTime, UNIX_EPOCH};
use tracing::warn;
@@ -32,15 +33,16 @@ pub(crate) async fn read_admin_provider_pool_cooldown_counts(
runtime: &RuntimeState,
provider_ids: &[String],
) -> BTreeMap<String, usize> {
let mut counts = BTreeMap::new();
for provider_id in provider_ids {
join_all(provider_ids.iter().map(|provider_id| async move {
let count = runtime
.set_len(&pool_cooldown_index_key(provider_id))
.await
.unwrap_or(0);
counts.insert(provider_id.clone(), count);
}
counts
(provider_id.clone(), count)
}))
.await
.into_iter()
.collect()
}
pub(crate) async fn read_admin_provider_pool_runtime_state(
@@ -170,11 +172,15 @@ pub(crate) async fn read_admin_provider_pool_runtime_state(
}
let now = current_unix_secs();
for (key_id, cost_key) in key_ids.iter().zip(cost_keys) {
let window_start = now.saturating_sub(pool_config.cost_window_seconds) as f64;
let total = runtime
.score_range_by_min(&cost_key, window_start)
.await
let cost_window_start = now.saturating_sub(pool_config.cost_window_seconds) as f64;
let cost_results = join_all(
cost_keys
.iter()
.map(|cost_key| runtime.score_range_by_min(cost_key, cost_window_start)),
)
.await;
for (key_id, members) in key_ids.iter().zip(cost_results) {
let total = members
.unwrap_or_default()
.iter()
.map(|member| parse_pool_cost_member(member))
@@ -184,11 +190,15 @@ pub(crate) async fn read_admin_provider_pool_runtime_state(
}
}
for (key_id, latency_key) in key_ids.iter().zip(latency_keys) {
let window_start = now.saturating_sub(pool_config.latency_window_seconds) as f64;
let samples = runtime
.score_range_by_min(&latency_key, window_start)
.await
let latency_window_start = now.saturating_sub(pool_config.latency_window_seconds) as f64;
let latency_results = join_all(
latency_keys
.iter()
.map(|latency_key| runtime.score_range_by_min(latency_key, latency_window_start)),
)
.await;
for (key_id, members) in key_ids.iter().zip(latency_results) {
let samples = members
.unwrap_or_default()
.iter()
.map(|member| parse_pool_latency_member(member))

View File

@@ -19,6 +19,7 @@ use axum::{
response::{IntoResponse, Response},
Json,
};
use futures_util::future::join_all;
use serde_json::{json, Value};
pub(super) async fn build_admin_pool_overview_response(
@@ -67,47 +68,66 @@ pub(super) async fn build_admin_pool_overview_response(
.collect::<BTreeMap<_, _>>();
let probe_config = PoolQuotaProbeWorkerConfig::from_env();
let mut runtime_metrics_by_provider = BTreeMap::new();
for (provider, pool_config) in &pool_enabled_providers {
let active_keys = key_stats_by_provider
.get(&provider.id)
.map(|item| item.active_keys as usize)
.unwrap_or(0);
let hot_count = if pool_config.probing_enabled {
state
.runtime_state()
.set_len(&admin_provider_pool_quota_probe_active_members_key(
&provider.id,
))
.await
.unwrap_or(0)
} else {
0
};
let demand_snapshot = read_provider_pool_demand_snapshot(
state.runtime_state(),
&provider.id,
active_keys,
probe_config.max_keys_per_provider,
)
.await;
let burst_pending = pool_config.probing_enabled
&& provider_pool_burst_pending(state.runtime_state(), &provider.id).await;
runtime_metrics_by_provider.insert(
provider.id.clone(),
json!({
"provider_hot_count": hot_count,
"provider_desired_hot": if pool_config.probing_enabled {
demand_snapshot.desired_hot
} else {
0
},
"provider_in_flight": demand_snapshot.in_flight,
"provider_ema_in_flight": demand_snapshot.ema_in_flight,
"provider_burst_pending": burst_pending,
}),
);
}
let runtime_metrics_by_provider = join_all(pool_enabled_providers.iter().map(
|(provider, pool_config)| {
let provider_id = provider.id.clone();
let probing_enabled = pool_config.probing_enabled;
let active_keys = key_stats_by_provider
.get(&provider.id)
.map(|item| item.active_keys as usize)
.unwrap_or(0);
let max_keys_per_provider = probe_config.max_keys_per_provider;
async move {
let hot_count_future = async {
if probing_enabled {
state
.runtime_state()
.set_len(&admin_provider_pool_quota_probe_active_members_key(
&provider_id,
))
.await
.unwrap_or(0)
} else {
0
}
};
let demand_snapshot_future = read_provider_pool_demand_snapshot(
state.runtime_state(),
&provider_id,
active_keys,
max_keys_per_provider,
);
let burst_pending_future = async {
probing_enabled
&& provider_pool_burst_pending(state.runtime_state(), &provider_id).await
};
let (hot_count, demand_snapshot, burst_pending) = tokio::join!(
hot_count_future,
demand_snapshot_future,
burst_pending_future
);
(
provider_id,
json!({
"provider_hot_count": hot_count,
"provider_desired_hot": if probing_enabled {
demand_snapshot.desired_hot
} else {
0
},
"provider_in_flight": demand_snapshot.in_flight,
"provider_ema_in_flight": demand_snapshot.ema_in_flight,
"provider_burst_pending": burst_pending,
}),
)
}
},
))
.await
.into_iter()
.collect::<BTreeMap<_, _>>();
let providers = pool_enabled_providers
.into_iter()