Persist ranking metadata through request traces

This commit is contained in:
fawney19
2026-04-27 14:05:19 +08:00
parent f52220a8ec
commit 91db4eefd0
19 changed files with 215 additions and 32 deletions

View File

@@ -4,6 +4,7 @@ use serde_json::Value;
use uuid::Uuid;
use crate::ai_pipeline::planner::candidate_affinity_cache::remember_scheduler_affinity_for_candidate;
use crate::ai_pipeline::planner::candidate_metadata::append_ranking_metadata_to_object;
use crate::ai_pipeline::planner::candidate_resolution::{
EligibleLocalExecutionCandidate, SkippedLocalExecutionCandidate,
};
@@ -177,34 +178,7 @@ fn local_candidate_extra_data_with_ranking(
}
None => serde_json::Map::new(),
};
object.insert(
"ranking_mode".to_string(),
Value::String(format!("{:?}", ranking.ranking_mode)),
);
object.insert(
"priority_mode".to_string(),
Value::String(format!("{:?}", ranking.priority_mode)),
);
object.insert(
"ranking_index".to_string(),
Value::Number(serde_json::Number::from(ranking.ranking_index as u64)),
);
object.insert(
"priority_slot".to_string(),
Value::Number(serde_json::Number::from(i64::from(ranking.priority_slot))),
);
if let Some(promoted_by) = ranking.promoted_by {
object.insert(
"promoted_by".to_string(),
Value::String(promoted_by.to_string()),
);
}
if let Some(demoted_by) = ranking.demoted_by {
object.insert(
"demoted_by".to_string(),
Value::String(demoted_by.to_string()),
);
}
append_ranking_metadata_to_object(&mut object, ranking);
Some(Value::Object(object))
}
@@ -372,7 +346,10 @@ pub(crate) async fn persist_skipped_local_execution_candidates(
&generated_candidate_id,
required_capabilities,
skipped_candidate.skip_reason,
skipped_candidate.extra_data,
local_candidate_extra_data_with_ranking(
skipped_candidate.extra_data,
skipped_candidate.ranking.as_ref(),
),
error_context,
record_runtime_miss_diagnostic,
)
@@ -641,13 +618,23 @@ mod tests {
"pool-skipped",
Some(json!({ "pool_advanced": {} })),
)),
ranking: None,
extra_data: None,
},
SkippedLocalExecutionCandidate {
candidate: sample_candidate("normal-skipped"),
skip_reason: "key_inactive",
transport: None,
extra_data: None,
ranking: Some(SchedulerRankingOutcome {
original_index: 2,
ranking_index: 1,
priority_mode: SchedulerPriorityMode::Provider,
ranking_mode: SchedulerRankingMode::CacheAffinity,
priority_slot: 9,
promoted_by: None,
demoted_by: Some("cross_format"),
}),
extra_data: Some(json!({ "existing": "value" })),
},
],
"persist skipped should not fail",
@@ -662,5 +649,19 @@ mod tests {
assert_eq!(stored.len(), 1);
assert_eq!(stored[0].key_id.as_deref(), Some("normal-skipped"));
assert_eq!(stored[0].candidate_index, 0);
let extra_data = stored[0]
.extra_data
.as_ref()
.and_then(serde_json::Value::as_object)
.expect("skipped ranking metadata should persist");
assert_eq!(extra_data.get("existing"), Some(&json!("value")));
assert_eq!(
extra_data.get("ranking_mode"),
Some(&json!("CacheAffinity"))
);
assert_eq!(extra_data.get("priority_mode"), Some(&json!("Provider")));
assert_eq!(extra_data.get("ranking_index"), Some(&json!(1)));
assert_eq!(extra_data.get("priority_slot"), Some(&json!(9)));
assert_eq!(extra_data.get("demoted_by"), Some(&json!("cross_format")));
}
}

View File

