mirror of
https://github.com/fawney19/Aether.git
synced 2026-09-02 01:10:23 +08:00
Clean up planner candidate modules
This commit is contained in:
@@ -0,0 +1,65 @@
|
||||
use aether_scheduler_core::{
|
||||
build_scheduler_affinity_cache_key_for_api_key_id, SchedulerAffinityTarget,
|
||||
SchedulerMinimalCandidateSelectionCandidate,
|
||||
};
|
||||
|
||||
use crate::ai_pipeline::{GatewayAuthApiKeySnapshot, PlannerAppState};
|
||||
use crate::scheduler::affinity::SCHEDULER_AFFINITY_TTL;
|
||||
|
||||
const PLANNER_SCHEDULER_AFFINITY_MAX_ENTRIES: usize = 10_000;
|
||||
|
||||
pub(crate) fn read_cached_scheduler_affinity_target(
|
||||
state: PlannerAppState<'_>,
|
||||
auth_snapshot: Option<&GatewayAuthApiKeySnapshot>,
|
||||
client_api_format: &str,
|
||||
requested_model: Option<&str>,
|
||||
) -> Option<SchedulerAffinityTarget> {
|
||||
let requested_model = requested_model
|
||||
.map(str::trim)
|
||||
.filter(|value| !value.is_empty())?;
|
||||
let api_key_id = auth_snapshot
|
||||
.map(|snapshot| snapshot.api_key_id.trim())
|
||||
.filter(|value| !value.is_empty())?;
|
||||
let cache_key = build_scheduler_affinity_cache_key_for_api_key_id(
|
||||
api_key_id,
|
||||
client_api_format,
|
||||
requested_model,
|
||||
)?;
|
||||
|
||||
state
|
||||
.app()
|
||||
.read_scheduler_affinity_target(&cache_key, SCHEDULER_AFFINITY_TTL)
|
||||
}
|
||||
|
||||
pub(crate) fn remember_scheduler_affinity_for_candidate(
|
||||
state: PlannerAppState<'_>,
|
||||
auth_snapshot: Option<&GatewayAuthApiKeySnapshot>,
|
||||
client_api_format: &str,
|
||||
requested_model: &str,
|
||||
candidate: &SchedulerMinimalCandidateSelectionCandidate,
|
||||
) {
|
||||
let Some(api_key_id) = auth_snapshot
|
||||
.map(|snapshot| snapshot.api_key_id.trim())
|
||||
.filter(|value| !value.is_empty())
|
||||
else {
|
||||
return;
|
||||
};
|
||||
let Some(cache_key) = build_scheduler_affinity_cache_key_for_api_key_id(
|
||||
api_key_id,
|
||||
client_api_format,
|
||||
requested_model,
|
||||
) else {
|
||||
return;
|
||||
};
|
||||
|
||||
state.app().remember_scheduler_affinity_target(
|
||||
&cache_key,
|
||||
SchedulerAffinityTarget {
|
||||
provider_id: candidate.provider_id.clone(),
|
||||
endpoint_id: candidate.endpoint_id.clone(),
|
||||
key_id: candidate.key_id.clone(),
|
||||
},
|
||||
SCHEDULER_AFFINITY_TTL,
|
||||
PLANNER_SCHEDULER_AFFINITY_MAX_ENTRIES,
|
||||
);
|
||||
}
|
||||
@@ -1 +0,0 @@
|
||||
pub(crate) use super::candidate_resolution::*;
|
||||
@@ -3,7 +3,7 @@ use aether_scheduler_core::SchedulerRankingOutcome;
|
||||
use serde_json::Value;
|
||||
use uuid::Uuid;
|
||||
|
||||
use crate::ai_pipeline::planner::candidate_affinity::remember_scheduler_affinity_for_candidate;
|
||||
use crate::ai_pipeline::planner::candidate_affinity_cache::remember_scheduler_affinity_for_candidate;
|
||||
use crate::ai_pipeline::planner::candidate_resolution::{
|
||||
EligibleLocalExecutionCandidate, SkippedLocalExecutionCandidate,
|
||||
};
|
||||
|
||||
@@ -3,35 +3,26 @@ use std::collections::BTreeMap;
|
||||
use tracing::warn;
|
||||
|
||||
use crate::ai_pipeline::{
|
||||
request_candidate_api_format_preference, GatewayAuthApiKeySnapshot,
|
||||
GatewayProviderTransportSnapshot, PlannerAppState,
|
||||
request_candidate_api_format_preference, GatewayAuthApiKeySnapshot, PlannerAppState,
|
||||
};
|
||||
use crate::handlers::shared::provider_pool::admin_provider_pool_config_from_config_value;
|
||||
use crate::scheduler::affinity::SCHEDULER_AFFINITY_TTL;
|
||||
use crate::scheduler::config::{
|
||||
read_scheduler_ordering_config, SchedulerOrderingConfig, SchedulerSchedulingMode,
|
||||
};
|
||||
use aether_scheduler_core::{
|
||||
apply_scheduler_candidate_ranking, build_scheduler_affinity_cache_key_for_api_key_id,
|
||||
matches_affinity_target, requested_capability_priority_for_candidate, SchedulerAffinityTarget,
|
||||
apply_scheduler_candidate_ranking, matches_affinity_target,
|
||||
requested_capability_priority_for_candidate, SchedulerAffinityTarget,
|
||||
SchedulerMinimalCandidateSelectionCandidate, SchedulerPriorityMode, SchedulerRankableCandidate,
|
||||
SchedulerRankingContext, SchedulerRankingMode, SchedulerTunnelAffinityBucket,
|
||||
};
|
||||
|
||||
use super::candidate_resolution::{
|
||||
read_candidate_transport_snapshot, EligibleLocalExecutionCandidate,
|
||||
use super::candidate_affinity_cache::read_cached_scheduler_affinity_target;
|
||||
use super::candidate_resolution::EligibleLocalExecutionCandidate;
|
||||
use super::candidate_transport_ordering::{
|
||||
resolve_cached_candidate_execution_ordering, resolve_cached_candidate_tunnel_owner_affinity,
|
||||
resolve_cached_transport_execution_ordering, CandidateExecutionOrdering,
|
||||
};
|
||||
|
||||
const PLANNER_SCHEDULER_AFFINITY_MAX_ENTRIES: usize = 10_000;
|
||||
|
||||
type CandidateTransportIdentity<'a> = (&'a str, &'a str, &'a str);
|
||||
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
struct CandidateExecutionOrdering {
|
||||
tunnel_bucket: SchedulerTunnelAffinityBucket,
|
||||
keep_priority_on_conversion: bool,
|
||||
}
|
||||
|
||||
pub(crate) async fn prefer_local_tunnel_owner_candidates(
|
||||
state: PlannerAppState<'_>,
|
||||
candidates: Vec<SchedulerMinimalCandidateSelectionCandidate>,
|
||||
@@ -137,21 +128,29 @@ pub(crate) async fn rank_eligible_local_execution_candidates(
|
||||
) -> Vec<EligibleLocalExecutionCandidate> {
|
||||
let ordering_config = read_scheduler_ordering_config_or_default(state).await;
|
||||
let mut candidates = candidates;
|
||||
let cached_affinity_target = read_cached_affinity_target_for_ranking(
|
||||
let affinity_requested_model = requested_model
|
||||
.map(str::trim)
|
||||
.filter(|value| !value.is_empty())
|
||||
.or_else(|| {
|
||||
candidates
|
||||
.first()
|
||||
.map(|candidate| candidate.candidate.global_model_name.as_str())
|
||||
});
|
||||
let cached_affinity_target = read_cached_scheduler_affinity_target(
|
||||
state,
|
||||
auth_snapshot,
|
||||
normalized_client_api_format,
|
||||
requested_model,
|
||||
&candidates,
|
||||
affinity_requested_model,
|
||||
);
|
||||
let mut rankables = Vec::with_capacity(candidates.len());
|
||||
let mut ordering_cache = BTreeMap::new();
|
||||
|
||||
for (original_index, eligible) in candidates.iter().enumerate() {
|
||||
let ordering = resolve_cached_eligible_candidate_execution_ordering(
|
||||
let ordering = resolve_cached_transport_execution_ordering(
|
||||
state,
|
||||
&mut ordering_cache,
|
||||
eligible,
|
||||
&eligible.candidate,
|
||||
eligible.transport.as_ref(),
|
||||
ordering_config,
|
||||
)
|
||||
.await;
|
||||
@@ -230,35 +229,6 @@ fn local_execution_candidate_uses_pool(eligible: &EligibleLocalExecutionCandidat
|
||||
.is_some()
|
||||
}
|
||||
|
||||
fn read_cached_affinity_target_for_ranking(
|
||||
state: PlannerAppState<'_>,
|
||||
auth_snapshot: Option<&GatewayAuthApiKeySnapshot>,
|
||||
client_api_format: &str,
|
||||
requested_model: Option<&str>,
|
||||
candidates: &[EligibleLocalExecutionCandidate],
|
||||
) -> Option<SchedulerAffinityTarget> {
|
||||
let requested_model = requested_model
|
||||
.map(str::trim)
|
||||
.filter(|value| !value.is_empty())
|
||||
.or_else(|| {
|
||||
candidates
|
||||
.first()
|
||||
.map(|candidate| candidate.candidate.global_model_name.as_str())
|
||||
})?;
|
||||
let api_key_id = auth_snapshot
|
||||
.map(|snapshot| snapshot.api_key_id.trim())
|
||||
.filter(|value| !value.is_empty())?;
|
||||
let cache_key = build_scheduler_affinity_cache_key_for_api_key_id(
|
||||
api_key_id,
|
||||
client_api_format,
|
||||
requested_model,
|
||||
)?;
|
||||
|
||||
state
|
||||
.app()
|
||||
.read_scheduler_affinity_target(&cache_key, SCHEDULER_AFFINITY_TTL)
|
||||
}
|
||||
|
||||
fn planner_ranking_context(ordering_config: SchedulerOrderingConfig) -> SchedulerRankingContext {
|
||||
SchedulerRankingContext {
|
||||
priority_mode: ordering_config.priority_mode,
|
||||
@@ -284,190 +254,6 @@ fn api_format_matches(left: &str, right: &str) -> bool {
|
||||
normalize_api_format_alias(left) == normalize_api_format_alias(right)
|
||||
}
|
||||
|
||||
pub(crate) fn remember_scheduler_affinity_for_candidate(
|
||||
state: PlannerAppState<'_>,
|
||||
auth_snapshot: Option<&GatewayAuthApiKeySnapshot>,
|
||||
client_api_format: &str,
|
||||
requested_model: &str,
|
||||
candidate: &SchedulerMinimalCandidateSelectionCandidate,
|
||||
) {
|
||||
let Some(api_key_id) = auth_snapshot
|
||||
.map(|snapshot| snapshot.api_key_id.trim())
|
||||
.filter(|value| !value.is_empty())
|
||||
else {
|
||||
return;
|
||||
};
|
||||
let Some(cache_key) = build_scheduler_affinity_cache_key_for_api_key_id(
|
||||
api_key_id,
|
||||
client_api_format,
|
||||
requested_model,
|
||||
) else {
|
||||
return;
|
||||
};
|
||||
|
||||
state.app().remember_scheduler_affinity_target(
|
||||
&cache_key,
|
||||
SchedulerAffinityTarget {
|
||||
provider_id: candidate.provider_id.clone(),
|
||||
endpoint_id: candidate.endpoint_id.clone(),
|
||||
key_id: candidate.key_id.clone(),
|
||||
},
|
||||
SCHEDULER_AFFINITY_TTL,
|
||||
PLANNER_SCHEDULER_AFFINITY_MAX_ENTRIES,
|
||||
);
|
||||
}
|
||||
|
||||
async fn resolve_candidate_tunnel_owner_affinity(
|
||||
state: PlannerAppState<'_>,
|
||||
candidate: &SchedulerMinimalCandidateSelectionCandidate,
|
||||
) -> SchedulerTunnelAffinityBucket {
|
||||
let Some(transport) = read_candidate_transport_snapshot(state, candidate).await else {
|
||||
return SchedulerTunnelAffinityBucket::Neutral;
|
||||
};
|
||||
|
||||
resolve_tunnel_owner_affinity_from_transport(state, &transport).await
|
||||
}
|
||||
|
||||
async fn resolve_cached_candidate_tunnel_owner_affinity<'a>(
|
||||
state: PlannerAppState<'_>,
|
||||
cache: &mut BTreeMap<CandidateTransportIdentity<'a>, SchedulerTunnelAffinityBucket>,
|
||||
candidate: &'a SchedulerMinimalCandidateSelectionCandidate,
|
||||
) -> SchedulerTunnelAffinityBucket {
|
||||
let identity = candidate_transport_identity(candidate);
|
||||
if let Some(bucket) = cache.get(&identity).copied() {
|
||||
return bucket;
|
||||
}
|
||||
|
||||
let bucket = resolve_candidate_tunnel_owner_affinity(state, candidate).await;
|
||||
cache.insert(identity, bucket);
|
||||
bucket
|
||||
}
|
||||
|
||||
async fn resolve_candidate_execution_ordering(
|
||||
state: PlannerAppState<'_>,
|
||||
candidate: &SchedulerMinimalCandidateSelectionCandidate,
|
||||
ordering_config: SchedulerOrderingConfig,
|
||||
) -> CandidateExecutionOrdering {
|
||||
let Some(transport) = read_candidate_transport_snapshot(state, candidate).await else {
|
||||
return CandidateExecutionOrdering {
|
||||
tunnel_bucket: SchedulerTunnelAffinityBucket::Neutral,
|
||||
keep_priority_on_conversion: ordering_config.keep_priority_on_conversion,
|
||||
};
|
||||
};
|
||||
|
||||
resolve_candidate_execution_ordering_from_transport(state, &transport, ordering_config).await
|
||||
}
|
||||
|
||||
async fn resolve_cached_candidate_execution_ordering<'a>(
|
||||
state: PlannerAppState<'_>,
|
||||
cache: &mut BTreeMap<CandidateTransportIdentity<'a>, CandidateExecutionOrdering>,
|
||||
candidate: &'a SchedulerMinimalCandidateSelectionCandidate,
|
||||
ordering_config: SchedulerOrderingConfig,
|
||||
) -> CandidateExecutionOrdering {
|
||||
let identity = candidate_transport_identity(candidate);
|
||||
if let Some(ordering) = cache.get(&identity).copied() {
|
||||
return ordering;
|
||||
}
|
||||
|
||||
let ordering = resolve_candidate_execution_ordering(state, candidate, ordering_config).await;
|
||||
cache.insert(identity, ordering);
|
||||
ordering
|
||||
}
|
||||
|
||||
async fn resolve_cached_eligible_candidate_execution_ordering<'a>(
|
||||
state: PlannerAppState<'_>,
|
||||
cache: &mut BTreeMap<CandidateTransportIdentity<'a>, CandidateExecutionOrdering>,
|
||||
eligible: &'a EligibleLocalExecutionCandidate,
|
||||
ordering_config: SchedulerOrderingConfig,
|
||||
) -> CandidateExecutionOrdering {
|
||||
let identity = candidate_transport_identity(&eligible.candidate);
|
||||
if let Some(ordering) = cache.get(&identity).copied() {
|
||||
return ordering;
|
||||
}
|
||||
|
||||
let ordering = resolve_candidate_execution_ordering_from_transport(
|
||||
state,
|
||||
&eligible.transport,
|
||||
ordering_config,
|
||||
)
|
||||
.await;
|
||||
cache.insert(identity, ordering);
|
||||
ordering
|
||||
}
|
||||
|
||||
async fn resolve_candidate_execution_ordering_from_transport(
|
||||
state: PlannerAppState<'_>,
|
||||
transport: &GatewayProviderTransportSnapshot,
|
||||
ordering_config: SchedulerOrderingConfig,
|
||||
) -> CandidateExecutionOrdering {
|
||||
CandidateExecutionOrdering {
|
||||
tunnel_bucket: resolve_tunnel_owner_affinity_from_transport(state, transport).await,
|
||||
keep_priority_on_conversion: ordering_config.keep_priority_on_conversion
|
||||
|| transport.provider.keep_priority_on_conversion,
|
||||
}
|
||||
}
|
||||
|
||||
async fn resolve_tunnel_owner_affinity_from_transport(
|
||||
state: PlannerAppState<'_>,
|
||||
transport: &GatewayProviderTransportSnapshot,
|
||||
) -> SchedulerTunnelAffinityBucket {
|
||||
let Some(proxy) = state
|
||||
.app()
|
||||
.resolve_transport_proxy_snapshot_with_tunnel_affinity(transport)
|
||||
.await
|
||||
else {
|
||||
return SchedulerTunnelAffinityBucket::Neutral;
|
||||
};
|
||||
if proxy.enabled == Some(false) {
|
||||
return SchedulerTunnelAffinityBucket::Neutral;
|
||||
}
|
||||
let Some(node_id) = proxy
|
||||
.node_id
|
||||
.as_deref()
|
||||
.map(str::trim)
|
||||
.filter(|value| !value.is_empty())
|
||||
else {
|
||||
return SchedulerTunnelAffinityBucket::Neutral;
|
||||
};
|
||||
|
||||
if state.app().tunnel.has_local_proxy(node_id) {
|
||||
return SchedulerTunnelAffinityBucket::LocalTunnel;
|
||||
}
|
||||
|
||||
match state
|
||||
.app()
|
||||
.tunnel
|
||||
.lookup_attachment_owner(state.app().data.as_ref(), node_id)
|
||||
.await
|
||||
{
|
||||
Ok(Some(owner)) if owner.gateway_instance_id == state.app().tunnel.local_instance_id() => {
|
||||
SchedulerTunnelAffinityBucket::LocalTunnel
|
||||
}
|
||||
Ok(Some(_)) => SchedulerTunnelAffinityBucket::RemoteTunnel,
|
||||
Ok(None) => SchedulerTunnelAffinityBucket::Neutral,
|
||||
Err(error) => {
|
||||
warn!(
|
||||
event_name = "candidate_affinity_tunnel_owner_lookup_failed",
|
||||
log_type = "event",
|
||||
node_id = node_id,
|
||||
error = %error,
|
||||
"failed to load tunnel attachment owner while evaluating scheduler affinity"
|
||||
);
|
||||
SchedulerTunnelAffinityBucket::Neutral
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn candidate_transport_identity(
|
||||
candidate: &SchedulerMinimalCandidateSelectionCandidate,
|
||||
) -> CandidateTransportIdentity<'_> {
|
||||
(
|
||||
candidate.provider_id.as_str(),
|
||||
candidate.endpoint_id.as_str(),
|
||||
candidate.key_id.as_str(),
|
||||
)
|
||||
}
|
||||
|
||||
fn candidate_api_format_preference(client_api_format: &str, provider_api_format: &str) -> (u8, u8) {
|
||||
request_candidate_api_format_preference(client_api_format, provider_api_format)
|
||||
.unwrap_or((u8::MAX, u8::MAX))
|
||||
@@ -501,9 +287,9 @@ mod tests {
|
||||
use aether_scheduler_core::RANKING_REASON_CACHED_AFFINITY;
|
||||
use serde_json::json;
|
||||
|
||||
use super::super::candidate_affinity_cache::remember_scheduler_affinity_for_candidate;
|
||||
use super::{
|
||||
prefer_local_tunnel_owner_candidates, rank_local_execution_candidates,
|
||||
remember_scheduler_affinity_for_candidate, PlannerAppState,
|
||||
prefer_local_tunnel_owner_candidates, rank_local_execution_candidates, PlannerAppState,
|
||||
SchedulerMinimalCandidateSelectionCandidate,
|
||||
};
|
||||
use crate::ai_pipeline::planner::candidate_resolution::filter_and_rank_local_execution_candidates;
|
||||
@@ -1466,7 +1252,7 @@ mod tests {
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn realtime_gate_skips_cross_format_candidates_when_conversion_is_disabled() {
|
||||
async fn realtime_gate_reports_cross_format_candidates_when_conversion_is_disabled() {
|
||||
let provider_catalog = InMemoryProviderCatalogReadRepository::seed(
|
||||
vec![
|
||||
sample_provider_with_options("provider-cross", true, 0),
|
||||
@@ -1533,11 +1319,13 @@ mod tests {
|
||||
|
||||
assert_eq!(ranked.len(), 1);
|
||||
assert_eq!(ranked[0].candidate.endpoint_id, "endpoint-same");
|
||||
assert!(skipped.is_empty());
|
||||
assert_eq!(skipped.len(), 1);
|
||||
assert_eq!(skipped[0].candidate.endpoint_id, "endpoint-cross");
|
||||
assert_eq!(skipped[0].skip_reason, "format_conversion_disabled");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn realtime_gate_hides_cross_format_disablement_when_same_key_has_exact_endpoint() {
|
||||
async fn realtime_gate_reports_cross_format_disablement_when_same_key_has_exact_endpoint() {
|
||||
let provider_catalog = InMemoryProviderCatalogReadRepository::seed(
|
||||
vec![sample_provider_with_options("provider-shared", false, 0)],
|
||||
vec![
|
||||
@@ -1591,7 +1379,9 @@ mod tests {
|
||||
|
||||
assert_eq!(ranked.len(), 1);
|
||||
assert_eq!(ranked[0].candidate.endpoint_id, "endpoint-exact");
|
||||
assert!(skipped.is_empty());
|
||||
assert_eq!(skipped.len(), 1);
|
||||
assert_eq!(skipped[0].candidate.endpoint_id, "endpoint-cross");
|
||||
assert_eq!(skipped[0].skip_reason, "format_conversion_disabled");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
@@ -10,7 +10,7 @@ use crate::ai_pipeline::{
|
||||
};
|
||||
use crate::orchestration::LocalExecutionCandidateMetadata;
|
||||
|
||||
use super::candidate_affinity::rank_eligible_local_execution_candidates;
|
||||
use super::candidate_ranking::rank_eligible_local_execution_candidates;
|
||||
use super::pool_scheduler::apply_local_execution_pool_scheduler;
|
||||
|
||||
#[derive(Debug, Clone, PartialEq)]
|
||||
|
||||
@@ -0,0 +1,168 @@
|
||||
use std::collections::BTreeMap;
|
||||
|
||||
use aether_scheduler_core::{
|
||||
SchedulerMinimalCandidateSelectionCandidate, SchedulerTunnelAffinityBucket,
|
||||
};
|
||||
use tracing::warn;
|
||||
|
||||
use crate::ai_pipeline::{GatewayProviderTransportSnapshot, PlannerAppState};
|
||||
use crate::scheduler::config::SchedulerOrderingConfig;
|
||||
|
||||
use super::candidate_resolution::read_candidate_transport_snapshot;
|
||||
|
||||
pub(super) type CandidateTransportIdentity<'a> = (&'a str, &'a str, &'a str);
|
||||
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
pub(super) struct CandidateExecutionOrdering {
|
||||
pub(super) tunnel_bucket: SchedulerTunnelAffinityBucket,
|
||||
pub(super) keep_priority_on_conversion: bool,
|
||||
}
|
||||
|
||||
pub(super) async fn resolve_cached_candidate_tunnel_owner_affinity<'a>(
|
||||
state: PlannerAppState<'_>,
|
||||
cache: &mut BTreeMap<CandidateTransportIdentity<'a>, SchedulerTunnelAffinityBucket>,
|
||||
candidate: &'a SchedulerMinimalCandidateSelectionCandidate,
|
||||
) -> SchedulerTunnelAffinityBucket {
|
||||
let identity = candidate_transport_identity(candidate);
|
||||
if let Some(bucket) = cache.get(&identity).copied() {
|
||||
return bucket;
|
||||
}
|
||||
|
||||
let bucket = resolve_candidate_tunnel_owner_affinity(state, candidate).await;
|
||||
cache.insert(identity, bucket);
|
||||
bucket
|
||||
}
|
||||
|
||||
pub(super) async fn resolve_cached_candidate_execution_ordering<'a>(
|
||||
state: PlannerAppState<'_>,
|
||||
cache: &mut BTreeMap<CandidateTransportIdentity<'a>, CandidateExecutionOrdering>,
|
||||
candidate: &'a SchedulerMinimalCandidateSelectionCandidate,
|
||||
ordering_config: SchedulerOrderingConfig,
|
||||
) -> CandidateExecutionOrdering {
|
||||
let identity = candidate_transport_identity(candidate);
|
||||
if let Some(ordering) = cache.get(&identity).copied() {
|
||||
return ordering;
|
||||
}
|
||||
|
||||
let ordering = resolve_candidate_execution_ordering(state, candidate, ordering_config).await;
|
||||
cache.insert(identity, ordering);
|
||||
ordering
|
||||
}
|
||||
|
||||
pub(super) async fn resolve_cached_transport_execution_ordering<'a>(
|
||||
state: PlannerAppState<'_>,
|
||||
cache: &mut BTreeMap<CandidateTransportIdentity<'a>, CandidateExecutionOrdering>,
|
||||
candidate: &'a SchedulerMinimalCandidateSelectionCandidate,
|
||||
transport: &GatewayProviderTransportSnapshot,
|
||||
ordering_config: SchedulerOrderingConfig,
|
||||
) -> CandidateExecutionOrdering {
|
||||
let identity = candidate_transport_identity(candidate);
|
||||
if let Some(ordering) = cache.get(&identity).copied() {
|
||||
return ordering;
|
||||
}
|
||||
|
||||
let ordering =
|
||||
resolve_candidate_execution_ordering_from_transport(state, transport, ordering_config)
|
||||
.await;
|
||||
cache.insert(identity, ordering);
|
||||
ordering
|
||||
}
|
||||
|
||||
async fn resolve_candidate_tunnel_owner_affinity(
|
||||
state: PlannerAppState<'_>,
|
||||
candidate: &SchedulerMinimalCandidateSelectionCandidate,
|
||||
) -> SchedulerTunnelAffinityBucket {
|
||||
let Some(transport) = read_candidate_transport_snapshot(state, candidate).await else {
|
||||
return SchedulerTunnelAffinityBucket::Neutral;
|
||||
};
|
||||
|
||||
resolve_tunnel_owner_affinity_from_transport(state, &transport).await
|
||||
}
|
||||
|
||||
async fn resolve_candidate_execution_ordering(
|
||||
state: PlannerAppState<'_>,
|
||||
candidate: &SchedulerMinimalCandidateSelectionCandidate,
|
||||
ordering_config: SchedulerOrderingConfig,
|
||||
) -> CandidateExecutionOrdering {
|
||||
let Some(transport) = read_candidate_transport_snapshot(state, candidate).await else {
|
||||
return CandidateExecutionOrdering {
|
||||
tunnel_bucket: SchedulerTunnelAffinityBucket::Neutral,
|
||||
keep_priority_on_conversion: ordering_config.keep_priority_on_conversion,
|
||||
};
|
||||
};
|
||||
|
||||
resolve_candidate_execution_ordering_from_transport(state, &transport, ordering_config).await
|
||||
}
|
||||
|
||||
async fn resolve_candidate_execution_ordering_from_transport(
|
||||
state: PlannerAppState<'_>,
|
||||
transport: &GatewayProviderTransportSnapshot,
|
||||
ordering_config: SchedulerOrderingConfig,
|
||||
) -> CandidateExecutionOrdering {
|
||||
CandidateExecutionOrdering {
|
||||
tunnel_bucket: resolve_tunnel_owner_affinity_from_transport(state, transport).await,
|
||||
keep_priority_on_conversion: ordering_config.keep_priority_on_conversion
|
||||
|| transport.provider.keep_priority_on_conversion,
|
||||
}
|
||||
}
|
||||
|
||||
async fn resolve_tunnel_owner_affinity_from_transport(
|
||||
state: PlannerAppState<'_>,
|
||||
transport: &GatewayProviderTransportSnapshot,
|
||||
) -> SchedulerTunnelAffinityBucket {
|
||||
let Some(proxy) = state
|
||||
.app()
|
||||
.resolve_transport_proxy_snapshot_with_tunnel_affinity(transport)
|
||||
.await
|
||||
else {
|
||||
return SchedulerTunnelAffinityBucket::Neutral;
|
||||
};
|
||||
if proxy.enabled == Some(false) {
|
||||
return SchedulerTunnelAffinityBucket::Neutral;
|
||||
}
|
||||
let Some(node_id) = proxy
|
||||
.node_id
|
||||
.as_deref()
|
||||
.map(str::trim)
|
||||
.filter(|value| !value.is_empty())
|
||||
else {
|
||||
return SchedulerTunnelAffinityBucket::Neutral;
|
||||
};
|
||||
|
||||
if state.app().tunnel.has_local_proxy(node_id) {
|
||||
return SchedulerTunnelAffinityBucket::LocalTunnel;
|
||||
}
|
||||
|
||||
match state
|
||||
.app()
|
||||
.tunnel
|
||||
.lookup_attachment_owner(state.app().data.as_ref(), node_id)
|
||||
.await
|
||||
{
|
||||
Ok(Some(owner)) if owner.gateway_instance_id == state.app().tunnel.local_instance_id() => {
|
||||
SchedulerTunnelAffinityBucket::LocalTunnel
|
||||
}
|
||||
Ok(Some(_)) => SchedulerTunnelAffinityBucket::RemoteTunnel,
|
||||
Ok(None) => SchedulerTunnelAffinityBucket::Neutral,
|
||||
Err(error) => {
|
||||
warn!(
|
||||
event_name = "candidate_transport_ordering_tunnel_owner_lookup_failed",
|
||||
log_type = "event",
|
||||
node_id = node_id,
|
||||
error = %error,
|
||||
"failed to load tunnel attachment owner while evaluating scheduler candidate ordering"
|
||||
);
|
||||
SchedulerTunnelAffinityBucket::Neutral
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn candidate_transport_identity(
|
||||
candidate: &SchedulerMinimalCandidateSelectionCandidate,
|
||||
) -> CandidateTransportIdentity<'_> {
|
||||
(
|
||||
candidate.provider_id.as_str(),
|
||||
candidate.endpoint_id.as_str(),
|
||||
candidate.key_id.as_str(),
|
||||
)
|
||||
}
|
||||
@@ -4,13 +4,14 @@ use crate::ai_pipeline::contracts::{
|
||||
use crate::ai_pipeline::GatewayControlDecision;
|
||||
use crate::{AppState, GatewayError};
|
||||
|
||||
mod candidate_affinity;
|
||||
mod candidate_eligibility;
|
||||
mod candidate_affinity_cache;
|
||||
mod candidate_materialization;
|
||||
mod candidate_metadata;
|
||||
mod candidate_preparation;
|
||||
mod candidate_ranking;
|
||||
mod candidate_resolution;
|
||||
mod candidate_source;
|
||||
mod candidate_transport_ordering;
|
||||
mod common;
|
||||
mod decision;
|
||||
mod decision_input;
|
||||
|
||||
@@ -11,7 +11,7 @@ use serde_json::{json, Value};
|
||||
use tracing::warn;
|
||||
use uuid::Uuid;
|
||||
|
||||
use crate::ai_pipeline::planner::candidate_affinity::prefer_local_tunnel_owner_candidates;
|
||||
use crate::ai_pipeline::planner::candidate_ranking::prefer_local_tunnel_owner_candidates;
|
||||
use crate::ai_pipeline::planner::common::{
|
||||
EXECUTION_RUNTIME_STREAM_DECISION_ACTION, EXECUTION_RUNTIME_SYNC_DECISION_ACTION,
|
||||
};
|
||||
|
||||
@@ -418,13 +418,14 @@ fn ai_pipeline_planner_gateway_state_seam_is_split_by_role() {
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn ai_pipeline_planner_separates_local_candidate_eligibility_from_affinity_ranking() {
|
||||
fn ai_pipeline_planner_separates_local_candidate_resolution_from_ranking() {
|
||||
let planner_mod = read_workspace_file("apps/aether-gateway/src/ai_pipeline/planner/mod.rs");
|
||||
for pattern in [
|
||||
"mod candidate_affinity;",
|
||||
"mod candidate_eligibility;",
|
||||
"mod candidate_affinity_cache;",
|
||||
"mod candidate_ranking;",
|
||||
"mod candidate_resolution;",
|
||||
"mod candidate_preparation;",
|
||||
"mod candidate_transport_ordering;",
|
||||
] {
|
||||
assert!(
|
||||
planner_mod.contains(pattern),
|
||||
@@ -445,31 +446,74 @@ fn ai_pipeline_planner_separates_local_candidate_eligibility_from_affinity_ranki
|
||||
);
|
||||
}
|
||||
|
||||
let candidate_eligibility =
|
||||
read_workspace_file("apps/aether-gateway/src/ai_pipeline/planner/candidate_eligibility.rs");
|
||||
assert!(
|
||||
candidate_eligibility.contains("pub(crate) use super::candidate_resolution::*;"),
|
||||
"planner/candidate_eligibility.rs should remain a compatibility shim"
|
||||
);
|
||||
assert!(
|
||||
!candidate_eligibility.contains("async fn filter_and_rank_local_execution_candidates("),
|
||||
"planner/candidate_eligibility.rs should not keep resolution implementation"
|
||||
!planner_mod.contains("mod candidate_eligibility;"),
|
||||
"planner should not wire the removed candidate eligibility compatibility shim"
|
||||
);
|
||||
|
||||
let candidate_affinity =
|
||||
read_workspace_file("apps/aether-gateway/src/ai_pipeline/planner/candidate_affinity.rs");
|
||||
let candidate_ranking =
|
||||
read_workspace_file("apps/aether-gateway/src/ai_pipeline/planner/candidate_ranking.rs");
|
||||
assert!(
|
||||
candidate_affinity.contains("#[cfg(test)]\nasync fn rank_local_execution_candidates("),
|
||||
"planner/candidate_affinity.rs should keep raw local ranking as a test-only helper"
|
||||
candidate_ranking.contains("#[cfg(test)]\nasync fn rank_local_execution_candidates("),
|
||||
"planner/candidate_ranking.rs should keep raw local ranking as a test-only helper"
|
||||
);
|
||||
for forbidden in [
|
||||
"struct SkippedLocalExecutionCandidate",
|
||||
"async fn current_local_execution_candidate_skip_reason(",
|
||||
"pub(crate) async fn filter_and_rank_local_execution_candidates(",
|
||||
"resolve_transport_proxy_snapshot_with_tunnel_affinity",
|
||||
] {
|
||||
assert!(
|
||||
!candidate_affinity.contains(forbidden),
|
||||
"planner/candidate_affinity.rs should not own local candidate eligibility helper {forbidden}"
|
||||
!candidate_ranking.contains(forbidden),
|
||||
"planner/candidate_ranking.rs should not own local candidate resolution or transport ordering helper {forbidden}"
|
||||
);
|
||||
}
|
||||
|
||||
let candidate_affinity_cache = read_workspace_file(
|
||||
"apps/aether-gateway/src/ai_pipeline/planner/candidate_affinity_cache.rs",
|
||||
);
|
||||
for pattern in [
|
||||
"pub(crate) fn read_cached_scheduler_affinity_target(",
|
||||
"pub(crate) fn remember_scheduler_affinity_for_candidate(",
|
||||
] {
|
||||
assert!(
|
||||
candidate_affinity_cache.contains(pattern),
|
||||
"planner/candidate_affinity_cache.rs should own {pattern}"
|
||||
);
|
||||
}
|
||||
for forbidden in [
|
||||
"apply_scheduler_candidate_ranking",
|
||||
"rank_eligible_local_execution_candidates",
|
||||
"filter_and_rank_local_execution_candidates",
|
||||
] {
|
||||
assert!(
|
||||
!candidate_affinity_cache.contains(forbidden),
|
||||
"planner/candidate_affinity_cache.rs should not own ranking or resolution helper {forbidden}"
|
||||
);
|
||||
}
|
||||
|
||||
let candidate_transport_ordering = read_workspace_file(
|
||||
"apps/aether-gateway/src/ai_pipeline/planner/candidate_transport_ordering.rs",
|
||||
);
|
||||
for pattern in [
|
||||
"pub(super) struct CandidateExecutionOrdering {",
|
||||
"resolve_cached_candidate_execution_ordering",
|
||||
"resolve_cached_transport_execution_ordering",
|
||||
"resolve_transport_proxy_snapshot_with_tunnel_affinity",
|
||||
] {
|
||||
assert!(
|
||||
candidate_transport_ordering.contains(pattern),
|
||||
"planner/candidate_transport_ordering.rs should own {pattern}"
|
||||
);
|
||||
}
|
||||
for forbidden in [
|
||||
"apply_scheduler_candidate_ranking",
|
||||
"rank_eligible_local_execution_candidates",
|
||||
"filter_and_rank_local_execution_candidates",
|
||||
] {
|
||||
assert!(
|
||||
!candidate_transport_ordering.contains(forbidden),
|
||||
"planner/candidate_transport_ordering.rs should not own ranking or resolution helper {forbidden}"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -565,17 +565,17 @@ fn scheduler_candidate_runtime_paths_depend_on_scheduler_core_and_state_trait()
|
||||
);
|
||||
}
|
||||
|
||||
let planner_candidate_affinity =
|
||||
read_workspace_file("apps/aether-gateway/src/ai_pipeline/planner/candidate_affinity.rs");
|
||||
let planner_candidate_ranking =
|
||||
read_workspace_file("apps/aether-gateway/src/ai_pipeline/planner/candidate_ranking.rs");
|
||||
assert!(
|
||||
planner_candidate_affinity.contains("use aether_scheduler_core::{")
|
||||
&& planner_candidate_affinity.contains("SchedulerMinimalCandidateSelectionCandidate"),
|
||||
"planner/candidate_affinity.rs should depend directly on core minimal candidate DTO"
|
||||
planner_candidate_ranking.contains("use aether_scheduler_core::{")
|
||||
&& planner_candidate_ranking.contains("SchedulerMinimalCandidateSelectionCandidate"),
|
||||
"planner/candidate_ranking.rs should depend directly on core minimal candidate DTO"
|
||||
);
|
||||
assert!(
|
||||
!planner_candidate_affinity
|
||||
!planner_candidate_ranking
|
||||
.contains("crate::scheduler::SchedulerMinimalCandidateSelectionCandidate"),
|
||||
"planner/candidate_affinity.rs should not depend on scheduler candidate DTO re-export"
|
||||
"planner/candidate_ranking.rs should not depend on scheduler candidate DTO re-export"
|
||||
);
|
||||
|
||||
let request_candidate_runtime =
|
||||
|
||||
Reference in New Issue
Block a user