mirror of
https://github.com/fawney19/Aether.git
synced 2026-09-02 17:30:23 +08:00
Improve pool score probing rules
This commit is contained in:
@@ -4,9 +4,10 @@ mod tests;
|
||||
|
||||
pub(crate) use runtime::{
|
||||
cancel_proxy_upgrade_rollout, clear_proxy_upgrade_rollout_conflicts,
|
||||
inspect_proxy_upgrade_rollout, list_admin_cleanup_run_records,
|
||||
perform_oauth_token_refresh_once, perform_pool_quota_probe_once, perform_provider_checkin_once,
|
||||
rebuild_admin_stats_once, record_completed_cleanup_run, record_proxy_upgrade_traffic_success,
|
||||
ensure_provider_key_pool_scores_for_keys, inspect_proxy_upgrade_rollout,
|
||||
list_admin_cleanup_run_records, perform_oauth_token_refresh_once,
|
||||
perform_pool_quota_probe_once, perform_provider_checkin_once, rebuild_admin_stats_once,
|
||||
record_completed_cleanup_run, record_proxy_upgrade_traffic_success,
|
||||
restore_proxy_upgrade_rollout_skipped_nodes, retry_proxy_upgrade_rollout_node,
|
||||
run_admin_system_cleanup_once, skip_proxy_upgrade_rollout_node, spawn_audit_cleanup_worker,
|
||||
spawn_db_maintenance_worker, spawn_gemini_file_mapping_cleanup_worker,
|
||||
|
||||
@@ -70,8 +70,9 @@ pub(crate) use pool_quota_probe::{
|
||||
PoolQuotaProbeWorkerConfig,
|
||||
};
|
||||
pub(crate) use pool_score_rebuild::{
|
||||
perform_pool_score_rebuild_once, perform_pool_score_rebuild_once_with_config,
|
||||
spawn_pool_score_rebuild_worker, PoolScoreRebuildRunSummary, PoolScoreRebuildWorkerConfig,
|
||||
ensure_provider_key_pool_scores_for_keys, perform_pool_score_rebuild_once,
|
||||
perform_pool_score_rebuild_once_with_config, spawn_pool_score_rebuild_worker,
|
||||
PoolScoreRebuildRunSummary, PoolScoreRebuildWorkerConfig,
|
||||
};
|
||||
pub(crate) use provider_checkin::{perform_provider_checkin_once, ProviderCheckinRunSummary};
|
||||
use proxy_node_metrics_cleanup::*;
|
||||
|
||||
@@ -21,6 +21,8 @@ use crate::admin_api::{
|
||||
};
|
||||
use crate::{AppState, GatewayError};
|
||||
|
||||
use super::pool_score_rebuild::ensure_provider_key_pool_scores_for_keys;
|
||||
|
||||
const POOL_QUOTA_PROBE_REDIS_PREFIX: &str = "ap:quota_probe:last";
|
||||
const POOL_QUOTA_PROBE_DEFAULT_SCAN_INTERVAL_SECONDS: u64 = 60;
|
||||
const POOL_QUOTA_PROBE_MIN_SCAN_INTERVAL_SECONDS: u64 = 15;
|
||||
@@ -495,7 +497,9 @@ async fn record_score_probe_results_from_payload(
|
||||
key_id,
|
||||
attempted_at,
|
||||
probe_result_succeeded(item),
|
||||
probe_result_hard_state(item),
|
||||
probe_result_hard_state(item).or_else(|| {
|
||||
(!probe_result_succeeded(item)).then_some(PoolMemberHardState::Cooldown)
|
||||
}),
|
||||
serde_json::json!({
|
||||
"last_probe": {
|
||||
"source": "pool_quota_probe",
|
||||
@@ -520,7 +524,7 @@ async fn record_score_probe_results_from_payload(
|
||||
key_id,
|
||||
attempted_at,
|
||||
false,
|
||||
None,
|
||||
Some(PoolMemberHardState::Cooldown),
|
||||
serde_json::json!({
|
||||
"last_probe": {
|
||||
"source": "pool_quota_probe",
|
||||
@@ -650,12 +654,7 @@ pub(crate) async fn perform_pool_quota_probe_once_with_config(
|
||||
.filter_map(|(provider, provider_type)| {
|
||||
let pool_config = admin_provider_pool_config(&provider)?;
|
||||
if pool_config.probing_enabled {
|
||||
Some((
|
||||
provider,
|
||||
provider_type,
|
||||
pool_config.probing_interval_minutes,
|
||||
pool_config.probe_concurrency.clamp(1, 64) as usize,
|
||||
))
|
||||
Some((provider, provider_type, pool_config))
|
||||
} else {
|
||||
None
|
||||
}
|
||||
@@ -668,7 +667,7 @@ pub(crate) async fn perform_pool_quota_probe_once_with_config(
|
||||
|
||||
let provider_ids = providers
|
||||
.iter()
|
||||
.map(|(provider, _, _, _)| provider.id.clone())
|
||||
.map(|(provider, _, _)| provider.id.clone())
|
||||
.collect::<Vec<_>>();
|
||||
let mut endpoints_by_provider = BTreeMap::<String, Vec<StoredProviderCatalogEndpoint>>::new();
|
||||
for endpoint in state
|
||||
@@ -688,7 +687,7 @@ pub(crate) async fn perform_pool_quota_probe_once_with_config(
|
||||
..PoolQuotaProbeRunSummary::empty()
|
||||
};
|
||||
|
||||
for (provider, provider_type, interval_minutes, probe_concurrency) in providers {
|
||||
for (provider, provider_type, pool_config) in providers {
|
||||
let endpoints = endpoints_by_provider
|
||||
.remove(&provider.id)
|
||||
.unwrap_or_default();
|
||||
@@ -702,6 +701,7 @@ pub(crate) async fn perform_pool_quota_probe_once_with_config(
|
||||
continue;
|
||||
};
|
||||
|
||||
let interval_minutes = pool_config.probing_interval_minutes;
|
||||
let interval_seconds = interval_minutes.clamp(1, 1440).saturating_mul(60);
|
||||
let keys = select_keys_for_provider(
|
||||
state,
|
||||
@@ -722,11 +722,44 @@ pub(crate) async fn perform_pool_quota_probe_once_with_config(
|
||||
summary.selected_keys += selected_count;
|
||||
|
||||
let selected_key_ids = keys.iter().map(|key| key.id.clone()).collect::<Vec<_>>();
|
||||
let score_ensure_budget = (pool_config.score_fallback_scan_limit as usize)
|
||||
.min(50_000)
|
||||
.max(selected_count.min(50_000));
|
||||
match ensure_provider_key_pool_scores_for_keys(
|
||||
state,
|
||||
&provider,
|
||||
&pool_config,
|
||||
&endpoints,
|
||||
&keys,
|
||||
now_ts,
|
||||
score_ensure_budget,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(upserted) if upserted > 0 => {
|
||||
debug!(
|
||||
provider_id = %provider.id,
|
||||
key_count = selected_count,
|
||||
scores_upserted = upserted,
|
||||
"gateway pool quota probe: ensured score rows for selected probe keys"
|
||||
);
|
||||
}
|
||||
Ok(_) => {}
|
||||
Err(err) => {
|
||||
warn!(
|
||||
provider_id = %provider.id,
|
||||
key_count = selected_count,
|
||||
error = ?err,
|
||||
"gateway pool quota probe: failed to ensure score rows for selected probe keys"
|
||||
);
|
||||
}
|
||||
}
|
||||
for key_id in &selected_key_ids {
|
||||
record_score_probe_in_progress_for_key(state, &provider.id, key_id, now_ts).await;
|
||||
}
|
||||
|
||||
let provider_short_id = provider.id.chars().take(8).collect::<String>();
|
||||
let probe_concurrency = pool_config.probe_concurrency.clamp(1, 64) as usize;
|
||||
let probe_concurrency = probe_concurrency.min(config.global_concurrency).max(1);
|
||||
let probe_results = stream::iter(keys.into_iter().map(|key| {
|
||||
let key_id = key.id.clone();
|
||||
|
||||
@@ -3,10 +3,14 @@ use std::time::{Duration, SystemTime, UNIX_EPOCH};
|
||||
|
||||
use aether_data_contracts::repository::global_models::AdminProviderModelListQuery;
|
||||
use aether_data_contracts::repository::pool_scores::GetPoolMemberScoresByIdsQuery;
|
||||
use aether_data_contracts::repository::provider_catalog::{
|
||||
StoredProviderCatalogEndpoint, StoredProviderCatalogKey, StoredProviderCatalogProvider,
|
||||
};
|
||||
use tracing::{debug, info, warn};
|
||||
|
||||
use crate::admin_api::admin_provider_pool_config;
|
||||
use crate::ai_serving::build_provider_key_pool_score_upsert;
|
||||
use crate::handlers::shared::provider_pool::AdminProviderPoolConfig;
|
||||
use crate::{AppState, GatewayError};
|
||||
|
||||
const POOL_SCORE_REBUILD_DEFAULT_INTERVAL_SECONDS: u64 = 300;
|
||||
@@ -130,6 +134,136 @@ struct ProviderScoreBuildItem {
|
||||
score_id: String,
|
||||
}
|
||||
|
||||
pub(crate) async fn ensure_provider_key_pool_scores_for_keys(
|
||||
state: &AppState,
|
||||
provider: &StoredProviderCatalogProvider,
|
||||
pool_config: &AdminProviderPoolConfig,
|
||||
endpoints: &[StoredProviderCatalogEndpoint],
|
||||
keys: &[StoredProviderCatalogKey],
|
||||
now_unix_secs: u64,
|
||||
max_upserts: usize,
|
||||
) -> Result<usize, GatewayError> {
|
||||
if max_upserts == 0
|
||||
|| keys.is_empty()
|
||||
|| !state.data.has_pool_score_reader()
|
||||
|| !state.data.has_pool_score_writer()
|
||||
{
|
||||
return Ok(0);
|
||||
}
|
||||
|
||||
let endpoints = endpoints
|
||||
.iter()
|
||||
.filter(|endpoint| endpoint.is_active && !endpoint.api_format.trim().is_empty())
|
||||
.collect::<Vec<_>>();
|
||||
let keys = keys
|
||||
.iter()
|
||||
.filter(|key| key.is_active && key.provider_id == provider.id)
|
||||
.collect::<Vec<_>>();
|
||||
if endpoints.is_empty() || keys.is_empty() {
|
||||
return Ok(0);
|
||||
}
|
||||
|
||||
let models = state
|
||||
.list_admin_provider_models(&AdminProviderModelListQuery {
|
||||
provider_id: provider.id.clone(),
|
||||
is_active: Some(true),
|
||||
offset: 0,
|
||||
limit: 10_000,
|
||||
})
|
||||
.await?
|
||||
.into_iter()
|
||||
.filter(|model| model.is_available)
|
||||
.collect::<Vec<_>>();
|
||||
if models.is_empty() {
|
||||
return Ok(0);
|
||||
}
|
||||
|
||||
let max_items = endpoints
|
||||
.len()
|
||||
.saturating_mul(models.len())
|
||||
.saturating_mul(keys.len())
|
||||
.min(max_upserts);
|
||||
let mut build_items = Vec::with_capacity(max_items);
|
||||
'outer: for (endpoint_index, endpoint) in endpoints.iter().enumerate() {
|
||||
for (model_index, model) in models.iter().enumerate() {
|
||||
for (key_index, key) in keys.iter().enumerate() {
|
||||
let draft = build_provider_key_pool_score_upsert(
|
||||
key,
|
||||
provider.provider_type.as_str(),
|
||||
endpoint.api_format.trim(),
|
||||
Some(model.id.as_str()),
|
||||
None,
|
||||
now_unix_secs,
|
||||
pool_config.score_rules,
|
||||
);
|
||||
build_items.push(ProviderScoreBuildItem {
|
||||
endpoint_index,
|
||||
model_index,
|
||||
key_index,
|
||||
score_id: draft.id,
|
||||
});
|
||||
if build_items.len() >= max_upserts {
|
||||
break 'outer;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
if build_items.is_empty() {
|
||||
return Ok(0);
|
||||
}
|
||||
|
||||
let existing_score_ids = state
|
||||
.data
|
||||
.get_pool_member_scores_by_ids(&GetPoolMemberScoresByIdsQuery {
|
||||
ids: build_items
|
||||
.iter()
|
||||
.map(|item| item.score_id.clone())
|
||||
.collect(),
|
||||
})
|
||||
.await
|
||||
.unwrap_or_else(|err| {
|
||||
debug!(
|
||||
provider_id = %provider.id,
|
||||
error = ?err,
|
||||
"gateway pool score ensure: failed to read existing scores by id"
|
||||
);
|
||||
Vec::new()
|
||||
})
|
||||
.into_iter()
|
||||
.map(|score| score.id)
|
||||
.collect::<std::collections::BTreeSet<_>>();
|
||||
|
||||
let mut upserted = 0usize;
|
||||
for item in &build_items {
|
||||
if existing_score_ids.contains(&item.score_id) {
|
||||
continue;
|
||||
}
|
||||
let endpoint = endpoints[item.endpoint_index];
|
||||
let model = &models[item.model_index];
|
||||
let key = keys[item.key_index];
|
||||
let upsert = build_provider_key_pool_score_upsert(
|
||||
key,
|
||||
provider.provider_type.as_str(),
|
||||
endpoint.api_format.trim(),
|
||||
Some(model.id.as_str()),
|
||||
None,
|
||||
now_unix_secs,
|
||||
pool_config.score_rules,
|
||||
);
|
||||
if state
|
||||
.data
|
||||
.upsert_pool_member_score(upsert)
|
||||
.await
|
||||
.map_err(|err| GatewayError::Internal(format!("{err:?}")))?
|
||||
.is_some()
|
||||
{
|
||||
upserted = upserted.saturating_add(1);
|
||||
}
|
||||
}
|
||||
|
||||
Ok(upserted)
|
||||
}
|
||||
|
||||
pub(crate) async fn perform_pool_score_rebuild_once_with_config(
|
||||
state: &AppState,
|
||||
config: PoolScoreRebuildWorkerConfig,
|
||||
@@ -145,16 +279,18 @@ pub(crate) async fn perform_pool_score_rebuild_once_with_config(
|
||||
.list_provider_catalog_providers(true)
|
||||
.await?
|
||||
.into_iter()
|
||||
.filter(|provider| admin_provider_pool_config(provider).is_some())
|
||||
.filter_map(|provider| {
|
||||
admin_provider_pool_config(&provider).map(|config| (provider, config))
|
||||
})
|
||||
.collect::<Vec<_>>();
|
||||
providers.sort_by(|left, right| left.id.cmp(&right.id));
|
||||
providers.sort_by(|left, right| left.0.id.cmp(&right.0.id));
|
||||
if providers.is_empty() {
|
||||
return Ok(PoolScoreRebuildRunSummary::empty());
|
||||
}
|
||||
|
||||
let provider_ids = providers
|
||||
.iter()
|
||||
.map(|provider| provider.id.clone())
|
||||
.map(|(provider, _)| provider.id.clone())
|
||||
.collect::<Vec<_>>();
|
||||
let mut endpoints_by_provider = BTreeMap::new();
|
||||
for endpoint in state
|
||||
@@ -196,7 +332,7 @@ pub(crate) async fn perform_pool_score_rebuild_once_with_config(
|
||||
break;
|
||||
}
|
||||
last_provider_index = Some(provider_index);
|
||||
let provider = providers[provider_index].clone();
|
||||
let (provider, pool_config) = providers[provider_index].clone();
|
||||
let endpoints = endpoints_by_provider
|
||||
.remove(&provider.id)
|
||||
.unwrap_or_default();
|
||||
@@ -252,6 +388,7 @@ pub(crate) async fn perform_pool_score_rebuild_once_with_config(
|
||||
Some(model.id.as_str()),
|
||||
None,
|
||||
now,
|
||||
pool_config.score_rules,
|
||||
);
|
||||
build_items.push(ProviderScoreBuildItem {
|
||||
endpoint_index,
|
||||
@@ -306,6 +443,7 @@ pub(crate) async fn perform_pool_score_rebuild_once_with_config(
|
||||
Some(model.id.as_str()),
|
||||
existing,
|
||||
now,
|
||||
pool_config.score_rules,
|
||||
);
|
||||
if state
|
||||
.data
|
||||
|
||||
Reference in New Issue
Block a user