@@ -1,5 +1,5 @@
use aether_contracts::ProxySnapshot;
use aether_scheduler_core::SchedulerMinimalCandidateSelectionCandidate;
use aether_scheduler_core::{SchedulerMinimalCandidateSelectionCandidate, SchedulerRankingOutcome};
use serde_json::{json, Map, Value};
use crate::ai_pipeline::planner::candidate_resolution::EligibleLocalExecutionCandidate;
@@ -26,6 +26,40 @@ pub(crate) struct LocalExecutionCandidateMetadataParts<'a> {
pub(crate) extra_fields: Map<String, Value>,
}
pub(crate) fn append_ranking_metadata_to_object(
object: &mut Map<String, Value>,
ranking: &SchedulerRankingOutcome,
) {
object.insert(
"ranking_mode".to_string(),
Value::String(format!("{:?}", ranking.ranking_mode)),
);
object.insert(
"priority_mode".to_string(),
Value::String(format!("{:?}", ranking.priority_mode)),
);
object.insert(
"ranking_index".to_string(),
Value::Number(serde_json::Number::from(ranking.ranking_index as u64)),
);
object.insert(
"priority_slot".to_string(),
Value::Number(serde_json::Number::from(i64::from(ranking.priority_slot))),
);
if let Some(promoted_by) = ranking.promoted_by {
object.insert(
"promoted_by".to_string(),
Value::String(promoted_by.to_string()),
);
}
if let Some(demoted_by) = ranking.demoted_by {
object.insert(
"demoted_by".to_string(),
Value::String(demoted_by.to_string()),
);
}
}
pub(crate) fn build_request_trace_proxy_value(
transport: Option<&GatewayProviderTransportSnapshot>,
resolved_proxy: Option<&ProxySnapshot>,

View File

@@ -27,6 +27,7 @@ pub(crate) struct SkippedLocalExecutionCandidate {
pub(crate) candidate: SchedulerMinimalCandidateSelectionCandidate,
pub(crate) skip_reason: &'static str,
pub(crate) transport: Option<Arc<GatewayProviderTransportSnapshot>>,
pub(crate) ranking: Option<SchedulerRankingOutcome>,
pub(crate) extra_data: Option<serde_json::Value>,
}
@@ -131,6 +132,7 @@ where
candidate,
skip_reason: "transport_snapshot_missing",
transport: None,
ranking: None,
extra_data: None,
});
continue;
@@ -145,6 +147,7 @@ where
candidate,
skip_reason,
transport: Some(transport),
ranking: None,
extra_data: None,
}),
None => selectable.push(EligibleLocalExecutionCandidate {

View File

@@ -123,6 +123,7 @@ pub(crate) async fn materialize_local_same_format_provider_candidate_attempts(
candidate: item.candidate,
skip_reason: item.skip_reason,
transport: None,
ranking: None,
extra_data: None,
})
.chain(skipped_candidates)

View File

@@ -102,6 +102,7 @@ pub(crate) async fn maybe_build_local_same_format_provider_decision_payload_for_
client_api_format: spec_metadata.api_format,
mapped_model: Some(&resolved.mapped_model),
candidate_group_id: eligible.orchestration.candidate_group_id.as_deref(),
ranking: eligible.ranking.as_ref(),
upstream_url: Some(&resolved.upstream_url),
header_rules: resolved.transport.endpoint.header_rules.as_ref(),
body_rules: resolved.transport.endpoint.body_rules.as_ref(),

View File

@@ -449,6 +449,7 @@ fn schedule_pool_group(
candidate,
skip_reason: POOL_ACCOUNT_BLOCKED_SKIP_REASON,
transport: Some(transport),
ranking,
extra_data: None,
});
continue;
@@ -459,6 +460,7 @@ fn schedule_pool_group(
candidate,
skip_reason: POOL_ACCOUNT_EXHAUSTED_SKIP_REASON,
transport: Some(transport),
ranking,
extra_data: None,
});
continue;
@@ -469,6 +471,7 @@ fn schedule_pool_group(
candidate,
skip_reason: POOL_COOLDOWN_SKIP_REASON,
transport: Some(transport),
ranking,
extra_data: None,
});
continue;
@@ -482,6 +485,7 @@ fn schedule_pool_group(
candidate,
skip_reason: POOL_COST_LIMIT_REACHED_SKIP_REASON,
transport: Some(transport),
ranking,
extra_data: None,
});
continue;

View File

