Rename transport ordering facts

This commit is contained in:
fawney19
2026-04-27 16:21:57 +08:00
parent 218e73e324
commit 98e4e91a98
3 changed files with 39 additions and 38 deletions

View File

@@ -18,8 +18,8 @@ use aether_scheduler_core::{
use super::candidate_affinity_cache::read_cached_scheduler_affinity_target; use super::candidate_affinity_cache::read_cached_scheduler_affinity_target;
use super::candidate_resolution::EligibleLocalExecutionCandidate; use super::candidate_resolution::EligibleLocalExecutionCandidate;
use super::candidate_transport_ordering::{ use super::candidate_transport_ranking_facts::{
resolve_cached_transport_execution_ordering, CandidateExecutionOrdering, resolve_cached_transport_ranking_facts, CandidateTransportRankingFacts,
}; };
pub(crate) async fn rank_eligible_local_execution_candidates( pub(crate) async fn rank_eligible_local_execution_candidates(
@@ -50,7 +50,7 @@ pub(crate) async fn rank_eligible_local_execution_candidates(
let mut ordering_cache = BTreeMap::new(); let mut ordering_cache = BTreeMap::new();
for (original_index, eligible) in candidates.iter().enumerate() { for (original_index, eligible) in candidates.iter().enumerate() {
let ordering = resolve_cached_transport_execution_ordering( let ranking_facts = resolve_cached_transport_ranking_facts(
state, state,
&mut ordering_cache, &mut ordering_cache,
&eligible.candidate, &eligible.candidate,
@@ -61,7 +61,7 @@ pub(crate) async fn rank_eligible_local_execution_candidates(
rankables.push(rankable_candidate_from_candidate( rankables.push(rankable_candidate_from_candidate(
&eligible.candidate, &eligible.candidate,
original_index, original_index,
ordering, ranking_facts,
normalized_client_api_format, normalized_client_api_format,
eligible.provider_api_format.as_str(), eligible.provider_api_format.as_str(),
required_capabilities, required_capabilities,
@@ -89,7 +89,7 @@ pub(crate) async fn rank_eligible_local_execution_candidates(
fn rankable_candidate_from_candidate( fn rankable_candidate_from_candidate(
candidate: &SchedulerMinimalCandidateSelectionCandidate, candidate: &SchedulerMinimalCandidateSelectionCandidate,
original_index: usize, original_index: usize,
ordering: CandidateExecutionOrdering, ranking_facts: CandidateTransportRankingFacts,
normalized_client_api_format: &str, normalized_client_api_format: &str,
provider_api_format: &str, provider_api_format: &str,
required_capabilities: Option<&serde_json::Value>, required_capabilities: Option<&serde_json::Value>,
@@ -109,9 +109,9 @@ fn rankable_candidate_from_candidate(
candidate, candidate,
)) ))
.with_cached_affinity_match(cached_affinity_match) .with_cached_affinity_match(cached_affinity_match)
.with_tunnel_bucket(ordering.tunnel_bucket) .with_tunnel_bucket(ranking_facts.tunnel_bucket)
.with_format_state( .with_format_state(
!is_same_format && !ordering.keep_priority_on_conversion, !is_same_format && !ranking_facts.keep_priority_on_conversion,
candidate_api_format_preference(normalized_client_api_format, provider_api_format), candidate_api_format_preference(normalized_client_api_format, provider_api_format),
) )
} }
@@ -195,7 +195,7 @@ mod tests {
use serde_json::json; use serde_json::json;
use super::super::candidate_affinity_cache::remember_scheduler_affinity_for_candidate; use super::super::candidate_affinity_cache::remember_scheduler_affinity_for_candidate;
use super::super::candidate_transport_ordering::resolve_cached_candidate_execution_ordering; use super::super::candidate_transport_ranking_facts::resolve_cached_candidate_transport_ranking_facts;
use super::{PlannerAppState, SchedulerMinimalCandidateSelectionCandidate}; use super::{PlannerAppState, SchedulerMinimalCandidateSelectionCandidate};
use crate::ai_pipeline::planner::candidate_resolution::filter_and_rank_local_execution_candidates; use crate::ai_pipeline::planner::candidate_resolution::filter_and_rank_local_execution_candidates;
use crate::data::auth::GatewayAuthApiKeySnapshot; use crate::data::auth::GatewayAuthApiKeySnapshot;
@@ -217,7 +217,7 @@ mod tests {
let mut ordering_cache = BTreeMap::new(); let mut ordering_cache = BTreeMap::new();
for (original_index, candidate) in candidates.iter().enumerate() { for (original_index, candidate) in candidates.iter().enumerate() {
let ordering = resolve_cached_candidate_execution_ordering( let ranking_facts = resolve_cached_candidate_transport_ranking_facts(
state, state,
&mut ordering_cache, &mut ordering_cache,
candidate, candidate,
@@ -227,7 +227,7 @@ mod tests {
rankables.push(super::rankable_candidate_from_candidate( rankables.push(super::rankable_candidate_from_candidate(
candidate, candidate,
original_index, original_index,
ordering, ranking_facts,
normalized_client_api_format.as_str(), normalized_client_api_format.as_str(),
candidate.endpoint_api_format.as_str(), candidate.endpoint_api_format.as_str(),
required_capabilities, required_capabilities,

View File

@@ -13,67 +13,68 @@ use super::candidate_resolution::read_candidate_transport_snapshot;
pub(super) type CandidateTransportIdentity<'a> = (&'a str, &'a str, &'a str); pub(super) type CandidateTransportIdentity<'a> = (&'a str, &'a str, &'a str);
#[derive(Debug, Clone, Copy, PartialEq, Eq)] #[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(super) struct CandidateExecutionOrdering { pub(super) struct CandidateTransportRankingFacts {
pub(super) tunnel_bucket: SchedulerTunnelAffinityBucket, pub(super) tunnel_bucket: SchedulerTunnelAffinityBucket,
pub(super) keep_priority_on_conversion: bool, pub(super) keep_priority_on_conversion: bool,
} }
pub(super) async fn resolve_cached_candidate_execution_ordering<'a>( pub(super) async fn resolve_cached_candidate_transport_ranking_facts<'a>(
state: PlannerAppState<'_>, state: PlannerAppState<'_>,
cache: &mut BTreeMap<CandidateTransportIdentity<'a>, CandidateExecutionOrdering>, cache: &mut BTreeMap<CandidateTransportIdentity<'a>, CandidateTransportRankingFacts>,
candidate: &'a SchedulerMinimalCandidateSelectionCandidate, candidate: &'a SchedulerMinimalCandidateSelectionCandidate,
ordering_config: SchedulerOrderingConfig, ordering_config: SchedulerOrderingConfig,
) -> CandidateExecutionOrdering { ) -> CandidateTransportRankingFacts {
let identity = candidate_transport_identity(candidate); let identity = candidate_transport_identity(candidate);
if let Some(ordering) = cache.get(&identity).copied() { if let Some(facts) = cache.get(&identity).copied() {
return ordering; return facts;
} }
let ordering = resolve_candidate_execution_ordering(state, candidate, ordering_config).await; let facts = resolve_candidate_transport_ranking_facts(state, candidate, ordering_config).await;
cache.insert(identity, ordering); cache.insert(identity, facts);
ordering facts
} }
pub(super) async fn resolve_cached_transport_execution_ordering<'a>( pub(super) async fn resolve_cached_transport_ranking_facts<'a>(
state: PlannerAppState<'_>, state: PlannerAppState<'_>,
cache: &mut BTreeMap<CandidateTransportIdentity<'a>, CandidateExecutionOrdering>, cache: &mut BTreeMap<CandidateTransportIdentity<'a>, CandidateTransportRankingFacts>,
candidate: &'a SchedulerMinimalCandidateSelectionCandidate, candidate: &'a SchedulerMinimalCandidateSelectionCandidate,
transport: &GatewayProviderTransportSnapshot, transport: &GatewayProviderTransportSnapshot,
ordering_config: SchedulerOrderingConfig, ordering_config: SchedulerOrderingConfig,
) -> CandidateExecutionOrdering { ) -> CandidateTransportRankingFacts {
let identity = candidate_transport_identity(candidate); let identity = candidate_transport_identity(candidate);
if let Some(ordering) = cache.get(&identity).copied() { if let Some(facts) = cache.get(&identity).copied() {
return ordering; return facts;
} }
let ordering = let facts =
resolve_candidate_execution_ordering_from_transport(state, transport, ordering_config) resolve_candidate_transport_ranking_facts_from_transport(state, transport, ordering_config)
.await; .await;
cache.insert(identity, ordering); cache.insert(identity, facts);
ordering facts
} }
async fn resolve_candidate_execution_ordering( async fn resolve_candidate_transport_ranking_facts(
state: PlannerAppState<'_>, state: PlannerAppState<'_>,
candidate: &SchedulerMinimalCandidateSelectionCandidate, candidate: &SchedulerMinimalCandidateSelectionCandidate,
ordering_config: SchedulerOrderingConfig, ordering_config: SchedulerOrderingConfig,
) -> CandidateExecutionOrdering { ) -> CandidateTransportRankingFacts {
let Some(transport) = read_candidate_transport_snapshot(state, candidate).await else { let Some(transport) = read_candidate_transport_snapshot(state, candidate).await else {
return CandidateExecutionOrdering { return CandidateTransportRankingFacts {
tunnel_bucket: SchedulerTunnelAffinityBucket::Neutral, tunnel_bucket: SchedulerTunnelAffinityBucket::Neutral,
keep_priority_on_conversion: ordering_config.keep_priority_on_conversion, keep_priority_on_conversion: ordering_config.keep_priority_on_conversion,
}; };
}; };
resolve_candidate_execution_ordering_from_transport(state, &transport, ordering_config).await resolve_candidate_transport_ranking_facts_from_transport(state, &transport, ordering_config)
.await
} }
async fn resolve_candidate_execution_ordering_from_transport( async fn resolve_candidate_transport_ranking_facts_from_transport(
state: PlannerAppState<'_>, state: PlannerAppState<'_>,
transport: &GatewayProviderTransportSnapshot, transport: &GatewayProviderTransportSnapshot,
ordering_config: SchedulerOrderingConfig, ordering_config: SchedulerOrderingConfig,
) -> CandidateExecutionOrdering { ) -> CandidateTransportRankingFacts {
CandidateExecutionOrdering { CandidateTransportRankingFacts {
tunnel_bucket: resolve_tunnel_owner_affinity_from_transport(state, transport).await, tunnel_bucket: resolve_tunnel_owner_affinity_from_transport(state, transport).await,
keep_priority_on_conversion: ordering_config.keep_priority_on_conversion keep_priority_on_conversion: ordering_config.keep_priority_on_conversion
|| transport.provider.keep_priority_on_conversion, || transport.provider.keep_priority_on_conversion,
@@ -120,11 +121,11 @@ async fn resolve_tunnel_owner_affinity_from_transport(
Ok(None) => SchedulerTunnelAffinityBucket::Neutral, Ok(None) => SchedulerTunnelAffinityBucket::Neutral,
Err(error) => { Err(error) => {
warn!( warn!(
event_name = "candidate_transport_ordering_tunnel_owner_lookup_failed", event_name = "candidate_transport_ranking_facts_tunnel_owner_lookup_failed",
log_type = "event", log_type = "event",
node_id = node_id, node_id = node_id,
error = %error, error = %error,
"failed to load tunnel attachment owner while evaluating scheduler candidate ordering" "failed to load tunnel attachment owner while evaluating candidate transport ranking facts"
); );
SchedulerTunnelAffinityBucket::Neutral SchedulerTunnelAffinityBucket::Neutral
} }

View File

@@ -11,7 +11,7 @@ mod candidate_preparation;
mod candidate_ranking; mod candidate_ranking;
mod candidate_resolution; mod candidate_resolution;
mod candidate_source; mod candidate_source;
mod candidate_transport_ordering; mod candidate_transport_ranking_facts;
mod common; mod common;
mod decision; mod decision;
mod decision_input; mod decision_input;