feat(codex): 按操作语义路由 Responses V2 压缩

(cherry picked from commit 2fc604e047)
This commit is contained in:
MMEXA
2026-07-16 00:34:42 +08:00
committed by elky
parent 598b2fb374
commit 20b27a13b2
75 changed files with 732 additions and 50 deletions
@@ -706,6 +706,7 @@ pub(crate) async fn build_lazy_requested_model_execution_candidate_attempt_sourc
trace_id: &str,
client_api_format: &str,
requested_model: &str,
request_operation: Option<&str>,
require_streaming: bool,
auth_snapshot: &GatewayAuthApiKeySnapshot,
client_session_affinity: Option<&ClientSessionAffinity>,
@@ -734,6 +735,7 @@ where
model_directive_policy,
client_api_format,
requested_model,
request_operation,
require_streaming,
required_capabilities,
auth_snapshot,
@@ -1141,6 +1143,7 @@ async fn resolve_priority_candidate_page_with_cache(
let key = CandidateResolvedPageCacheKey::new(
&cursor.requested_model,
cursor.page_cursor.resolved_page_cache_request_operation(),
&cursor.client_api_format,
true,
&cursor.auth_snapshot,
@@ -2190,6 +2193,7 @@ mod tests {
&model_directive_policy,
"openai:chat",
"gpt-5",
None,
true,
None,
&auth_snapshot,
@@ -2245,6 +2249,7 @@ mod tests {
&model_directive_policy,
"openai:chat",
"gpt-5",
None,
true,
None,
&auth_snapshot,
@@ -2282,6 +2287,7 @@ mod tests {
&model_directive_policy,
"openai:chat",
"gpt-5",
None,
true,
None,
&auth_snapshot,
@@ -6,9 +6,10 @@ use aether_routing_core::ResolvedRoutingPolicy;
use aether_runtime::ConcurrencyPermit;
use aether_scheduler_core::{
enumerate_minimal_candidate_selection_with_model_directives, normalize_api_format,
resolve_requested_global_model_name_with_model_directives,
row_supports_requested_model_with_model_directives, ClientSessionAffinity,
EnumerateMinimalCandidateSelectionInput, SchedulerMinimalCandidateSelectionCandidate,
resolve_requested_global_model_name_with_model_directives_and_request_operation,
row_supports_requested_model_with_model_directives_and_request_operation,
ClientSessionAffinity, EnumerateMinimalCandidateSelectionInput,
SchedulerMinimalCandidateSelectionCandidate,
};
use async_trait::async_trait;
use std::collections::{BTreeMap, BTreeSet, VecDeque};
@@ -56,6 +57,7 @@ struct GatewayLocalCandidatePreselectionPort<'a> {
state: PlannerAppState<'a>,
client_api_format: &'a str,
requested_model: &'a str,
request_operation: Option<&'a str>,
require_streaming: bool,
required_capabilities: Option<&'a serde_json::Value>,
auth_snapshot: &'a GatewayAuthApiKeySnapshot,
@@ -112,7 +114,7 @@ impl AiCandidatePreselectionPort for GatewayLocalCandidatePreselectionPort<'_> {
let auth_snapshot = matches_client_format.then_some(self.auth_snapshot);
let (candidates, skipped_candidates) = self
.state
.list_selectable_candidates_with_skip_reasons(
.list_selectable_candidates_with_skip_reasons_for_request_operation(
candidate_api_format,
self.routing_model(candidate_api_format),
self.require_streaming,
@@ -121,6 +123,7 @@ impl AiCandidatePreselectionPort for GatewayLocalCandidatePreselectionPort<'_> {
self.client_session_affinity,
self.ranking_seed,
false,
self.request_operation,
)
.await?;
@@ -197,6 +200,7 @@ pub(crate) async fn preselect_local_execution_candidates_with_serving(
model_directive_policy: &crate::system_features::ModelDirectivePolicySnapshot,
client_api_format: &str,
requested_model: &str,
request_operation: Option<&str>,
require_streaming: bool,
required_capabilities: Option<&serde_json::Value>,
auth_snapshot: &GatewayAuthApiKeySnapshot,
@@ -221,6 +225,7 @@ pub(crate) async fn preselect_local_execution_candidates_with_serving(
model_directive_policy,
client_api_format,
requested_model,
request_operation,
require_streaming,
required_capabilities,
auth_snapshot,
@@ -239,6 +244,7 @@ pub(crate) async fn preselect_local_execution_candidates_for_api_formats_with_se
model_directive_policy: &crate::system_features::ModelDirectivePolicySnapshot,
client_api_format: &str,
requested_model: &str,
request_operation: Option<&str>,
require_streaming: bool,
required_capabilities: Option<&serde_json::Value>,
auth_snapshot: &GatewayAuthApiKeySnapshot,
@@ -263,6 +269,7 @@ pub(crate) async fn preselect_local_execution_candidates_for_api_formats_with_se
state,
client_api_format,
requested_model,
request_operation,
require_streaming,
required_capabilities,
auth_snapshot,
@@ -283,6 +290,7 @@ pub(crate) struct LocalCandidatePreselectionPageCursor<'a> {
trace_id: String,
client_api_format: String,
requested_model: String,
request_operation: Option<String>,
require_streaming: bool,
required_capabilities: Option<serde_json::Value>,
auth_snapshot: GatewayAuthApiKeySnapshot,
@@ -336,6 +344,7 @@ impl<'a> LocalCandidatePreselectionPageCursor<'a> {
model_directive_policy: &crate::system_features::ModelDirectivePolicySnapshot,
client_api_format: &str,
requested_model: &str,
request_operation: Option<&str>,
require_streaming: bool,
required_capabilities: Option<&serde_json::Value>,
auth_snapshot: &GatewayAuthApiKeySnapshot,
@@ -370,6 +379,7 @@ impl<'a> LocalCandidatePreselectionPageCursor<'a> {
trace_id: trace_id.unwrap_or_default().to_string(),
client_api_format: client_api_format.to_string(),
requested_model: requested_model.to_string(),
request_operation: request_operation.map(str::to_string),
require_streaming,
required_capabilities: required_capabilities.cloned(),
auth_snapshot: auth_snapshot.clone(),
@@ -452,6 +462,10 @@ impl<'a> LocalCandidatePreselectionPageCursor<'a> {
self.key_mode.cache_key_name()
}
pub(crate) fn resolved_page_cache_request_operation(&self) -> Option<&str> {
self.request_operation.as_deref()
}
pub(crate) fn resolved_page_cache_use_api_format_alias_match(&self) -> bool {
self.use_api_format_alias_match
}
@@ -515,6 +529,7 @@ impl<'a> LocalCandidatePreselectionPageCursor<'a> {
> {
let key = CandidatePageCacheKey::new(
&self.requested_model,
self.request_operation.as_deref(),
&self.client_api_format,
self.require_streaming,
&self.auth_snapshot,
@@ -965,11 +980,12 @@ impl<'a> LocalCandidatePreselectionPageCursor<'a> {
.map_err(|err| GatewayError::Internal(err.to_string()))?
.into_iter()
.filter(|row| {
row_supports_requested_model_with_model_directives(
row_supports_requested_model_with_model_directives_and_request_operation(
row,
&routing_model,
normalized_api_format,
false,
self.request_operation.as_deref(),
)
})
.collect::<Vec<_>>();
@@ -1009,12 +1025,15 @@ impl<'a> LocalCandidatePreselectionPageCursor<'a> {
if let Some(value) = self.resolved_global_model_names.get(normalized_api_format) {
value.clone()
} else {
let Some(value) = resolve_requested_global_model_name_with_model_directives(
&rows,
&routing_model,
normalized_api_format,
false,
) else {
let Some(value) =
resolve_requested_global_model_name_with_model_directives_and_request_operation(
&rows,
&routing_model,
normalized_api_format,
false,
self.request_operation.as_deref(),
)
else {
return Ok(None);
};
self.resolved_global_model_names
@@ -1037,6 +1056,7 @@ impl<'a> LocalCandidatePreselectionPageCursor<'a> {
EnumerateMinimalCandidateSelectionInput {
rows,
normalized_api_format,
request_operation: self.request_operation.as_deref(),
requested_model_name: &routing_model,
resolved_global_model_name: resolved_global_model_name.as_str(),
require_streaming: self.require_streaming,
@@ -1316,6 +1336,7 @@ mod tests {
&model_directive_policy,
"openai:chat",
"gpt-5",
None,
true,
None,
&auth_snapshot,
@@ -1531,6 +1552,7 @@ mod tests {
priority: 1,
api_formats: None,
endpoint_ids: Some(vec!["endpoint-opg-openai".to_string()]),
operations: None,
}]),
model_supports_streaming: Some(true),
model_is_active: true,
@@ -1561,6 +1583,7 @@ mod tests {
&model_directive_policy,
"claude:messages",
"gpt-5.5-xhigh",
None,
false,
None,
&auth_snapshot,
@@ -1590,6 +1613,70 @@ mod tests {
);
}
#[tokio::test]
async fn paged_preselection_prefers_operation_scoped_mapping_for_compaction() {
let mut row = openai_responses_mapping_row();
row.global_model_mappings = None;
row.global_model_name = "gpt-5.6-sol".to_string();
row.model_provider_model_name = "gpt-5.6-sol".to_string();
row.model_provider_model_mappings = Some(vec![
StoredProviderModelMapping {
name: "gpt-5.6-sol".to_string(),
priority: 1,
api_formats: Some(vec!["openai:responses".to_string()]),
endpoint_ids: None,
operations: None,
},
StoredProviderModelMapping {
name: "gpt-5.6-terra".to_string(),
priority: 1,
api_formats: Some(vec!["openai:responses".to_string()]),
endpoint_ids: None,
operations: Some(vec!["compact".to_string()]),
},
]);
let repository: Arc<dyn MinimalCandidateSelectionReadRepository> =
Arc::new(InMemoryMinimalCandidateSelectionReadRepository::seed([row]));
let data_state =
GatewayDataState::with_minimal_candidate_selection_reader_for_tests(repository);
let app = AppState::new()
.expect("gateway state should build")
.with_data_state_for_tests(data_state);
let auth_snapshot = unrestricted_auth_snapshot();
let model_directive_policy =
crate::system_features::ModelDirectivePolicySnapshot::load(&app).await;
let mut cursor = LocalCandidatePreselectionPageCursor::new(
PlannerAppState::new(&app),
&model_directive_policy,
"openai:responses",
"gpt-5.6-sol",
Some("compact"),
false,
None,
&auth_snapshot,
None,
None,
None,
true,
LocalCandidatePreselectionKeyMode::ProviderEndpointKeyModelAndApiFormat,
true,
None,
)
.await;
let page = cursor
.next_page()
.await
.expect("preselection should succeed")
.expect("compact mapping should find a provider");
assert_eq!(page.candidates.len(), 1);
assert_eq!(
page.candidates[0].selected_provider_model_name,
"gpt-5.6-terra"
);
}
#[tokio::test]
async fn custom_policy_suffix_uses_the_same_base_model_for_candidate_selection() {
let mut row = openai_responses_mapping_row();
@@ -1634,6 +1721,7 @@ mod tests {
&model_directive_policy,
"openai:responses",
"deployment-alias-VendorFuture",
None,
false,
None,
&auth_snapshot,
@@ -1695,6 +1783,7 @@ mod tests {
&model_directive_policy,
"claude:messages",
"deepseek-v4-pro",
None,
false,
None,
&auth_snapshot,
@@ -1773,6 +1862,7 @@ mod tests {
&model_directive_policy,
"claude:messages",
"gpt-5",
None,
false,
None,
&auth_snapshot,
@@ -124,6 +124,7 @@ pub(super) async fn materialize_local_standard_candidate_attempts(
&input.model_directive_policy,
spec_metadata.api_format,
&input.requested_model,
None,
false,
input.required_capabilities.as_ref(),
&input.auth_snapshot,
@@ -250,6 +251,7 @@ pub(super) async fn build_local_standard_candidate_attempt_source<'a>(
trace_id,
spec_metadata.api_format,
&input.requested_model,
None,
spec_metadata.require_streaming,
&input.auth_snapshot,
input.client_session_affinity.as_ref(),
@@ -345,6 +347,7 @@ async fn maybe_append_gemini_image_openai_image_preselection(
&input.model_directive_policy,
spec_metadata.api_format,
&input.requested_model,
None,
spec_metadata.require_streaming,
input.required_capabilities.as_ref(),
&input.auth_snapshot,
@@ -297,6 +297,7 @@ pub(crate) async fn build_lazy_local_openai_chat_candidate_attempt_source<'a>(
trace_id,
"openai:chat",
&input.requested_model,
None,
require_streaming,
&input.auth_snapshot,
input.client_session_affinity.as_ref(),
@@ -24,6 +24,7 @@ pub(crate) async fn list_local_openai_chat_candidates(
&input.model_directive_policy,
"openai:chat",
&input.requested_model,
None,
require_streaming,
input.required_capabilities.as_ref(),
&input.auth_snapshot,
@@ -31,8 +31,8 @@ use crate::ai_serving::planner::spec_metadata::local_openai_responses_spec_metad
use crate::ai_serving::planner::CandidateFailureDiagnostic;
use crate::ai_serving::{
ai_local_execution_contract_for_formats, extract_pool_sticky_session_token,
resolve_local_decision_execution_runtime_auth_context, ExecutionRuntimeAuthContext,
GatewayControlDecision, PlannerAppState,
openai_responses_request_operation, resolve_local_decision_execution_runtime_auth_context,
ExecutionRuntimeAuthContext, GatewayControlDecision, PlannerAppState,
};
use crate::client_session_affinity::client_session_affinity_from_parts;
use crate::{AppState, GatewayError};
@@ -163,6 +163,7 @@ pub(crate) async fn materialize_local_openai_responses_candidate_attempts(
spec: LocalOpenAiResponsesSpec,
) -> Result<(Vec<LocalOpenAiResponsesCandidateAttempt>, usize), GatewayError> {
let spec_metadata = local_openai_responses_spec_metadata(spec);
let request_operation = openai_responses_request_operation(spec_metadata.api_format, body_json);
let planner_state = PlannerAppState::new(state);
let sticky_session_token = extract_pool_sticky_session_token(body_json);
let auth_context: &ExecutionRuntimeAuthContext = &input.auth_context;
@@ -176,6 +177,7 @@ pub(crate) async fn materialize_local_openai_responses_candidate_attempts(
&input.model_directive_policy,
spec_metadata.api_format,
&input.requested_model,
request_operation,
spec_metadata.require_streaming,
input.required_capabilities.as_ref(),
&input.auth_snapshot,
@@ -262,6 +264,7 @@ pub(crate) async fn build_local_openai_responses_candidate_attempt_source<'a>(
spec: LocalOpenAiResponsesSpec,
) -> Result<(LocalOpenAiResponsesCandidateAttemptSource<'a>, usize), GatewayError> {
let spec_metadata = local_openai_responses_spec_metadata(spec);
let request_operation = openai_responses_request_operation(spec_metadata.api_format, body_json);
let planner_state = PlannerAppState::new(state);
let sticky_session_token = extract_pool_sticky_session_token(body_json);
let auth_context: &ExecutionRuntimeAuthContext = &input.auth_context;
@@ -287,6 +290,7 @@ pub(crate) async fn build_local_openai_responses_candidate_attempt_source<'a>(
trace_id,
spec_metadata.api_format,
&input.requested_model,
request_operation,
spec_metadata.require_streaming,
&input.auth_snapshot,
input.client_session_affinity.as_ref(),
@@ -372,6 +376,7 @@ pub(crate) async fn build_local_openai_responses_image_candidate_attempt_source<
&input.model_directive_policy,
spec_metadata.api_format,
&input.requested_model,
None,
false,
input.required_capabilities.as_ref(),
&input.auth_snapshot,
@@ -53,13 +53,45 @@ impl<'a> PlannerAppState<'a> {
Vec<SchedulerSkippedCandidate>,
),
GatewayError,
> {
self.list_selectable_candidates_with_skip_reasons_for_request_operation(
api_format,
global_model_name,
require_streaming,
required_capabilities,
auth_snapshot,
client_session_affinity,
now_unix_secs,
enable_model_directives,
None,
)
.await
}
pub(crate) async fn list_selectable_candidates_with_skip_reasons_for_request_operation(
self,
api_format: &str,
global_model_name: &str,
require_streaming: bool,
required_capabilities: Option<&serde_json::Value>,
auth_snapshot: Option<&GatewayAuthApiKeySnapshot>,
client_session_affinity: Option<&ClientSessionAffinity>,
now_unix_secs: u64,
enable_model_directives: bool,
request_operation: Option<&str>,
) -> Result<
(
Vec<SchedulerMinimalCandidateSelectionCandidate>,
Vec<SchedulerSkippedCandidate>,
),
GatewayError,
> {
let wait_timeout = Duration::from_millis(API_KEY_CONCURRENCY_WAIT_TIMEOUT_MS);
let wait_interval = Duration::from_millis(API_KEY_CONCURRENCY_WAIT_POLL_INTERVAL_MS.max(1));
let wait_deadline = Instant::now() + wait_timeout;
let mut attempt_now_unix_secs = now_unix_secs;
loop {
let result = crate::scheduler::candidate::list_selectable_candidates_with_skip_reasons(
let result = crate::scheduler::candidate::list_selectable_candidates_with_skip_reasons_for_request_operation(
self.app().data.as_ref(),
self.app(),
api_format,
@@ -70,6 +102,7 @@ impl<'a> PlannerAppState<'a> {
client_session_affinity,
attempt_now_unix_secs,
enable_model_directives,
request_operation,
)
.await?;
@@ -165,5 +165,5 @@ pub(crate) use aether_ai_formats::api::{
pub(crate) use aether_ai_formats::{
api_format_defaults_to_client_error_failover, api_format_defaults_to_non_stream,
api_format_permission_covers, intersect_api_format_allowed_lists, is_embedding_api_format,
is_rerank_api_format,
is_rerank_api_format, openai_responses_request_operation,
};