@@ -1,8 +1,10 @@
use std::collections::BTreeMap;
use aether_scheduler_core::SchedulerRankingOutcome;
use serde_json::{Map, Value};
use crate::ai_pipeline::contracts::ExecutionRuntimeAuthContext;
use crate::ai_pipeline::planner::candidate_metadata::append_ranking_metadata_to_object;
use crate::orchestration::ExecutionAttemptIdentity;
pub(crate) struct LocalExecutionReportContextParts<'a> {
@@ -23,6 +25,7 @@ pub(crate) struct LocalExecutionReportContextParts<'a> {
pub(crate) client_api_format: &'a str,
pub(crate) mapped_model: Option<&'a str>,
pub(crate) candidate_group_id: Option<&'a str>,
pub(crate) ranking: Option<&'a SchedulerRankingOutcome>,
pub(crate) upstream_url: Option<&'a str>,
pub(crate) header_rules: Option<&'a Value>,
pub(crate) body_rules: Option<&'a Value>,
@@ -172,6 +175,9 @@ pub(crate) fn build_local_execution_report_context(
Value::String(candidate_group_id.to_string()),
);
}
if let Some(ranking) = parts.ranking {
append_ranking_metadata_to_object(&mut object, ranking);
}
if let Some(upstream_url) = parts.upstream_url {
object.insert(
"upstream_url".to_string(),

View File

@@ -83,6 +83,7 @@ pub(super) async fn maybe_build_local_gemini_files_decision_payload_for_candidat
client_api_format: GEMINI_FILES_CLIENT_API_FORMAT,
mapped_model: None,
candidate_group_id: eligible.orchestration.candidate_group_id.as_deref(),
ranking: eligible.ranking.as_ref(),
upstream_url: None,
header_rules: transport.endpoint.header_rules.as_ref(),
body_rules: transport.endpoint.body_rules.as_ref(),

View File

@@ -80,6 +80,7 @@ pub(super) async fn maybe_build_local_openai_image_decision_payload_for_candidat
client_api_format: spec_metadata.api_format,
mapped_model: Some(&resolved.mapped_model),
candidate_group_id: eligible.orchestration.candidate_group_id.as_deref(),
ranking: eligible.ranking.as_ref(),
upstream_url: Some(&resolved.upstream_url),
header_rules: transport.endpoint.header_rules.as_ref(),
body_rules: transport.endpoint.body_rules.as_ref(),

View File

@@ -128,6 +128,7 @@ pub(super) async fn list_local_openai_image_candidate_attempts(
candidate: item.candidate,
skip_reason: item.skip_reason,
transport: None,
ranking: None,
extra_data: None,
})
.collect(),

View File

@@ -66,6 +66,7 @@ pub(super) async fn maybe_build_local_video_create_decision_payload_for_candidat
client_api_format: spec_metadata.api_format,
mapped_model: Some(&resolved.mapped_model),
candidate_group_id: eligible.orchestration.candidate_group_id.as_deref(),
ranking: eligible.ranking.as_ref(),
upstream_url: None,
header_rules: transport.endpoint.header_rules.as_ref(),
body_rules: transport.endpoint.body_rules.as_ref(),

View File

@@ -140,6 +140,7 @@ pub(super) async fn list_local_video_create_candidate_attempts(
candidate: item.candidate,
skip_reason: item.skip_reason,
transport: None,
ranking: None,
extra_data: None,
})
.collect(),

View File

@@ -154,6 +154,7 @@ pub(super) async fn materialize_local_standard_candidate_attempts(
candidate: skipped_candidate.candidate,
skip_reason: skipped_candidate.skip_reason,
transport: None,
ranking: None,
extra_data: None,
});
}

View File

