Merge upstream/main

This commit is contained in:
ZheFox
2026-05-18 13:11:11 +08:00
124 changed files with 11698 additions and 581 deletions
@@ -2,6 +2,7 @@ use serde_json::json;
use tracing::debug;
use crate::ai_serving::build_request_trace_proxy_value;
use crate::ai_serving::planner::decision_input::apply_provider_request_routing_policy_to_decision;
use crate::ai_serving::planner::report_context::{
build_local_execution_report_context, insert_provider_stream_event_api_format,
LocalExecutionReportContextParts,
@@ -15,7 +16,7 @@ use crate::ai_serving::transport::{
};
use crate::{
append_execution_contract_fields_to_value, append_local_failover_policy_to_value,
AiExecutionDecision, AppState,
AiExecutionDecision, AppState, GatewayError,
};
use super::request::resolve_local_openai_responses_candidate_payload_parts;
@@ -30,7 +31,7 @@ pub(crate) async fn maybe_build_local_openai_responses_decision_payload_for_cand
input: &LocalOpenAiResponsesDecisionInput,
attempt: LocalOpenAiResponsesCandidateAttempt,
spec: LocalOpenAiResponsesSpec,
) -> Option<AiExecutionDecision> {
) -> Result<Option<AiExecutionDecision>, GatewayError> {
let spec_metadata = local_openai_responses_spec_metadata(spec);
let attempt_identity = attempt.attempt_identity();
let LocalOpenAiResponsesCandidateAttempt {
@@ -39,7 +40,7 @@ pub(crate) async fn maybe_build_local_openai_responses_decision_payload_for_cand
candidate_id,
..
} = attempt;
let resolved = resolve_local_openai_responses_candidate_payload_parts(
let Some(resolved) = resolve_local_openai_responses_candidate_payload_parts(
state,
parts,
trace_id,
@@ -50,7 +51,10 @@ pub(crate) async fn maybe_build_local_openai_responses_decision_payload_for_cand
&candidate_id,
spec,
)
.await?;
.await
else {
return Ok(None);
};
let candidate = &eligible.candidate;
let prompt_cache_key = resolved
@@ -102,6 +106,7 @@ pub(crate) async fn maybe_build_local_openai_responses_decision_payload_for_cand
&mut extra_fields,
resolved.transport.provider.provider_type.as_str(),
);
let effective_headers = input.effective_headers(&parts.headers);
let report_context = append_local_failover_policy_to_value(
append_execution_contract_fields_to_value(
build_local_execution_report_context(LocalExecutionReportContextParts {
@@ -129,7 +134,7 @@ pub(crate) async fn maybe_build_local_openai_responses_decision_payload_for_cand
body_rules: resolved.transport.endpoint.body_rules.as_ref(),
provider_request_method: Some(serde_json::Value::Null),
provider_request_headers: Some(&resolved.provider_request_headers),
original_headers: &parts.headers,
original_headers: effective_headers,
request_path: Some(parts.uri.path()),
request_query_string: parts.uri.query(),
request_origin: Some(crate::ai_serving::request_origin_from_parts(parts)),
@@ -197,39 +202,39 @@ pub(crate) async fn maybe_build_local_openai_responses_decision_payload_for_cand
image_request_summary: _,
} = resolved;
Some(build_ai_execution_decision_response(
AiExecutionDecisionResponseParts {
decision_is_stream: spec_metadata.require_streaming,
decision_kind: spec_metadata.decision_kind.to_string(),
execution_strategy,
conversion_mode,
request_id: trace_id.to_string(),
candidate_id: candidate_id.clone(),
provider_name: transport.provider.name.clone(),
provider_id: candidate.provider_id.clone(),
endpoint_id: candidate.endpoint_id.clone(),
key_id: candidate.key_id.clone(),
upstream_base_url: transport.endpoint.base_url.clone(),
upstream_url,
provider_request_method: None,
auth_header: Some(auth_header),
auth_value: Some(auth_value),
provider_api_format,
client_api_format: spec_metadata.api_format.to_string(),
model_name: input.requested_model.clone(),
mapped_model,
prompt_cache_key,
provider_request_headers,
provider_request_body: Some(provider_request_body),
provider_request_body_base64: None,
content_type: Some("application/json".to_string()),
proxy,
transport_profile,
timeouts,
upstream_is_stream,
report_kind: spec_metadata.report_kind.map(ToOwned::to_owned),
report_context: Some(report_context),
auth_context: input.auth_context.clone(),
},
))
let mut decision = build_ai_execution_decision_response(AiExecutionDecisionResponseParts {
decision_is_stream: spec_metadata.require_streaming,
decision_kind: spec_metadata.decision_kind.to_string(),
execution_strategy,
conversion_mode,
request_id: trace_id.to_string(),
candidate_id: candidate_id.clone(),
provider_name: transport.provider.name.clone(),
provider_id: candidate.provider_id.clone(),
endpoint_id: candidate.endpoint_id.clone(),
key_id: candidate.key_id.clone(),
upstream_base_url: transport.endpoint.base_url.clone(),
upstream_url,
provider_request_method: None,
auth_header: Some(auth_header),
auth_value: Some(auth_value),
provider_api_format,
client_api_format: spec_metadata.api_format.to_string(),
model_name: input.requested_model.clone(),
mapped_model,
prompt_cache_key,
provider_request_headers,
provider_request_body: Some(provider_request_body),
provider_request_body_base64: None,
content_type: Some("application/json".to_string()),
proxy,
transport_profile,
timeouts,
upstream_is_stream,
report_kind: spec_metadata.report_kind.map(ToOwned::to_owned),
report_context: Some(report_context),
auth_context: input.auth_context.clone(),
});
apply_provider_request_routing_policy_to_decision(input, &mut decision)?;
Ok(Some(decision))
}
@@ -256,6 +256,7 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
);
let force_body_stream_field =
endpoint_config_forces_body_stream_field(transport.endpoint.config.as_ref());
let effective_headers = input.effective_headers(&parts.headers);
let Some(mut base_provider_request_body) = (if needs_bidirectional_conversion {
build_cross_format_openai_responses_request_body(
body_json,
@@ -271,7 +272,7 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
transport.endpoint.body_rules.as_ref()
},
Some(input.auth_context.api_key_id.as_str()),
&parts.headers,
effective_headers,
enable_model_directives,
)
} else {
@@ -288,7 +289,7 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
transport.endpoint.body_rules.as_ref()
},
Some(input.auth_context.api_key_id.as_str()),
&parts.headers,
effective_headers,
enable_model_directives,
)
}) else {
@@ -452,7 +453,7 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
transport,
provider_api_format,
same_format,
headers: &parts.headers,
headers: effective_headers,
auth_header: &auth_header,
auth_value: &auth_value,
extra_headers: &extra_headers,
@@ -483,7 +484,7 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
apply_codex_openai_responses_special_headers(
&mut provider_request_headers,
&provider_request_body,
&parts.headers,
effective_headers,
transport.provider.provider_type.as_str(),
provider_api_format,
Some(trace_id),
@@ -1022,12 +1023,13 @@ async fn build_kiro_openai_responses_payload_parts(
kiro_auth: &KiroRequestAuth,
) -> Option<LocalOpenAiResponsesCandidatePayloadParts> {
let candidate = &eligible.candidate;
let effective_headers = input.effective_headers(&parts.headers);
let provider_request_body = match build_kiro_provider_request_body(
&claude_request_body,
&mapped_model,
&kiro_auth.auth_config,
transport.endpoint.body_rules.as_ref(),
Some(&parts.headers),
Some(effective_headers),
) {
Some(body) => body,
None => {
@@ -1078,7 +1080,7 @@ async fn build_kiro_openai_responses_payload_parts(
}
};
let provider_request_headers = match build_kiro_provider_headers(KiroProviderHeadersInput {
headers: &parts.headers,
headers: effective_headers,
provider_request_body: &provider_request_body,
original_request_body: original_body_json,
header_rules: transport.endpoint.header_rules.as_ref(),
@@ -20,6 +20,7 @@ use crate::ai_serving::planner::candidate_source::{
};
use crate::ai_serving::planner::common::extract_standard_requested_model;
use crate::ai_serving::planner::decision_input::{
attach_routing_policy_to_local_requested_model_input,
build_local_requested_model_decision_input, resolve_local_authenticated_decision_input,
};
use crate::ai_serving::planner::materialization_policy::{
@@ -50,7 +51,7 @@ pub(crate) async fn resolve_local_openai_responses_decision_input(
decision: &GatewayControlDecision,
body_json: &serde_json::Value,
plan_kind: &str,
) -> Option<LocalOpenAiResponsesDecisionInput> {
) -> Result<Option<LocalOpenAiResponsesDecisionInput>, GatewayError> {
let Some(auth_context) = resolve_local_decision_execution_runtime_auth_context(decision) else {
warn!(
trace_id = %trace_id,
@@ -67,7 +68,7 @@ pub(crate) async fn resolve_local_openai_responses_decision_input(
extract_standard_requested_model(body_json).as_deref(),
"missing_auth_context",
);
return None;
return Ok(None);
};
let Some(requested_model) = extract_standard_requested_model(body_json) else {
@@ -83,7 +84,7 @@ pub(crate) async fn resolve_local_openai_responses_decision_input(
None,
"missing_requested_model",
);
return None;
return Ok(None);
};
let resolved_input = match resolve_local_authenticated_decision_input(
@@ -110,7 +111,7 @@ pub(crate) async fn resolve_local_openai_responses_decision_input(
Some(requested_model.as_str()),
"auth_snapshot_missing",
);
return None;
return Ok(None);
}
Err(err) => {
warn!(
@@ -126,14 +127,30 @@ pub(crate) async fn resolve_local_openai_responses_decision_input(
Some(requested_model.as_str()),
"auth_snapshot_read_failed",
);
return None;
return Err(err);
}
};
let mut input = build_local_requested_model_decision_input(resolved_input, requested_model);
input.request_auth_channel = decision.request_auth_channel.clone();
input.client_session_affinity = client_session_affinity_from_parts(parts, Some(body_json));
Some(input)
if let Err(err) = attach_routing_policy_to_local_requested_model_input(
state,
parts,
&mut input,
body_json,
"openai:responses",
)
.await
{
warn!(
trace_id = %trace_id,
error = ?err,
"gateway local openai responses decision routing profile resolution failed"
);
return Err(err);
}
Ok(Some(input))
}
pub(crate) async fn materialize_local_openai_responses_candidate_attempts(
@@ -159,6 +176,7 @@ pub(crate) async fn materialize_local_openai_responses_candidate_attempts(
spec_metadata.require_streaming,
input.required_capabilities.as_ref(),
&input.auth_snapshot,
input.routing_policy.as_ref(),
input.client_session_affinity.as_ref(),
true,
LocalCandidatePreselectionKeyMode::ProviderEndpointKeyModelAndApiFormat,
@@ -172,6 +190,7 @@ pub(crate) async fn materialize_local_openai_responses_candidate_attempts(
Some(&input.auth_snapshot),
input.client_session_affinity.as_ref(),
input.required_capabilities.as_ref(),
input.routing_policy.as_ref(),
sticky_session_token.as_deref(),
input.request_auth_channel.as_deref(),
persistence_policy,
@@ -268,6 +287,7 @@ pub(crate) async fn build_local_openai_responses_candidate_attempt_source<'a>(
&input.auth_snapshot,
input.client_session_affinity.as_ref(),
input.required_capabilities.as_ref(),
input.routing_policy.as_ref(),
sticky_session_token.as_deref(),
input.request_auth_channel.as_deref(),
persistence_policy,
@@ -350,6 +370,7 @@ pub(crate) async fn build_local_openai_responses_image_candidate_attempt_source<
false,
input.required_capabilities.as_ref(),
&input.auth_snapshot,
input.routing_policy.as_ref(),
input.client_session_affinity.as_ref(),
true,
LocalCandidatePreselectionKeyMode::ProviderEndpointKeyModelAndApiFormat,
@@ -365,6 +386,7 @@ pub(crate) async fn build_local_openai_responses_image_candidate_attempt_source<
Some(&input.auth_snapshot),
input.client_session_affinity.as_ref(),
input.required_capabilities.as_ref(),
input.routing_policy.as_ref(),
sticky_session_token.as_deref(),
input.request_auth_channel.as_deref(),
persistence_policy,
@@ -103,10 +103,11 @@ pub(crate) async fn maybe_build_sync_local_openai_responses_decision_payload(
let Some(input) = resolve_local_openai_responses_decision_input(
state, parts, trace_id, decision, body_json, plan_kind,
)
.await
.await?
else {
return Ok(None);
};
let body_json = input.effective_body_json(body_json);
let (mut source, _) = build_local_openai_responses_candidate_attempt_source(
state, trace_id, &input, body_json, spec,
@@ -117,7 +118,7 @@ pub(crate) async fn maybe_build_sync_local_openai_responses_decision_payload(
if let Some(payload) = maybe_build_local_openai_responses_decision_payload_for_candidate(
state, parts, trace_id, body_json, &input, attempt, spec,
)
.await
.await?
{
return Ok(Some(payload));
}
@@ -141,10 +142,11 @@ pub(crate) async fn maybe_build_stream_local_openai_responses_decision_payload(
let Some(input) = resolve_local_openai_responses_decision_input(
state, parts, trace_id, decision, body_json, plan_kind,
)
.await
.await?
else {
return Ok(None);
};
let body_json = input.effective_body_json(body_json);
let (mut source, _) = build_local_openai_responses_candidate_attempt_source(
state, trace_id, &input, body_json, spec,
@@ -155,7 +157,7 @@ pub(crate) async fn maybe_build_stream_local_openai_responses_decision_payload(
if let Some(payload) = maybe_build_local_openai_responses_decision_payload_for_candidate(
state, parts, trace_id, body_json, &input, attempt, spec,
)
.await
.await?
{
return Ok(Some(payload));
}
@@ -29,7 +29,7 @@ pub(crate) struct LocalOpenAiResponsesSyncAttemptSource<'a> {
state: &'a AppState,
parts: &'a http::request::Parts,
trace_id: &'a str,
body_json: &'a serde_json::Value,
body_json: serde_json::Value,
input: LocalOpenAiResponsesDecisionInput,
spec: LocalOpenAiResponsesSpec,
candidates: LocalOpenAiResponsesCandidateAttemptSource<'a>,
@@ -39,7 +39,7 @@ pub(crate) struct LocalOpenAiResponsesStreamAttemptSource<'a> {
state: &'a AppState,
parts: &'a http::request::Parts,
trace_id: &'a str,
body_json: &'a serde_json::Value,
body_json: serde_json::Value,
input: LocalOpenAiResponsesDecisionInput,
spec: LocalOpenAiResponsesSpec,
candidates: LocalOpenAiResponsesCandidateAttemptSource<'a>,
@@ -62,7 +62,7 @@ pub(super) async fn build_local_sync_attempt_source<'a>(
body_json,
spec_metadata.decision_kind,
)
.await
.await?
else {
return Ok(None);
};
@@ -74,8 +74,13 @@ pub(super) async fn build_local_sync_attempt_source<'a>(
Some(input.requested_model.as_str()),
"candidate_evaluation_incomplete",
);
let effective_body_json = input.effective_body_json(body_json).clone();
let (candidates, candidate_count) = build_local_openai_responses_candidate_attempt_source(
state, trace_id, &input, body_json, spec,
state,
trace_id,
&input,
&effective_body_json,
spec,
)
.await?;
apply_local_runtime_candidate_evaluation_progress(state, trace_id, candidate_count);
@@ -88,7 +93,7 @@ pub(super) async fn build_local_sync_attempt_source<'a>(
state,
parts,
trace_id,
body_json,
body_json: effective_body_json,
input,
spec,
candidates,
@@ -114,7 +119,7 @@ pub(super) async fn build_local_stream_attempt_source<'a>(
body_json,
spec_metadata.decision_kind,
)
.await
.await?
else {
return Ok(None);
};
@@ -126,8 +131,13 @@ pub(super) async fn build_local_stream_attempt_source<'a>(
Some(input.requested_model.as_str()),
"candidate_evaluation_incomplete",
);
let effective_body_json = input.effective_body_json(body_json).clone();
let (candidates, candidate_count) = build_local_openai_responses_candidate_attempt_source(
state, trace_id, &input, body_json, spec,
state,
trace_id,
&input,
&effective_body_json,
spec,
)
.await?;
apply_local_runtime_candidate_evaluation_progress(state, trace_id, candidate_count);
@@ -140,7 +150,7 @@ pub(super) async fn build_local_stream_attempt_source<'a>(
state,
parts,
trace_id,
body_json,
body_json: effective_body_json,
input,
spec,
candidates,
@@ -214,19 +224,19 @@ impl LocalOpenAiResponsesSyncAttemptSource<'_> {
self.state,
self.parts,
self.trace_id,
self.body_json,
&self.body_json,
&self.input,
attempt,
self.spec,
)
.await
.await?
else {
return Ok(None);
};
match build_openai_responses_sync_plan_from_decision(
self.parts,
self.body_json,
&self.body_json,
payload,
self.spec.compact,
) {
@@ -252,19 +262,19 @@ impl LocalOpenAiResponsesStreamAttemptSource<'_> {
self.state,
self.parts,
self.trace_id,
self.body_json,
&self.body_json,
&self.input,
attempt,
self.spec,
)
.await
.await?
else {
return Ok(None);
};
match build_openai_responses_stream_plan_from_decision(
self.parts,
self.body_json,
&self.body_json,
payload,
self.spec.compact,
) {
@@ -298,7 +308,7 @@ pub(super) async fn build_local_sync_plan_and_reports(
body_json,
spec_metadata.decision_kind,
)
.await
.await?
else {
return Ok(Vec::new());
};
@@ -325,7 +335,7 @@ pub(super) async fn build_local_sync_plan_and_reports(
let Some(payload) = maybe_build_local_openai_responses_decision_payload_for_candidate(
state, parts, trace_id, body_json, &input, attempt, spec,
)
.await
.await?
else {
continue;
};
@@ -370,7 +380,7 @@ pub(super) async fn build_local_stream_plan_and_reports(
body_json,
spec_metadata.decision_kind,
)
.await
.await?
else {
return Ok(Vec::new());
};
@@ -397,7 +407,7 @@ pub(super) async fn build_local_stream_plan_and_reports(
let Some(payload) = maybe_build_local_openai_responses_decision_payload_for_candidate(
state, parts, trace_id, body_json, &input, attempt, spec,
)
.await
.await?
else {
continue;
};