Merge remote-tracking branch 'origin/main' into codex/concurrency-hardening

This commit is contained in:
elky
2026-09-10 08:16:50 +08:00
15 changed files with 539 additions and 139 deletions
@@ -103,10 +103,11 @@ pub(crate) async fn maybe_build_local_admin_provider_reads_response(
.build_admin_provider_summary_payload(&provider_id)
.await
{
Some(payload) => Json(payload).into_response(),
None => build_admin_provider_not_found_response(format!(
Ok(Some(payload)) => Json(payload).into_response(),
Ok(None) => build_admin_provider_not_found_response(format!(
"Provider {provider_id} 不存在"
)),
Err(_) => build_admin_providers_data_unavailable_response(),
},
));
}
@@ -173,16 +173,17 @@ pub(crate) async fn maybe_build_local_admin_provider_writes_response(
.build_admin_provider_summary_payload(&provider_id)
.await
{
Some(payload) => attach_admin_audit_response(
Ok(Some(payload)) => attach_admin_audit_response(
Json(payload).into_response(),
"admin_provider_updated",
"update_provider",
"provider",
&provider_id,
),
None => build_admin_provider_not_found_response(format!(
Ok(None) => build_admin_provider_not_found_response(format!(
"Provider {provider_id} 不存在"
)),
Err(_) => build_admin_providers_data_unavailable_response(),
},
));
}
@@ -1,5 +1,6 @@
use super::value::build_admin_provider_summary_value;
use crate::handlers::admin::request::AdminAppState;
use crate::GatewayError;
use aether_data_contracts::repository::provider_catalog::{
StoredProviderCatalogEndpoint, StoredProviderCatalogKey,
};
@@ -10,19 +11,23 @@ use std::time::{SystemTime, UNIX_EPOCH};
pub(crate) async fn build_admin_provider_summary_payload(
state: &AdminAppState<'_>,
provider_id: &str,
) -> Option<serde_json::Value> {
) -> Result<Option<serde_json::Value>, GatewayError> {
let state = state.as_ref();
if !state.has_provider_catalog_data_reader() {
return None;
return Err(GatewayError::Internal(
"Admin provider catalog data unavailable".to_string(),
));
}
let provider_ids = vec![provider_id.to_string()];
let provider = state
let Some(provider) = state
.read_provider_catalog_providers_by_ids(&provider_ids)
.await
.ok()?
.await?
.into_iter()
.next()?;
.next()
else {
return Ok(None);
};
let (
endpoints_result,
keys_result,
@@ -36,8 +41,8 @@ pub(crate) async fn build_admin_provider_summary_payload(
state.list_provider_model_stats(&provider_ids),
state.list_active_global_model_ids_by_provider_ids(&provider_ids),
);
let endpoints = endpoints_result.ok().unwrap_or_default();
let keys = keys_result.ok().unwrap_or_default();
let endpoints = endpoints_result?;
let keys = keys_result?;
let quota_snapshot = quota_snapshot_result.ok().flatten();
let model_stats = model_stats_result
.ok()
@@ -57,7 +62,7 @@ pub(crate) async fn build_admin_provider_summary_payload(
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_secs();
Some(build_admin_provider_summary_value(
Ok(Some(build_admin_provider_summary_value(
&provider,
&endpoints,
&keys,
@@ -65,7 +70,7 @@ pub(crate) async fn build_admin_provider_summary_payload(
model_stats.as_ref(),
active_global_model_ids,
now_unix_secs,
))
)))
}
pub(crate) async fn build_admin_providers_summary_payload(
@@ -94,11 +99,7 @@ pub(crate) async fn build_admin_providers_summary_payload(
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)
.await
.ok()
.unwrap_or_default();
let mut providers = state.list_provider_catalog_providers(false).await.ok()?;
let all_provider_ids = providers
.iter()
.map(|provider| provider.id.clone())
@@ -109,8 +110,7 @@ pub(crate) async fn build_admin_providers_summary_payload(
state
.list_provider_catalog_endpoints_by_provider_ids(&all_provider_ids)
.await
.ok()
.unwrap_or_default()
.ok()?
};
let active_global_model_refs = if !requires_model_filter || all_provider_ids.is_empty() {
Vec::new()
@@ -202,8 +202,8 @@ pub(crate) async fn build_admin_providers_summary_payload(
state.list_active_global_model_ids_by_provider_ids(&provider_ids),
);
(
endpoints_result.ok().unwrap_or_default(),
keys_result.ok().unwrap_or_default(),
endpoints_result.ok()?,
keys_result.ok()?,
model_stats_result.ok().unwrap_or_default(),
active_global_model_refs_result.ok().unwrap_or_default(),
)
@@ -95,8 +95,11 @@ pub(crate) fn build_admin_provider_summary_value(
let scores = endpoint_keys
.iter()
.filter(|key| endpoint.is_active && key.is_active)
.filter_map(|key| provider_key_health_score(key, &endpoint.api_format))
.filter(|score| score.is_finite())
.map(|key| {
provider_key_health_score(key, &endpoint.api_format)
.filter(|score| score.is_finite())
.unwrap_or(1.0)
})
.collect::<Vec<_>>();
let health_score =
(!scores.is_empty()).then(|| scores.iter().sum::<f64>() / scores.len() as f64);
@@ -141,7 +141,7 @@ impl<'a> AdminAppState<'a> {
pub(crate) async fn build_admin_provider_summary_payload(
&self,
provider_id: &str,
) -> Option<serde_json::Value> {
) -> Result<Option<serde_json::Value>, GatewayError> {
crate::handlers::admin::provider::summary::build_admin_provider_summary_payload(
self,
provider_id,
@@ -3074,6 +3074,45 @@ mod tests {
.expect("key transport should build")
}
#[test]
fn admin_provider_key_health_response_preserves_v0_7_13_defaults() {
let state = AppState::new().expect("gateway should build");
for (health, expected_score) in [
(None, json!(1.0)),
(Some(json!({})), json!(1.0)),
(
Some(json!({"openai:chat": {"consecutive_failures": 0}})),
json!(1.0),
),
(
Some(json!({"openai:chat": {"health_score": 0.0}})),
json!(0.0),
),
(
Some(json!({"openai:chat": {"health_score": 1.0}})),
json!(1.0),
),
(
Some(json!({
"openai:chat": {"health_score": 0.25},
"openai:responses": {"health_score": 0.75},
})),
json!(0.25),
),
] {
let mut key = sample_catalog_key();
key.health_by_format = health;
let payload = build_admin_provider_key_response(
&state,
&key,
"openai",
&["openai:chat".to_string()],
1_000,
);
assert_eq!(payload["health_score"], expected_score);
}
}
#[test]
fn responses_key_scope_covers_search_in_one_direction() {
let mut responses_key = sample_catalog_key();
@@ -40,6 +40,8 @@ use crate::data::GatewayDataState;
const ADMIN_PROVIDERS_DATA_UNAVAILABLE_DETAIL: &str = "Admin provider catalog data unavailable";
mod health;
async fn provider_health_summary(
endpoints: &[StoredProviderCatalogEndpoint],
keys: &[StoredProviderCatalogKey],
@@ -170,7 +172,7 @@ async fn admin_provider_summary_health_ignores_disabled_keys() {
}
#[tokio::test]
async fn admin_provider_summary_health_does_not_inflate_observed_scores_with_missing_data() {
async fn admin_provider_summary_health_averages_v0_7_13_defaults_with_observed_scores() {
let endpoint = sample_endpoint(
"endpoint-chat",
"provider-openai",
@@ -183,21 +185,52 @@ async fn admin_provider_summary_health_does_not_inflate_observed_scores_with_mis
sample_key("key-unobserved", "provider-openai", "openai:chat", "test"),
sample_key("key-other-format", "provider-openai", "openai:chat", "test")
.with_health_fields(
Some(json!({"openai:responses": {"health_score": 1.0}})),
Some(json!({"openai:responses": {"health_score": 0.0}})),
None,
),
];
let payload = provider_health_summary(&[endpoint], &keys).await;
assert_eq!(payload["endpoint_health_details"][0]["health_score"], 0.2);
let expected_score = (0.2 + 1.0 + 1.0) / 3.0;
assert_eq!(
payload["endpoint_health_details"][0]["health_score"],
expected_score
);
assert_eq!(payload["endpoint_health_details"][0]["active_keys"], 3);
assert_eq!(payload["avg_health_score"], 0.2);
assert_eq!(payload["unhealthy_endpoints"], 1);
assert_eq!(payload["avg_health_score"], expected_score);
assert_eq!(payload["unhealthy_endpoints"], 0);
}
#[tokio::test]
async fn admin_provider_summary_health_is_unknown_without_active_observations() {
async fn admin_provider_summary_health_defaults_unobserved_enabled_keys_to_one() {
let endpoint = sample_endpoint(
"endpoint-chat",
"provider-openai",
"openai:chat",
"https://api.openai.example",
);
for health in [
None,
Some(json!({})),
Some(json!({"openai:chat": {"consecutive_failures": 0}})),
Some(json!({"openai:chat": {"health_score": null}})),
Some(json!({"openai:chat": {"health_score": "invalid"}})),
Some(json!({"openai:responses": {"health_score": 0.0}})),
] {
let key = sample_key("key-unobserved", "provider-openai", "openai:chat", "test")
.with_health_fields(health, None);
let payload = provider_health_summary(std::slice::from_ref(&endpoint), &[key]).await;
assert_eq!(payload["endpoint_health_details"][0]["health_score"], 1.0);
assert_eq!(payload["endpoint_health_details"][0]["active_keys"], 1);
assert_eq!(payload["avg_health_score"], 1.0);
assert_eq!(payload["unhealthy_endpoints"], 0);
}
}
#[tokio::test]
async fn admin_provider_summary_health_is_unknown_without_enabled_keys() {
let endpoint = sample_endpoint(
"endpoint-chat",
"provider-openai",
@@ -207,16 +240,7 @@ async fn admin_provider_summary_health_is_unknown_without_active_observations()
let mut disabled_key = sample_key("key-disabled", "provider-openai", "openai:chat", "test")
.with_health_fields(Some(json!({"openai:chat": {"health_score": 0.2}})), None);
disabled_key.is_active = false;
for keys in [
Vec::new(),
vec![sample_key(
"key-unobserved",
"provider-openai",
"openai:chat",
"test",
)],
vec![disabled_key],
] {
for keys in [Vec::new(), vec![disabled_key]] {
let payload = provider_health_summary(std::slice::from_ref(&endpoint), &keys).await;
assert_eq!(
@@ -233,7 +257,7 @@ async fn admin_provider_summary_health_is_unknown_without_active_observations()
}
#[tokio::test]
async fn admin_provider_summary_health_excludes_disabled_and_unobserved_endpoints() {
async fn admin_provider_summary_health_excludes_disabled_endpoints_and_defaults_unobserved_ones() {
let mut disabled_endpoint = sample_endpoint(
"endpoint-disabled",
"provider-openai",
@@ -279,16 +303,22 @@ async fn admin_provider_summary_health_excludes_disabled_and_unobserved_endpoint
let payload = provider_health_summary(&endpoints, &keys).await;
assert_eq!(payload["endpoint_health_details"][0]["health_score"], 0.8);
assert_eq!(
payload["endpoint_health_details"][1]["health_score"],
json!(null)
);
assert_eq!(
payload["endpoint_health_details"][2]["health_score"],
json!(null)
);
assert_eq!(payload["avg_health_score"], 0.8);
let details = payload["endpoint_health_details"]
.as_array()
.expect("endpoint health details should be an array");
for (api_format, health_score, is_active) in [
("openai:chat", json!(0.8), true),
("openai:responses", json!(null), false),
("openai:embedding", json!(1.0), true),
] {
let detail = details
.iter()
.find(|detail| detail["api_format"] == api_format)
.expect("endpoint health detail should exist");
assert_eq!(detail["health_score"], health_score, "{api_format}");
assert_eq!(detail["is_active"], is_active, "{api_format}");
}
assert_eq!(payload["avg_health_score"], 0.9);
assert_eq!(payload["unhealthy_endpoints"], 0);
}
@@ -0,0 +1,269 @@
use super::*;
use aether_data_contracts::repository::provider_catalog::{
ProviderCatalogKeyListQuery, StoredProviderCatalogKeyMaintenanceSummary,
StoredProviderCatalogKeyPage, StoredProviderCatalogKeyStats,
};
use aether_data_contracts::DataLayerError;
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum FailedRead {
Providers,
Endpoints,
Keys,
}
struct FailingSummaryRepository {
inner: InMemoryProviderCatalogReadRepository,
failed_read: FailedRead,
}
impl FailingSummaryRepository {
fn check(&self, operation: FailedRead) -> Result<(), DataLayerError> {
if self.failed_read == operation {
return Err(DataLayerError::InvalidConfiguration(
"injected summary read failure".to_string(),
));
}
Ok(())
}
}
#[async_trait::async_trait]
impl ProviderCatalogReadRepository for FailingSummaryRepository {
async fn list_providers(
&self,
active_only: bool,
) -> Result<Vec<StoredProviderCatalogProvider>, DataLayerError> {
self.check(FailedRead::Providers)?;
self.inner.list_providers(active_only).await
}
async fn list_providers_by_ids(
&self,
provider_ids: &[String],
) -> Result<Vec<StoredProviderCatalogProvider>, DataLayerError> {
self.check(FailedRead::Providers)?;
self.inner.list_providers_by_ids(provider_ids).await
}
async fn list_endpoints_by_ids(
&self,
endpoint_ids: &[String],
) -> Result<Vec<StoredProviderCatalogEndpoint>, DataLayerError> {
self.inner.list_endpoints_by_ids(endpoint_ids).await
}
async fn list_endpoints_by_provider_ids(
&self,
provider_ids: &[String],
) -> Result<Vec<StoredProviderCatalogEndpoint>, DataLayerError> {
self.check(FailedRead::Endpoints)?;
self.inner
.list_endpoints_by_provider_ids(provider_ids)
.await
}
async fn list_keys_by_ids(
&self,
key_ids: &[String],
) -> Result<Vec<StoredProviderCatalogKey>, DataLayerError> {
self.inner.list_keys_by_ids(key_ids).await
}
async fn list_keys_by_provider_ids(
&self,
provider_ids: &[String],
) -> Result<Vec<StoredProviderCatalogKey>, DataLayerError> {
self.inner.list_keys_by_provider_ids(provider_ids).await
}
async fn list_key_summaries_by_provider_ids(
&self,
provider_ids: &[String],
) -> Result<Vec<StoredProviderCatalogKey>, DataLayerError> {
self.check(FailedRead::Keys)?;
self.inner
.list_key_summaries_by_provider_ids(provider_ids)
.await
}
async fn list_key_maintenance_summaries_by_provider_ids(
&self,
provider_ids: &[String],
) -> Result<Vec<StoredProviderCatalogKeyMaintenanceSummary>, DataLayerError> {
self.inner
.list_key_maintenance_summaries_by_provider_ids(provider_ids)
.await
}
async fn list_keys_page(
&self,
query: &ProviderCatalogKeyListQuery,
) -> Result<StoredProviderCatalogKeyPage, DataLayerError> {
self.inner.list_keys_page(query).await
}
async fn list_key_stats_by_provider_ids(
&self,
provider_ids: &[String],
) -> Result<Vec<StoredProviderCatalogKeyStats>, DataLayerError> {
self.inner
.list_key_stats_by_provider_ids(provider_ids)
.await
}
}
#[tokio::test]
async fn admin_provider_summary_health_read_errors_do_not_look_like_empty_accounts() {
for failed_read in [
FailedRead::Providers,
FailedRead::Endpoints,
FailedRead::Keys,
] {
let repository = Arc::new(FailingSummaryRepository {
inner: InMemoryProviderCatalogReadRepository::seed(
vec![sample_provider("provider-openai", "openai", 10)],
vec![sample_endpoint(
"endpoint-chat",
"provider-openai",
"openai:chat",
"https://api.openai.example",
)],
vec![
sample_key("key-observed", "provider-openai", "openai:chat", "test")
.with_health_fields(
Some(json!({"openai:chat": {"health_score": 1.0}})),
None,
),
],
),
failed_read,
});
let state = AppState::new()
.expect("gateway should build")
.with_data_state_for_tests(GatewayDataState::with_provider_catalog_reader_for_tests(
repository,
));
for uri in [
"/api/admin/providers/summary",
"/api/admin/providers/summary?api_format=openai%3Achat",
"/api/admin/providers/provider-openai/summary",
] {
let response =
local_admin_providers_response(&state, http::Method::GET, uri, None).await;
assert_eq!(
response.status(),
StatusCode::SERVICE_UNAVAILABLE,
"{failed_read:?}: {uri}"
);
let body = axum::body::to_bytes(response.into_body(), 1024 * 1024)
.await
.expect("error body should read");
let payload: serde_json::Value =
serde_json::from_slice(&body).expect("error should parse");
assert_eq!(payload["detail"], ADMIN_PROVIDERS_DATA_UNAVAILABLE_DETAIL);
assert!(payload.get("items").is_none());
assert!(payload.get("endpoint_health_details").is_none());
}
}
}
#[tokio::test]
async fn admin_provider_summary_health_missing_provider_remains_not_found() {
let repository = Arc::new(InMemoryProviderCatalogReadRepository::seed(
vec![],
vec![],
vec![],
));
let state = AppState::new()
.expect("gateway should build")
.with_data_state_for_tests(GatewayDataState::with_provider_catalog_reader_for_tests(
repository,
));
let response = local_admin_providers_response(
&state,
http::Method::GET,
"/api/admin/providers/provider-missing/summary",
None,
)
.await;
assert_eq!(response.status(), StatusCode::NOT_FOUND);
}
#[tokio::test]
async fn admin_provider_summary_health_counts_inherited_reverse_proxy_accounts() {
for (provider_type, auth_type, api_format) in [
("codex", "oauth", "openai:responses"),
("kiro", "oauth", "claude:messages"),
("kiro", "bearer", "claude:messages"),
("gemini_cli", "oauth", "gemini:generate_content"),
("antigravity", "oauth", "gemini:generate_content"),
("vertex_ai", "service_account", "gemini:generate_content"),
("chatgpt_web", "oauth", "openai:chat"),
("chatgpt_web", "bearer", "openai:chat"),
("windsurf", "oauth", "openai:chat"),
] {
for configured_formats in [None, Some(json!([]))] {
for score in [None, Some(0.0), Some(0.75)] {
let mut provider = sample_provider("provider-reverse", provider_type, 10);
provider.provider_type = provider_type.to_string();
let endpoint = sample_endpoint(
"endpoint-reverse",
&provider.id,
api_format,
"https://reverse.example",
);
let mut key = sample_key("key-reverse", &provider.id, api_format, "test");
key.auth_type = auth_type.to_string();
key.api_formats = configured_formats.clone();
key.encrypted_api_key = Some("summary".to_string());
key.encrypted_auth_config = Some("{}".to_string());
key.health_by_format =
score.map(|score| json!({api_format: {"health_score": score}}));
let repository = Arc::new(InMemoryProviderCatalogReadRepository::seed(
vec![provider],
vec![endpoint],
vec![key],
));
let state = AppState::new()
.expect("gateway should build")
.with_data_state_for_tests(
GatewayDataState::with_provider_catalog_reader_for_tests(repository),
);
for uri in [
"/api/admin/providers/summary",
"/api/admin/providers/provider-reverse/summary",
] {
let response =
local_admin_providers_response(&state, http::Method::GET, uri, None).await;
assert_eq!(
response.status(),
StatusCode::OK,
"{provider_type}/{auth_type}: {uri}"
);
let body = axum::body::to_bytes(response.into_body(), 1024 * 1024)
.await
.expect("summary body should read");
let payload: serde_json::Value =
serde_json::from_slice(&body).expect("summary should parse");
let summary = payload
.get("items")
.map(|items| &items[0])
.unwrap_or(&payload);
let detail = &summary["endpoint_health_details"][0];
assert_eq!(detail["total_keys"], 1, "{provider_type}/{auth_type}");
assert_eq!(detail["active_keys"], 1, "{provider_type}/{auth_type}");
assert_eq!(
detail["health_score"],
json!(score.unwrap_or(1.0)),
"{provider_type}/{auth_type}"
);
assert_eq!(summary["avg_health_score"], json!(score.unwrap_or(1.0)));
}
}
}
}
}