@@ -82,6 +82,7 @@ pub(super) async fn maybe_build_local_standard_decision_payload_for_candidate(
client_api_format: spec_metadata.api_format,
mapped_model: Some(&resolved.mapped_model),
candidate_group_id: eligible.orchestration.candidate_group_id.as_deref(),
ranking: eligible.ranking.as_ref(),
upstream_url: Some(&resolved.upstream_url),
header_rules: resolved.transport.endpoint.header_rules.as_ref(),
body_rules: resolved.transport.endpoint.body_rules.as_ref(),

View File

@@ -100,6 +100,7 @@ pub(crate) async fn maybe_build_local_openai_chat_decision_payload_for_candidate
client_api_format: "openai:chat",
mapped_model: Some(&resolved.mapped_model),
candidate_group_id: eligible.orchestration.candidate_group_id.as_deref(),
ranking: eligible.ranking.as_ref(),
upstream_url: Some(&resolved.upstream_url),
header_rules: resolved.transport.endpoint.header_rules.as_ref(),
body_rules: resolved.transport.endpoint.body_rules.as_ref(),

View File

@@ -78,6 +78,7 @@ pub(crate) async fn list_local_openai_chat_candidates(
candidate: skipped_candidate.candidate,
skip_reason: skipped_candidate.skip_reason,
transport: None,
ranking: None,
extra_data: None,
});
}

View File

@@ -98,6 +98,7 @@ pub(crate) async fn maybe_build_local_openai_responses_decision_payload_for_cand
client_api_format: spec_metadata.api_format,
mapped_model: Some(&resolved.mapped_model),
candidate_group_id: eligible.orchestration.candidate_group_id.as_deref(),
ranking: eligible.ranking.as_ref(),
upstream_url: Some(&resolved.upstream_url),
header_rules: resolved.transport.endpoint.header_rules.as_ref(),
body_rules: resolved.transport.endpoint.body_rules.as_ref(),

View File

@@ -210,6 +210,7 @@ pub(crate) async fn materialize_local_openai_responses_candidate_attempts(
candidate: skipped_candidate.candidate,
skip_reason: skipped_candidate.skip_reason,
transport: None,
ranking: None,
extra_data: None,
});
}

View File

