feat(codex): add optional minimum quota reserve for pool scheduling

This commit is contained in:
ZheFox
2026-09-23 14:44:29 +08:00
parent 2a9d8d3b25
commit 1a4eba1005
16 changed files with 951 additions and 47 deletions
@@ -131,8 +131,13 @@ async fn schedule_pool_page_candidates(
entry.1.insert(candidate.candidate.key_id.clone());
}
let key_context_by_id =
read_pool_catalog_key_contexts_by_id(state, &candidates, provider_model_name).await;
let key_context_by_id = read_pool_catalog_key_contexts_by_id(
state,
&candidates,
provider_model_name,
effective_pool_config,
)
.await;
let mut runtime_by_provider = BTreeMap::new();
let mut pool_config_by_provider = BTreeMap::new();
@@ -1029,6 +1034,26 @@ impl<'a> PoolKeyCursor<'a> {
return None;
}
if pool_config.reserve_minimum_quota
&& admin_provider_pool_pure::admin_pool_key_minimum_quota_reached(
&key,
self.group.candidate.provider_type.as_str(),
Some(self.group.candidate.selected_provider_model_name.as_str()),
)
{
self.seen_key_ids.insert(key.id.clone());
self.record_skip_reason(POOL_ACCOUNT_EXHAUSTED_SKIP_REASON);
self.skipped_candidates
.push(SkippedLocalExecutionCandidate {
candidate: pool_candidate_from_catalog_key(&self.group, key),
skip_reason: POOL_ACCOUNT_EXHAUSTED_SKIP_REASON,
transport: None,
ranking: self.group.ranking.clone(),
extra_data: None,
});
return None;
}
let candidate = pool_candidate_from_catalog_key(&self.group, key);
self.build_eligible_candidate(candidate).await
}
@@ -1429,15 +1454,23 @@ async fn read_pool_catalog_key_contexts_by_id(
state: PlannerAppState<'_>,
candidates: &[EligibleLocalExecutionCandidate],
provider_model_name: Option<&str>,
effective_pool_config: Option<&AdminProviderPoolConfig>,
) -> BTreeMap<String, PoolCatalogKeyContext> {
let mut key_ids = Vec::new();
let mut provider_type_by_key_id = BTreeMap::<String, String>::new();
let mut reserve_minimum_quota_key_ids = BTreeSet::new();
for candidate in candidates {
if pool_config_for_candidate(candidate).is_none() {
let Some(pool_config) = effective_pool_config
.cloned()
.or_else(|| pool_config_for_candidate(candidate))
else {
continue;
}
};
let key_id = candidate.candidate.key_id.clone();
if pool_config.reserve_minimum_quota {
reserve_minimum_quota_key_ids.insert(key_id.clone());
}
if let Entry::Vacant(entry) = provider_type_by_key_id.entry(key_id.clone()) {
entry.insert(candidate.transport.provider.provider_type.clone());
key_ids.push(key_id);
@@ -1489,16 +1522,20 @@ async fn read_pool_catalog_key_contexts_by_id(
.get(&key.id)
.map(String::as_str)
.unwrap_or_default();
(
key.id.clone(),
build_pool_catalog_key_context(
state,
&provider_pool_service,
let mut context = build_pool_catalog_key_context(
state,
&provider_pool_service,
&key,
provider_type,
provider_model_name,
);
context.quota_exhausted |= reserve_minimum_quota_key_ids.contains(&key.id)
&& admin_provider_pool_pure::admin_pool_key_minimum_quota_reached(
&key,
provider_type,
provider_model_name,
),
)
);
(key.id.clone(), context)
})
.collect::<BTreeMap<_, _>>();
// A key can disappear between the candidate-row and catalog reads. Keep
@@ -3964,6 +4001,91 @@ mod tests {
}));
}
#[tokio::test]
async fn pool_key_cursor_reserve_minimum_quota_filters_pages_and_sticky_hits() {
for reserve_enabled in [false, true] {
for sticky in [false, true] {
for used_percent in [99.0, 98.0] {
let provider_config = Some(json!({
"pool_advanced": {
"reserve_minimum_quota": reserve_enabled,
"skip_exhausted_accounts": false
}
}));
let provider =
sample_codex_pool_provider("provider-pool", 0, provider_config.clone());
let endpoint = sample_codex_pool_endpoint("provider-pool", "endpoint-1");
let mut reserved = sample_codex_pool_key("provider-pool", "key-low");
reserved.upstream_metadata = Some(json!({
"codex": {"primary_used_percent": used_percent}
}));
let ready = sample_codex_pool_key("provider-pool", "key-ready");
let rows = vec![
sample_codex_pool_row("provider-pool", "endpoint-1", "key-low", 0),
sample_codex_pool_row("provider-pool", "endpoint-1", "key-ready", 0),
];
let data_state = GatewayDataState::with_provider_catalog_and_minimal_candidate_selection_for_tests(
Arc::new(InMemoryProviderCatalogReadRepository::seed(
vec![provider], vec![endpoint], vec![reserved, ready],
)),
Arc::new(InMemoryMinimalCandidateSelectionReadRepository::seed(rows)),
)
.with_encryption_key_for_tests(aether_crypto::DEVELOPMENT_ENCRYPTION_KEY);
let app = AppState::new()
.expect("state should build")
.with_data_state_for_tests(data_state);
let group =
sample_codex_pool_group("provider-pool", "endpoint-1", 0, provider_config);
let pool_config =
pool_config_for_candidate(&group).expect("pool config should parse");
let sticky_token = sticky.then_some("reserve-session");
if sticky {
record_admin_provider_pool_success(
app.runtime_state.as_ref(),
"provider-pool",
"key-low",
&pool_config,
sticky_token,
0,
None,
)
.await;
}
let mut cursor = PoolKeyCursor::new(
PlannerAppState::new(&app),
group,
sticky_token,
None,
None,
);
cursor.window_size = 1;
cursor.page_size = 1;
let mut returned = Vec::new();
while let Some(candidate) = cursor.next_key().await {
returned.push(candidate.candidate.key_id);
}
let reserve_reached = reserve_enabled && used_percent >= 99.0;
assert_eq!(
returned.contains(&"key-low".to_string()),
!reserve_reached,
"reserve={reserve_enabled}, sticky={sticky}, used={used_percent}"
);
assert!(returned.contains(&"key-ready".to_string()));
if reserve_reached {
assert_eq!(
cursor
.skip_reason_counts
.get(POOL_ACCOUNT_EXHAUSTED_SKIP_REASON),
Some(&1)
);
} else if sticky {
assert_eq!(returned.first().map(String::as_str), Some("key-low"));
}
}
}
}
}
#[tokio::test]
async fn pool_key_cursor_does_not_spend_effective_scan_budget_on_exhausted_accounts() {
let provider_config = Some(json!({
@@ -411,6 +411,7 @@ pub(crate) fn admin_provider_pool_config_from_config_value(
unschedulable_rules: Vec::new(),
lru_enabled: false,
skip_exhausted_accounts: false,
reserve_minimum_quota: false,
sticky_session_ttl_seconds: 3600,
latency_window_seconds: 3600,
latency_sample_limit: 50,
@@ -446,6 +447,10 @@ pub(crate) fn admin_provider_pool_config_from_config_value(
.get("skip_exhausted_accounts")
.and_then(Value::as_bool)
.unwrap_or(false),
reserve_minimum_quota: pool_advanced
.get("reserve_minimum_quota")
.and_then(Value::as_bool)
.unwrap_or(false),
sticky_session_ttl_seconds: pool_advanced
.get("sticky_session_ttl_seconds")
.and_then(json_u64)
@@ -574,6 +579,22 @@ mod tests {
let config = admin_provider_pool_config(&provider).expect("pool config should exist");
assert!(!config.skip_exhausted_accounts);
assert!(!config.reserve_minimum_quota);
}
#[test]
fn parses_reserve_minimum_quota_independently_of_skip_exhausted_accounts() {
for enabled in [false, true] {
let provider = sample_provider(json!({
"pool_advanced": {
"reserve_minimum_quota": enabled,
"skip_exhausted_accounts": false
}
}));
let config = admin_provider_pool_config(&provider).expect("pool config should exist");
assert_eq!(config.reserve_minimum_quota, enabled);
assert!(!config.skip_exhausted_accounts);
}
}
#[test]
@@ -651,6 +651,7 @@ mod tests {
unschedulable_rules: Vec::new(),
lru_enabled: true,
skip_exhausted_accounts: false,
reserve_minimum_quota: false,
sticky_session_ttl_seconds: 120,
latency_window_seconds: 600,
latency_sample_limit: 10,
@@ -1117,10 +1117,16 @@ pub(super) fn build_admin_pool_key_payload(
let health_score = admin_pool_health_score(key);
let circuit_breaker_open = false;
let auth_semantics = provider_key_auth_semantics(key, provider_type);
let account_quota_exhausted = pool_config
.as_ref()
.is_some_and(|config| config.skip_exhausted_accounts)
&& admin_provider_pool_pure::admin_pool_key_account_quota_exhausted(key, provider_type);
let account_quota_exhausted = pool_config.as_ref().is_some_and(|config| {
(config.skip_exhausted_accounts
&& admin_provider_pool_pure::admin_pool_key_account_quota_exhausted(key, provider_type))
|| (config.reserve_minimum_quota
&& admin_provider_pool_pure::admin_pool_key_minimum_quota_reached(
key,
provider_type,
None,
))
});
let auth_config = state.parse_catalog_auth_config_json(key);
let oauth_expires_at =
admin_pool_derive_oauth_expires_at(provider_type, key, auth_config.as_ref());
@@ -488,9 +488,16 @@ pub(super) fn admin_pool_key_visible_status_filter(
) {
return status;
}
if pool_config.is_some_and(|config| config.skip_exhausted_accounts)
&& admin_provider_pool_pure::admin_pool_key_account_quota_exhausted(key, provider_type)
{
if pool_config.is_some_and(|config| {
(config.skip_exhausted_accounts
&& admin_provider_pool_pure::admin_pool_key_account_quota_exhausted(key, provider_type))
|| (config.reserve_minimum_quota
&& admin_provider_pool_pure::admin_pool_key_minimum_quota_reached(
key,
provider_type,
None,
))
}) {
return "quota_exhausted";
}
if !key.is_active {
@@ -1340,6 +1340,7 @@ async fn provider_query_apply_pool_scheduler_to_test_candidates(
provider.id.clone(),
provider_query_ai_pool_runtime_state(&runtime),
);
let reserve_minimum_quota = pool_config.reserve_minimum_quota;
let pool_config =
provider_query_ai_pool_scheduling_config(pool_config, provider.provider_type.as_str());
let inputs = keys
@@ -1351,6 +1352,14 @@ async fn provider_query_apply_pool_scheduler_to_test_candidates(
effective_model: effective_model.to_string(),
scheduler_skip_reason: None,
};
let mut key_context =
provider_query_pool_catalog_key_context(state, &key, &provider.provider_type);
key_context.quota_exhausted |= reserve_minimum_quota
&& admin_provider_pool_pure::admin_pool_key_minimum_quota_reached(
&key,
&provider.provider_type,
Some(effective_model),
);
AiPoolCandidateInput {
facts: AiPoolCandidateFacts {
provider_id: provider.id.clone(),
@@ -1362,11 +1371,7 @@ async fn provider_query_apply_pool_scheduler_to_test_candidates(
key_internal_priority: key.internal_priority,
},
pool_config: Some(pool_config.clone()),
key_context: provider_query_pool_catalog_key_context(
state,
&key,
&provider.provider_type,
),
key_context,
candidate,
}
})
@@ -65,6 +65,7 @@ pub(crate) struct AdminProviderPoolConfig {
pub(crate) unschedulable_rules: Vec<AdminProviderPoolUnschedulableRule>,
pub(crate) lru_enabled: bool,
pub(crate) skip_exhausted_accounts: bool,
pub(crate) reserve_minimum_quota: bool,
pub(crate) sticky_session_ttl_seconds: u64,
pub(crate) latency_window_seconds: u64,
pub(crate) latency_sample_limit: u64,
@@ -40,10 +40,6 @@ pub(super) async fn read_candidate_runtime_selection_snapshot(
) -> Result<CandidateRuntimeSelectionSnapshot, GatewayError> {
let provider_concurrent_limits = read_provider_concurrent_limits(state, candidates).await?;
let provider_pool_state = read_provider_pool_state_map(state, candidates).await?;
let provider_skip_exhausted_accounts = provider_pool_state
.iter()
.map(|(provider_id, state)| (provider_id.clone(), state.skip_exhausted_accounts))
.collect::<BTreeMap<_, _>>();
let pool_provider_ids = provider_pool_state
.iter()
.filter_map(|(provider_id, state)| state.pool_enabled.then_some(provider_id.clone()))
@@ -62,7 +58,7 @@ pub(super) async fn read_candidate_runtime_selection_snapshot(
let key_account_quota_exhausted = read_key_account_quota_exhaustion_map(
candidates,
&provider_key_rpm_states,
&provider_skip_exhausted_accounts,
&provider_pool_state,
);
let key_oauth_invalid =
read_key_oauth_invalid_map(candidates, &provider_key_rpm_states, now_unix_secs);
@@ -360,6 +356,7 @@ async fn read_provider_quota_block_map(
struct ProviderPoolState {
pool_enabled: bool,
skip_exhausted_accounts: bool,
reserve_minimum_quota: bool,
}
async fn read_provider_pool_state_map(
@@ -391,11 +388,17 @@ async fn read_provider_pool_state_map(
.and_then(|value| value.get("skip_exhausted_accounts"))
.and_then(serde_json::Value::as_bool)
.unwrap_or(false);
let reserve_minimum_quota = pool_advanced
.and_then(serde_json::Value::as_object)
.and_then(|value| value.get("reserve_minimum_quota"))
.and_then(serde_json::Value::as_bool)
.unwrap_or(false);
(
provider.id,
ProviderPoolState {
pool_enabled: pool_advanced.is_some(),
skip_exhausted_accounts,
reserve_minimum_quota,
},
)
})
@@ -405,7 +408,7 @@ async fn read_provider_pool_state_map(
fn read_key_account_quota_exhaustion_map(
candidates: &[SchedulerMinimalCandidateSelectionCandidate],
provider_key_rpm_states: &BTreeMap<String, StoredProviderCatalogKey>,
provider_skip_exhausted_accounts: &BTreeMap<String, bool>,
provider_pool_state: &BTreeMap<String, ProviderPoolState>,
) -> BTreeMap<String, bool> {
candidates
.iter()
@@ -431,11 +434,19 @@ fn read_key_account_quota_exhaustion_map(
candidate.provider_type.as_str(),
candidate.selected_provider_model_name.as_str(),
);
let skip_configured = provider_skip_exhausted_accounts
let pool_state = provider_pool_state
.get(candidate.provider_id.as_str())
.copied()
.unwrap_or(false);
hard_blocked || (skip_configured && account_exhausted)
.unwrap_or_default();
let reserve_reached = pool_state.reserve_minimum_quota
&& admin_provider_pool_pure::admin_pool_key_minimum_quota_reached(
key,
candidate.provider_type.as_str(),
Some(candidate.selected_provider_model_name.as_str()),
);
hard_blocked
|| reserve_reached
|| (pool_state.skip_exhausted_accounts && account_exhausted)
});
(candidate.key_id.clone(), exhausted)
})
@@ -566,3 +577,64 @@ fn read_provider_key_rpm_reset_at_map(
})
.collect::<BTreeMap<_, _>>()
}
#[cfg(test)]
mod reserve_minimum_quota_tests {
use super::*;
use serde_json::json;
#[test]
fn reserve_minimum_quota_is_independent_of_skip_exhausted_accounts() {
let candidate = SchedulerMinimalCandidateSelectionCandidate {
provider_id: "provider-codex".to_string(),
provider_name: "codex".to_string(),
provider_type: "codex".to_string(),
provider_priority: 0,
endpoint_id: "endpoint-codex".to_string(),
endpoint_api_format: "openai:responses".to_string(),
key_id: "key-codex".to_string(),
key_name: "codex".to_string(),
key_auth_type: "oauth".to_string(),
key_internal_priority: 0,
key_global_priority_for_format: None,
key_capabilities: None,
model_id: "model-codex".to_string(),
global_model_id: "global-model-codex".to_string(),
global_model_name: "gpt-5".to_string(),
selected_provider_model_name: "gpt-5".to_string(),
supports_streaming: true,
mapping_matched_model: None,
};
let mut key = StoredProviderCatalogKey::new(
candidate.key_id.clone(),
candidate.provider_id.clone(),
"codex".to_string(),
"oauth".to_string(),
None,
true,
)
.expect("key should build");
for reserve_enabled in [false, true] {
for used_percent in [99.0, 98.0] {
key.upstream_metadata =
Some(json!({"codex": {"primary_used_percent": used_percent}}));
let exhausted = read_key_account_quota_exhaustion_map(
std::slice::from_ref(&candidate),
&BTreeMap::from([(key.id.clone(), key.clone())]),
&BTreeMap::from([(
candidate.provider_id.clone(),
ProviderPoolState {
pool_enabled: true,
reserve_minimum_quota: reserve_enabled,
skip_exhausted_accounts: false,
},
)]),
);
assert_eq!(
exhausted.get(&key.id),
Some(&(reserve_enabled && used_percent >= 99.0))
);
}
}
}
}
@@ -2364,6 +2364,78 @@ async fn gateway_marks_exhausted_codex_pool_key_as_blocked_when_flag_enabled() {
assert_eq!(keys[0]["account_quota"], json!("5H剩余 0.0%"));
}
#[tokio::test]
async fn gateway_reserve_minimum_quota_marks_and_filters_codex_pool_keys() {
for (reserve_enabled, used_percent, exhausted) in [
(true, 99.0, true),
(false, 99.0, false),
(true, 98.0, false),
] {
let mut provider = sample_provider("provider-codex", "codex", 10);
provider.provider_type = "codex".to_string();
provider.config = Some(json!({
"pool_advanced": {
"reserve_minimum_quota": reserve_enabled,
"skip_exhausted_accounts": false,
"auto_remove_quota_exhausted_keys": true
}
}));
let mut key = sample_key(
"key-codex-reserved",
"provider-codex",
"openai:responses",
"oauth-placeholder",
);
key.auth_type = "oauth".to_string();
key.upstream_metadata = Some(json!({
"codex": {"primary_used_percent": used_percent}
}));
let state = AppState::new()
.expect("gateway should build")
.with_data_state_for_tests(GatewayDataState::with_provider_catalog_reader_for_tests(
Arc::new(InMemoryProviderCatalogReadRepository::seed(
vec![provider],
Vec::new(),
vec![key],
)),
));
for status in ["all", "quota_exhausted", "available"] {
let response = local_admin_pool_response(
&state,
http::Method::GET,
&format!("/api/admin/pool/provider-codex/keys?page=1&page_size=50&status={status}"),
None,
)
.await;
assert_eq!(response.status(), StatusCode::OK);
let payload: serde_json::Value = serde_json::from_slice(
&to_bytes(response.into_body(), usize::MAX)
.await
.expect("body should read"),
)
.expect("json body should parse");
let keys = payload["keys"].as_array().expect("keys should be array");
let visible = status == "all" || (status == "quota_exhausted") == exhausted;
assert_eq!(
keys.len(),
usize::from(visible),
"reserve={reserve_enabled}, used={used_percent}, status={status}"
);
if visible {
assert_eq!(keys[0]["key_id"], "key-codex-reserved");
assert_eq!(
keys[0]["scheduling_reason"] == "account_quota_exhausted",
exhausted
);
if exhausted {
assert_eq!(keys[0]["scheduling_status"], "blocked");
assert_eq!(keys[0]["scheduling_label"], "额度耗尽");
}
}
}
}
}
#[tokio::test]
async fn gateway_lists_inherited_fixed_provider_api_formats_for_pool_keys() {
let mut provider = sample_provider("provider-codex", "codex", 10).with_transport_fields(