mirror of
https://github.com/fawney19/Aether.git
synced 2026-10-08 10:27:46 +08:00
Merge upstream/main into main
This commit is contained in:
@@ -41,7 +41,8 @@ use super::{
|
||||
LocalSameFormatProviderSpec,
|
||||
};
|
||||
use crate::ai_serving::planner::standard::{
|
||||
codex_model_capabilities_for_transport, same_format_provider_request_body_failure_extra_data,
|
||||
codex_model_capabilities_for_transport, openai_provider_request_contract_failure_extra_data,
|
||||
same_format_provider_request_body_failure_extra_data,
|
||||
};
|
||||
|
||||
pub(crate) fn resolve_same_format_provider_transport_unsupported_reason_for_trace(
|
||||
@@ -267,33 +268,39 @@ pub(crate) async fn resolve_local_same_format_provider_candidate_payload_parts(
|
||||
prepared.mapped_model.as_str(),
|
||||
source_model,
|
||||
);
|
||||
if crate::ai_serving::finalize_openai_provider_request_with_codex_model_capabilities(
|
||||
&mut base_provider_request_body,
|
||||
crate::ai_serving::OpenAiProviderRequestFinalization {
|
||||
source_api_format: spec.api_format,
|
||||
provider_api_format: prepared.provider_api_format.as_str(),
|
||||
provider_type: transport.provider.provider_type.as_str(),
|
||||
provider_model: prepared.mapped_model.as_str(),
|
||||
source_model,
|
||||
body_rules: transport.endpoint.body_rules.as_ref(),
|
||||
upstream_is_stream: prepared.upstream_is_stream,
|
||||
require_body_stream_field: request_requires_body_stream_field(
|
||||
body_json,
|
||||
prepared.force_body_stream_field,
|
||||
),
|
||||
},
|
||||
codex_model_capabilities.as_ref(),
|
||||
)
|
||||
.is_err()
|
||||
if let Err(violation) =
|
||||
crate::ai_serving::finalize_openai_provider_request_with_codex_model_capabilities(
|
||||
&mut base_provider_request_body,
|
||||
crate::ai_serving::OpenAiProviderRequestFinalization {
|
||||
source_api_format: spec.api_format,
|
||||
provider_api_format: prepared.provider_api_format.as_str(),
|
||||
provider_type: transport.provider.provider_type.as_str(),
|
||||
provider_model: prepared.mapped_model.as_str(),
|
||||
source_model,
|
||||
body_rules: transport.endpoint.body_rules.as_ref(),
|
||||
upstream_is_stream: prepared.upstream_is_stream,
|
||||
require_body_stream_field: request_requires_body_stream_field(
|
||||
body_json,
|
||||
prepared.force_body_stream_field,
|
||||
),
|
||||
},
|
||||
codex_model_capabilities.as_ref(),
|
||||
)
|
||||
{
|
||||
mark_skipped_local_same_format_provider_candidate(
|
||||
mark_skipped_local_same_format_provider_candidate_with_extra_data(
|
||||
state,
|
||||
input,
|
||||
trace_id,
|
||||
candidate,
|
||||
attempt.candidate_index,
|
||||
&attempt.candidate_id,
|
||||
"provider_request_body_missing",
|
||||
"provider_request_body_build_failed",
|
||||
Some(openai_provider_request_contract_failure_extra_data(
|
||||
&violation,
|
||||
spec.api_format,
|
||||
prepared.provider_api_format.as_str(),
|
||||
"same_format_provider_request_finalization",
|
||||
)),
|
||||
)
|
||||
.await;
|
||||
return Ok(None);
|
||||
|
||||
@@ -23,7 +23,8 @@ use crate::ai_serving::planner::spec_metadata::local_standard_spec_metadata;
|
||||
use crate::ai_serving::planner::standard::{
|
||||
apply_codex_openai_special_headers, apply_deepseek_tool_call_thinking_compat,
|
||||
codex_model_capabilities_for_transport, is_deepseek_provider,
|
||||
request_body_build_failure_extra_data, request_conversion_failure_extra_data,
|
||||
openai_provider_request_contract_failure_extra_data, request_body_build_failure_extra_data,
|
||||
request_conversion_failure_extra_data,
|
||||
};
|
||||
use crate::ai_serving::transport::kiro::{
|
||||
build_kiro_provider_headers, build_kiro_provider_request_body,
|
||||
@@ -748,24 +749,24 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
|
||||
prepared_candidate.mapped_model.as_str(),
|
||||
source_model,
|
||||
);
|
||||
if crate::ai_serving::finalize_openai_provider_request_with_codex_model_capabilities(
|
||||
&mut provider_request_body,
|
||||
crate::ai_serving::OpenAiProviderRequestFinalization {
|
||||
source_api_format: spec_metadata.api_format,
|
||||
provider_api_format,
|
||||
provider_type: transport.provider.provider_type.as_str(),
|
||||
provider_model: prepared_candidate.mapped_model.as_str(),
|
||||
source_model,
|
||||
body_rules: transport.endpoint.body_rules.as_ref(),
|
||||
upstream_is_stream,
|
||||
require_body_stream_field: request_requires_body_stream_field(
|
||||
body_json,
|
||||
force_body_stream_field,
|
||||
),
|
||||
},
|
||||
codex_model_capabilities.as_ref(),
|
||||
)
|
||||
.is_err()
|
||||
if let Err(violation) =
|
||||
crate::ai_serving::finalize_openai_provider_request_with_codex_model_capabilities(
|
||||
&mut provider_request_body,
|
||||
crate::ai_serving::OpenAiProviderRequestFinalization {
|
||||
source_api_format: spec_metadata.api_format,
|
||||
provider_api_format,
|
||||
provider_type: transport.provider.provider_type.as_str(),
|
||||
provider_model: prepared_candidate.mapped_model.as_str(),
|
||||
source_model,
|
||||
body_rules: transport.endpoint.body_rules.as_ref(),
|
||||
upstream_is_stream,
|
||||
require_body_stream_field: request_requires_body_stream_field(
|
||||
body_json,
|
||||
force_body_stream_field,
|
||||
),
|
||||
},
|
||||
codex_model_capabilities.as_ref(),
|
||||
)
|
||||
{
|
||||
mark_skipped_local_standard_candidate_with_extra_data(
|
||||
state,
|
||||
@@ -775,15 +776,12 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
|
||||
attempt.candidate_index,
|
||||
&attempt.candidate_id,
|
||||
"provider_request_body_build_failed",
|
||||
request_conversion_failure_extra_data(
|
||||
body_json,
|
||||
Some(openai_provider_request_contract_failure_extra_data(
|
||||
&violation,
|
||||
spec_metadata.api_format,
|
||||
provider_api_format,
|
||||
Some(prepared_candidate.mapped_model.as_str()),
|
||||
Some(parts.uri.path()),
|
||||
upstream_is_stream,
|
||||
"standard_family_request_finalization",
|
||||
),
|
||||
)),
|
||||
)
|
||||
.await;
|
||||
return Ok(None);
|
||||
|
||||
@@ -65,8 +65,8 @@ pub(crate) use crate::ai_serving::{
|
||||
normalize_openai_responses_request_to_openai_chat_request, parse_openai_tool_result_content,
|
||||
};
|
||||
pub(crate) use aether_ai_serving::{
|
||||
request_body_build_failure_extra_data, request_conversion_failure_extra_data,
|
||||
same_format_provider_request_body_failure_extra_data,
|
||||
openai_provider_request_contract_failure_extra_data, request_body_build_failure_extra_data,
|
||||
request_conversion_failure_extra_data, same_format_provider_request_body_failure_extra_data,
|
||||
};
|
||||
|
||||
pub(crate) fn build_standard_upstream_url(
|
||||
|
||||
@@ -81,6 +81,10 @@ pub(crate) fn build_local_openai_responses_request_body_with_codex_model_capabil
|
||||
&mut provider_request_body,
|
||||
provider_api_format,
|
||||
);
|
||||
crate::ai_serving::strip_incompatible_openai_responses_reasoning_items(
|
||||
&mut provider_request_body,
|
||||
provider_api_format,
|
||||
);
|
||||
enforce_provider_body_stream_policy(
|
||||
&mut provider_request_body,
|
||||
provider_api_format,
|
||||
@@ -172,6 +176,10 @@ pub(crate) fn build_cross_format_openai_responses_request_body_with_codex_model_
|
||||
&mut provider_request_body,
|
||||
provider_api_format,
|
||||
);
|
||||
crate::ai_serving::strip_incompatible_openai_responses_reasoning_items(
|
||||
&mut provider_request_body,
|
||||
provider_api_format,
|
||||
);
|
||||
enforce_provider_body_stream_policy(
|
||||
&mut provider_request_body,
|
||||
provider_api_format,
|
||||
|
||||
@@ -152,6 +152,43 @@ fn local_openai_responses_wrapper_preserves_body_order_after_edits() {
|
||||
assert!(provider_request_body.get("instructions").is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn local_openai_responses_wrapper_strips_foreign_reasoning_item_ids() {
|
||||
let body_json = json!({
|
||||
"model": "gpt-5.4",
|
||||
"input": [
|
||||
{"type": "reasoning", "id": "rs_provider_123", "summary": []},
|
||||
{
|
||||
"type": "reasoning",
|
||||
"id": "item_72d3bd8d367d01977ace23f1",
|
||||
"summary": []
|
||||
},
|
||||
{"type": "message", "role": "user", "content": "continue"}
|
||||
]
|
||||
});
|
||||
|
||||
let provider_request_body = build_local_openai_responses_request_body(
|
||||
&body_json,
|
||||
"gpt-5.4",
|
||||
false,
|
||||
false,
|
||||
"codex",
|
||||
"openai:responses",
|
||||
None,
|
||||
None,
|
||||
&http::HeaderMap::new(),
|
||||
false,
|
||||
)
|
||||
.expect("local OpenAI Responses body should build");
|
||||
|
||||
let input = provider_request_body["input"]
|
||||
.as_array()
|
||||
.expect("input array");
|
||||
assert_eq!(input.len(), 2);
|
||||
assert_eq!(input[0]["id"], "rs_provider_123");
|
||||
assert_eq!(input[1]["type"], "message");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn local_openai_responses_compact_wrapper_strips_store_for_same_format_requests() {
|
||||
let body_json = json!({
|
||||
|
||||
+18
-25
@@ -29,7 +29,8 @@ use crate::ai_serving::planner::standard::{
|
||||
apply_deepseek_tool_call_thinking_compat, build_cross_format_openai_chat_request_body,
|
||||
build_cross_format_openai_chat_upstream_url, build_local_openai_chat_request_body,
|
||||
build_local_openai_chat_upstream_url, codex_model_capabilities_for_transport,
|
||||
request_body_build_failure_extra_data, request_conversion_failure_extra_data,
|
||||
openai_provider_request_contract_failure_extra_data, request_body_build_failure_extra_data,
|
||||
request_conversion_failure_extra_data,
|
||||
};
|
||||
use crate::ai_serving::transport::antigravity::is_antigravity_provider_transport;
|
||||
use crate::ai_serving::transport::auth::resolve_local_openai_bearer_auth;
|
||||
@@ -109,7 +110,7 @@ fn finalize_openai_chat_provider_request_body(
|
||||
original_body: &Value,
|
||||
transport: &GatewayProviderTransportSnapshot,
|
||||
mapped_model: &str,
|
||||
) -> bool {
|
||||
) -> Option<Value> {
|
||||
if let Some(mapping) = custom_directive_mapping {
|
||||
crate::ai_serving::apply_model_directive_mapping_patch(provider_request_body, mapping);
|
||||
}
|
||||
@@ -156,7 +157,15 @@ fn finalize_openai_chat_provider_request_body(
|
||||
},
|
||||
codex_model_capabilities.as_ref(),
|
||||
)
|
||||
.is_ok()
|
||||
.err()
|
||||
.map(|violation| {
|
||||
openai_provider_request_contract_failure_extra_data(
|
||||
&violation,
|
||||
"openai:chat",
|
||||
provider_api_format,
|
||||
"openai_chat_request_finalization",
|
||||
)
|
||||
})
|
||||
}
|
||||
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
@@ -288,7 +297,7 @@ pub(crate) async fn resolve_local_openai_chat_candidate_payload_parts(
|
||||
.await;
|
||||
return Ok(None);
|
||||
};
|
||||
if !finalize_openai_chat_provider_request_body(
|
||||
if let Some(extra_data) = finalize_openai_chat_provider_request_body(
|
||||
&mut provider_request_body,
|
||||
model_directive_mapping.as_ref(),
|
||||
provider_api_format,
|
||||
@@ -306,11 +315,7 @@ pub(crate) async fn resolve_local_openai_chat_candidate_payload_parts(
|
||||
candidate_index,
|
||||
candidate_id,
|
||||
"provider_request_body_build_failed",
|
||||
request_body_build_failure_extra_data(
|
||||
body_json,
|
||||
"openai:chat",
|
||||
provider_api_format,
|
||||
),
|
||||
Some(extra_data),
|
||||
)
|
||||
.await;
|
||||
return Ok(None);
|
||||
@@ -501,7 +506,7 @@ pub(crate) async fn resolve_local_openai_chat_candidate_payload_parts(
|
||||
"openai_chat_payload_body_build",
|
||||
body_build_started_at.elapsed().as_millis() as u64,
|
||||
);
|
||||
if !finalize_openai_chat_provider_request_body(
|
||||
if let Some(extra_data) = finalize_openai_chat_provider_request_body(
|
||||
&mut provider_request_body,
|
||||
model_directive_mapping.as_ref(),
|
||||
"openai:chat",
|
||||
@@ -519,11 +524,7 @@ pub(crate) async fn resolve_local_openai_chat_candidate_payload_parts(
|
||||
candidate_index,
|
||||
candidate_id,
|
||||
"provider_request_body_build_failed",
|
||||
request_body_build_failure_extra_data(
|
||||
body_json,
|
||||
"openai:chat",
|
||||
provider_api_format,
|
||||
),
|
||||
Some(extra_data),
|
||||
)
|
||||
.await;
|
||||
return Ok(None);
|
||||
@@ -813,7 +814,7 @@ pub(crate) async fn resolve_local_openai_chat_candidate_payload_parts(
|
||||
.await;
|
||||
return Ok(None);
|
||||
};
|
||||
if !finalize_openai_chat_provider_request_body(
|
||||
if let Some(extra_data) = finalize_openai_chat_provider_request_body(
|
||||
&mut provider_request_body,
|
||||
model_directive_mapping.as_ref(),
|
||||
provider_api_format.as_str(),
|
||||
@@ -831,15 +832,7 @@ pub(crate) async fn resolve_local_openai_chat_candidate_payload_parts(
|
||||
candidate_index,
|
||||
candidate_id,
|
||||
"provider_request_body_build_failed",
|
||||
request_conversion_failure_extra_data(
|
||||
body_json,
|
||||
"openai:chat",
|
||||
provider_api_format.as_str(),
|
||||
Some(prepared_candidate.mapped_model.as_str()),
|
||||
Some(parts.uri.path()),
|
||||
upstream_is_stream,
|
||||
"openai_chat_request_conversion",
|
||||
),
|
||||
Some(extra_data),
|
||||
)
|
||||
.await;
|
||||
return Ok(None);
|
||||
|
||||
+24
-26
@@ -32,7 +32,8 @@ use crate::ai_serving::planner::standard::{
|
||||
build_cross_format_openai_responses_upstream_url,
|
||||
build_local_openai_responses_request_body_with_codex_model_capabilities,
|
||||
build_local_openai_responses_upstream_url, codex_model_capabilities_for_transport,
|
||||
request_body_build_failure_extra_data, request_conversion_failure_extra_data,
|
||||
openai_provider_request_contract_failure_extra_data, request_body_build_failure_extra_data,
|
||||
request_conversion_failure_extra_data,
|
||||
};
|
||||
use crate::ai_serving::transport::antigravity::is_antigravity_provider_transport;
|
||||
use crate::ai_serving::transport::auth::{
|
||||
@@ -530,24 +531,24 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
|
||||
provider_api_format,
|
||||
Some(body_json),
|
||||
);
|
||||
if crate::ai_serving::finalize_openai_provider_request_with_codex_model_capabilities(
|
||||
&mut base_provider_request_body,
|
||||
crate::ai_serving::OpenAiProviderRequestFinalization {
|
||||
source_api_format: spec_metadata.api_format,
|
||||
provider_api_format,
|
||||
provider_type: transport.provider.provider_type.as_str(),
|
||||
provider_model: mapped_model.as_str(),
|
||||
source_model,
|
||||
body_rules: transport.endpoint.body_rules.as_ref(),
|
||||
upstream_is_stream,
|
||||
require_body_stream_field: request_requires_body_stream_field(
|
||||
body_json,
|
||||
force_body_stream_field,
|
||||
),
|
||||
},
|
||||
codex_model_capabilities.as_ref(),
|
||||
)
|
||||
.is_err()
|
||||
if let Err(violation) =
|
||||
crate::ai_serving::finalize_openai_provider_request_with_codex_model_capabilities(
|
||||
&mut base_provider_request_body,
|
||||
crate::ai_serving::OpenAiProviderRequestFinalization {
|
||||
source_api_format: spec_metadata.api_format,
|
||||
provider_api_format,
|
||||
provider_type: transport.provider.provider_type.as_str(),
|
||||
provider_model: mapped_model.as_str(),
|
||||
source_model,
|
||||
body_rules: transport.endpoint.body_rules.as_ref(),
|
||||
upstream_is_stream,
|
||||
require_body_stream_field: request_requires_body_stream_field(
|
||||
body_json,
|
||||
force_body_stream_field,
|
||||
),
|
||||
},
|
||||
codex_model_capabilities.as_ref(),
|
||||
)
|
||||
{
|
||||
mark_skipped_local_openai_responses_candidate_with_extra_data(
|
||||
state,
|
||||
@@ -557,15 +558,12 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
|
||||
candidate_index,
|
||||
candidate_id,
|
||||
"provider_request_body_build_failed",
|
||||
request_conversion_failure_extra_data(
|
||||
body_json,
|
||||
Some(openai_provider_request_contract_failure_extra_data(
|
||||
&violation,
|
||||
spec_metadata.api_format,
|
||||
provider_api_format,
|
||||
Some(mapped_model.as_str()),
|
||||
Some(parts.uri.path()),
|
||||
upstream_is_stream,
|
||||
"openai_responses_request_conversion",
|
||||
),
|
||||
"openai_responses_request_finalization",
|
||||
)),
|
||||
)
|
||||
.await;
|
||||
return Ok(None);
|
||||
|
||||
@@ -172,7 +172,9 @@ 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, openai_responses_request_operation, ApiOperation, ClientSurface,
|
||||
is_rerank_api_format, openai_responses_request_operation,
|
||||
openai_responses_synthetic_reasoning_item_id,
|
||||
strip_incompatible_openai_responses_reasoning_items, ApiOperation, ClientSurface,
|
||||
};
|
||||
|
||||
pub(crate) fn plan_kind_matches_api_operation(
|
||||
|
||||
@@ -1,6 +1,8 @@
|
||||
use http::Uri;
|
||||
|
||||
use super::{classify_control_route, headers};
|
||||
use crate::control::GatewayPublicRequestContext;
|
||||
use crate::handlers::shared::local_proxy_route_requires_buffered_body;
|
||||
|
||||
#[test]
|
||||
fn classifies_admin_pool_overview_as_admin_proxy_route() {
|
||||
@@ -92,6 +94,16 @@ fn classifies_admin_pool_provider_key_routes_as_admin_proxy_route() {
|
||||
batch_update.route_kind.as_deref(),
|
||||
Some("batch_update_keys")
|
||||
);
|
||||
let batch_update_context = GatewayPublicRequestContext::from_request_parts(
|
||||
"trace-admin-pool-batch-update",
|
||||
&http::Method::PATCH,
|
||||
&batch_update_uri,
|
||||
&headers,
|
||||
Some(batch_update),
|
||||
);
|
||||
assert!(local_proxy_route_requires_buffered_body(
|
||||
&batch_update_context
|
||||
));
|
||||
|
||||
let resolve_selection_uri: Uri = "/api/admin/pool/provider-1/keys/resolve-selection"
|
||||
.parse()
|
||||
|
||||
@@ -26,6 +26,7 @@ use crate::ai_serving::api::{
|
||||
CanonicalContentPart, CanonicalStreamEvent, CanonicalStreamFrame, ClaudeClientEmitter,
|
||||
OpenAIChatClientEmitter, OpenAIResponsesClientEmitter, StreamingCanonicalUsage,
|
||||
};
|
||||
use crate::ai_serving::openai_responses_synthetic_reasoning_item_id;
|
||||
use crate::clock::current_unix_secs;
|
||||
use crate::execution_runtime::ndjson::encode_stream_frame_ndjson;
|
||||
use crate::execution_runtime::transport::{
|
||||
@@ -2705,7 +2706,7 @@ fn openai_responses_body(
|
||||
let mut output = Vec::new();
|
||||
if !collected.thinking.trim().is_empty() {
|
||||
output.push(json!({
|
||||
"id": format!("{response_id}_rs_0"),
|
||||
"id": openai_responses_synthetic_reasoning_item_id(&response_id, 0),
|
||||
"type": "reasoning",
|
||||
"status": "completed",
|
||||
"summary": [{
|
||||
|
||||
@@ -8,6 +8,9 @@ use aether_data_contracts::repository::candidates::{
|
||||
use aether_data_contracts::repository::provider_catalog::{
|
||||
StoredProviderCatalogEndpoint, StoredProviderCatalogKey, StoredProviderCatalogProvider,
|
||||
};
|
||||
use aether_data_contracts::repository::usage::{
|
||||
ROUTING_CANDIDATE_SKIP_REASON_METADATA_KEY, ROUTING_FAILURE_DIAGNOSTIC_METADATA_KEY,
|
||||
};
|
||||
use aether_usage_runtime::{
|
||||
build_usage_event_data_seed, UsageEvent, UsageEventData, UsageEventType,
|
||||
};
|
||||
@@ -195,6 +198,7 @@ pub(crate) async fn build_local_execution_exhaustion(
|
||||
data.provider_api_key_id = data
|
||||
.provider_api_key_id
|
||||
.or_else(|| candidate.key_id.clone());
|
||||
attach_runtime_miss_candidate_usage_metadata(&mut data, candidate);
|
||||
}
|
||||
|
||||
exhaustion.data = data;
|
||||
@@ -385,23 +389,25 @@ pub(crate) async fn record_failed_usage_for_runtime_miss_request(
|
||||
|
||||
let selected_candidate =
|
||||
select_last_runtime_miss_executed_candidate(&context.candidate_contexts);
|
||||
let api_format = selected_candidate
|
||||
let routing_candidate = selected_candidate
|
||||
.or_else(|| select_last_runtime_miss_routing_candidate(&context.candidate_contexts));
|
||||
let api_format = routing_candidate
|
||||
.and_then(|value| value.client_api_format.clone())
|
||||
.or_else(|| {
|
||||
trimmed_non_empty(decision.and_then(|value| value.auth_endpoint_signature.as_deref()))
|
||||
});
|
||||
let provider_api_format = selected_candidate
|
||||
let provider_api_format = routing_candidate
|
||||
.and_then(|value| value.provider_api_format.clone())
|
||||
.or_else(|| api_format.clone());
|
||||
let provider_name = selected_candidate
|
||||
let provider_name = routing_candidate
|
||||
.and_then(|value| value.provider_name.clone())
|
||||
.or_else(|| selected_candidate.and_then(|value| value.candidate.provider_id.clone()))
|
||||
.or_else(|| routing_candidate.and_then(|value| value.candidate.provider_id.clone()))
|
||||
.unwrap_or_else(|| "unknown".to_string());
|
||||
let model = trimmed_non_empty(diagnostic.and_then(|value| value.requested_model.as_deref()))
|
||||
.or_else(|| selected_candidate.and_then(|value| value.global_model_name.clone()))
|
||||
.or_else(|| selected_candidate.and_then(|value| value.selected_provider_model_name.clone()))
|
||||
.or_else(|| routing_candidate.and_then(|value| value.global_model_name.clone()))
|
||||
.or_else(|| routing_candidate.and_then(|value| value.selected_provider_model_name.clone()))
|
||||
.unwrap_or_else(|| "unknown".to_string());
|
||||
let target_model = selected_candidate
|
||||
let target_model = routing_candidate
|
||||
.and_then(|value| value.selected_provider_model_name.clone())
|
||||
.filter(|value| !value.eq_ignore_ascii_case(model.as_str()));
|
||||
|
||||
@@ -436,10 +442,10 @@ pub(crate) async fn record_failed_usage_for_runtime_miss_request(
|
||||
provider_name,
|
||||
model,
|
||||
target_model,
|
||||
provider_id: selected_candidate.and_then(|value| value.candidate.provider_id.clone()),
|
||||
provider_endpoint_id: selected_candidate
|
||||
provider_id: routing_candidate.and_then(|value| value.candidate.provider_id.clone()),
|
||||
provider_endpoint_id: routing_candidate
|
||||
.and_then(|value| value.candidate.endpoint_id.clone()),
|
||||
provider_api_key_id: selected_candidate.and_then(|value| value.candidate.key_id.clone()),
|
||||
provider_api_key_id: routing_candidate.and_then(|value| value.candidate.key_id.clone()),
|
||||
request_type: Some(infer_request_type(api_format.as_deref())),
|
||||
api_format: api_format.clone(),
|
||||
api_family: api_format
|
||||
@@ -459,7 +465,7 @@ pub(crate) async fn record_failed_usage_for_runtime_miss_request(
|
||||
.as_deref()
|
||||
.and_then(infer_endpoint_kind)
|
||||
.map(ToOwned::to_owned),
|
||||
has_format_conversion: selected_candidate.and_then(|value| {
|
||||
has_format_conversion: routing_candidate.and_then(|value| {
|
||||
value
|
||||
.client_api_format
|
||||
.as_deref()
|
||||
@@ -478,13 +484,16 @@ pub(crate) async fn record_failed_usage_for_runtime_miss_request(
|
||||
client_response_body: Some(client_body),
|
||||
..UsageEventData::default()
|
||||
};
|
||||
if let Some(candidate) = routing_candidate {
|
||||
insert_runtime_miss_candidate_usage_metadata(&mut request_metadata, &candidate.candidate);
|
||||
}
|
||||
apply_runtime_miss_usage_routing(
|
||||
&mut data,
|
||||
&mut request_metadata,
|
||||
execution_path,
|
||||
selected_candidate.map(|value| value.candidate.id.as_str()),
|
||||
selected_candidate.map(|value| value.candidate.candidate_index),
|
||||
selected_candidate.and_then(|value| value.key_name.as_deref()),
|
||||
routing_candidate.map(|value| value.candidate.id.as_str()),
|
||||
routing_candidate.map(|value| value.candidate.candidate_index),
|
||||
routing_candidate.and_then(|value| value.key_name.as_deref()),
|
||||
diagnostic,
|
||||
decision.and_then(|value| value.route_family.as_deref()),
|
||||
decision.and_then(|value| value.route_kind.as_deref()),
|
||||
@@ -637,6 +646,22 @@ fn select_last_runtime_miss_executed_candidate(
|
||||
})
|
||||
}
|
||||
|
||||
fn select_last_runtime_miss_routing_candidate(
|
||||
candidates: &[RuntimeMissCandidateContext],
|
||||
) -> Option<&RuntimeMissCandidateContext> {
|
||||
candidates.iter().max_by_key(|candidate| {
|
||||
(
|
||||
candidate.candidate.retry_index,
|
||||
candidate.candidate.candidate_index,
|
||||
candidate
|
||||
.candidate
|
||||
.finished_at_unix_ms
|
||||
.or(candidate.candidate.started_at_unix_ms)
|
||||
.unwrap_or(candidate.candidate.created_at_unix_ms),
|
||||
)
|
||||
})
|
||||
}
|
||||
|
||||
fn request_candidate_represents_provider_execution(candidate: &StoredRequestCandidate) -> bool {
|
||||
matches!(
|
||||
candidate.status,
|
||||
@@ -935,6 +960,66 @@ fn candidate_extra_data_string(candidate: &StoredRequestCandidate, key: &str) ->
|
||||
.map(ToOwned::to_owned)
|
||||
}
|
||||
|
||||
fn attach_runtime_miss_candidate_usage_metadata(
|
||||
data: &mut UsageEventData,
|
||||
candidate: &StoredRequestCandidate,
|
||||
) {
|
||||
let mut metadata = match data.request_metadata.take() {
|
||||
Some(Value::Object(object)) => object,
|
||||
Some(other) => Map::from_iter([("seed".to_string(), other)]),
|
||||
None => Map::new(),
|
||||
};
|
||||
insert_runtime_miss_candidate_usage_metadata(&mut metadata, candidate);
|
||||
data.request_metadata = (!metadata.is_empty()).then_some(Value::Object(metadata));
|
||||
}
|
||||
|
||||
fn insert_runtime_miss_candidate_usage_metadata(
|
||||
metadata: &mut Map<String, Value>,
|
||||
candidate: &StoredRequestCandidate,
|
||||
) {
|
||||
if let Some(skip_reason) = candidate
|
||||
.skip_reason
|
||||
.as_deref()
|
||||
.map(str::trim)
|
||||
.filter(|value| !value.is_empty())
|
||||
{
|
||||
metadata.insert(
|
||||
ROUTING_CANDIDATE_SKIP_REASON_METADATA_KEY.to_string(),
|
||||
Value::String(skip_reason.to_string()),
|
||||
);
|
||||
}
|
||||
|
||||
let diagnostic = candidate
|
||||
.extra_data
|
||||
.as_ref()
|
||||
.and_then(Value::as_object)
|
||||
.and_then(|extra_data| {
|
||||
extra_data
|
||||
.get("failure_diagnostic")
|
||||
.filter(|value| {
|
||||
value.as_object().is_some_and(|diagnostic| {
|
||||
diagnostic.get("safe_to_show") != Some(&Value::Bool(false))
|
||||
})
|
||||
})
|
||||
.or_else(|| {
|
||||
extra_data
|
||||
.get("request_conversion_error")
|
||||
.filter(|v| v.is_object())
|
||||
})
|
||||
.or_else(|| {
|
||||
extra_data
|
||||
.get("request_body_build_error")
|
||||
.filter(|v| v.is_object())
|
||||
})
|
||||
});
|
||||
if let Some(diagnostic) = diagnostic {
|
||||
metadata.insert(
|
||||
ROUTING_FAILURE_DIAGNOSTIC_METADATA_KEY.to_string(),
|
||||
diagnostic.clone(),
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
fn runtime_miss_candidate_failure_diagnostic(
|
||||
candidate: &RuntimeMissCandidateContext,
|
||||
) -> Option<RuntimeMissFailureDiagnostic> {
|
||||
@@ -1170,9 +1255,10 @@ fn trimmed_non_empty(value: Option<&str>) -> Option<String> {
|
||||
mod tests {
|
||||
use super::{
|
||||
apply_runtime_miss_usage_routing, beautify_local_execution_client_error_message,
|
||||
insert_runtime_miss_candidate_usage_metadata,
|
||||
request_candidate_represents_provider_execution, runtime_miss_client_error_body,
|
||||
select_last_runtime_miss_executed_candidate, LocalExecutionRuntimeMissContext,
|
||||
RuntimeMissCandidateContext,
|
||||
select_last_runtime_miss_executed_candidate, select_last_runtime_miss_routing_candidate,
|
||||
LocalExecutionRuntimeMissContext, RuntimeMissCandidateContext,
|
||||
};
|
||||
use crate::constants::EXECUTION_PATH_LOCAL_EXECUTION_RUNTIME_MISS;
|
||||
use crate::state::LocalExecutionRuntimeMissDiagnostic;
|
||||
@@ -1263,7 +1349,7 @@ mod tests {
|
||||
|
||||
#[test]
|
||||
fn runtime_miss_executed_candidate_selection_ignores_skipped_only_histories() {
|
||||
let skipped_candidate = StoredRequestCandidate::new(
|
||||
let mut skipped_candidate = StoredRequestCandidate::new(
|
||||
"cand-skipped".to_string(),
|
||||
"req-1".to_string(),
|
||||
Some("user-1".to_string()),
|
||||
@@ -1290,6 +1376,26 @@ mod tests {
|
||||
None,
|
||||
)
|
||||
.expect("candidate should build");
|
||||
skipped_candidate.skip_reason = Some("provider_request_body_build_failed".to_string());
|
||||
skipped_candidate.extra_data = Some(json!({
|
||||
"failure_diagnostic": {
|
||||
"kind": "request_body_build",
|
||||
"path": "$.reasoning.summary",
|
||||
"message": "invalid reasoning summary",
|
||||
"safe_to_show": true
|
||||
}
|
||||
}));
|
||||
|
||||
let mut request_metadata = Map::new();
|
||||
insert_runtime_miss_candidate_usage_metadata(&mut request_metadata, &skipped_candidate);
|
||||
assert_eq!(
|
||||
request_metadata["routing_candidate_skip_reason"],
|
||||
"provider_request_body_build_failed"
|
||||
);
|
||||
assert_eq!(
|
||||
request_metadata["routing_failure_diagnostic"]["path"],
|
||||
"$.reasoning.summary"
|
||||
);
|
||||
|
||||
assert!(!request_candidate_represents_provider_execution(
|
||||
&skipped_candidate
|
||||
@@ -1307,6 +1413,11 @@ mod tests {
|
||||
}];
|
||||
|
||||
assert!(select_last_runtime_miss_executed_candidate(&contexts).is_none());
|
||||
assert_eq!(
|
||||
select_last_runtime_miss_routing_candidate(&contexts)
|
||||
.map(|candidate| candidate.candidate.id.as_str()),
|
||||
Some("cand-skipped")
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
||||
@@ -236,6 +236,15 @@ async fn admin_monitoring_trace_request_falls_back_to_usage_routing_snapshot() {
|
||||
usage.provider_api_key_id = Some("provider-key-1".to_string());
|
||||
usage.error_message = Some("no local stream plans".to_string());
|
||||
usage.response_time_ms = Some(45);
|
||||
usage.request_metadata = Some(json!({
|
||||
"routing_candidate_skip_reason": "provider_request_body_build_failed",
|
||||
"routing_failure_diagnostic": {
|
||||
"kind": "request_body_build",
|
||||
"path": "$.reasoning.summary",
|
||||
"message": "上游请求体语义校验失败",
|
||||
"safe_to_show": true
|
||||
}
|
||||
}));
|
||||
|
||||
let usage_repository = Arc::new(InMemoryUsageReadRepository::seed(vec![usage]));
|
||||
let data_state =
|
||||
@@ -267,6 +276,10 @@ async fn admin_monitoring_trace_request_falls_back_to_usage_routing_snapshot() {
|
||||
assert_eq!(payload["final_status"], json!("failed"));
|
||||
assert_eq!(payload["candidates"][0]["id"], json!("routing-cand-1"));
|
||||
assert_eq!(payload["candidates"][0]["status"], json!("failed"));
|
||||
assert_eq!(
|
||||
payload["candidates"][0]["skip_reason"],
|
||||
json!("provider_request_body_build_failed")
|
||||
);
|
||||
assert_eq!(
|
||||
payload["candidates"][0]["error_type"],
|
||||
json!("no_local_stream_plans")
|
||||
@@ -279,6 +292,10 @@ async fn admin_monitoring_trace_request_falls_back_to_usage_routing_snapshot() {
|
||||
payload["candidates"][0]["extra_data"]["execution_path"],
|
||||
json!("local_execution_runtime_miss")
|
||||
);
|
||||
assert_eq!(
|
||||
payload["candidates"][0]["extra_data"]["failure_diagnostic"]["path"],
|
||||
json!("$.reasoning.summary")
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
|
||||
@@ -203,7 +203,7 @@ fn build_admin_monitoring_usage_routing_snapshot_trace(
|
||||
endpoint_id: usage.provider_endpoint_id.clone(),
|
||||
key_id: usage.provider_api_key_id.clone(),
|
||||
status,
|
||||
skip_reason: None,
|
||||
skip_reason: usage.routing_candidate_skip_reason().map(ToOwned::to_owned),
|
||||
is_cached: false,
|
||||
status_code: usage.status_code,
|
||||
error_type: usage
|
||||
@@ -356,6 +356,9 @@ fn build_admin_monitoring_usage_routing_snapshot_extra_data(
|
||||
if let Some(candidate_index) = usage.routing_candidate_index() {
|
||||
object.insert("candidate_index".to_string(), json!(candidate_index));
|
||||
}
|
||||
if let Some(diagnostic) = usage.routing_failure_diagnostic() {
|
||||
object.insert("failure_diagnostic".to_string(), diagnostic.clone());
|
||||
}
|
||||
Some(Value::Object(object))
|
||||
}
|
||||
|
||||
|
||||
@@ -377,6 +377,7 @@ pub(crate) fn admin_proxy_local_requires_buffered_body(
|
||||
| (Some("users_manage"), http::Method::PATCH, Some("lock_user_api_key"))
|
||||
| (Some("pool_manage"), http::Method::POST, Some("batch_import_keys"))
|
||||
| (Some("pool_manage"), http::Method::POST, Some("batch_action_keys"))
|
||||
| (Some("pool_manage"), http::Method::PATCH, Some("batch_update_keys"))
|
||||
| (Some("pool_manage"), http::Method::POST, Some("resolve_selection"))
|
||||
| (Some("usage_manage"), http::Method::POST, Some("replay"))
|
||||
| (Some("wallets_manage"), http::Method::POST, Some("adjust_balance"))
|
||||
|
||||
@@ -3615,19 +3615,26 @@ async fn gateway_batch_updates_shared_pool_key_configuration() {
|
||||
Vec::new(),
|
||||
vec![first_key, second_key],
|
||||
));
|
||||
let state = AppState::new()
|
||||
.expect("gateway should build")
|
||||
.with_data_state_for_tests(
|
||||
GatewayDataState::with_provider_catalog_repository_for_tests(Arc::clone(
|
||||
&provider_catalog_repository,
|
||||
)),
|
||||
);
|
||||
let gateway = build_router_with_state(
|
||||
AppState::new()
|
||||
.expect("gateway should build")
|
||||
.with_data_state_for_tests(
|
||||
GatewayDataState::with_provider_catalog_repository_for_tests(Arc::clone(
|
||||
&provider_catalog_repository,
|
||||
)),
|
||||
),
|
||||
);
|
||||
let (gateway_url, gateway_handle) = start_server(gateway).await;
|
||||
|
||||
let response = local_admin_pool_response(
|
||||
&state,
|
||||
http::Method::PATCH,
|
||||
"/api/admin/pool/provider-openai/keys/batch-update",
|
||||
Some(json!({
|
||||
let response = reqwest::Client::new()
|
||||
.patch(format!(
|
||||
"{gateway_url}/api/admin/pool/provider-openai/keys/batch-update"
|
||||
))
|
||||
.header(crate::constants::GATEWAY_HEADER, "rust-phase3b")
|
||||
.header(TRUSTED_ADMIN_USER_ID_HEADER, "admin-user-123")
|
||||
.header(TRUSTED_ADMIN_USER_ROLE_HEADER, "admin")
|
||||
.header(TRUSTED_ADMIN_SESSION_ID_HEADER, "session-123")
|
||||
.json(&json!({
|
||||
"key_ids": ["key-openai-b", "key-openai-a", "key-openai-a"],
|
||||
"patch": {
|
||||
"api_formats": ["openai:responses"],
|
||||
@@ -3638,17 +3645,13 @@ async fn gateway_batch_updates_shared_pool_key_configuration() {
|
||||
"locked_models": [],
|
||||
"note": null
|
||||
}
|
||||
})),
|
||||
)
|
||||
.await;
|
||||
}))
|
||||
.send()
|
||||
.await
|
||||
.expect("request should succeed");
|
||||
|
||||
assert_eq!(response.status(), StatusCode::OK);
|
||||
let payload: serde_json::Value = serde_json::from_slice(
|
||||
&to_bytes(response.into_body(), usize::MAX)
|
||||
.await
|
||||
.expect("body should read"),
|
||||
)
|
||||
.expect("json body should parse");
|
||||
let payload: serde_json::Value = response.json().await.expect("json body should parse");
|
||||
assert_eq!(payload["affected"], json!(2));
|
||||
assert_eq!(payload["model_sync"], serde_json::Value::Null);
|
||||
|
||||
@@ -3670,6 +3673,8 @@ async fn gateway_batch_updates_shared_pool_key_configuration() {
|
||||
assert_eq!(key.locked_models, None);
|
||||
assert_eq!(key.note, None);
|
||||
}
|
||||
|
||||
gateway_handle.abort();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
|
||||
@@ -1942,14 +1942,14 @@ async fn gateway_records_failed_usage_when_all_local_claude_cli_candidates_are_s
|
||||
stored_usage.user_id.as_deref(),
|
||||
Some("user-claude-cli-usage-local-miss-1")
|
||||
);
|
||||
assert_eq!(stored_usage.provider_name, "unknown");
|
||||
assert_eq!(stored_usage.provider_name, "RightCode");
|
||||
assert_eq!(stored_usage.model, "gpt-5.4");
|
||||
assert_eq!(stored_usage.api_format.as_deref(), Some("claude:messages"));
|
||||
assert_eq!(
|
||||
stored_usage.endpoint_api_format.as_deref(),
|
||||
Some("claude:messages")
|
||||
Some("openai:responses")
|
||||
);
|
||||
assert_eq!(stored_usage.routing_key_name(), None);
|
||||
assert_eq!(stored_usage.routing_key_name(), Some("codex"));
|
||||
assert_eq!(stored_usage.routing_planner_kind(), Some("claude_cli_sync"));
|
||||
assert_eq!(stored_usage.routing_route_family(), Some("claude"));
|
||||
assert_eq!(stored_usage.routing_route_kind(), Some("messages"));
|
||||
@@ -2011,7 +2011,14 @@ async fn gateway_records_failed_usage_when_all_local_claude_cli_candidates_are_s
|
||||
stored_candidates[0].skip_reason.as_deref(),
|
||||
Some("format_conversion_disabled")
|
||||
);
|
||||
assert_eq!(stored_usage.routing_candidate_id(), None);
|
||||
assert_eq!(
|
||||
stored_usage.routing_candidate_id(),
|
||||
Some(stored_candidates[0].id.as_str())
|
||||
);
|
||||
assert_eq!(
|
||||
stored_usage.routing_candidate_skip_reason(),
|
||||
Some("format_conversion_disabled")
|
||||
);
|
||||
assert_eq!(*public_hits.lock().expect("mutex should lock"), 0);
|
||||
|
||||
gateway_handle.abort();
|
||||
|
||||
Reference in New Issue
Block a user