@@ -24,6 +24,12 @@ pub struct SchedulerRequestCandidateReportContext {
pub body_rules: Option<Value>,
pub proxy: Option<Value>,
pub error_flow: Option<Value>,
pub ranking_mode: Option<String>,
pub priority_mode: Option<String>,
pub ranking_index: Option<u32>,
pub priority_slot: Option<i32>,
pub promoted_by: Option<String>,
pub demoted_by: Option<String>,
}
#[derive(Debug, Clone, PartialEq)]
@@ -59,6 +65,12 @@ struct ReportCandidateExtraDataInput {
body_rules: Option<Value>,
proxy: Option<Value>,
error_flow: Option<Value>,
ranking_mode: Option<String>,
priority_mode: Option<String>,
ranking_index: Option<u32>,
priority_slot: Option<i32>,
promoted_by: Option<String>,
demoted_by: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
@@ -141,6 +153,12 @@ pub fn parse_request_candidate_report_context(
.get("error_flow")
.cloned()
.filter(|value| !value.is_null()),
ranking_mode: string_field(report_context, "ranking_mode"),
priority_mode: string_field(report_context, "priority_mode"),
ranking_index: u32_field(report_context, "ranking_index"),
priority_slot: i32_field(report_context, "priority_slot"),
promoted_by: string_field(report_context, "promoted_by"),
demoted_by: string_field(report_context, "demoted_by"),
})
}
@@ -170,6 +188,12 @@ pub fn resolve_report_request_candidate_slot(
body_rules,
proxy,
error_flow,
ranking_mode,
priority_mode,
ranking_index,
priority_slot,
promoted_by,
demoted_by,
} = metadata;
let request_id = request_id?;
let synthesized_extra_data = build_report_candidate_extra_data(ReportCandidateExtraDataInput {
@@ -182,6 +206,12 @@ pub fn resolve_report_request_candidate_slot(
body_rules,
proxy,
error_flow,
ranking_mode,
priority_mode,
ranking_index,
priority_slot,
promoted_by,
demoted_by,
});
let created_at_unix_ms = matched_candidate
.as_ref()
@@ -346,6 +376,12 @@ pub fn build_local_request_candidate_status_record(
body_rules: metadata.body_rules.clone(),
proxy: metadata.proxy.clone(),
error_flow: metadata.error_flow.clone(),
ranking_mode: metadata.ranking_mode.clone(),
priority_mode: metadata.priority_mode.clone(),
ranking_index: metadata.ranking_index,
priority_slot: metadata.priority_slot,
promoted_by: metadata.promoted_by.clone(),
demoted_by: metadata.demoted_by.clone(),
});
let created_at_unix_ms = started_at_unix_ms.or(finished_at_unix_ms);
@@ -500,6 +536,14 @@ fn u32_field_from_object(object: &Map<String, Value>, key: &str) -> Option<u32>
.and_then(|value| u32::try_from(value).ok())
}
fn i32_field(value: &Value, key: &str) -> Option<i32> {
value
.as_object()
.and_then(|object| object.get(key))
.and_then(Value::as_i64)
.and_then(|value| i32::try_from(value).ok())
}
fn match_existing_report_candidate<'a>(
candidates: &'a [StoredRequestCandidate],
metadata: &SchedulerRequestCandidateReportContext,
@@ -558,6 +602,12 @@ fn build_report_candidate_extra_data(input: ReportCandidateExtraDataInput) -> Op
body_rules,
proxy,
error_flow,
ranking_mode,
priority_mode,
ranking_index,
priority_slot,
promoted_by,
demoted_by,
} = input;
let mut extra_data = Map::with_capacity(8);
extra_data.insert("gateway_execution_runtime".to_string(), Value::Bool(true));
@@ -595,6 +645,30 @@ fn build_report_candidate_extra_data(input: ReportCandidateExtraDataInput) -> Op
if let Some(error_flow) = error_flow {
extra_data.insert("error_flow".to_string(), error_flow);
}
if let Some(ranking_mode) = ranking_mode {
extra_data.insert("ranking_mode".to_string(), Value::String(ranking_mode));
}
if let Some(priority_mode) = priority_mode {
extra_data.insert("priority_mode".to_string(), Value::String(priority_mode));
}
if let Some(ranking_index) = ranking_index {
extra_data.insert(
"ranking_index".to_string(),
Value::Number(ranking_index.into()),
);
}
if let Some(priority_slot) = priority_slot {
extra_data.insert(
"priority_slot".to_string(),
Value::Number(priority_slot.into()),
);
}
if let Some(promoted_by) = promoted_by {
extra_data.insert("promoted_by".to_string(), Value::String(promoted_by));
}
if let Some(demoted_by) = demoted_by {
extra_data.insert("demoted_by".to_string(), Value::String(demoted_by));
}
(!extra_data.is_empty()).then_some(Value::Object(extra_data))
}
@@ -881,7 +955,13 @@ mod tests {
"provider_api_format": "openai:responses",
"upstream_url": "https://example.com/v1/responses",
"mapped_model": "gpt-5-upstream",
"key_name": "primary"
"key_name": "primary",
"ranking_mode": "CacheAffinity",
"priority_mode": "Provider",
"ranking_index": 2,
"priority_slot": 7,
"promoted_by": "cached_affinity",
"demoted_by": "cross_format"
})),
status_update: SchedulerRequestCandidateStatusUpdate {
status: RequestCandidateStatus::Failed,
@@ -915,6 +995,48 @@ mod tests {
.and_then(|value| value.get("mapped_model")),
Some(&json!("gpt-5-upstream"))
);
assert_eq!(
record
.extra_data
.as_ref()
.and_then(|value| value.get("ranking_mode")),
Some(&json!("CacheAffinity"))
);
assert_eq!(
record
.extra_data
.as_ref()
.and_then(|value| value.get("priority_mode")),
Some(&json!("Provider"))
);
assert_eq!(
record
.extra_data
.as_ref()
.and_then(|value| value.get("ranking_index")),
Some(&json!(2))
);
assert_eq!(
record
.extra_data
.as_ref()
.and_then(|value| value.get("priority_slot")),
Some(&json!(7))
);
assert_eq!(
record
.extra_data
.as_ref()
.and_then(|value| value.get("promoted_by")),
Some(&json!("cached_affinity"))
);
assert_eq!(
record
.extra_data
.as_ref()
.and_then(|value| value.get("demoted_by")),
Some(&json!("cross_format"))
);
}
#[test]