Files
Aether/apps/aether-gateway/src/ai_serving/planner/pool_scheduler.rs
T

2471 lines
85 KiB
Rust
Raw Normal View History

2026-05-03 20:14:29 +08:00
use std::collections::{btree_map::Entry, BTreeMap, BTreeSet, VecDeque};
use std::sync::atomic::{AtomicU64, Ordering as AtomicOrdering};
2026-05-11 00:12:05 +08:00
use aether_admin::provider::pool as admin_provider_pool_pure;
2026-05-02 13:23:54 +08:00
use aether_ai_serving::{
2026-05-11 00:12:05 +08:00
normalize_enabled_ai_pool_presets, run_ai_pool_scheduler, AiPoolCandidateFacts,
AiPoolCandidateInput, AiPoolCandidateOrchestration, AiPoolCatalogKeyContext,
AiPoolRuntimeState, AiPoolSchedulingConfig, AiPoolSchedulingPreset,
2026-05-02 13:23:54 +08:00
};
2026-05-03 20:14:29 +08:00
use aether_data_contracts::repository::candidate_selection::{
2026-05-11 00:12:05 +08:00
StoredMinimalCandidateSelectionRow, StoredPoolKeyCandidateOrder,
StoredPoolKeyCandidateRowsQuery,
2026-05-03 20:14:29 +08:00
};
use aether_data_contracts::repository::provider_catalog::StoredProviderCatalogKey;
use serde_json::{Map, Value};
use tracing::warn;
2026-05-02 13:23:54 +08:00
use crate::ai_serving::planner::candidate_resolution::{
2026-05-03 20:14:29 +08:00
candidate_auth_channel_skip_reason, read_candidate_transport_snapshot,
EligibleLocalExecutionCandidate, LocalExecutionCandidateKind, SkippedLocalExecutionCandidate,
};
use crate::ai_serving::{
candidate_common_transport_skip_reason, CandidateTransportPolicyFacts, PlannerAppState,
};
use crate::clock::current_unix_ms;
use crate::handlers::shared::provider_pool::admin_provider_pool_config_from_config_value;
use crate::handlers::shared::provider_pool::read_admin_provider_pool_runtime_state;
use crate::handlers::shared::provider_pool::{
2026-05-11 00:12:05 +08:00
try_claim_admin_provider_pool_key, AdminProviderPoolConfig, AdminProviderPoolRuntimeState,
};
use crate::handlers::shared::{
parse_catalog_auth_config_json, provider_key_health_summary,
provider_key_status_snapshot_payload,
};
use crate::orchestration::LocalExecutionCandidateMetadata;
use crate::provider_key_auth::provider_key_auth_semantics;
static LOAD_BALANCE_SEQUENCE: AtomicU64 = AtomicU64::new(0);
2026-05-03 20:14:29 +08:00
const DEFAULT_POOL_KEY_PAGE_SIZE: u32 = 128;
const DEFAULT_POOL_MAX_SCANNED_KEYS: u32 = 1024;
2026-05-02 13:23:54 +08:00
type PoolCatalogKeyContext = AiPoolCatalogKeyContext;
pub(crate) async fn apply_local_execution_pool_scheduler(
state: PlannerAppState<'_>,
candidates: Vec<EligibleLocalExecutionCandidate>,
sticky_session_token: Option<&str>,
2026-05-03 20:14:29 +08:00
requested_model: Option<&str>,
request_auth_channel: Option<&str>,
) -> (
Vec<EligibleLocalExecutionCandidate>,
Vec<SkippedLocalExecutionCandidate>,
) {
if candidates.is_empty() {
return (Vec::new(), Vec::new());
}
let sticky_session_token = sticky_session_token
.map(str::trim)
.filter(|value| !value.is_empty());
2026-05-03 20:14:29 +08:00
let mut scheduled = Vec::new();
let mut skipped = Vec::new();
for candidate in candidates {
if candidate.kind == LocalExecutionCandidateKind::PoolGroup {
let mut expanded = expand_pool_group_candidate(
state,
candidate,
sticky_session_token,
requested_model,
request_auth_channel,
)
.await;
scheduled.append(&mut expanded.0);
skipped.append(&mut expanded.1);
} else {
scheduled.push(candidate);
}
}
(scheduled, skipped)
}
async fn schedule_pool_page_candidates(
state: PlannerAppState<'_>,
candidates: Vec<EligibleLocalExecutionCandidate>,
sticky_session_token: Option<&str>,
) -> (
Vec<EligibleLocalExecutionCandidate>,
Vec<SkippedLocalExecutionCandidate>,
) {
if candidates.is_empty() {
return (Vec::new(), Vec::new());
}
let mut provider_runtime_requirements =
BTreeMap::<String, (AdminProviderPoolConfig, BTreeSet<String>)>::new();
for candidate in &candidates {
let Some(pool_config) = pool_config_for_candidate(candidate) else {
continue;
};
let entry = provider_runtime_requirements
.entry(candidate.candidate.provider_id.clone())
.or_insert_with(|| (pool_config.clone(), BTreeSet::new()));
entry.1.insert(candidate.candidate.key_id.clone());
}
let key_context_by_id = read_pool_catalog_key_contexts_by_id(state, &candidates).await;
let mut runtime_by_provider = BTreeMap::new();
for (provider_id, (pool_config, key_ids)) in provider_runtime_requirements {
let key_ids = key_ids.into_iter().collect::<Vec<_>>();
2026-05-08 00:18:12 +08:00
let runtime = if key_ids.is_empty() {
AdminProviderPoolRuntimeState::default()
} else {
read_admin_provider_pool_runtime_state(
state.app().runtime_state.as_ref(),
provider_id.as_str(),
&key_ids,
&pool_config,
sticky_session_token,
)
.await
};
runtime_by_provider.insert(provider_id, runtime);
}
apply_local_execution_pool_scheduler_with_runtime_map(
candidates,
&runtime_by_provider,
&key_context_by_id,
)
}
2026-05-03 20:14:29 +08:00
async fn expand_pool_group_candidate(
state: PlannerAppState<'_>,
group: EligibleLocalExecutionCandidate,
sticky_session_token: Option<&str>,
requested_model: Option<&str>,
request_auth_channel: Option<&str>,
) -> (
Vec<EligibleLocalExecutionCandidate>,
Vec<SkippedLocalExecutionCandidate>,
) {
let mut cursor = PoolKeyCursor::new(
state,
group,
sticky_session_token,
requested_model,
request_auth_channel,
);
let mut scheduled = Vec::new();
let mut skipped = Vec::new();
while let Some(candidate) = cursor.next_key().await {
scheduled.push(candidate);
}
skipped.append(&mut cursor.take_skipped_candidates());
if scheduled.is_empty() {
cursor.log_exhausted();
}
(scheduled, skipped)
}
pub(crate) struct PoolKeyCursor<'a> {
state: PlannerAppState<'a>,
group: EligibleLocalExecutionCandidate,
sticky_session_token: Option<String>,
requested_model: Option<String>,
request_auth_channel: Option<String>,
2026-05-11 00:12:05 +08:00
pool_key_order: StoredPoolKeyCandidateOrder,
2026-05-03 20:14:29 +08:00
next_offset: u32,
scanned_keys: u32,
page_size: u32,
max_scanned_keys: u32,
skip_reason_counts: BTreeMap<&'static str, u32>,
next_pool_key_index: u32,
2026-05-11 00:12:05 +08:00
sticky_candidate_loaded: bool,
seen_key_ids: BTreeSet<String>,
2026-05-03 20:14:29 +08:00
queued_candidates: VecDeque<EligibleLocalExecutionCandidate>,
skipped_candidates: Vec<SkippedLocalExecutionCandidate>,
exhausted_logged: bool,
}
impl<'a> PoolKeyCursor<'a> {
pub(crate) fn new(
state: PlannerAppState<'a>,
group: EligibleLocalExecutionCandidate,
sticky_session_token: Option<&str>,
requested_model: Option<&str>,
request_auth_channel: Option<&str>,
) -> Self {
2026-05-11 00:12:05 +08:00
let pool_key_order = pool_key_candidate_order_for_group(&group);
2026-05-03 20:14:29 +08:00
Self {
state,
group,
sticky_session_token: sticky_session_token.map(str::to_string),
requested_model: requested_model.map(str::to_string),
request_auth_channel: request_auth_channel.map(str::to_string),
2026-05-11 00:12:05 +08:00
pool_key_order,
2026-05-03 20:14:29 +08:00
next_offset: 0,
scanned_keys: 0,
page_size: DEFAULT_POOL_KEY_PAGE_SIZE,
max_scanned_keys: DEFAULT_POOL_MAX_SCANNED_KEYS,
skip_reason_counts: BTreeMap::new(),
next_pool_key_index: 0,
2026-05-11 00:12:05 +08:00
sticky_candidate_loaded: false,
seen_key_ids: BTreeSet::new(),
2026-05-03 20:14:29 +08:00
queued_candidates: VecDeque::new(),
skipped_candidates: Vec::new(),
exhausted_logged: false,
}
}
pub(crate) async fn next_key(&mut self) -> Option<EligibleLocalExecutionCandidate> {
loop {
2026-05-11 00:12:05 +08:00
if let Some(candidate) = self.next_queued_candidate().await {
2026-05-03 20:14:29 +08:00
return Some(candidate);
}
2026-05-11 00:12:05 +08:00
if !self.sticky_candidate_loaded {
self.sticky_candidate_loaded = true;
if let Some(candidate) = self.sticky_candidate().await {
self.queued_candidates.push_back(candidate);
continue;
}
2026-05-03 20:14:29 +08:00
}
2026-05-11 00:12:05 +08:00
if !self.refill_queued_candidates().await {
return None;
2026-05-03 20:14:29 +08:00
}
}
}
pub(crate) fn take_skipped_candidates(&mut self) -> Vec<SkippedLocalExecutionCandidate> {
std::mem::take(&mut self.skipped_candidates)
}
pub(crate) fn log_exhausted(&mut self) {
if self.exhausted_logged {
return;
}
self.exhausted_logged = true;
warn!(
event_name = "pool_group_exhausted",
log_type = "event",
provider_id = %self.group.candidate.provider_id,
endpoint_id = %self.group.candidate.endpoint_id,
model_id = %self.group.candidate.model_id,
scanned_keys = self.scanned_keys,
skip_reason_counts = ?self.skip_reason_counts,
"gateway pool scheduler exhausted pool group without a schedulable key"
);
}
async fn next_page_candidates(&mut self) -> Option<Vec<EligibleLocalExecutionCandidate>> {
if self.scanned_keys >= self.max_scanned_keys {
return None;
}
let limit = self
.page_size
.min(self.max_scanned_keys - self.scanned_keys);
let query = StoredPoolKeyCandidateRowsQuery {
api_format: self.group.candidate.endpoint_api_format.clone(),
provider_id: self.group.candidate.provider_id.clone(),
endpoint_id: self.group.candidate.endpoint_id.clone(),
model_id: self.group.candidate.model_id.clone(),
selected_provider_model_name: self.group.candidate.selected_provider_model_name.clone(),
2026-05-11 00:12:05 +08:00
order: self.pool_key_order.clone(),
2026-05-03 20:14:29 +08:00
offset: self.next_offset,
limit,
};
let rows = match self
.state
.app()
.list_pool_key_candidate_rows_for_group(&query)
.await
{
Ok(rows) => rows,
Err(err) => {
warn!(
event_name = "pool_group_key_page_load_failed",
log_type = "event",
provider_id = %self.group.candidate.provider_id,
endpoint_id = %self.group.candidate.endpoint_id,
model_id = %self.group.candidate.model_id,
selected_provider_model_name = %self.group.candidate.selected_provider_model_name,
offset = self.next_offset,
limit,
error = ?err,
"gateway pool scheduler failed to read pool key page"
);
return None;
}
};
if rows.is_empty() {
return None;
}
self.scanned_keys += rows.len() as u32;
self.next_offset = self.next_offset.saturating_add(rows.len() as u32);
Some(self.build_page_eligible_candidates(rows).await)
}
2026-05-11 00:12:05 +08:00
async fn sticky_candidate(&mut self) -> Option<EligibleLocalExecutionCandidate> {
let pool_config = pool_config_for_candidate(&self.group)?;
let runtime = read_admin_provider_pool_runtime_state(
self.state.app().runtime_state.as_ref(),
self.group.candidate.provider_id.as_str(),
&[],
&pool_config,
self.sticky_session_token.as_deref(),
)
.await;
let sticky_key_id = runtime.sticky_bound_key_id?;
if self.seen_key_ids.contains(&sticky_key_id) {
return None;
}
let key = match self
.state
.app()
.read_provider_catalog_keys_by_ids(std::slice::from_ref(&sticky_key_id))
.await
{
Ok(mut keys) => keys.pop()?,
Err(err) => {
warn!(
event_name = "pool_group_sticky_key_load_failed",
log_type = "event",
provider_id = %self.group.candidate.provider_id,
endpoint_id = %self.group.candidate.endpoint_id,
model_id = %self.group.candidate.model_id,
key_id = %sticky_key_id,
error = ?err,
"gateway pool scheduler failed to read sticky pool key"
);
return None;
}
};
if key.provider_id != self.group.candidate.provider_id {
return None;
}
let candidate = pool_candidate_from_catalog_key(&self.group, key);
self.build_eligible_candidate(candidate).await
}
async fn refill_queued_candidates(&mut self) -> bool {
let mut candidates = Vec::new();
let refill_target = self.page_size.max(1) as usize;
// Keep pool expansion bounded; the cursor schedules one page-sized window at a time.
while candidates.len() < refill_target {
let Some(mut page_candidates) = self.next_page_candidates().await else {
break;
};
candidates.append(&mut page_candidates);
}
if candidates.is_empty() {
return false;
}
let (mut scheduled, mut skipped) = schedule_pool_page_candidates(
self.state,
candidates,
self.sticky_session_token.as_deref(),
)
.await;
self.record_skipped_candidates(&skipped);
self.queued_candidates.extend(scheduled.drain(..));
self.skipped_candidates.append(&mut skipped);
!self.queued_candidates.is_empty()
}
async fn next_queued_candidate(&mut self) -> Option<EligibleLocalExecutionCandidate> {
if self.queued_candidates.len() > 1 {
self.reschedule_queued_candidates().await;
}
while let Some(candidate) = self.queued_candidates.pop_front() {
let (mut scheduled, mut skipped) = schedule_pool_page_candidates(
self.state,
vec![candidate],
self.sticky_session_token.as_deref(),
)
.await;
self.record_skipped_candidates(&skipped);
self.skipped_candidates.append(&mut skipped);
let Some(mut candidate) = scheduled.pop() else {
continue;
};
if !self.attach_pool_key_lease(&mut candidate).await {
continue;
}
candidate.orchestration.pool_key_index = Some(self.next_pool_key_index);
self.next_pool_key_index = self.next_pool_key_index.saturating_add(1);
return Some(candidate);
}
None
}
async fn attach_pool_key_lease(
&mut self,
candidate: &mut EligibleLocalExecutionCandidate,
) -> bool {
if !pool_key_order_requires_exclusive_lease(&self.pool_key_order) {
return true;
}
let owner = self.state.app().tunnel.local_instance_id();
match try_claim_admin_provider_pool_key(
self.state.app().runtime_state.as_ref(),
candidate.candidate.provider_id.as_str(),
candidate.candidate.key_id.as_str(),
owner,
)
.await
{
Ok(Some(lease)) => {
candidate.orchestration.pool_key_lease = Some(lease);
true
}
Ok(None) => {
self.record_skip_reason("pool_key_lease_busy");
false
}
Err(err) => {
warn!(
event_name = "pool_key_lease_claim_failed",
log_type = "event",
provider_id = %candidate.candidate.provider_id,
endpoint_id = %candidate.candidate.endpoint_id,
model_id = %candidate.candidate.model_id,
key_id = %candidate.candidate.key_id,
error = ?err,
"gateway pool scheduler failed to claim pool key lease; scheduling key without lease"
);
true
}
}
}
async fn reschedule_queued_candidates(&mut self) {
let candidates = self.queued_candidates.drain(..).collect::<Vec<_>>();
if candidates.is_empty() {
return;
}
let (mut scheduled, mut skipped) = schedule_pool_page_candidates(
self.state,
candidates,
self.sticky_session_token.as_deref(),
)
.await;
self.record_skipped_candidates(&skipped);
self.queued_candidates.extend(scheduled.drain(..));
self.skipped_candidates.append(&mut skipped);
}
2026-05-03 20:14:29 +08:00
async fn build_page_eligible_candidates(
&mut self,
rows: Vec<StoredMinimalCandidateSelectionRow>,
) -> Vec<EligibleLocalExecutionCandidate> {
let mut candidates = Vec::with_capacity(rows.len());
for row in rows {
let candidate = pool_candidate_from_row(&self.group, row);
2026-05-11 00:12:05 +08:00
if let Some(candidate) = self.build_eligible_candidate(candidate).await {
candidates.push(candidate);
2026-05-03 20:14:29 +08:00
}
}
candidates
}
2026-05-11 00:12:05 +08:00
async fn build_eligible_candidate(
&mut self,
candidate: aether_scheduler_core::SchedulerMinimalCandidateSelectionCandidate,
) -> Option<EligibleLocalExecutionCandidate> {
if !self.seen_key_ids.insert(candidate.key_id.clone()) {
return None;
}
let Some(transport) = read_candidate_transport_snapshot(self.state, &candidate).await
else {
self.record_skip_reason("transport_snapshot_missing");
return None;
};
if let Some(skip_reason) =
candidate_auth_channel_skip_reason(&transport, self.request_auth_channel.as_deref())
{
self.record_skip_reason(skip_reason);
return None;
}
if let Some(skip_reason) = candidate_common_transport_skip_reason(
&transport,
pool_candidate_transport_policy_facts(&candidate),
self.requested_model.as_deref(),
) {
self.record_skip_reason(skip_reason);
return None;
}
Some(EligibleLocalExecutionCandidate {
kind: LocalExecutionCandidateKind::SingleKey,
candidate,
provider_api_format: transport.endpoint.api_format.trim().to_ascii_lowercase(),
transport: std::sync::Arc::new(transport),
orchestration: LocalExecutionCandidateMetadata::default(),
ranking: self.group.ranking.clone(),
})
}
2026-05-03 20:14:29 +08:00
fn record_skip_reason(&mut self, reason: &'static str) {
*self.skip_reason_counts.entry(reason).or_insert(0) += 1;
}
2026-05-11 00:12:05 +08:00
fn record_skipped_candidates(&mut self, skipped_candidates: &[SkippedLocalExecutionCandidate]) {
for skipped_candidate in skipped_candidates {
self.record_skip_reason(skipped_candidate.skip_reason);
}
}
}
fn pool_key_order_requires_exclusive_lease(order: &StoredPoolKeyCandidateOrder) -> bool {
!matches!(
order,
StoredPoolKeyCandidateOrder::CacheAffinity | StoredPoolKeyCandidateOrder::SingleAccount
)
2026-05-03 20:14:29 +08:00
}
fn pool_candidate_transport_policy_facts(
candidate: &aether_scheduler_core::SchedulerMinimalCandidateSelectionCandidate,
) -> CandidateTransportPolicyFacts<'_> {
CandidateTransportPolicyFacts {
endpoint_api_format: candidate.endpoint_api_format.as_str(),
global_model_name: candidate.global_model_name.as_str(),
selected_provider_model_name: candidate.selected_provider_model_name.as_str(),
mapping_matched_model: candidate.mapping_matched_model.as_deref(),
}
}
fn pool_candidate_from_row(
group: &EligibleLocalExecutionCandidate,
row: StoredMinimalCandidateSelectionRow,
) -> aether_scheduler_core::SchedulerMinimalCandidateSelectionCandidate {
let mut candidate = group.candidate.clone();
candidate.key_id = row.key_id;
candidate.key_name = row.key_name;
candidate.key_auth_type = row.key_auth_type;
candidate.key_internal_priority = row.key_internal_priority;
candidate.key_global_priority_for_format =
aether_scheduler_core::extract_global_priority_for_format(
row.key_global_priority_by_format.as_ref(),
group.candidate.endpoint_api_format.as_str(),
)
.ok()
.flatten();
candidate.key_capabilities = row.key_capabilities;
candidate
}
2026-05-11 00:12:05 +08:00
fn pool_candidate_from_catalog_key(
group: &EligibleLocalExecutionCandidate,
key: StoredProviderCatalogKey,
) -> aether_scheduler_core::SchedulerMinimalCandidateSelectionCandidate {
let mut candidate = group.candidate.clone();
candidate.key_id = key.id;
candidate.key_name = key.name;
candidate.key_auth_type = key.auth_type;
candidate.key_internal_priority = key.internal_priority;
candidate.key_global_priority_for_format =
aether_scheduler_core::extract_global_priority_for_format(
key.global_priority_by_format.as_ref(),
group.candidate.endpoint_api_format.as_str(),
)
.ok()
.flatten();
candidate.key_capabilities = key.capabilities;
candidate
}
async fn read_pool_catalog_key_contexts_by_id(
state: PlannerAppState<'_>,
candidates: &[EligibleLocalExecutionCandidate],
) -> BTreeMap<String, PoolCatalogKeyContext> {
let mut key_ids = Vec::new();
let mut provider_type_by_key_id = BTreeMap::<String, String>::new();
for candidate in candidates {
if pool_config_for_candidate(candidate).is_none() {
continue;
}
let key_id = candidate.candidate.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);
}
}
if key_ids.is_empty() {
return BTreeMap::new();
}
let keys = match state
.app()
.read_provider_catalog_keys_by_ids(&key_ids)
.await
{
Ok(keys) => keys,
Err(err) => {
warn!(
error = ?err,
key_count = key_ids.len(),
"gateway pool scheduler: failed to read catalog key metadata"
);
return BTreeMap::new();
}
};
keys.into_iter()
.map(|key| {
let provider_type = provider_type_by_key_id
.get(&key.id)
.map(String::as_str)
.unwrap_or_default();
(
key.id.clone(),
build_pool_catalog_key_context(state, &key, provider_type),
)
})
.collect()
}
fn build_pool_catalog_key_context(
state: PlannerAppState<'_>,
key: &StoredProviderCatalogKey,
provider_type: &str,
) -> PoolCatalogKeyContext {
let status_snapshot = provider_key_status_snapshot_payload(key, provider_type);
let quota_snapshot = status_snapshot
.as_object()
.and_then(|snapshot| snapshot.get("quota"))
.and_then(Value::as_object);
let account_snapshot = status_snapshot
.as_object()
.and_then(|snapshot| snapshot.get("account"))
.and_then(Value::as_object);
let (health_score, _, _, _, _) = provider_key_health_summary(key);
let health_score = key
.health_by_format
.as_ref()
.and_then(Value::as_object)
.filter(|payload| !payload.is_empty())
.map(|_| health_score);
let latency_avg_ms = key
.success_count
.filter(|count| *count > 0)
.zip(key.total_response_time_ms)
.map(|(success_count, total_response_time_ms)| {
f64::from(total_response_time_ms) / f64::from(success_count)
})
.filter(|value| value.is_finite() && *value >= 0.0);
PoolCatalogKeyContext {
oauth_plan_type: quota_snapshot
.and_then(|quota| quota.get("plan_type"))
.and_then(Value::as_str)
.and_then(|value| normalize_pool_plan_type(value, provider_type))
.or_else(|| derive_pool_oauth_plan_type(state, key, provider_type)),
quota_usage_ratio: quota_snapshot
.and_then(|quota| quota.get("usage_ratio"))
.and_then(json_f64)
.map(|value| value.clamp(0.0, 1.0)),
quota_reset_seconds: quota_snapshot
.and_then(|quota| quota.get("reset_seconds"))
.and_then(json_f64)
.filter(|value| *value >= 0.0),
account_blocked: account_snapshot
.and_then(|account| account.get("blocked"))
.and_then(Value::as_bool)
2026-05-11 00:12:05 +08:00
.unwrap_or(false)
|| admin_provider_pool_pure::admin_pool_key_is_known_banned(key),
quota_exhausted: pool_catalog_key_quota_exhausted(key, provider_type, quota_snapshot),
health_score,
latency_avg_ms,
catalog_lru_score: Some(key.last_used_at_unix_secs.unwrap_or(0) as f64),
}
}
2026-05-11 00:12:05 +08:00
fn pool_catalog_key_quota_exhausted(
key: &StoredProviderCatalogKey,
provider_type: &str,
quota_snapshot: Option<&Map<String, Value>>,
) -> bool {
match provider_type.trim().to_ascii_lowercase().as_str() {
"codex" | "kiro" | "chatgpt_web" => {
admin_provider_pool_pure::admin_pool_key_account_quota_exhausted(key, provider_type)
}
_ => quota_snapshot
.and_then(|quota| quota.get("exhausted"))
.and_then(Value::as_bool)
.unwrap_or(false),
}
}
fn derive_pool_oauth_plan_type(
state: PlannerAppState<'_>,
key: &StoredProviderCatalogKey,
provider_type: &str,
) -> Option<String> {
if !provider_key_auth_semantics(key, provider_type).oauth_managed() {
return None;
}
let provider_type_key = provider_type.trim().to_ascii_lowercase();
if let Some(upstream_metadata) = key.upstream_metadata.as_ref().and_then(Value::as_object) {
let provider_bucket = upstream_metadata
.get(&provider_type_key)
.and_then(Value::as_object);
for source in provider_bucket
.into_iter()
.chain(std::iter::once(upstream_metadata))
{
if let Some(plan_type) = pool_plan_type_from_source(
source,
provider_type,
&[
"plan_type",
"tier",
"subscription_title",
"subscription_plan",
"plan",
],
) {
return Some(plan_type);
}
}
}
parse_catalog_auth_config_json(state.app(), key).and_then(|auth_config| {
pool_plan_type_from_source(
&auth_config,
provider_type,
&["plan_type", "tier", "plan", "subscription_plan"],
)
})
}
fn pool_plan_type_from_source(
source: &Map<String, Value>,
provider_type: &str,
fields: &[&str],
) -> Option<String> {
for field in fields {
let Some(value) = source.get(*field).and_then(Value::as_str) else {
continue;
};
if let Some(normalized) = normalize_pool_plan_type(value, provider_type) {
return Some(normalized);
}
}
None
}
fn normalize_pool_plan_type(value: &str, provider_type: &str) -> Option<String> {
let mut normalized = value.trim().to_string();
if normalized.is_empty() {
return None;
}
let provider_type = provider_type.trim().to_ascii_lowercase();
if !provider_type.is_empty() && normalized.to_ascii_lowercase().starts_with(&provider_type) {
normalized = normalized[provider_type.len()..]
.trim_matches(|ch: char| [' ', ':', '-', '_'].contains(&ch))
.to_string();
}
let normalized = normalized.trim().to_ascii_lowercase();
(!normalized.is_empty()).then_some(normalized)
}
fn json_f64(value: &Value) -> Option<f64> {
match value {
Value::Number(number) => number.as_f64(),
Value::String(text) => text.trim().parse::<f64>().ok(),
_ => None,
}
.filter(|value| value.is_finite())
}
fn apply_local_execution_pool_scheduler_with_runtime_map(
candidates: Vec<EligibleLocalExecutionCandidate>,
runtime_by_provider: &BTreeMap<String, AdminProviderPoolRuntimeState>,
key_context_by_id: &BTreeMap<String, PoolCatalogKeyContext>,
) -> (
Vec<EligibleLocalExecutionCandidate>,
Vec<SkippedLocalExecutionCandidate>,
) {
2026-05-02 13:23:54 +08:00
let runtime_by_provider = runtime_by_provider
.iter()
.map(|(provider_id, runtime)| (provider_id.clone(), ai_pool_runtime_state(runtime)))
.collect::<BTreeMap<_, _>>();
let inputs = candidates
.into_iter()
.map(|candidate| {
let key_context = key_context_by_id
.get(&candidate.candidate.key_id)
.cloned()
.unwrap_or_default();
AiPoolCandidateInput {
facts: ai_pool_candidate_facts(&candidate),
pool_config: pool_config_for_candidate(&candidate).map(ai_pool_scheduling_config),
key_context,
candidate,
}
2026-05-02 13:23:54 +08:00
})
.collect::<Vec<_>>();
let outcome = run_ai_pool_scheduler(inputs, &runtime_by_provider, pool_sort_seed().as_str());
2026-05-02 13:23:54 +08:00
let candidates = outcome
.candidates
.into_iter()
.map(|scheduled| apply_ai_pool_orchestration(scheduled.candidate, scheduled.orchestration))
.collect::<Vec<_>>();
let skipped_candidates = outcome
.skipped_candidates
.into_iter()
.map(|skipped| SkippedLocalExecutionCandidate {
candidate: skipped.candidate.candidate,
skip_reason: skipped.skip_reason,
transport: Some(skipped.candidate.transport),
ranking: skipped.candidate.ranking,
extra_data: None,
})
.collect::<Vec<_>>();
2026-05-02 13:23:54 +08:00
(candidates, skipped_candidates)
}
fn pool_config_for_candidate(
candidate: &EligibleLocalExecutionCandidate,
) -> Option<AdminProviderPoolConfig> {
admin_provider_pool_config_from_config_value(candidate.transport.provider.config.as_ref())
}
2026-05-11 00:12:05 +08:00
fn pool_key_candidate_order_for_group(
group: &EligibleLocalExecutionCandidate,
) -> StoredPoolKeyCandidateOrder {
let Some(pool_config) = pool_config_for_candidate(group) else {
return StoredPoolKeyCandidateOrder::InternalPriority;
};
let presets = pool_config
.scheduling_presets
.iter()
.map(|preset| AiPoolSchedulingPreset {
preset: preset.preset.clone(),
enabled: preset.enabled,
mode: preset.mode.clone(),
})
.collect::<Vec<_>>();
let active_presets = normalize_enabled_ai_pool_presets(
&presets,
group.transport.provider.provider_type.as_str(),
);
if let Some(distribution_mode) = active_presets
.iter()
.find(|preset| pool_distribution_mode_preset(preset.as_str()))
.map(String::as_str)
{
return match distribution_mode {
"cache_affinity" => StoredPoolKeyCandidateOrder::CacheAffinity,
"load_balance" => StoredPoolKeyCandidateOrder::LoadBalance {
seed: pool_sort_seed(),
},
"single_account" => StoredPoolKeyCandidateOrder::SingleAccount,
_ => StoredPoolKeyCandidateOrder::InternalPriority,
};
}
if pool_config.lru_enabled {
return StoredPoolKeyCandidateOrder::Lru;
}
StoredPoolKeyCandidateOrder::InternalPriority
}
fn pool_distribution_mode_preset(preset: &str) -> bool {
matches!(preset, "cache_affinity" | "load_balance" | "single_account")
}
2026-05-02 13:23:54 +08:00
fn pool_sort_seed() -> String {
let now_ms = current_unix_ms();
let sequence = LOAD_BALANCE_SEQUENCE.fetch_add(1, AtomicOrdering::Relaxed);
2026-05-02 13:23:54 +08:00
format!("{now_ms}:{sequence}")
}
fn ai_pool_candidate_facts(candidate: &EligibleLocalExecutionCandidate) -> AiPoolCandidateFacts {
AiPoolCandidateFacts {
provider_id: candidate.candidate.provider_id.clone(),
endpoint_id: candidate.candidate.endpoint_id.clone(),
model_id: candidate.candidate.model_id.clone(),
selected_provider_model_name: candidate.candidate.selected_provider_model_name.clone(),
provider_api_format: candidate.provider_api_format.clone(),
provider_type: candidate.transport.provider.provider_type.clone(),
key_id: candidate.candidate.key_id.clone(),
key_internal_priority: candidate.candidate.key_internal_priority,
}
}
2026-05-02 13:23:54 +08:00
fn ai_pool_scheduling_config(config: AdminProviderPoolConfig) -> AiPoolSchedulingConfig {
AiPoolSchedulingConfig {
scheduling_presets: config
.scheduling_presets
.into_iter()
.map(|preset| AiPoolSchedulingPreset {
preset: preset.preset,
enabled: preset.enabled,
mode: preset.mode,
})
.collect(),
lru_enabled: config.lru_enabled,
skip_exhausted_accounts: config.skip_exhausted_accounts,
cost_limit_per_key_tokens: config.cost_limit_per_key_tokens,
}
}
2026-05-02 13:23:54 +08:00
fn ai_pool_runtime_state(runtime: &AdminProviderPoolRuntimeState) -> AiPoolRuntimeState {
AiPoolRuntimeState {
sticky_bound_key_id: runtime.sticky_bound_key_id.clone(),
cooldown_reason_by_key: runtime.cooldown_reason_by_key.clone(),
cost_window_usage_by_key: runtime.cost_window_usage_by_key.clone(),
latency_avg_ms_by_key: runtime.latency_avg_ms_by_key.clone(),
lru_score_by_key: runtime.lru_score_by_key.clone(),
}
}
2026-05-02 13:23:54 +08:00
fn apply_ai_pool_orchestration(
mut candidate: EligibleLocalExecutionCandidate,
orchestration: AiPoolCandidateOrchestration,
) -> EligibleLocalExecutionCandidate {
let scheduler_affinity_epoch = candidate.orchestration.scheduler_affinity_epoch;
2026-05-02 13:23:54 +08:00
candidate.orchestration = LocalExecutionCandidateMetadata {
candidate_group_id: orchestration.candidate_group_id,
pool_key_index: orchestration.pool_key_index,
2026-05-11 00:12:05 +08:00
pool_key_lease: None,
scheduler_affinity_epoch,
2026-05-02 13:23:54 +08:00
};
candidate
}
#[cfg(test)]
mod tests {
use super::{
apply_local_execution_pool_scheduler_with_runtime_map, build_pool_catalog_key_context,
2026-05-11 00:12:05 +08:00
pool_config_for_candidate, pool_key_order_requires_exclusive_lease, PoolCatalogKeyContext,
PoolKeyCursor,
};
2026-05-03 20:14:29 +08:00
use crate::ai_serving::planner::candidate_resolution::{
EligibleLocalExecutionCandidate, LocalExecutionCandidateKind,
};
2026-05-02 13:23:54 +08:00
use crate::ai_serving::PlannerAppState;
use crate::data::GatewayDataState;
2026-05-11 00:12:05 +08:00
use crate::handlers::shared::provider_pool::{
record_admin_provider_pool_error, release_admin_provider_pool_key_lease,
try_claim_admin_provider_pool_key, AdminProviderPoolRuntimeState,
};
use crate::orchestration::LocalExecutionCandidateMetadata;
use crate::AppState;
2026-05-02 13:23:54 +08:00
use aether_ai_serving::{normalize_enabled_ai_pool_presets, AiPoolSchedulingPreset};
2026-05-11 00:12:05 +08:00
use aether_data::repository::candidate_selection::InMemoryMinimalCandidateSelectionReadRepository;
use aether_data::repository::provider_catalog::InMemoryProviderCatalogReadRepository;
2026-05-11 00:12:05 +08:00
use aether_data_contracts::repository::candidate_selection::{
StoredMinimalCandidateSelectionRow, StoredPoolKeyCandidateOrder,
};
use aether_data_contracts::repository::provider_catalog::{
StoredProviderCatalogEndpoint, StoredProviderCatalogKey, StoredProviderCatalogProvider,
};
use aether_provider_transport::snapshot::{
GatewayProviderTransportEndpoint, GatewayProviderTransportKey,
GatewayProviderTransportProvider,
};
use aether_scheduler_core::SchedulerMinimalCandidateSelectionCandidate;
use serde_json::json;
2026-05-11 00:12:05 +08:00
use std::collections::{BTreeMap, VecDeque};
use std::sync::Arc;
#[test]
fn pool_scheduler_groups_interleaved_candidates_and_reorders_internal_keys() {
let pool_first = sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"key-pool-a",
10,
Some(json!({ "pool_advanced": { "lru_enabled": true } })),
);
let other =
sample_eligible_candidate("provider-other", "endpoint-2", "key-other", 10, None);
let pool_second = sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"key-pool-b",
10,
Some(json!({ "pool_advanced": { "lru_enabled": true } })),
);
let mut runtime_by_provider = BTreeMap::new();
runtime_by_provider.insert(
"provider-pool".to_string(),
AdminProviderPoolRuntimeState {
lru_score_by_key: BTreeMap::from([
("key-pool-a".to_string(), 20.0),
("key-pool-b".to_string(), 10.0),
]),
..AdminProviderPoolRuntimeState::default()
},
);
let (reordered, skipped) = apply_local_execution_pool_scheduler_with_runtime_map(
vec![pool_first, other, pool_second],
&runtime_by_provider,
&BTreeMap::new(),
);
assert!(skipped.is_empty());
assert_eq!(
reordered
.iter()
.map(|item| item.candidate.key_id.as_str())
.collect::<Vec<_>>(),
vec!["key-pool-b", "key-pool-a", "key-other"]
);
}
#[test]
fn pool_scheduler_uses_catalog_last_used_when_runtime_lru_is_missing() {
let recent_key = sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"key-recent",
10,
Some(json!({ "pool_advanced": { "lru_enabled": true } })),
);
let older_key = sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"key-older",
10,
Some(json!({ "pool_advanced": { "lru_enabled": true } })),
);
let key_context_by_id = BTreeMap::from([
(
"key-recent".to_string(),
PoolCatalogKeyContext {
catalog_lru_score: Some(200.0),
..PoolCatalogKeyContext::default()
},
),
(
"key-older".to_string(),
PoolCatalogKeyContext {
catalog_lru_score: Some(100.0),
..PoolCatalogKeyContext::default()
},
),
]);
let (reordered, skipped) = apply_local_execution_pool_scheduler_with_runtime_map(
vec![recent_key, older_key],
&BTreeMap::new(),
&key_context_by_id,
);
assert!(skipped.is_empty());
assert_eq!(
reordered
.iter()
.map(|item| item.candidate.key_id.as_str())
.collect::<Vec<_>>(),
vec!["key-older", "key-recent"]
);
}
#[test]
fn pool_scheduler_attaches_group_and_pool_metadata_to_ranked_candidates() {
let pool_first = sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"key-pool-a",
10,
Some(json!({ "pool_advanced": { "lru_enabled": true } })),
);
let other =
sample_eligible_candidate("provider-other", "endpoint-2", "key-other", 10, None);
let pool_second = sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"key-pool-b",
10,
Some(json!({ "pool_advanced": { "lru_enabled": true } })),
);
let mut runtime_by_provider = BTreeMap::new();
runtime_by_provider.insert(
"provider-pool".to_string(),
AdminProviderPoolRuntimeState {
lru_score_by_key: BTreeMap::from([
("key-pool-a".to_string(), 20.0),
("key-pool-b".to_string(), 10.0),
]),
..AdminProviderPoolRuntimeState::default()
},
);
let (reordered, skipped) = apply_local_execution_pool_scheduler_with_runtime_map(
vec![pool_first, other, pool_second],
&runtime_by_provider,
&BTreeMap::new(),
);
assert!(skipped.is_empty());
assert_eq!(reordered.len(), 3);
assert_eq!(
reordered[0].orchestration,
LocalExecutionCandidateMetadata {
candidate_group_id: Some(
"provider=provider-pool|endpoint=endpoint-1|model=model-1|selected_model=gpt-5|api_format=openai:chat|singleton_key=*"
.to_string(),
),
pool_key_index: Some(0),
2026-05-11 00:12:05 +08:00
pool_key_lease: None,
scheduler_affinity_epoch: None,
}
);
assert_eq!(reordered[1].orchestration.pool_key_index, Some(1));
assert_eq!(
reordered[1].orchestration.candidate_group_id,
reordered[0].orchestration.candidate_group_id
);
assert_eq!(
reordered[2].orchestration,
LocalExecutionCandidateMetadata {
candidate_group_id: Some(
"provider=provider-other|endpoint=endpoint-2|model=model-1|selected_model=gpt-5|api_format=openai:chat|singleton_key=key-other"
.to_string(),
),
pool_key_index: None,
2026-05-11 00:12:05 +08:00
pool_key_lease: None,
scheduler_affinity_epoch: None,
}
);
}
#[test]
fn pool_scheduler_promotes_sticky_hit_before_other_sorted_keys() {
let key_a = sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"key-a",
10,
Some(json!({
"pool_advanced": {
"scheduling_presets": [{"preset": "cache_affinity", "enabled": true}]
}
})),
);
let key_b = sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"key-b",
10,
Some(json!({
"pool_advanced": {
"scheduling_presets": [{"preset": "cache_affinity", "enabled": true}]
}
})),
);
let mut runtime_by_provider = BTreeMap::new();
runtime_by_provider.insert(
"provider-pool".to_string(),
AdminProviderPoolRuntimeState {
sticky_bound_key_id: Some("key-a".to_string()),
lru_score_by_key: BTreeMap::from([
("key-a".to_string(), 50.0),
("key-b".to_string(), 10.0),
]),
..AdminProviderPoolRuntimeState::default()
},
);
let (reordered, skipped) = apply_local_execution_pool_scheduler_with_runtime_map(
vec![key_a, key_b],
&runtime_by_provider,
&BTreeMap::new(),
);
assert!(skipped.is_empty());
assert_eq!(
reordered
.iter()
.map(|item| item.candidate.key_id.as_str())
.collect::<Vec<_>>(),
vec!["key-a", "key-b"]
);
}
#[test]
fn pool_scheduler_promotes_sticky_hit_regardless_distribution_mode() {
let key_a = sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"key-a",
10,
Some(json!({
"pool_advanced": {
"scheduling_presets": [{"preset": "quota_balanced", "enabled": true}]
}
})),
);
let key_b = sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"key-b",
10,
Some(json!({
"pool_advanced": {
"scheduling_presets": [{"preset": "quota_balanced", "enabled": true}]
}
})),
);
let runtime_by_provider = BTreeMap::from([(
"provider-pool".to_string(),
AdminProviderPoolRuntimeState {
sticky_bound_key_id: Some("key-a".to_string()),
..AdminProviderPoolRuntimeState::default()
},
)]);
let (reordered, skipped) = apply_local_execution_pool_scheduler_with_runtime_map(
vec![key_b, key_a],
&runtime_by_provider,
&BTreeMap::new(),
);
assert!(skipped.is_empty());
assert_eq!(
reordered
.iter()
.map(|item| item.candidate.key_id.as_str())
.collect::<Vec<_>>(),
vec!["key-a", "key-b"]
);
}
#[test]
fn pool_scheduler_skips_cooldown_and_cost_exhausted_keys() {
let key_ready = sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"key-ready",
10,
Some(json!({
"pool_advanced": {
"cost_limit_per_key_tokens": 100
}
})),
);
let key_cooldown = sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"key-cooldown",
10,
Some(json!({
"pool_advanced": {
"cost_limit_per_key_tokens": 100
}
})),
);
let key_cost = sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"key-cost",
10,
Some(json!({
"pool_advanced": {
"cost_limit_per_key_tokens": 100
}
})),
);
let mut runtime_by_provider = BTreeMap::new();
runtime_by_provider.insert(
"provider-pool".to_string(),
AdminProviderPoolRuntimeState {
cooldown_reason_by_key: BTreeMap::from([(
"key-cooldown".to_string(),
"429".to_string(),
)]),
cost_window_usage_by_key: BTreeMap::from([("key-cost".to_string(), 100)]),
..AdminProviderPoolRuntimeState::default()
},
);
let (reordered, skipped) = apply_local_execution_pool_scheduler_with_runtime_map(
vec![key_ready, key_cooldown, key_cost],
&runtime_by_provider,
&BTreeMap::new(),
);
assert_eq!(
reordered
.iter()
.map(|item| item.candidate.key_id.as_str())
.collect::<Vec<_>>(),
vec!["key-ready"]
);
assert_eq!(
skipped
.iter()
.map(|item| (item.candidate.key_id.as_str(), item.skip_reason))
.collect::<Vec<_>>(),
vec![
("key-cooldown", "pool_cooldown"),
("key-cost", "pool_cost_limit_reached"),
]
);
}
#[test]
2026-05-11 00:12:05 +08:00
fn pool_scheduler_applies_distribution_mode_before_strategy_presets() {
let key_a = sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"key-a",
50,
Some(json!({
"pool_advanced": {
"scheduling_presets": [
2026-05-11 00:12:05 +08:00
{"preset": "cache_affinity", "enabled": true},
{"preset": "priority_first", "enabled": true}
]
}
})),
);
let key_b = sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"key-b",
10,
Some(json!({
"pool_advanced": {
"scheduling_presets": [
2026-05-11 00:12:05 +08:00
{"preset": "cache_affinity", "enabled": true},
{"preset": "priority_first", "enabled": true}
]
}
})),
);
let mut runtime_by_provider = BTreeMap::new();
runtime_by_provider.insert(
"provider-pool".to_string(),
AdminProviderPoolRuntimeState {
lru_score_by_key: BTreeMap::from([
2026-05-11 00:12:05 +08:00
("key-a".to_string(), 100.0),
("key-b".to_string(), 5.0),
]),
..AdminProviderPoolRuntimeState::default()
},
);
let (reordered, skipped) = apply_local_execution_pool_scheduler_with_runtime_map(
vec![key_a, key_b],
&runtime_by_provider,
&BTreeMap::new(),
);
assert!(skipped.is_empty());
assert_eq!(
reordered
.iter()
.map(|item| item.candidate.key_id.as_str())
.collect::<Vec<_>>(),
2026-05-11 00:12:05 +08:00
vec!["key-a", "key-b"]
);
}
#[test]
fn pool_scheduler_uses_plan_preset_with_catalog_context() {
let key_free = sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"key-free",
10,
Some(json!({
"pool_advanced": {
"scheduling_presets": [{"preset": "plus_first", "enabled": true}]
}
})),
);
let key_plus = sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"key-plus",
10,
Some(json!({
"pool_advanced": {
"scheduling_presets": [{"preset": "plus_first", "enabled": true}]
}
})),
);
let key_context_by_id = BTreeMap::from([
(
"key-free".to_string(),
PoolCatalogKeyContext {
oauth_plan_type: Some("free".to_string()),
..PoolCatalogKeyContext::default()
},
),
(
"key-plus".to_string(),
PoolCatalogKeyContext {
oauth_plan_type: Some("plus".to_string()),
..PoolCatalogKeyContext::default()
},
),
]);
let (reordered, skipped) = apply_local_execution_pool_scheduler_with_runtime_map(
vec![key_free, key_plus],
&BTreeMap::new(),
&key_context_by_id,
);
assert!(skipped.is_empty());
assert_eq!(
reordered
.iter()
.map(|item| item.candidate.key_id.as_str())
.collect::<Vec<_>>(),
vec!["key-plus", "key-free"]
);
}
#[test]
fn pool_scheduler_plus_first_treats_plus_and_pro_as_top_tier() {
let key_plus = sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"key-plus",
10,
Some(json!({
"pool_advanced": {
"scheduling_presets": [{"preset": "plus_first", "enabled": true}]
}
})),
);
let key_pro = sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"key-pro",
10,
Some(json!({
"pool_advanced": {
"scheduling_presets": [{"preset": "plus_first", "enabled": true}]
}
})),
);
let key_team = sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"key-team",
10,
Some(json!({
"pool_advanced": {
"scheduling_presets": [{"preset": "plus_first", "enabled": true}]
}
})),
);
let key_context_by_id = BTreeMap::from([
(
"key-plus".to_string(),
PoolCatalogKeyContext {
oauth_plan_type: Some("plus".to_string()),
catalog_lru_score: Some(300.0),
..PoolCatalogKeyContext::default()
},
),
(
"key-pro".to_string(),
PoolCatalogKeyContext {
oauth_plan_type: Some("pro".to_string()),
catalog_lru_score: Some(100.0),
..PoolCatalogKeyContext::default()
},
),
(
"key-team".to_string(),
PoolCatalogKeyContext {
oauth_plan_type: Some("team".to_string()),
catalog_lru_score: Some(50.0),
..PoolCatalogKeyContext::default()
},
),
]);
let (reordered, skipped) = apply_local_execution_pool_scheduler_with_runtime_map(
vec![key_plus, key_pro, key_team],
&BTreeMap::new(),
&key_context_by_id,
);
assert!(skipped.is_empty());
assert_eq!(
reordered
.iter()
.map(|item| item.candidate.key_id.as_str())
.collect::<Vec<_>>(),
vec!["key-plus", "key-pro", "key-team"]
);
}
#[test]
fn pool_scheduler_supports_pro_first_plan_preset() {
let key_plus = sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"key-plus",
10,
Some(json!({
"pool_advanced": {
"scheduling_presets": [{"preset": "pro_first", "enabled": true}]
}
})),
);
let key_pro = sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"key-pro",
10,
Some(json!({
"pool_advanced": {
"scheduling_presets": [{"preset": "pro_first", "enabled": true}]
}
})),
);
let key_team = sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"key-team",
10,
Some(json!({
"pool_advanced": {
"scheduling_presets": [{"preset": "pro_first", "enabled": true}]
}
})),
);
let key_context_by_id = BTreeMap::from([
(
"key-plus".to_string(),
PoolCatalogKeyContext {
oauth_plan_type: Some("plus".to_string()),
..PoolCatalogKeyContext::default()
},
),
(
"key-pro".to_string(),
PoolCatalogKeyContext {
oauth_plan_type: Some("pro".to_string()),
..PoolCatalogKeyContext::default()
},
),
(
"key-team".to_string(),
PoolCatalogKeyContext {
oauth_plan_type: Some("team".to_string()),
..PoolCatalogKeyContext::default()
},
),
]);
let (reordered, skipped) = apply_local_execution_pool_scheduler_with_runtime_map(
vec![key_plus, key_team, key_pro],
&BTreeMap::new(),
&key_context_by_id,
);
assert!(skipped.is_empty());
assert_eq!(
reordered
.iter()
.map(|item| item.candidate.key_id.as_str())
.collect::<Vec<_>>(),
vec!["key-pro", "key-plus", "key-team"]
);
}
#[test]
fn pool_scheduler_defaults_empty_pool_advanced_to_cache_affinity() {
let key_a = sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"key-a",
10,
Some(json!({ "pool_advanced": {} })),
);
let key_b = sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"key-b",
10,
Some(json!({ "pool_advanced": {} })),
);
let runtime_by_provider = BTreeMap::from([(
"provider-pool".to_string(),
AdminProviderPoolRuntimeState {
lru_score_by_key: BTreeMap::from([
("key-a".to_string(), 10.0),
("key-b".to_string(), 200.0),
]),
..AdminProviderPoolRuntimeState::default()
},
)]);
let (reordered, skipped) = apply_local_execution_pool_scheduler_with_runtime_map(
vec![key_a, key_b],
&runtime_by_provider,
&BTreeMap::new(),
);
assert!(skipped.is_empty());
assert_eq!(
reordered
.iter()
.map(|item| item.candidate.key_id.as_str())
.collect::<Vec<_>>(),
vec!["key-b", "key-a"]
);
}
#[test]
2026-05-11 00:12:05 +08:00
fn normalizes_distribution_mode_before_strategy_presets() {
2026-05-02 13:23:54 +08:00
let presets = normalize_enabled_ai_pool_presets(
&[
2026-05-02 13:23:54 +08:00
AiPoolSchedulingPreset {
preset: "lru".to_string(),
enabled: false,
mode: None,
},
2026-05-02 13:23:54 +08:00
AiPoolSchedulingPreset {
preset: "single_account".to_string(),
enabled: true,
mode: None,
},
2026-05-02 13:23:54 +08:00
AiPoolSchedulingPreset {
preset: "cache_affinity".to_string(),
enabled: true,
mode: None,
},
2026-05-02 13:23:54 +08:00
AiPoolSchedulingPreset {
preset: "priority_first".to_string(),
enabled: true,
mode: None,
},
],
"openai",
);
2026-05-02 13:23:54 +08:00
assert_eq!(presets, ["single_account", "priority_first"]);
}
2026-05-11 00:12:05 +08:00
#[test]
fn pool_key_cursor_uses_distribution_order_for_page_queries() {
let app = AppState::new().expect("state should build");
let load_balance_group = sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"pool-group",
10,
Some(json!({
"pool_advanced": {
"scheduling_presets": [
{"preset": "load_balance", "enabled": true},
{"preset": "priority_first", "enabled": true}
]
}
})),
);
let load_balance_cursor = PoolKeyCursor::new(
PlannerAppState::new(&app),
load_balance_group,
None,
None,
None,
);
assert!(matches!(
load_balance_cursor.pool_key_order,
StoredPoolKeyCandidateOrder::LoadBalance { ref seed } if !seed.is_empty()
));
let lru_group = sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"pool-group",
10,
Some(json!({ "pool_advanced": { "lru_enabled": true } })),
);
let lru_cursor =
PoolKeyCursor::new(PlannerAppState::new(&app), lru_group, None, None, None);
assert_eq!(lru_cursor.pool_key_order, StoredPoolKeyCandidateOrder::Lru);
}
#[test]
fn pool_key_lease_exclusivity_follows_distribution_mode() {
assert!(!pool_key_order_requires_exclusive_lease(
&StoredPoolKeyCandidateOrder::CacheAffinity
));
assert!(!pool_key_order_requires_exclusive_lease(
&StoredPoolKeyCandidateOrder::SingleAccount
));
assert!(pool_key_order_requires_exclusive_lease(
&StoredPoolKeyCandidateOrder::Lru
));
assert!(pool_key_order_requires_exclusive_lease(
&StoredPoolKeyCandidateOrder::LoadBalance {
seed: "seed".to_string()
}
));
}
#[tokio::test]
async fn pool_key_cursor_revalidates_queued_candidates_before_return() {
let app = AppState::new().expect("state should build");
let provider_config = Some(json!({
"pool_advanced": {
"rate_limit_cooldown_seconds": 300
}
}));
let group = sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"pool-group",
10,
provider_config.clone(),
);
let mut cursor =
PoolKeyCursor::new(PlannerAppState::new(&app), group.clone(), None, None, None);
cursor.queued_candidates = VecDeque::from([
sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"key-a",
10,
provider_config.clone(),
),
sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"key-b",
10,
provider_config.clone(),
),
sample_eligible_candidate("provider-pool", "endpoint-1", "key-c", 10, provider_config),
]);
let first = cursor.next_key().await.expect("first key should schedule");
assert_eq!(first.candidate.key_id, "key-a");
let pool_config = pool_config_for_candidate(&group).expect("pool config should parse");
record_admin_provider_pool_error(
app.runtime_state.as_ref(),
"provider-pool",
"key-b",
&pool_config,
429,
None,
None,
)
.await;
let second = cursor
.next_key()
.await
.expect("second schedulable key should skip cooled-down key");
assert_eq!(second.candidate.key_id, "key-c");
assert_eq!(second.orchestration.pool_key_index, Some(1));
assert!(cursor.next_key().await.is_none());
let skipped = cursor.take_skipped_candidates();
assert!(skipped.iter().any(|candidate| {
candidate.candidate.key_id == "key-b" && candidate.skip_reason == "pool_cooldown"
}));
}
#[tokio::test]
async fn pool_key_cursor_skips_busy_lease_and_claims_returned_key() {
let app = AppState::new().expect("state should build");
let provider_config = Some(json!({ "pool_advanced": { "lru_enabled": true } }));
let held_lease = try_claim_admin_provider_pool_key(
app.runtime_state.as_ref(),
"provider-pool",
"key-a",
"test-owner",
)
.await
.expect("lease claim should not fail")
.expect("first key should be claimed by test");
let group = sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"pool-group",
10,
provider_config.clone(),
);
let mut cursor = PoolKeyCursor::new(PlannerAppState::new(&app), group, None, None, None);
cursor.queued_candidates = VecDeque::from([
sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"key-a",
10,
provider_config.clone(),
),
sample_eligible_candidate("provider-pool", "endpoint-1", "key-b", 10, provider_config),
]);
let candidate = cursor
.next_key()
.await
.expect("cursor should skip busy key and return next key");
assert_eq!(candidate.candidate.key_id, "key-b");
assert_eq!(candidate.orchestration.pool_key_index, Some(0));
assert_eq!(
cursor.skip_reason_counts.get("pool_key_lease_busy"),
Some(&1)
);
let returned_lease = candidate
.orchestration
.pool_key_lease
.clone()
.expect("returned key should carry its lease");
assert!(
try_claim_admin_provider_pool_key(
app.runtime_state.as_ref(),
"provider-pool",
"key-b",
"other-owner",
)
.await
.expect("second claim should not fail")
.is_none(),
"returned key should remain leased until execution releases it"
);
release_admin_provider_pool_key_lease(app.runtime_state.as_ref(), &held_lease)
.await
.expect("held lease release should not fail");
release_admin_provider_pool_key_lease(app.runtime_state.as_ref(), &returned_lease)
.await
.expect("returned lease release should not fail");
}
#[tokio::test]
async fn pool_key_cursor_allows_busy_lease_for_cache_affinity_mode() {
let app = AppState::new().expect("state should build");
let provider_config = Some(json!({
"pool_advanced": {
"scheduling_presets": [
{"preset": "cache_affinity", "enabled": true}
]
}
}));
let held_lease = try_claim_admin_provider_pool_key(
app.runtime_state.as_ref(),
"provider-pool",
"key-a",
"test-owner",
)
.await
.expect("lease claim should not fail")
.expect("first key should be claimed by test");
let group = sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"pool-group",
10,
provider_config.clone(),
);
let mut cursor = PoolKeyCursor::new(PlannerAppState::new(&app), group, None, None, None);
cursor.queued_candidates = VecDeque::from([
sample_eligible_candidate(
"provider-pool",
"endpoint-1",
"key-a",
10,
provider_config.clone(),
),
sample_eligible_candidate("provider-pool", "endpoint-1", "key-b", 10, provider_config),
]);
let candidate = cursor
.next_key()
.await
.expect("cache-affinity cursor should allow repeated busy key");
assert_eq!(candidate.candidate.key_id, "key-a");
assert_eq!(candidate.orchestration.pool_key_index, Some(0));
assert!(candidate.orchestration.pool_key_lease.is_none());
assert!(!cursor
.skip_reason_counts
.contains_key("pool_key_lease_busy"));
release_admin_provider_pool_key_lease(app.runtime_state.as_ref(), &held_lease)
.await
.expect("held lease release should not fail");
}
#[tokio::test]
async fn pool_key_cursor_simulates_large_lru_pool_with_lazy_pages_and_dynamic_skips() {
const KEY_COUNT: usize = 2048;
let provider_config = Some(json!({
"pool_advanced": {
"lru_enabled": true,
"rate_limit_cooldown_seconds": 300
}
}));
let (provider, endpoint, keys, rows) =
large_pool_fixture(KEY_COUNT, provider_config.clone());
let data_state =
GatewayDataState::with_provider_catalog_and_minimal_candidate_selection_for_tests(
Arc::new(InMemoryProviderCatalogReadRepository::seed(
vec![provider],
vec![endpoint],
keys,
)),
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_eligible_candidate(
"provider-pool",
"endpoint-1",
"pool-group",
10,
provider_config.clone(),
);
let pool_config = pool_config_for_candidate(&group).expect("pool config should parse");
let mut held_leases = Vec::new();
for key_id in ["key-00000", "key-00001"] {
held_leases.push(
try_claim_admin_provider_pool_key(
app.runtime_state.as_ref(),
"provider-pool",
key_id,
"large-pool-test",
)
.await
.expect("lease claim should not fail")
.expect("test key should claim"),
);
}
for key_id in ["key-00002", "key-00003"] {
record_admin_provider_pool_error(
app.runtime_state.as_ref(),
"provider-pool",
key_id,
&pool_config,
429,
None,
None,
)
.await;
}
let mut cursor = PoolKeyCursor::new(PlannerAppState::new(&app), group, None, None, None);
cursor.page_size = 64;
cursor.max_scanned_keys = KEY_COUNT as u32;
let mut returned_ids = Vec::new();
let mut returned_leases = Vec::new();
for _ in 0..10 {
let candidate = cursor
.next_key()
.await
.expect("large pool should return first page candidates");
returned_ids.push(candidate.candidate.key_id.clone());
returned_leases.push(
candidate
.orchestration
.pool_key_lease
.expect("lru pool candidate should carry exclusive lease"),
);
}
assert_eq!(returned_ids.first().map(String::as_str), Some("key-00004"));
assert_eq!(returned_ids.last().map(String::as_str), Some("key-00013"));
assert_eq!(cursor.scanned_keys, 64);
assert!(
cursor.queued_candidates.len() <= cursor.page_size as usize,
"cursor should only retain the current page window"
);
record_admin_provider_pool_error(
app.runtime_state.as_ref(),
"provider-pool",
"key-00014",
&pool_config,
429,
None,
None,
)
.await;
held_leases.push(
try_claim_admin_provider_pool_key(
app.runtime_state.as_ref(),
"provider-pool",
"key-00015",
"large-pool-test",
)
.await
.expect("dynamic lease claim should not fail")
.expect("dynamic busy key should claim"),
);
let candidate = cursor
.next_key()
.await
.expect("cursor should skip dynamically unhealthy and busy keys");
assert_eq!(candidate.candidate.key_id, "key-00016");
returned_ids.push(candidate.candidate.key_id.clone());
returned_leases.push(
candidate
.orchestration
.pool_key_lease
.expect("lru pool candidate should carry exclusive lease"),
);
while let Some(candidate) = cursor.next_key().await {
returned_ids.push(candidate.candidate.key_id.clone());
returned_leases.push(
candidate
.orchestration
.pool_key_lease
.expect("lru pool candidate should carry exclusive lease"),
);
}
assert_eq!(returned_ids.len(), KEY_COUNT - 6);
assert_eq!(returned_ids.last().map(String::as_str), Some("key-02047"));
assert_eq!(cursor.scanned_keys, KEY_COUNT as u32);
assert_eq!(
cursor.skip_reason_counts.get("pool_key_lease_busy"),
Some(&3)
);
assert_eq!(cursor.skip_reason_counts.get("pool_cooldown"), Some(&3));
for skipped in [
"key-00000",
"key-00001",
"key-00002",
"key-00003",
"key-00014",
"key-00015",
] {
assert!(
!returned_ids.iter().any(|key_id| key_id == skipped),
"{skipped} should have been skipped"
);
}
for lease in held_leases.into_iter().chain(returned_leases) {
release_admin_provider_pool_key_lease(app.runtime_state.as_ref(), &lease)
.await
.expect("lease release should not fail");
}
}
#[test]
fn builds_pool_catalog_context_from_status_snapshot_and_auth_config() {
let mut key = StoredProviderCatalogKey::new(
"key-1".to_string(),
"provider-1".to_string(),
"key-1".to_string(),
"oauth".to_string(),
None,
true,
)
.expect("key should build")
.with_transport_fields(
None,
"secret".to_string(),
None,
None,
None,
None,
None,
None,
None,
)
.expect("transport fields should build");
key.status_snapshot = Some(json!({
"account": {"blocked": false},
"quota": {
"usage_ratio": 0.25,
"reset_seconds": 3600,
"exhausted": false,
"plan_type": "team"
}
}));
key.success_count = Some(4);
key.total_response_time_ms = Some(200);
key.last_used_at_unix_secs = Some(1_711_000_123);
let app = AppState::new()
.expect("state should build")
.with_data_state_for_tests(GatewayDataState::with_provider_catalog_reader_for_tests(
Arc::new(InMemoryProviderCatalogReadRepository::seed(
Vec::new(),
Vec::new(),
vec![key.clone()],
)),
));
let context = build_pool_catalog_key_context(PlannerAppState::new(&app), &key, "codex");
assert_eq!(context.oauth_plan_type.as_deref(), Some("team"));
assert_eq!(context.quota_usage_ratio, Some(0.25));
assert_eq!(context.quota_reset_seconds, Some(3600.0));
assert_eq!(context.latency_avg_ms, Some(50.0));
assert_eq!(context.catalog_lru_score, Some(1_711_000_123.0));
}
2026-05-11 00:12:05 +08:00
#[test]
fn pool_catalog_context_ignores_stale_codex_exhausted_snapshot_when_windows_have_capacity() {
let mut key = sample_catalog_oauth_key("key-stale-exhausted");
key.upstream_metadata = Some(json!({
"codex": {
"primary_used_percent": 100.0
}
}));
key.status_snapshot = Some(json!({
"quota": {
"version": 2,
"provider_type": "codex",
"code": "exhausted",
"exhausted": true,
"usage_ratio": 0.0,
"windows": [
{
"code": "weekly",
"used_ratio": 0.0,
"remaining_ratio": 1.0
},
{
"code": "5h",
"used_ratio": 0.0,
"remaining_ratio": 1.0
}
]
}
}));
let app = app_state_with_catalog_key(key.clone());
let context = build_pool_catalog_key_context(PlannerAppState::new(&app), &key, "codex");
assert!(!context.quota_exhausted);
}
#[test]
fn pool_catalog_context_marks_codex_metadata_exhausted() {
let mut key = sample_catalog_oauth_key("key-metadata-exhausted");
key.upstream_metadata = Some(json!({
"codex": {
"secondary_used_percent": 100.0
}
}));
let app = app_state_with_catalog_key(key.clone());
let context = build_pool_catalog_key_context(PlannerAppState::new(&app), &key, "codex");
assert!(context.quota_exhausted);
}
#[test]
fn pool_catalog_context_preserves_snapshot_exhaustion_for_snapshot_only_providers() {
let mut key = sample_catalog_oauth_key("key-antigravity-exhausted");
key.status_snapshot = Some(json!({
"quota": {
"version": 2,
"provider_type": "antigravity",
"code": "exhausted",
"exhausted": true,
"windows": [
{
"code": "gemini-2.5-pro",
"used_ratio": 1.0,
"remaining_ratio": 0.0
}
]
}
}));
let app = app_state_with_catalog_key(key.clone());
let context =
build_pool_catalog_key_context(PlannerAppState::new(&app), &key, "antigravity");
assert!(context.quota_exhausted);
}
#[test]
fn pool_catalog_context_marks_known_banned_account_from_metadata() {
let mut key = sample_catalog_oauth_key("key-account-banned");
key.upstream_metadata = Some(json!({
"codex": {
"account_disabled": true,
"reason": "deactivated_workspace"
}
}));
let app = app_state_with_catalog_key(key.clone());
let context = build_pool_catalog_key_context(PlannerAppState::new(&app), &key, "codex");
assert!(context.account_blocked);
}
fn sample_catalog_oauth_key(key_id: &str) -> StoredProviderCatalogKey {
StoredProviderCatalogKey::new(
key_id.to_string(),
"provider-1".to_string(),
key_id.to_string(),
"oauth".to_string(),
None,
true,
)
.expect("key should build")
.with_transport_fields(
None,
"secret".to_string(),
None,
None,
None,
None,
None,
None,
None,
)
.expect("transport fields should build")
}
fn app_state_with_catalog_key(key: StoredProviderCatalogKey) -> AppState {
AppState::new()
.expect("state should build")
.with_data_state_for_tests(GatewayDataState::with_provider_catalog_reader_for_tests(
Arc::new(InMemoryProviderCatalogReadRepository::seed(
Vec::new(),
Vec::new(),
vec![key],
)),
))
}
fn large_pool_fixture(
key_count: usize,
provider_config: Option<serde_json::Value>,
) -> (
StoredProviderCatalogProvider,
StoredProviderCatalogEndpoint,
Vec<StoredProviderCatalogKey>,
Vec<StoredMinimalCandidateSelectionRow>,
) {
let provider = StoredProviderCatalogProvider::new(
"provider-pool".to_string(),
"provider-pool".to_string(),
Some("https://example.com".to_string()),
"openai".to_string(),
)
.expect("provider should build")
.with_routing_fields(0)
.with_transport_fields(
true,
false,
false,
None,
None,
None,
None,
None,
provider_config,
);
let endpoint = StoredProviderCatalogEndpoint::new(
"endpoint-1".to_string(),
"provider-pool".to_string(),
"openai:chat".to_string(),
Some("openai".to_string()),
Some("chat".to_string()),
true,
)
.expect("endpoint should build")
.with_health_score(1.0)
.with_transport_fields(
"https://example.com/v1/chat/completions".to_string(),
None,
None,
None,
None,
None,
None,
None,
)
.expect("endpoint transport should build");
let mut keys = Vec::with_capacity(key_count);
let mut rows = Vec::with_capacity(key_count);
for index in 0..key_count {
let key_id = format!("key-{index:05}");
let mut key = StoredProviderCatalogKey::new(
key_id.clone(),
"provider-pool".to_string(),
key_id.clone(),
"api_key".to_string(),
None,
true,
)
.expect("key should build")
.with_transport_fields(
Some(json!(["openai:chat"])),
Some(format!("secret-{index}")),
None,
None,
None,
None,
None,
None,
None,
)
.expect("key transport should build");
key.internal_priority = 10;
key.last_used_at_unix_secs = Some(index as u64);
keys.push(key);
rows.push(StoredMinimalCandidateSelectionRow {
provider_id: "provider-pool".to_string(),
provider_name: "provider-pool".to_string(),
provider_type: "openai".to_string(),
provider_priority: 0,
provider_is_active: true,
endpoint_id: "endpoint-1".to_string(),
endpoint_api_format: "openai:chat".to_string(),
endpoint_api_family: Some("openai".to_string()),
endpoint_kind: Some("chat".to_string()),
endpoint_is_active: true,
key_id: key_id.clone(),
key_name: key_id,
key_auth_type: "api_key".to_string(),
key_is_active: true,
key_api_formats: Some(vec!["openai:chat".to_string()]),
key_allowed_models: None,
key_capabilities: None,
key_internal_priority: 10,
key_global_priority_by_format: None,
model_id: "model-1".to_string(),
global_model_id: "global-model-1".to_string(),
global_model_name: "gpt-5".to_string(),
global_model_mappings: None,
global_model_supports_streaming: Some(true),
model_provider_model_name: "gpt-5".to_string(),
model_provider_model_mappings: None,
model_supports_streaming: Some(true),
model_is_active: true,
model_is_available: true,
});
}
(provider, endpoint, keys, rows)
}
fn sample_eligible_candidate(
provider_id: &str,
endpoint_id: &str,
key_id: &str,
internal_priority: i32,
provider_config: Option<serde_json::Value>,
) -> EligibleLocalExecutionCandidate {
EligibleLocalExecutionCandidate {
2026-05-03 20:14:29 +08:00
kind: if provider_config.is_some() {
LocalExecutionCandidateKind::PoolGroup
} else {
LocalExecutionCandidateKind::SingleKey
},
candidate: SchedulerMinimalCandidateSelectionCandidate {
provider_id: provider_id.to_string(),
provider_name: provider_id.to_string(),
provider_type: "codex".to_string(),
provider_priority: 10,
endpoint_id: endpoint_id.to_string(),
endpoint_api_format: "openai:chat".to_string(),
key_id: key_id.to_string(),
key_name: key_id.to_string(),
key_auth_type: "api_key".to_string(),
key_internal_priority: internal_priority,
key_global_priority_for_format: Some(1),
key_capabilities: None,
model_id: "model-1".to_string(),
global_model_id: "global-model-1".to_string(),
global_model_name: "gpt-5".to_string(),
selected_provider_model_name: "gpt-5".to_string(),
mapping_matched_model: None,
},
provider_api_format: "openai:chat".to_string(),
orchestration: LocalExecutionCandidateMetadata::default(),
2026-04-27 12:34:03 +08:00
ranking: None,
2026-05-02 13:23:54 +08:00
transport: Arc::new(crate::ai_serving::GatewayProviderTransportSnapshot {
provider: GatewayProviderTransportProvider {
id: provider_id.to_string(),
name: provider_id.to_string(),
provider_type: "codex".to_string(),
website: None,
is_active: true,
keep_priority_on_conversion: false,
enable_format_conversion: false,
concurrent_limit: None,
max_retries: None,
proxy: None,
request_timeout_secs: None,
stream_first_byte_timeout_secs: None,
config: provider_config,
},
endpoint: GatewayProviderTransportEndpoint {
id: endpoint_id.to_string(),
provider_id: provider_id.to_string(),
api_format: "openai:chat".to_string(),
api_family: Some("openai".to_string()),
endpoint_kind: Some("chat".to_string()),
is_active: true,
base_url: "https://example.com".to_string(),
header_rules: None,
body_rules: None,
max_retries: None,
custom_path: None,
config: None,
format_acceptance_config: None,
proxy: None,
},
key: GatewayProviderTransportKey {
id: key_id.to_string(),
provider_id: provider_id.to_string(),
name: key_id.to_string(),
auth_type: "api_key".to_string(),
is_active: true,
api_formats: Some(vec!["openai:chat".to_string()]),
2026-04-29 15:46:50 +08:00
auth_type_by_format: None,
allow_auth_channel_mismatch_formats: None,
2026-04-29 15:46:50 +08:00
allowed_models: None,
capabilities: None,
rate_multipliers: None,
global_priority_by_format: None,
expires_at_unix_secs: None,
proxy: None,
fingerprint: None,
decrypted_api_key: "secret".to_string(),
decrypted_auth_config: None,
},
}),
}
}
}