Add unified candidate failure diagnostics

This commit is contained in:
fawney19
2026-04-26 01:38:37 +08:00
parent d784c540b6
commit 1d3ea3232d
21 changed files with 1671 additions and 80 deletions

View File

@@ -41,8 +41,9 @@ pub(crate) use self::planner::{
extract_pool_sticky_session_token, maybe_build_stream_decision_payload,
maybe_build_stream_plan_payload, maybe_build_sync_decision_payload,
maybe_build_sync_plan_payload, planner_is_matching_stream_request,
set_local_openai_chat_execution_exhausted_diagnostic, GatewayAuthApiKeySnapshot,
GatewayProviderTransportSnapshot, LocalResolvedOAuthRequestAuth, PlannerAppState,
set_local_openai_chat_execution_exhausted_diagnostic, CandidateFailureDiagnostic,
CandidateFailureDiagnosticKind, GatewayAuthApiKeySnapshot, GatewayProviderTransportSnapshot,
LocalResolvedOAuthRequestAuth, PlannerAppState,
};
pub(crate) use self::pure::*;
pub(crate) use crate::control::GatewayControlDecision;

View File

@@ -6,6 +6,7 @@ use crate::ai_pipeline::planner::candidate_affinity::remember_scheduler_affinity
use crate::ai_pipeline::planner::candidate_eligibility::{
EligibleLocalExecutionCandidate, SkippedLocalExecutionCandidate,
};
use crate::ai_pipeline::planner::failure_diagnostic::CandidateFailureDiagnostic;
use crate::ai_pipeline::planner::runtime_miss::record_local_runtime_candidate_skip_reason;
use crate::ai_pipeline::{GatewayAuthApiKeySnapshot, PlannerAppState};
use crate::clock::current_unix_ms;
@@ -238,6 +239,56 @@ pub(crate) async fn mark_skipped_local_execution_candidate(
.await;
}
pub(crate) async fn mark_skipped_local_execution_candidate_with_extra_data(
state: &AppState,
trace_id: &str,
context: LocalSkippedCandidatePersistenceContext<'_>,
candidate: &SchedulerMinimalCandidateSelectionCandidate,
candidate_index: u32,
candidate_id: &str,
skip_reason: &'static str,
extra_data: Option<Value>,
) {
persist_skipped_local_execution_candidate(
state,
trace_id,
context.user_id,
context.api_key_id,
candidate,
candidate_index,
candidate_id,
context.required_capabilities,
skip_reason,
extra_data,
context.error_context,
context.record_runtime_miss_diagnostic,
)
.await;
}
pub(crate) async fn mark_skipped_local_execution_candidate_with_failure_diagnostic(
state: &AppState,
trace_id: &str,
context: LocalSkippedCandidatePersistenceContext<'_>,
candidate: &SchedulerMinimalCandidateSelectionCandidate,
candidate_index: u32,
candidate_id: &str,
skip_reason: &'static str,
diagnostic: CandidateFailureDiagnostic,
) {
mark_skipped_local_execution_candidate_with_extra_data(
state,
trace_id,
context,
candidate,
candidate_index,
candidate_id,
skip_reason,
Some(diagnostic.to_extra_data()),
)
.await;
}
#[allow(clippy::too_many_arguments)]
pub(crate) async fn persist_skipped_local_execution_candidates(
state: &AppState,

View File

@@ -0,0 +1,190 @@
use serde_json::{json, Value};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum CandidateFailureDiagnosticKind {
RequestBodyBuild,
RequestConversion,
BodyRules,
HeaderRules,
UrlBuild,
TransportAuth,
EnvelopeBuild,
}
impl CandidateFailureDiagnosticKind {
pub(crate) fn as_str(self) -> &'static str {
match self {
Self::RequestBodyBuild => "request_body_build",
Self::RequestConversion => "request_conversion",
Self::BodyRules => "body_rules",
Self::HeaderRules => "header_rules",
Self::UrlBuild => "url_build",
Self::TransportAuth => "transport_auth",
Self::EnvelopeBuild => "envelope_build",
}
}
}
#[derive(Debug, Clone)]
pub(crate) struct CandidateFailureDiagnostic {
pub(crate) kind: CandidateFailureDiagnosticKind,
pub(crate) path: String,
pub(crate) message: String,
pub(crate) source: Option<String>,
pub(crate) client_api_format: Option<String>,
pub(crate) provider_api_format: Option<String>,
pub(crate) safe_to_show: bool,
}
impl CandidateFailureDiagnostic {
pub(crate) fn new(
kind: CandidateFailureDiagnosticKind,
path: impl Into<String>,
message: impl Into<String>,
) -> Self {
Self {
kind,
path: path.into(),
message: message.into(),
source: None,
client_api_format: None,
provider_api_format: None,
safe_to_show: true,
}
}
pub(crate) fn source(mut self, source: impl Into<String>) -> Self {
self.source = Some(source.into());
self
}
pub(crate) fn formats(
mut self,
client_api_format: impl Into<String>,
provider_api_format: impl Into<String>,
) -> Self {
self.client_api_format = Some(client_api_format.into());
self.provider_api_format = Some(provider_api_format.into());
self
}
pub(crate) fn to_extra_data(&self) -> Value {
let diagnostic = self.to_value();
let mut extra_data = json!({
"failure_diagnostic": diagnostic,
});
// Compatibility for current usage UI and already persisted trace readers.
if self.kind == CandidateFailureDiagnosticKind::RequestBodyBuild {
if let Some(object) = extra_data.as_object_mut() {
object.insert(
"request_body_build_error".to_string(),
json!({
"path": self.path,
"message": self.message,
"client_api_format": self.client_api_format,
"provider_api_format": self.provider_api_format,
}),
);
}
}
extra_data
}
pub(crate) fn upstream_url_missing(
client_api_format: impl Into<String>,
provider_api_format: impl Into<String>,
source: impl Into<String>,
) -> Self {
Self::new(
CandidateFailureDiagnosticKind::UrlBuild,
"$.endpoint",
"无法构建上游请求地址;请检查 base_url、custom_path、API 格式和模型映射",
)
.formats(client_api_format, provider_api_format)
.source(source)
}
pub(crate) fn header_rules_apply_failed(
client_api_format: impl Into<String>,
provider_api_format: impl Into<String>,
source: impl Into<String>,
) -> Self {
Self::new(
CandidateFailureDiagnosticKind::HeaderRules,
"$.endpoint.header_rules",
"Header 规则应用失败;请检查规则格式、条件配置,或是否试图覆盖受保护认证头",
)
.formats(client_api_format, provider_api_format)
.source(source)
}
pub(crate) fn body_rules_apply_failed(
client_api_format: impl Into<String>,
provider_api_format: impl Into<String>,
source: impl Into<String>,
) -> Self {
Self::new(
CandidateFailureDiagnosticKind::BodyRules,
"$.endpoint.body_rules",
"Body 规则应用失败;请检查规则格式、条件配置,或规则输出是否仍是当前上游支持的请求体",
)
.formats(client_api_format, provider_api_format)
.source(source)
}
pub(crate) fn body_rules_unsupported_for_binary_upload(
client_api_format: impl Into<String>,
provider_api_format: impl Into<String>,
source: impl Into<String>,
) -> Self {
Self::new(
CandidateFailureDiagnosticKind::BodyRules,
"$.endpoint.body_rules",
"二进制上传暂不支持本地应用 Body 规则;请移除该 Endpoint 的 Body 规则或改用 JSON 请求体",
)
.formats(client_api_format, provider_api_format)
.source(source)
}
pub(crate) fn provider_request_body_missing(
client_api_format: impl Into<String>,
provider_api_format: impl Into<String>,
source: impl Into<String>,
) -> Self {
Self::new(
CandidateFailureDiagnosticKind::RequestBodyBuild,
"$",
"无法构建上游请求体;请检查请求体是否为支持的 JSON object以及该任务类型必需字段是否存在且取值受支持",
)
.formats(client_api_format, provider_api_format)
.source(source)
}
pub(crate) fn envelope_build_failed(
client_api_format: impl Into<String>,
provider_api_format: impl Into<String>,
source: impl Into<String>,
) -> Self {
Self::new(
CandidateFailureDiagnosticKind::EnvelopeBuild,
"$",
"无法构建上游请求封装;请检查该 Provider 的认证配置、模型映射、Endpoint Body 规则和当前请求体是否兼容",
)
.formats(client_api_format, provider_api_format)
.source(source)
}
fn to_value(&self) -> Value {
json!({
"kind": self.kind.as_str(),
"path": self.path,
"message": self.message,
"source": self.source,
"client_api_format": self.client_api_format,
"provider_api_format": self.provider_api_format,
"safe_to_show": self.safe_to_show,
})
}
}

View File

@@ -13,6 +13,7 @@ mod candidate_source;
mod common;
mod decision;
mod decision_input;
mod failure_diagnostic;
mod materialization_policy;
mod passthrough;
mod payload_metadata;
@@ -27,6 +28,9 @@ mod standard;
mod state;
pub(crate) use self::candidate_eligibility::extract_pool_sticky_session_token;
pub(crate) use self::failure_diagnostic::{
CandidateFailureDiagnostic, CandidateFailureDiagnosticKind,
};
pub(crate) use self::passthrough::{
build_local_same_format_stream_plan_and_reports, build_local_same_format_sync_plan_and_reports,
};

View File

@@ -1,6 +1,9 @@
use serde_json::json;
use crate::ai_pipeline::planner::candidate_materialization::mark_skipped_local_execution_candidate;
use crate::ai_pipeline::planner::candidate_materialization::{
mark_skipped_local_execution_candidate, mark_skipped_local_execution_candidate_with_extra_data,
mark_skipped_local_execution_candidate_with_failure_diagnostic,
};
use crate::ai_pipeline::planner::candidate_metadata::build_request_trace_proxy_value;
use crate::ai_pipeline::planner::materialization_policy::{
build_local_candidate_persistence_policy, LocalCandidatePersistencePolicyKind,
@@ -12,6 +15,7 @@ use crate::ai_pipeline::planner::report_context::{
build_local_execution_report_context, LocalExecutionReportContextParts,
};
use crate::ai_pipeline::planner::spec_metadata::local_same_format_provider_spec_metadata;
use crate::ai_pipeline::planner::CandidateFailureDiagnostic;
use crate::ai_pipeline::transport::{
resolve_transport_execution_timeouts, resolve_transport_tls_profile,
};
@@ -198,3 +202,61 @@ pub(super) async fn mark_skipped_local_same_format_provider_candidate(
)
.await;
}
#[allow(clippy::too_many_arguments)]
pub(super) async fn mark_skipped_local_same_format_provider_candidate_with_extra_data(
state: &AppState,
input: &LocalSameFormatProviderDecisionInput,
trace_id: &str,
candidate: &SchedulerMinimalCandidateSelectionCandidate,
candidate_index: u32,
candidate_id: &str,
skip_reason: &'static str,
extra_data: Option<serde_json::Value>,
) {
let persistence_policy = build_local_candidate_persistence_policy(
&input.auth_context,
input.required_capabilities.as_ref(),
LocalCandidatePersistencePolicyKind::SameFormatProviderDecision,
);
mark_skipped_local_execution_candidate_with_extra_data(
state,
trace_id,
persistence_policy.skipped,
candidate,
candidate_index,
candidate_id,
skip_reason,
extra_data,
)
.await;
}
#[allow(clippy::too_many_arguments)]
pub(super) async fn mark_skipped_local_same_format_provider_candidate_with_failure_diagnostic(
state: &AppState,
input: &LocalSameFormatProviderDecisionInput,
trace_id: &str,
candidate: &SchedulerMinimalCandidateSelectionCandidate,
candidate_index: u32,
candidate_id: &str,
skip_reason: &'static str,
diagnostic: CandidateFailureDiagnostic,
) {
let persistence_policy = build_local_candidate_persistence_policy(
&input.auth_context,
input.required_capabilities.as_ref(),
LocalCandidatePersistencePolicyKind::SameFormatProviderDecision,
);
mark_skipped_local_execution_candidate_with_failure_diagnostic(
state,
trace_id,
persistence_policy.skipped,
candidate,
candidate_index,
candidate_id,
skip_reason,
diagnostic,
)
.await;
}

View File

@@ -14,18 +14,23 @@ use crate::ai_pipeline::transport::auth::{
use crate::ai_pipeline::transport::claude_code::build_claude_code_passthrough_headers;
use crate::ai_pipeline::transport::kiro::{build_kiro_provider_headers, KiroProviderHeadersInput};
use crate::ai_pipeline::transport::{apply_local_header_rules, ensure_upstream_auth_header};
use crate::ai_pipeline::GatewayProviderTransportSnapshot;
use crate::ai_pipeline::{CandidateFailureDiagnostic, GatewayProviderTransportSnapshot};
use crate::AppState;
mod policy;
mod prepare;
use self::prepare::prepare_local_same_format_provider_candidate;
use super::payload::mark_skipped_local_same_format_provider_candidate;
use super::payload::{
mark_skipped_local_same_format_provider_candidate,
mark_skipped_local_same_format_provider_candidate_with_extra_data,
mark_skipped_local_same_format_provider_candidate_with_failure_diagnostic,
};
use super::{
LocalSameFormatProviderCandidateAttempt, LocalSameFormatProviderDecisionInput,
LocalSameFormatProviderSpec,
};
use crate::ai_pipeline::planner::standard::same_format_provider_request_body_failure_extra_data;
pub(crate) fn resolve_same_format_provider_transport_unsupported_reason_for_trace(
transport: &GatewayProviderTransportSnapshot,
@@ -118,7 +123,7 @@ pub(crate) async fn resolve_local_same_format_provider_candidate_payload_parts(
prepared.is_claude_code,
)
else {
mark_skipped_local_same_format_provider_candidate(
mark_skipped_local_same_format_provider_candidate_with_extra_data(
state,
input,
trace_id,
@@ -126,6 +131,16 @@ pub(crate) async fn resolve_local_same_format_provider_candidate_payload_parts(
attempt.candidate_index,
&attempt.candidate_id,
"provider_request_body_missing",
same_format_provider_request_body_failure_extra_data(
body_json,
attempt.eligible.provider_api_format.as_str(),
prepared.transport.endpoint.body_rules.as_ref(),
if prepared.kiro_auth.is_some() {
"kiro_envelope"
} else {
"same_format"
},
),
)
.await;
return None;
@@ -165,7 +180,7 @@ pub(crate) async fn resolve_local_same_format_provider_candidate_payload_parts(
) {
AntigravityRequestEnvelopeSupport::Supported(envelope) => envelope,
AntigravityRequestEnvelopeSupport::Unsupported(_) => {
mark_skipped_local_same_format_provider_candidate(
mark_skipped_local_same_format_provider_candidate_with_extra_data(
state,
input,
trace_id,
@@ -173,6 +188,12 @@ pub(crate) async fn resolve_local_same_format_provider_candidate_payload_parts(
attempt.candidate_index,
&attempt.candidate_id,
"provider_request_body_missing",
same_format_provider_request_body_failure_extra_data(
body_json,
attempt.eligible.provider_api_format.as_str(),
prepared.transport.endpoint.body_rules.as_ref(),
"antigravity_envelope",
),
)
.await;
return None;
@@ -190,7 +211,7 @@ pub(crate) async fn resolve_local_same_format_provider_candidate_payload_parts(
prepared.upstream_is_stream,
prepared.kiro_auth.as_ref(),
) else {
mark_skipped_local_same_format_provider_candidate(
mark_skipped_local_same_format_provider_candidate_with_failure_diagnostic(
state,
input,
trace_id,
@@ -198,6 +219,11 @@ pub(crate) async fn resolve_local_same_format_provider_candidate_payload_parts(
attempt.candidate_index,
&attempt.candidate_id,
"upstream_url_missing",
CandidateFailureDiagnostic::upstream_url_missing(
attempt.eligible.provider_api_format.as_str(),
attempt.eligible.provider_api_format.as_str(),
"same_format_provider_url",
),
)
.await;
return None;
@@ -271,7 +297,7 @@ pub(crate) async fn resolve_local_same_format_provider_candidate_payload_parts(
Some(provider_request_headers)
}
}) else {
mark_skipped_local_same_format_provider_candidate(
mark_skipped_local_same_format_provider_candidate_with_failure_diagnostic(
state,
input,
trace_id,
@@ -279,6 +305,11 @@ pub(crate) async fn resolve_local_same_format_provider_candidate_payload_parts(
attempt.candidate_index,
&attempt.candidate_id,
"transport_header_rules_apply_failed",
CandidateFailureDiagnostic::header_rules_apply_failed(
attempt.eligible.provider_api_format.as_str(),
attempt.eligible.provider_api_format.as_str(),
"same_format_provider_headers",
),
)
.await;
return None;

View File

@@ -11,12 +11,14 @@ use crate::ai_pipeline::transport::auth::{
use crate::ai_pipeline::transport::local_gemini_transport_unsupported_reason_with_network;
use crate::ai_pipeline::transport::url::build_gemini_files_passthrough_url;
use crate::ai_pipeline::transport::{apply_local_body_rules, apply_local_header_rules};
use crate::ai_pipeline::GatewayProviderTransportSnapshot;
use crate::ai_pipeline::{CandidateFailureDiagnostic, GatewayProviderTransportSnapshot};
use crate::AppState;
use super::support::{
mark_skipped_local_gemini_files_candidate, LocalGeminiFilesCandidateAttempt,
LocalGeminiFilesDecisionInput, GEMINI_FILES_CANDIDATE_API_FORMAT,
mark_skipped_local_gemini_files_candidate,
mark_skipped_local_gemini_files_candidate_with_failure_diagnostic,
LocalGeminiFilesCandidateAttempt, LocalGeminiFilesDecisionInput,
GEMINI_FILES_CANDIDATE_API_FORMAT, GEMINI_FILES_CLIENT_API_FORMAT,
};
use super::LocalGeminiFilesSpec;
@@ -90,7 +92,7 @@ pub(super) async fn resolve_local_gemini_files_candidate_payload_parts(
passthrough_path,
parts.uri.query(),
) else {
mark_skipped_local_gemini_files_candidate(
mark_skipped_local_gemini_files_candidate_with_failure_diagnostic(
state,
input,
trace_id,
@@ -98,6 +100,11 @@ pub(super) async fn resolve_local_gemini_files_candidate_payload_parts(
attempt.candidate_index,
&attempt.candidate_id,
"upstream_url_missing",
CandidateFailureDiagnostic::upstream_url_missing(
GEMINI_FILES_CLIENT_API_FORMAT,
GEMINI_FILES_CANDIDATE_API_FORMAT,
"gemini_files_passthrough_url",
),
)
.await;
return None;
@@ -120,7 +127,7 @@ pub(super) async fn resolve_local_gemini_files_candidate_payload_parts(
None
};
if provider_request_body_base64.is_some() && transport.endpoint.body_rules.is_some() {
mark_skipped_local_gemini_files_candidate(
mark_skipped_local_gemini_files_candidate_with_failure_diagnostic(
state,
input,
trace_id,
@@ -128,6 +135,11 @@ pub(super) async fn resolve_local_gemini_files_candidate_payload_parts(
attempt.candidate_index,
&attempt.candidate_id,
"transport_body_rules_unsupported_for_binary_upload",
CandidateFailureDiagnostic::body_rules_unsupported_for_binary_upload(
GEMINI_FILES_CLIENT_API_FORMAT,
GEMINI_FILES_CANDIDATE_API_FORMAT,
"gemini_files_binary_upload",
),
)
.await;
return None;
@@ -138,7 +150,7 @@ pub(super) async fn resolve_local_gemini_files_candidate_payload_parts(
transport.endpoint.body_rules.as_ref(),
Some(body_json),
) {
mark_skipped_local_gemini_files_candidate(
mark_skipped_local_gemini_files_candidate_with_failure_diagnostic(
state,
input,
trace_id,
@@ -146,6 +158,11 @@ pub(super) async fn resolve_local_gemini_files_candidate_payload_parts(
attempt.candidate_index,
&attempt.candidate_id,
"transport_body_rules_apply_failed",
CandidateFailureDiagnostic::body_rules_apply_failed(
GEMINI_FILES_CLIENT_API_FORMAT,
GEMINI_FILES_CANDIDATE_API_FORMAT,
"gemini_files_body_rules",
),
)
.await;
return None;
@@ -175,7 +192,7 @@ pub(super) async fn resolve_local_gemini_files_candidate_payload_parts(
.unwrap_or(original_request_body),
Some(original_request_body),
) {
mark_skipped_local_gemini_files_candidate(
mark_skipped_local_gemini_files_candidate_with_failure_diagnostic(
state,
input,
trace_id,
@@ -183,6 +200,11 @@ pub(super) async fn resolve_local_gemini_files_candidate_payload_parts(
attempt.candidate_index,
&attempt.candidate_id,
"transport_header_rules_apply_failed",
CandidateFailureDiagnostic::header_rules_apply_failed(
GEMINI_FILES_CLIENT_API_FORMAT,
GEMINI_FILES_CANDIDATE_API_FORMAT,
"gemini_files_header_rules",
),
)
.await;
return None;

View File

@@ -6,6 +6,7 @@ use crate::ai_pipeline::contracts::ExecutionRuntimeAuthContext;
use crate::ai_pipeline::planner::candidate_eligibility::filter_and_rank_local_execution_candidates_without_transport_pair_gate;
use crate::ai_pipeline::planner::candidate_materialization::{
mark_skipped_local_execution_candidate,
mark_skipped_local_execution_candidate_with_failure_diagnostic,
persist_available_local_execution_candidates_with_context,
persist_skipped_local_execution_candidates_with_context,
remember_first_local_candidate_affinity,
@@ -22,7 +23,8 @@ use crate::ai_pipeline::planner::materialization_policy::{
};
use crate::ai_pipeline::PlannerAppState;
use crate::ai_pipeline::{
resolve_local_decision_execution_runtime_auth_context, GatewayControlDecision,
resolve_local_decision_execution_runtime_auth_context, CandidateFailureDiagnostic,
GatewayControlDecision,
};
use crate::clock::current_unix_secs;
use crate::{AppState, GatewayError};
@@ -184,3 +186,31 @@ pub(super) async fn mark_skipped_local_gemini_files_candidate(
)
.await;
}
pub(super) async fn mark_skipped_local_gemini_files_candidate_with_failure_diagnostic(
state: &AppState,
input: &LocalGeminiFilesDecisionInput,
trace_id: &str,
candidate: &SchedulerMinimalCandidateSelectionCandidate,
candidate_index: u32,
candidate_id: &str,
skip_reason: &'static str,
diagnostic: CandidateFailureDiagnostic,
) {
let persistence_policy = build_local_candidate_persistence_policy(
&input.auth_context,
input.required_capabilities.as_ref(),
LocalCandidatePersistencePolicyKind::GeminiFilesDecision,
);
mark_skipped_local_execution_candidate_with_failure_diagnostic(
state,
trace_id,
persistence_policy.skipped,
candidate,
candidate_index,
candidate_id,
skip_reason,
diagnostic,
)
.await;
}

View File

@@ -17,14 +17,15 @@ use crate::ai_pipeline::transport::{
};
use crate::ai_pipeline::{
apply_codex_openai_cli_special_body_edits, apply_codex_openai_cli_special_headers,
GatewayProviderTransportSnapshot, PlannerAppState, CODEX_OPENAI_IMAGE_DEFAULT_MODEL,
CODEX_OPENAI_IMAGE_DEFAULT_VARIATION_MODEL,
CandidateFailureDiagnostic, GatewayProviderTransportSnapshot, PlannerAppState,
CODEX_OPENAI_IMAGE_DEFAULT_MODEL, CODEX_OPENAI_IMAGE_DEFAULT_VARIATION_MODEL,
};
use crate::AppState;
use super::support::{
mark_skipped_local_openai_image_candidate, LocalOpenAiImageCandidateAttempt,
LocalOpenAiImageDecisionInput,
mark_skipped_local_openai_image_candidate,
mark_skipped_local_openai_image_candidate_with_failure_diagnostic,
LocalOpenAiImageCandidateAttempt, LocalOpenAiImageDecisionInput,
};
use super::LocalOpenAiImageSpec;
@@ -190,7 +191,7 @@ pub(super) async fn resolve_local_openai_image_candidate_payload_parts(
let Some(normalized_request) = normalize_openai_image_request(parts, body_json, body_base64)
else {
mark_skipped_local_openai_image_candidate(
mark_skipped_local_openai_image_candidate_with_failure_diagnostic(
state,
input,
trace_id,
@@ -198,6 +199,11 @@ pub(super) async fn resolve_local_openai_image_candidate_payload_parts(
attempt.candidate_index,
&attempt.candidate_id,
"provider_request_body_missing",
CandidateFailureDiagnostic::provider_request_body_missing(
spec_metadata.api_format,
spec_metadata.api_format,
"openai_image_request_normalize",
),
)
.await;
return None;
@@ -228,7 +234,7 @@ pub(super) async fn resolve_local_openai_image_candidate_payload_parts(
&provider_request_body,
Some(body_json),
) {
mark_skipped_local_openai_image_candidate(
mark_skipped_local_openai_image_candidate_with_failure_diagnostic(
state,
input,
trace_id,
@@ -236,6 +242,11 @@ pub(super) async fn resolve_local_openai_image_candidate_payload_parts(
attempt.candidate_index,
&attempt.candidate_id,
"transport_header_rules_apply_failed",
CandidateFailureDiagnostic::header_rules_apply_failed(
spec_metadata.api_format,
spec_metadata.api_format,
"openai_image_header_rules",
),
)
.await;
return None;

View File

@@ -7,6 +7,7 @@ use crate::ai_pipeline::planner::candidate_eligibility::{
};
use crate::ai_pipeline::planner::candidate_materialization::{
mark_skipped_local_execution_candidate,
mark_skipped_local_execution_candidate_with_failure_diagnostic,
persist_available_local_execution_candidates_with_context,
persist_skipped_local_execution_candidates_with_context,
remember_first_local_candidate_affinity,
@@ -24,7 +25,8 @@ use crate::ai_pipeline::planner::materialization_policy::{
use crate::ai_pipeline::planner::spec_metadata::local_openai_image_spec_metadata;
use crate::ai_pipeline::PlannerAppState;
use crate::ai_pipeline::{
resolve_local_decision_execution_runtime_auth_context, GatewayControlDecision,
resolve_local_decision_execution_runtime_auth_context, CandidateFailureDiagnostic,
GatewayControlDecision,
};
use crate::clock::current_unix_secs;
use crate::AppState;
@@ -239,3 +241,31 @@ pub(super) async fn mark_skipped_local_openai_image_candidate(
)
.await;
}
pub(super) async fn mark_skipped_local_openai_image_candidate_with_failure_diagnostic(
state: &AppState,
input: &LocalOpenAiImageDecisionInput,
trace_id: &str,
candidate: &SchedulerMinimalCandidateSelectionCandidate,
candidate_index: u32,
candidate_id: &str,
skip_reason: &'static str,
diagnostic: CandidateFailureDiagnostic,
) {
let persistence_policy = build_local_candidate_persistence_policy(
&input.auth_context,
input.required_capabilities.as_ref(),
LocalCandidatePersistencePolicyKind::ImageDecision,
);
mark_skipped_local_execution_candidate_with_failure_diagnostic(
state,
trace_id,
persistence_policy.skipped,
candidate,
candidate_index,
candidate_id,
skip_reason,
diagnostic,
)
.await;
}

View File

@@ -17,12 +17,12 @@ use crate::ai_pipeline::transport::{
local_gemini_transport_unsupported_reason_with_network,
local_standard_transport_unsupported_reason_with_network,
};
use crate::ai_pipeline::GatewayProviderTransportSnapshot;
use crate::ai_pipeline::{CandidateFailureDiagnostic, GatewayProviderTransportSnapshot};
use crate::AppState;
use super::support::{
mark_skipped_local_video_candidate, LocalVideoCreateCandidateAttempt,
LocalVideoCreateDecisionInput,
mark_skipped_local_video_candidate, mark_skipped_local_video_candidate_with_failure_diagnostic,
LocalVideoCreateCandidateAttempt, LocalVideoCreateDecisionInput,
};
use super::{LocalVideoCreateFamily, LocalVideoCreateSpec};
@@ -110,7 +110,7 @@ pub(super) async fn resolve_local_video_create_candidate_payload_parts(
let Some(upstream_url) = build_video_upstream_url(parts, transport, &mapped_model, spec.family)
else {
mark_skipped_local_video_candidate(
mark_skipped_local_video_candidate_with_failure_diagnostic(
state,
input,
trace_id,
@@ -118,25 +118,35 @@ pub(super) async fn resolve_local_video_create_candidate_payload_parts(
attempt.candidate_index,
&attempt.candidate_id,
"upstream_url_missing",
CandidateFailureDiagnostic::upstream_url_missing(
spec_metadata.api_format,
spec_metadata.api_format,
"video_upstream_url",
),
)
.await;
return None;
};
let Some(provider_request_body) = build_provider_request_body(
let Ok(provider_request_body) = build_provider_request_body(
body_json,
spec.family,
&mapped_model,
transport.endpoint.body_rules.as_ref(),
) else {
mark_skipped_local_video_candidate(
mark_skipped_local_video_candidate_with_failure_diagnostic(
state,
input,
trace_id,
candidate,
attempt.candidate_index,
&attempt.candidate_id,
"provider_request_body_missing",
"transport_body_rules_apply_failed",
CandidateFailureDiagnostic::body_rules_apply_failed(
spec_metadata.api_format,
spec_metadata.api_format,
"video_body_rules",
),
)
.await;
return None;
@@ -155,7 +165,7 @@ pub(super) async fn resolve_local_video_create_candidate_payload_parts(
&provider_request_body,
Some(body_json),
) {
mark_skipped_local_video_candidate(
mark_skipped_local_video_candidate_with_failure_diagnostic(
state,
input,
trace_id,
@@ -163,6 +173,11 @@ pub(super) async fn resolve_local_video_create_candidate_payload_parts(
attempt.candidate_index,
&attempt.candidate_id,
"transport_header_rules_apply_failed",
CandidateFailureDiagnostic::header_rules_apply_failed(
spec_metadata.api_format,
spec_metadata.api_format,
"video_header_rules",
),
)
.await;
return None;
@@ -184,7 +199,7 @@ fn build_provider_request_body(
family: LocalVideoCreateFamily,
mapped_model: &str,
body_rules: Option<&serde_json::Value>,
) -> Option<serde_json::Value> {
) -> Result<serde_json::Value, ()> {
let mut provider_request_body = match family {
LocalVideoCreateFamily::OpenAi => {
let mut provider_request_body = body_json.as_object().cloned().unwrap_or_default();
@@ -195,9 +210,9 @@ fn build_provider_request_body(
LocalVideoCreateFamily::Gemini => body_json.clone(),
};
if !apply_local_body_rules(&mut provider_request_body, body_rules, Some(body_json)) {
return None;
return Err(());
}
Some(provider_request_body)
Ok(provider_request_body)
}
fn build_video_upstream_url(

View File

@@ -9,6 +9,7 @@ use crate::ai_pipeline::planner::candidate_eligibility::{
};
use crate::ai_pipeline::planner::candidate_materialization::{
mark_skipped_local_execution_candidate,
mark_skipped_local_execution_candidate_with_failure_diagnostic,
persist_available_local_execution_candidates_with_context,
persist_skipped_local_execution_candidates_with_context,
remember_first_local_candidate_affinity,
@@ -27,7 +28,8 @@ use crate::ai_pipeline::planner::materialization_policy::{
use crate::ai_pipeline::planner::spec_metadata::local_video_create_spec_metadata;
use crate::ai_pipeline::PlannerAppState;
use crate::ai_pipeline::{
resolve_local_decision_execution_runtime_auth_context, GatewayControlDecision,
resolve_local_decision_execution_runtime_auth_context, CandidateFailureDiagnostic,
GatewayControlDecision,
};
use crate::clock::current_unix_secs;
use crate::AppState;
@@ -251,3 +253,31 @@ pub(super) async fn mark_skipped_local_video_candidate(
)
.await;
}
pub(super) async fn mark_skipped_local_video_candidate_with_failure_diagnostic(
state: &AppState,
input: &LocalVideoCreateDecisionInput,
trace_id: &str,
candidate: &SchedulerMinimalCandidateSelectionCandidate,
candidate_index: u32,
candidate_id: &str,
skip_reason: &'static str,
diagnostic: CandidateFailureDiagnostic,
) {
let persistence_policy = build_local_candidate_persistence_policy(
&input.auth_context,
input.required_capabilities.as_ref(),
LocalCandidatePersistencePolicyKind::VideoDecision,
);
mark_skipped_local_execution_candidate_with_failure_diagnostic(
state,
trace_id,
persistence_policy.skipped,
candidate,
candidate_index,
candidate_id,
skip_reason,
diagnostic,
)
.await;
}

View File

@@ -1,4 +1,7 @@
use crate::ai_pipeline::planner::candidate_materialization::mark_skipped_local_execution_candidate;
use crate::ai_pipeline::planner::candidate_materialization::{
mark_skipped_local_execution_candidate, mark_skipped_local_execution_candidate_with_extra_data,
mark_skipped_local_execution_candidate_with_failure_diagnostic,
};
use crate::ai_pipeline::planner::candidate_metadata::build_request_trace_proxy_value;
use crate::ai_pipeline::planner::materialization_policy::{
build_local_candidate_persistence_policy, LocalCandidatePersistencePolicyKind,
@@ -10,6 +13,7 @@ use crate::ai_pipeline::planner::report_context::{
build_local_execution_report_context, LocalExecutionReportContextParts,
};
use crate::ai_pipeline::planner::spec_metadata::local_standard_spec_metadata;
use crate::ai_pipeline::planner::CandidateFailureDiagnostic;
use crate::ai_pipeline::transport::{
resolve_transport_execution_timeouts, resolve_transport_tls_profile,
};
@@ -179,3 +183,61 @@ pub(super) async fn mark_skipped_local_standard_candidate(
)
.await;
}
#[allow(clippy::too_many_arguments)]
pub(super) async fn mark_skipped_local_standard_candidate_with_extra_data(
state: &AppState,
input: &LocalStandardDecisionInput,
trace_id: &str,
candidate: &aether_scheduler_core::SchedulerMinimalCandidateSelectionCandidate,
candidate_index: u32,
candidate_id: &str,
skip_reason: &'static str,
extra_data: Option<serde_json::Value>,
) {
let persistence_policy = build_local_candidate_persistence_policy(
&input.auth_context,
input.required_capabilities.as_ref(),
LocalCandidatePersistencePolicyKind::StandardDecision,
);
mark_skipped_local_execution_candidate_with_extra_data(
state,
trace_id,
persistence_policy.skipped,
candidate,
candidate_index,
candidate_id,
skip_reason,
extra_data,
)
.await;
}
#[allow(clippy::too_many_arguments)]
pub(super) async fn mark_skipped_local_standard_candidate_with_failure_diagnostic(
state: &AppState,
input: &LocalStandardDecisionInput,
trace_id: &str,
candidate: &aether_scheduler_core::SchedulerMinimalCandidateSelectionCandidate,
candidate_index: u32,
candidate_id: &str,
skip_reason: &'static str,
diagnostic: CandidateFailureDiagnostic,
) {
let persistence_policy = build_local_candidate_persistence_policy(
&input.auth_context,
input.required_capabilities.as_ref(),
LocalCandidatePersistencePolicyKind::StandardDecision,
);
mark_skipped_local_execution_candidate_with_failure_diagnostic(
state,
trace_id,
persistence_policy.skipped,
candidate,
candidate_index,
candidate_id,
skip_reason,
diagnostic,
)
.await;
}

View File

@@ -8,7 +8,9 @@ use crate::ai_pipeline::planner::candidate_preparation::{
};
use crate::ai_pipeline::planner::common::force_upstream_streaming_for_provider;
use crate::ai_pipeline::planner::spec_metadata::local_standard_spec_metadata;
use crate::ai_pipeline::planner::standard::apply_codex_openai_cli_special_headers;
use crate::ai_pipeline::planner::standard::{
apply_codex_openai_cli_special_headers, request_body_build_failure_extra_data,
};
use crate::ai_pipeline::transport::apply_local_header_rules;
use crate::ai_pipeline::transport::auth::{
build_claude_passthrough_headers, build_openai_passthrough_headers, ensure_upstream_auth_header,
@@ -18,10 +20,15 @@ use crate::ai_pipeline::transport::kiro::{
KiroRequestAuth, KIRO_ENVELOPE_NAME,
};
use crate::ai_pipeline::transport::vertex::uses_vertex_api_key_query_auth;
use crate::ai_pipeline::{GatewayProviderTransportSnapshot, LocalResolvedOAuthRequestAuth};
use crate::ai_pipeline::{
CandidateFailureDiagnostic, GatewayProviderTransportSnapshot, LocalResolvedOAuthRequestAuth,
};
use crate::AppState;
use super::payload::mark_skipped_local_standard_candidate;
use super::payload::{
mark_skipped_local_standard_candidate, mark_skipped_local_standard_candidate_with_extra_data,
mark_skipped_local_standard_candidate_with_failure_diagnostic,
};
use super::{LocalStandardCandidateAttempt, LocalStandardDecisionInput, LocalStandardSpec};
pub(crate) struct LocalStandardCandidatePayloadParts {
@@ -193,14 +200,19 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
) {
Some(body) => body,
None => {
mark_skipped_local_standard_candidate(
mark_skipped_local_standard_candidate_with_extra_data(
state,
input,
trace_id,
candidate,
attempt.candidate_index,
&attempt.candidate_id,
"provider_request_body_missing",
"provider_request_body_build_failed",
request_body_build_failure_extra_data(
body_json,
spec_metadata.api_format,
provider_api_format,
),
)
.await;
return None;
@@ -236,7 +248,7 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
) {
Some(url) => url,
None => {
mark_skipped_local_standard_candidate(
mark_skipped_local_standard_candidate_with_failure_diagnostic(
state,
input,
trace_id,
@@ -244,6 +256,11 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
attempt.candidate_index,
&attempt.candidate_id,
"upstream_url_missing",
CandidateFailureDiagnostic::upstream_url_missing(
spec_metadata.api_format,
provider_api_format,
"standard_family_url",
),
)
.await;
return None;
@@ -280,7 +297,7 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
&provider_request_body,
Some(body_json),
) {
mark_skipped_local_standard_candidate(
mark_skipped_local_standard_candidate_with_failure_diagnostic(
state,
input,
trace_id,
@@ -288,6 +305,11 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
attempt.candidate_index,
&attempt.candidate_id,
"transport_header_rules_apply_failed",
CandidateFailureDiagnostic::header_rules_apply_failed(
spec_metadata.api_format,
provider_api_format,
"standard_family_headers",
),
)
.await;
return None;
@@ -361,14 +383,19 @@ async fn build_kiro_cross_format_payload_parts(
) {
Some(body) => body,
None => {
mark_skipped_local_standard_candidate(
mark_skipped_local_standard_candidate_with_extra_data(
state,
input,
trace_id,
candidate,
attempt.candidate_index,
&attempt.candidate_id,
"provider_request_body_missing",
"provider_request_body_build_failed",
request_body_build_failure_extra_data(
&claude_request_body,
provider_api_format,
provider_api_format,
),
)
.await;
return None;
@@ -384,7 +411,7 @@ async fn build_kiro_cross_format_payload_parts(
) {
Some(url) => url,
None => {
mark_skipped_local_standard_candidate(
mark_skipped_local_standard_candidate_with_failure_diagnostic(
state,
input,
trace_id,
@@ -392,6 +419,11 @@ async fn build_kiro_cross_format_payload_parts(
attempt.candidate_index,
&attempt.candidate_id,
"upstream_url_missing",
CandidateFailureDiagnostic::upstream_url_missing(
provider_api_format,
provider_api_format,
"standard_family_kiro_url",
),
)
.await;
return None;
@@ -409,7 +441,7 @@ async fn build_kiro_cross_format_payload_parts(
}) {
Some(headers) => headers,
None => {
mark_skipped_local_standard_candidate(
mark_skipped_local_standard_candidate_with_failure_diagnostic(
state,
input,
trace_id,
@@ -417,6 +449,11 @@ async fn build_kiro_cross_format_payload_parts(
attempt.candidate_index,
&attempt.candidate_id,
"transport_header_rules_apply_failed",
CandidateFailureDiagnostic::header_rules_apply_failed(
provider_api_format,
provider_api_format,
"standard_family_kiro_headers",
),
)
.await;
return None;

View File

@@ -12,6 +12,7 @@ mod family;
mod gemini;
mod normalize;
mod openai;
mod request_body_diagnostics;
pub(crate) use self::codex::apply_codex_openai_cli_special_headers;
pub(crate) use self::family::{
@@ -35,6 +36,9 @@ pub(crate) use self::openai::{
resolve_openai_chat_max_tokens, set_local_openai_chat_execution_exhausted_diagnostic,
value_as_u64,
};
pub(crate) use self::request_body_diagnostics::{
request_body_build_failure_extra_data, same_format_provider_request_body_failure_extra_data,
};
pub(crate) use crate::ai_pipeline::conversion::{
build_core_error_body_for_client_format, request_conversion_kind,
request_conversion_transport_supported, sync_chat_response_conversion_kind,

View File

@@ -12,7 +12,7 @@ use crate::ai_pipeline::planner::common::OPENAI_CHAT_STREAM_PLAN_KIND;
use crate::ai_pipeline::planner::standard::{
apply_codex_openai_cli_special_headers, 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,
build_local_openai_chat_upstream_url, request_body_build_failure_extra_data,
};
use crate::ai_pipeline::transport::apply_local_header_rules;
use crate::ai_pipeline::transport::auth::{
@@ -27,12 +27,16 @@ use crate::ai_pipeline::transport::kiro::{
use crate::ai_pipeline::transport::local_openai_chat_transport_unsupported_reason;
use crate::ai_pipeline::transport::vertex::uses_vertex_api_key_query_auth;
use crate::ai_pipeline::{
ConversionMode, ExecutionStrategy, GatewayProviderTransportSnapshot,
LocalResolvedOAuthRequestAuth,
CandidateFailureDiagnostic, ConversionMode, ExecutionStrategy,
GatewayProviderTransportSnapshot, LocalResolvedOAuthRequestAuth,
};
use crate::AppState;
use super::support::{mark_skipped_local_openai_chat_candidate, LocalOpenAiChatDecisionInput};
use super::support::{
mark_skipped_local_openai_chat_candidate,
mark_skipped_local_openai_chat_candidate_with_extra_data,
mark_skipped_local_openai_chat_candidate_with_failure_diagnostic, LocalOpenAiChatDecisionInput,
};
pub(crate) struct LocalOpenAiChatCandidatePayloadParts {
pub(super) auth_header: String,
@@ -118,21 +122,26 @@ pub(crate) async fn resolve_local_openai_chat_candidate_payload_parts(
upstream_is_stream,
transport.endpoint.body_rules.as_ref(),
) else {
mark_skipped_local_openai_chat_candidate(
mark_skipped_local_openai_chat_candidate_with_extra_data(
state,
input,
trace_id,
candidate,
candidate_index,
candidate_id,
"provider_request_body_missing",
"provider_request_body_build_failed",
request_body_build_failure_extra_data(
body_json,
"openai:chat",
provider_api_format,
),
)
.await;
return None;
};
let Some(upstream_url) = build_local_openai_chat_upstream_url(parts, transport) else {
mark_skipped_local_openai_chat_candidate(
mark_skipped_local_openai_chat_candidate_with_failure_diagnostic(
state,
input,
trace_id,
@@ -140,6 +149,11 @@ pub(crate) async fn resolve_local_openai_chat_candidate_payload_parts(
candidate_index,
candidate_id,
"upstream_url_missing",
CandidateFailureDiagnostic::upstream_url_missing(
"openai:chat",
provider_api_format,
"openai_chat_same_format_url",
),
)
.await;
return None;
@@ -159,7 +173,7 @@ pub(crate) async fn resolve_local_openai_chat_candidate_payload_parts(
&provider_request_body,
Some(body_json),
) {
mark_skipped_local_openai_chat_candidate(
mark_skipped_local_openai_chat_candidate_with_failure_diagnostic(
state,
input,
trace_id,
@@ -167,6 +181,11 @@ pub(crate) async fn resolve_local_openai_chat_candidate_payload_parts(
candidate_index,
candidate_id,
"transport_header_rules_apply_failed",
CandidateFailureDiagnostic::header_rules_apply_failed(
"openai:chat",
provider_api_format,
"openai_chat_same_format_headers",
),
)
.await;
return None;
@@ -343,14 +362,19 @@ pub(crate) async fn resolve_local_openai_chat_candidate_payload_parts(
},
Some(input.auth_context.api_key_id.as_str()),
) else {
mark_skipped_local_openai_chat_candidate(
mark_skipped_local_openai_chat_candidate_with_extra_data(
state,
input,
trace_id,
candidate,
candidate_index,
candidate_id,
"provider_request_body_missing",
"provider_request_body_build_failed",
request_body_build_failure_extra_data(
body_json,
"openai:chat",
provider_api_format.as_str(),
),
)
.await;
return None;
@@ -386,7 +410,7 @@ pub(crate) async fn resolve_local_openai_chat_candidate_payload_parts(
provider_api_format.as_str(),
upstream_is_stream,
) else {
mark_skipped_local_openai_chat_candidate(
mark_skipped_local_openai_chat_candidate_with_failure_diagnostic(
state,
input,
trace_id,
@@ -394,6 +418,11 @@ pub(crate) async fn resolve_local_openai_chat_candidate_payload_parts(
candidate_index,
candidate_id,
"upstream_url_missing",
CandidateFailureDiagnostic::upstream_url_missing(
"openai:chat",
provider_api_format.as_str(),
"openai_chat_cross_format_url",
),
)
.await;
return None;
@@ -430,7 +459,7 @@ pub(crate) async fn resolve_local_openai_chat_candidate_payload_parts(
&provider_request_body,
Some(body_json),
) {
mark_skipped_local_openai_chat_candidate(
mark_skipped_local_openai_chat_candidate_with_failure_diagnostic(
state,
input,
trace_id,
@@ -438,6 +467,11 @@ pub(crate) async fn resolve_local_openai_chat_candidate_payload_parts(
candidate_index,
candidate_id,
"transport_header_rules_apply_failed",
CandidateFailureDiagnostic::header_rules_apply_failed(
"openai:chat",
provider_api_format.as_str(),
"openai_chat_cross_format_headers",
),
)
.await;
return None;
@@ -522,14 +556,19 @@ async fn build_kiro_openai_chat_cross_format_payload_parts(
) {
Some(body) => body,
None => {
mark_skipped_local_openai_chat_candidate(
mark_skipped_local_openai_chat_candidate_with_failure_diagnostic(
state,
input,
trace_id,
candidate,
candidate_index,
candidate_id,
"provider_request_body_missing",
"provider_request_body_build_failed",
CandidateFailureDiagnostic::envelope_build_failed(
"openai:chat",
provider_api_format,
"openai_chat_kiro_envelope",
),
)
.await;
return None;
@@ -545,7 +584,7 @@ async fn build_kiro_openai_chat_cross_format_payload_parts(
) {
Some(url) => url,
None => {
mark_skipped_local_openai_chat_candidate(
mark_skipped_local_openai_chat_candidate_with_failure_diagnostic(
state,
input,
trace_id,
@@ -553,6 +592,11 @@ async fn build_kiro_openai_chat_cross_format_payload_parts(
candidate_index,
candidate_id,
"upstream_url_missing",
CandidateFailureDiagnostic::upstream_url_missing(
"openai:chat",
provider_api_format,
"openai_chat_kiro_url",
),
)
.await;
return None;
@@ -570,7 +614,7 @@ async fn build_kiro_openai_chat_cross_format_payload_parts(
}) {
Some(headers) => headers,
None => {
mark_skipped_local_openai_chat_candidate(
mark_skipped_local_openai_chat_candidate_with_failure_diagnostic(
state,
input,
trace_id,
@@ -578,6 +622,11 @@ async fn build_kiro_openai_chat_cross_format_payload_parts(
candidate_index,
candidate_id,
"transport_header_rules_apply_failed",
CandidateFailureDiagnostic::header_rules_apply_failed(
"openai:chat",
provider_api_format,
"openai_chat_kiro_headers",
),
)
.await;
return None;

View File

@@ -6,7 +6,8 @@ use crate::ai_pipeline::planner::candidate_eligibility::{
SkippedLocalExecutionCandidate,
};
use crate::ai_pipeline::planner::candidate_materialization::{
mark_skipped_local_execution_candidate,
mark_skipped_local_execution_candidate, mark_skipped_local_execution_candidate_with_extra_data,
mark_skipped_local_execution_candidate_with_failure_diagnostic,
persist_available_local_execution_candidates_with_context,
persist_skipped_local_execution_candidates_with_context,
remember_first_local_candidate_affinity,
@@ -19,6 +20,7 @@ use crate::ai_pipeline::planner::candidate_metadata::{
use crate::ai_pipeline::planner::materialization_policy::{
build_local_candidate_persistence_policy, LocalCandidatePersistencePolicyKind,
};
use crate::ai_pipeline::planner::CandidateFailureDiagnostic;
use crate::ai_pipeline::{ConversionMode, ExecutionStrategy, PlannerAppState};
use crate::AppState;
@@ -52,6 +54,66 @@ pub(crate) async fn mark_skipped_local_openai_chat_candidate(
.await;
}
#[allow(clippy::too_many_arguments)]
pub(crate) async fn mark_skipped_local_openai_chat_candidate_with_extra_data(
state: &AppState,
input: &LocalOpenAiChatDecisionInput,
trace_id: &str,
candidate: &SchedulerMinimalCandidateSelectionCandidate,
candidate_index: u32,
candidate_id: &str,
skip_reason: &'static str,
extra_data: Option<serde_json::Value>,
) {
let auth_context: &ExecutionRuntimeAuthContext = &input.auth_context;
let persistence_policy = build_local_candidate_persistence_policy(
auth_context,
input.required_capabilities.as_ref(),
LocalCandidatePersistencePolicyKind::OpenAiChatDecision,
);
mark_skipped_local_execution_candidate_with_extra_data(
state,
trace_id,
persistence_policy.skipped,
candidate,
candidate_index,
candidate_id,
skip_reason,
extra_data,
)
.await;
}
#[allow(clippy::too_many_arguments)]
pub(crate) async fn mark_skipped_local_openai_chat_candidate_with_failure_diagnostic(
state: &AppState,
input: &LocalOpenAiChatDecisionInput,
trace_id: &str,
candidate: &SchedulerMinimalCandidateSelectionCandidate,
candidate_index: u32,
candidate_id: &str,
skip_reason: &'static str,
diagnostic: CandidateFailureDiagnostic,
) {
let auth_context: &ExecutionRuntimeAuthContext = &input.auth_context;
let persistence_policy = build_local_candidate_persistence_policy(
auth_context,
input.required_capabilities.as_ref(),
LocalCandidatePersistencePolicyKind::OpenAiChatDecision,
);
mark_skipped_local_execution_candidate_with_failure_diagnostic(
state,
trace_id,
persistence_policy.skipped,
candidate,
candidate_index,
candidate_id,
skip_reason,
diagnostic,
)
.await;
}
pub(crate) async fn materialize_local_openai_chat_candidate_attempts(
state: &AppState,
trace_id: &str,

View File

@@ -14,7 +14,7 @@ use crate::ai_pipeline::planner::spec_metadata::local_openai_cli_spec_metadata;
use crate::ai_pipeline::planner::standard::{
apply_codex_openai_cli_special_headers, build_cross_format_openai_cli_request_body,
build_cross_format_openai_cli_upstream_url, build_local_openai_cli_request_body,
build_local_openai_cli_upstream_url,
build_local_openai_cli_upstream_url, request_body_build_failure_extra_data,
};
use crate::ai_pipeline::transport::antigravity::{
build_antigravity_safe_v1internal_request, build_antigravity_static_identity_headers,
@@ -34,13 +34,17 @@ use crate::ai_pipeline::transport::kiro::{
};
use crate::ai_pipeline::transport::local_standard_transport_unsupported_reason_with_network;
use crate::ai_pipeline::transport::vertex::uses_vertex_api_key_query_auth;
use crate::ai_pipeline::{ConversionMode, ExecutionStrategy};
use crate::ai_pipeline::{CandidateFailureDiagnostic, ConversionMode, ExecutionStrategy};
use crate::ai_pipeline::{
GatewayProviderTransportSnapshot, LocalResolvedOAuthRequestAuth, PlannerAppState,
};
use crate::AppState;
use super::support::{mark_skipped_local_openai_cli_candidate, LocalOpenAiCliDecisionInput};
use super::support::{
mark_skipped_local_openai_cli_candidate,
mark_skipped_local_openai_cli_candidate_with_extra_data,
mark_skipped_local_openai_cli_candidate_with_failure_diagnostic, LocalOpenAiCliDecisionInput,
};
use super::LocalOpenAiCliSpec;
const ANTIGRAVITY_ENVELOPE_NAME: &str = "antigravity:v1internal";
@@ -258,14 +262,19 @@ pub(crate) async fn resolve_local_openai_cli_candidate_payload_parts(
Some(input.auth_context.api_key_id.as_str()),
)
}) else {
mark_skipped_local_openai_cli_candidate(
mark_skipped_local_openai_cli_candidate_with_extra_data(
state,
input,
trace_id,
candidate,
candidate_index,
candidate_id,
"provider_request_body_missing",
"provider_request_body_build_failed",
request_body_build_failure_extra_data(
body_json,
spec_metadata.api_format,
provider_api_format,
),
)
.await;
return None;
@@ -304,14 +313,19 @@ pub(crate) async fn resolve_local_openai_cli_candidate_payload_parts(
) {
AntigravityRequestEnvelopeSupport::Supported(envelope) => envelope,
AntigravityRequestEnvelopeSupport::Unsupported(_) => {
mark_skipped_local_openai_cli_candidate(
mark_skipped_local_openai_cli_candidate_with_failure_diagnostic(
state,
input,
trace_id,
candidate,
candidate_index,
candidate_id,
"provider_request_body_missing",
"provider_request_body_build_failed",
CandidateFailureDiagnostic::envelope_build_failed(
spec_metadata.api_format,
provider_api_format,
"openai_cli_antigravity_envelope",
),
)
.await;
return None;
@@ -361,7 +375,7 @@ pub(crate) async fn resolve_local_openai_cli_candidate_payload_parts(
provider_api_format == "openai:compact",
)
}) else {
mark_skipped_local_openai_cli_candidate(
mark_skipped_local_openai_cli_candidate_with_failure_diagnostic(
state,
input,
trace_id,
@@ -369,6 +383,11 @@ pub(crate) async fn resolve_local_openai_cli_candidate_payload_parts(
candidate_index,
candidate_id,
"upstream_url_missing",
CandidateFailureDiagnostic::upstream_url_missing(
spec_metadata.api_format,
provider_api_format,
"openai_cli_url",
),
)
.await;
return None;
@@ -416,7 +435,7 @@ pub(crate) async fn resolve_local_openai_cli_candidate_payload_parts(
&provider_request_body,
Some(body_json),
) {
mark_skipped_local_openai_cli_candidate(
mark_skipped_local_openai_cli_candidate_with_failure_diagnostic(
state,
input,
trace_id,
@@ -424,6 +443,11 @@ pub(crate) async fn resolve_local_openai_cli_candidate_payload_parts(
candidate_index,
candidate_id,
"transport_header_rules_apply_failed",
CandidateFailureDiagnostic::header_rules_apply_failed(
spec_metadata.api_format,
provider_api_format,
"openai_cli_headers",
),
)
.await;
return None;
@@ -537,14 +561,19 @@ async fn build_kiro_openai_cli_payload_parts(
) {
Some(body) => body,
None => {
mark_skipped_local_openai_cli_candidate(
mark_skipped_local_openai_cli_candidate_with_failure_diagnostic(
state,
input,
trace_id,
candidate,
candidate_index,
candidate_id,
"provider_request_body_missing",
"provider_request_body_build_failed",
CandidateFailureDiagnostic::envelope_build_failed(
client_api_format,
provider_api_format,
"openai_cli_kiro_envelope",
),
)
.await;
return None;
@@ -560,7 +589,7 @@ async fn build_kiro_openai_cli_payload_parts(
) {
Some(url) => url,
None => {
mark_skipped_local_openai_cli_candidate(
mark_skipped_local_openai_cli_candidate_with_failure_diagnostic(
state,
input,
trace_id,
@@ -568,6 +597,11 @@ async fn build_kiro_openai_cli_payload_parts(
candidate_index,
candidate_id,
"upstream_url_missing",
CandidateFailureDiagnostic::upstream_url_missing(
client_api_format,
provider_api_format,
"openai_cli_kiro_url",
),
)
.await;
return None;
@@ -585,7 +619,7 @@ async fn build_kiro_openai_cli_payload_parts(
}) {
Some(headers) => headers,
None => {
mark_skipped_local_openai_cli_candidate(
mark_skipped_local_openai_cli_candidate_with_failure_diagnostic(
state,
input,
trace_id,
@@ -593,6 +627,11 @@ async fn build_kiro_openai_cli_payload_parts(
candidate_index,
candidate_id,
"transport_header_rules_apply_failed",
CandidateFailureDiagnostic::header_rules_apply_failed(
client_api_format,
provider_api_format,
"openai_cli_kiro_headers",
),
)
.await;
return None;

View File

@@ -10,7 +10,8 @@ use crate::ai_pipeline::planner::candidate_eligibility::{
SkippedLocalExecutionCandidate,
};
use crate::ai_pipeline::planner::candidate_materialization::{
mark_skipped_local_execution_candidate,
mark_skipped_local_execution_candidate, mark_skipped_local_execution_candidate_with_extra_data,
mark_skipped_local_execution_candidate_with_failure_diagnostic,
persist_available_local_execution_candidates_with_context,
persist_skipped_local_execution_candidates_with_context,
remember_first_local_candidate_affinity,
@@ -30,6 +31,7 @@ use crate::ai_pipeline::planner::materialization_policy::{
};
use crate::ai_pipeline::planner::runtime_miss::set_local_runtime_miss_diagnostic_reason;
use crate::ai_pipeline::planner::spec_metadata::local_openai_cli_spec_metadata;
use crate::ai_pipeline::planner::CandidateFailureDiagnostic;
use crate::ai_pipeline::PlannerAppState;
use crate::ai_pipeline::{
resolve_local_decision_execution_runtime_auth_context, ConversionMode, ExecutionStrategy,
@@ -361,3 +363,63 @@ pub(crate) async fn mark_skipped_local_openai_cli_candidate(
)
.await;
}
#[allow(clippy::too_many_arguments)]
pub(crate) async fn mark_skipped_local_openai_cli_candidate_with_extra_data(
state: &AppState,
input: &LocalOpenAiCliDecisionInput,
trace_id: &str,
candidate: &SchedulerMinimalCandidateSelectionCandidate,
candidate_index: u32,
candidate_id: &str,
skip_reason: &'static str,
extra_data: Option<serde_json::Value>,
) {
let auth_context: &ExecutionRuntimeAuthContext = &input.auth_context;
let persistence_policy = build_local_candidate_persistence_policy(
auth_context,
input.required_capabilities.as_ref(),
LocalCandidatePersistencePolicyKind::OpenAiCliDecision,
);
mark_skipped_local_execution_candidate_with_extra_data(
state,
trace_id,
persistence_policy.skipped,
candidate,
candidate_index,
candidate_id,
skip_reason,
extra_data,
)
.await;
}
#[allow(clippy::too_many_arguments)]
pub(crate) async fn mark_skipped_local_openai_cli_candidate_with_failure_diagnostic(
state: &AppState,
input: &LocalOpenAiCliDecisionInput,
trace_id: &str,
candidate: &SchedulerMinimalCandidateSelectionCandidate,
candidate_index: u32,
candidate_id: &str,
skip_reason: &'static str,
diagnostic: CandidateFailureDiagnostic,
) {
let auth_context: &ExecutionRuntimeAuthContext = &input.auth_context;
let persistence_policy = build_local_candidate_persistence_policy(
auth_context,
input.required_capabilities.as_ref(),
LocalCandidatePersistencePolicyKind::OpenAiCliDecision,
);
mark_skipped_local_execution_candidate_with_failure_diagnostic(
state,
trace_id,
persistence_policy.skipped,
candidate,
candidate_index,
candidate_id,
skip_reason,
diagnostic,
)
.await;
}

View File

@@ -0,0 +1,736 @@
use serde_json::Value;
use crate::ai_pipeline::planner::{CandidateFailureDiagnostic, CandidateFailureDiagnosticKind};
pub(crate) fn request_body_build_failure_extra_data(
body_json: &Value,
client_api_format: &str,
provider_api_format: &str,
) -> Option<Value> {
let diagnostic =
diagnose_request_body_build_failure(body_json, client_api_format, provider_api_format)?;
Some(
diagnostic
.formats(client_api_format, provider_api_format)
.source(request_body_build_source(
client_api_format,
provider_api_format,
))
.to_extra_data(),
)
}
pub(crate) fn same_format_provider_request_body_failure_extra_data(
body_json: &Value,
provider_api_format: &str,
body_rules: Option<&Value>,
context: &str,
) -> Option<Value> {
let diagnostic =
diagnose_same_format_provider_request_body_failure(body_json, body_rules, context)?;
Some(
diagnostic
.formats(provider_api_format, provider_api_format)
.source(context)
.to_extra_data(),
)
}
type RequestBodyBuildDiagnostic = CandidateFailureDiagnostic;
fn diagnose_request_body_build_failure(
body_json: &Value,
client_api_format: &str,
provider_api_format: &str,
) -> Option<RequestBodyBuildDiagnostic> {
if !body_json.is_object() {
return Some(diagnostic("$", "请求体必须是 JSON object"));
}
if client_api_format == "openai:cli" {
if let Some(diagnostic) = diagnose_openai_cli_request(body_json) {
return Some(diagnostic);
}
return Some(diagnostic(
"$",
"OpenAI CLI 请求体初步结构检查通过;失败可能发生在后续跨格式转换或 Body 规则应用",
));
}
if client_api_format == "openai:chat"
&& (provider_api_format.starts_with("claude:")
|| provider_api_format.starts_with("gemini:"))
{
return diagnose_openai_chat_cross_format_request(body_json, provider_api_format);
}
Some(diagnostic(
"$",
"请求体转换失败;当前转换器未返回更细的字段路径",
))
}
fn diagnose_same_format_provider_request_body_failure(
body_json: &Value,
body_rules: Option<&Value>,
context: &str,
) -> Option<RequestBodyBuildDiagnostic> {
if !body_json.is_object() {
return Some(diagnostic("$", "反代请求体必须是 JSON object"));
}
if body_rules.is_some_and(|rules| !rules.is_array()) {
return Some(diagnostic(
"$.endpoint.body_rules",
"Endpoint Body 规则必须是数组,本地反代无法应用该配置",
));
}
match context {
"kiro_envelope" => Some(diagnostic(
"$",
"Kiro 反代请求体包装失败;请检查 Kiro auth_config 与 Endpoint Body 规则",
)),
"antigravity_envelope" => Some(diagnostic(
"$",
"Antigravity 反代请求体包装失败;请检查请求体是否满足该传输封装要求",
)),
_ => Some(diagnostic(
"$",
"反代请求体构建失败;当前路径未返回更细的字段信息",
)),
}
}
fn diagnose_openai_chat_cross_format_request(
body_json: &Value,
provider_api_format: &str,
) -> Option<RequestBodyBuildDiagnostic> {
let request = body_json.as_object()?;
if let Some(messages) = request.get("messages") {
let Some(messages) = messages.as_array() else {
return Some(diagnostic(
"$.messages",
"OpenAI Chat 的 messages 必须是数组",
));
};
for (message_index, message) in messages.iter().enumerate() {
let Some(message_object) = message.as_object() else {
return Some(diagnostic(
format!("$.messages[{message_index}]"),
"message 必须是 object",
));
};
let role = message_object
.get("role")
.and_then(Value::as_str)
.unwrap_or_default()
.trim()
.to_ascii_lowercase();
match role.as_str() {
"system" | "developer" => {
if let Some(diagnostic) = diagnose_openai_text_content(
message_object.get("content"),
format!("$.messages[{message_index}].content"),
) {
return Some(diagnostic);
}
}
"user" | "assistant" => {
if let Some(diagnostic) = diagnose_openai_content_blocks(
message_object.get("content"),
format!("$.messages[{message_index}].content"),
role.as_str(),
) {
return Some(diagnostic);
}
if role == "assistant" {
if let Some(diagnostic) = diagnose_openai_assistant_tool_calls(
message_object.get("tool_calls"),
format!("$.messages[{message_index}].tool_calls"),
) {
return Some(diagnostic);
}
}
}
"tool" => {
let valid_tool_call_id = message_object
.get("tool_call_id")
.and_then(Value::as_str)
.map(str::trim)
.is_some_and(|value| !value.is_empty());
if !valid_tool_call_id {
return Some(diagnostic(
format!("$.messages[{message_index}].tool_call_id"),
"tool 消息必须包含非空 tool_call_id",
));
}
}
_ => {}
}
}
}
if let Some(diagnostic) = diagnose_openai_tools(request.get("tools"), provider_api_format) {
return Some(diagnostic);
}
diagnose_openai_tool_choice(request.get("tool_choice"))
}
fn diagnose_openai_cli_request(body_json: &Value) -> Option<RequestBodyBuildDiagnostic> {
let request = body_json.as_object()?;
if let Some(diagnostic) =
diagnose_openai_cli_text_content(request.get("instructions"), "$.instructions".to_string())
{
return Some(diagnostic);
}
if let Some(diagnostic) = diagnose_openai_cli_input(request.get("input")) {
return Some(diagnostic);
}
if let Some(diagnostic) = diagnose_openai_cli_tools(request.get("tools")) {
return Some(diagnostic);
}
diagnose_openai_cli_tool_choice(request.get("tool_choice"))
}
fn diagnose_openai_cli_input(input: Option<&Value>) -> Option<RequestBodyBuildDiagnostic> {
let Some(input) = input else {
return None;
};
match input {
Value::Null | Value::String(_) => None,
Value::Array(items) => {
for (item_index, item) in items.iter().enumerate() {
if item.is_string() {
continue;
}
let item_path = format!("$.input[{item_index}]");
let Some(item_object) = item.as_object() else {
return Some(diagnostic(
item_path,
"OpenAI CLI input 数组项必须是 string 或 object",
));
};
let item_type = item_object
.get("type")
.and_then(Value::as_str)
.unwrap_or("message")
.trim()
.to_ascii_lowercase();
match item_type.as_str() {
"message" => {
let role = item_object
.get("role")
.and_then(Value::as_str)
.unwrap_or("user")
.trim()
.to_ascii_lowercase();
if role == "system" || role == "developer" {
if let Some(diagnostic) = diagnose_openai_cli_text_content(
item_object.get("content"),
format!("{item_path}.content"),
) {
return Some(diagnostic);
}
} else if let Some(diagnostic) = diagnose_openai_cli_message_content(
item_object.get("content"),
format!("{item_path}.content"),
) {
return Some(diagnostic);
}
}
"function_call" => {
let valid_name = item_object
.get("name")
.and_then(Value::as_str)
.map(str::trim)
.is_some_and(|value| !value.is_empty());
if !valid_name {
return Some(diagnostic(
format!("{item_path}.name"),
"function_call 必须包含非空 name",
));
}
}
_ => {}
}
}
None
}
_ => Some(diagnostic(
"$.input",
"OpenAI CLI input 必须是 string、array 或 null",
)),
}
}
fn diagnose_openai_cli_text_content(
content: Option<&Value>,
path: String,
) -> Option<RequestBodyBuildDiagnostic> {
match content {
None | Some(Value::Null) | Some(Value::String(_)) => None,
Some(Value::Array(parts)) => {
for (part_index, part) in parts.iter().enumerate() {
if !part.is_object() {
return Some(diagnostic(
format!("{path}[{part_index}]"),
"文本 content 数组项必须是 object",
));
}
}
None
}
Some(_) => Some(diagnostic(
path,
"文本 content 必须是 string、array 或 null",
)),
}
}
fn diagnose_openai_cli_message_content(
content: Option<&Value>,
path: String,
) -> Option<RequestBodyBuildDiagnostic> {
match content {
None | Some(Value::Null) | Some(Value::String(_)) => None,
Some(Value::Array(parts)) => {
for (part_index, part) in parts.iter().enumerate() {
let part_path = format!("{path}[{part_index}]");
let Some(part_object) = part.as_object() else {
return Some(diagnostic(part_path, "message content 数组项必须是 object"));
};
let part_type = part_object
.get("type")
.and_then(Value::as_str)
.unwrap_or_default()
.trim()
.to_ascii_lowercase();
if matches!(
part_type.as_str(),
"input_image" | "output_image" | "image_url"
) && image_part_url(part_object).is_none()
{
return Some(diagnostic(
part_path,
"图片 content 缺少 image_url/url无法规范化为 OpenAI Chat 图片内容",
));
}
}
None
}
Some(_) => None,
}
}
fn diagnose_openai_text_content(
content: Option<&Value>,
path: String,
) -> Option<RequestBodyBuildDiagnostic> {
match content {
None | Some(Value::Null) | Some(Value::String(_)) => None,
Some(Value::Array(parts)) => {
for (part_index, part) in parts.iter().enumerate() {
if !part.is_object() {
return Some(diagnostic(
format!("{path}[{part_index}]"),
"content 数组项必须是 object",
));
}
}
None
}
Some(_) => Some(diagnostic(path, "content 必须是 string、array 或 null")),
}
}
fn diagnose_openai_content_blocks(
content: Option<&Value>,
path: String,
role: &str,
) -> Option<RequestBodyBuildDiagnostic> {
match content {
None | Some(Value::Null) | Some(Value::String(_)) => None,
Some(Value::Array(parts)) => {
for (part_index, part) in parts.iter().enumerate() {
let part_path = format!("{path}[{part_index}]");
let Some(part_object) = part.as_object() else {
return Some(diagnostic(part_path, "content 数组项必须是 object"));
};
let part_type = part_object
.get("type")
.and_then(Value::as_str)
.unwrap_or_default();
if matches!(part_type, "image_url" | "input_image" | "output_image")
&& role == "user"
&& image_part_url(part_object).is_none()
{
return Some(diagnostic(
part_path,
"图片 content 缺少 image_url/url无法转换为 Claude image block",
));
}
}
None
}
Some(_) => Some(diagnostic(path, "content 必须是 string、array 或 null")),
}
}
fn diagnose_openai_assistant_tool_calls(
tool_calls: Option<&Value>,
path: String,
) -> Option<RequestBodyBuildDiagnostic> {
let Some(tool_calls) = tool_calls else {
return None;
};
let Some(tool_calls) = tool_calls.as_array() else {
return Some(diagnostic(path, "assistant.tool_calls 必须是数组"));
};
for (tool_call_index, tool_call) in tool_calls.iter().enumerate() {
let tool_call_path = format!("{path}[{tool_call_index}]");
let Some(tool_call_object) = tool_call.as_object() else {
return Some(diagnostic(tool_call_path, "tool_call 必须是 object"));
};
let Some(function) = tool_call_object.get("function").and_then(Value::as_object) else {
return Some(diagnostic(
format!("{tool_call_path}.function"),
"tool_call 必须包含 function object",
));
};
let valid_name = function
.get("name")
.and_then(Value::as_str)
.map(str::trim)
.is_some_and(|value| !value.is_empty());
if !valid_name {
return Some(diagnostic(
format!("{tool_call_path}.function.name"),
"tool_call.function.name 必须是非空字符串",
));
}
}
None
}
fn diagnose_openai_tools(
tools: Option<&Value>,
provider_api_format: &str,
) -> Option<RequestBodyBuildDiagnostic> {
let Some(tools) = tools else {
return None;
};
let Some(tools) = tools.as_array() else {
return Some(diagnostic("$.tools", "OpenAI Chat 的 tools 必须是数组"));
};
for (tool_index, tool) in tools.iter().enumerate() {
let tool_path = format!("$.tools[{tool_index}]");
let Some(tool_object) = tool.as_object() else {
return Some(diagnostic(tool_path, "tool 必须是 object"));
};
if tool_object
.get("type")
.and_then(Value::as_str)
.is_some_and(|value| value != "function")
{
continue;
}
let Some(function) = tool_object.get("function").and_then(Value::as_object) else {
let native_tool_hint = if provider_api_format.starts_with("claude:") {
";如果这是 Claude 原生 tool请改为 OpenAI function tool 格式"
} else if provider_api_format.starts_with("gemini:") {
";如果这是 Gemini 原生 tool请改为 OpenAI function tool 格式"
} else {
""
};
return Some(diagnostic(
format!("{tool_path}.function"),
format!("OpenAI tool 必须包含 function object{native_tool_hint}"),
));
};
let valid_name = function
.get("name")
.and_then(Value::as_str)
.map(str::trim)
.is_some_and(|value| !value.is_empty());
if !valid_name {
return Some(diagnostic(
format!("{tool_path}.function.name"),
"OpenAI tool 的 function.name 必须是非空字符串",
));
}
}
None
}
fn diagnose_openai_cli_tools(tools: Option<&Value>) -> Option<RequestBodyBuildDiagnostic> {
let Some(tools) = tools else {
return None;
};
let Some(tool_values) = tools.as_array() else {
return None;
};
for (tool_index, tool) in tool_values.iter().enumerate() {
let tool_path = format!("$.tools[{tool_index}]");
let Some(tool_object) = tool.as_object() else {
return Some(diagnostic(tool_path, "OpenAI CLI tool 必须是 object"));
};
let tool_type = tool_object
.get("type")
.and_then(Value::as_str)
.unwrap_or("function")
.trim()
.to_ascii_lowercase();
if tool_type.starts_with("web_search")
|| tool_object.get("function").is_some()
|| tool_type != "function"
{
continue;
}
let valid_name = tool_object
.get("name")
.and_then(Value::as_str)
.map(str::trim)
.is_some_and(|value| !value.is_empty());
if !valid_name {
return Some(diagnostic(
format!("{tool_path}.name"),
"OpenAI CLI function tool 必须包含非空 name",
));
}
}
None
}
fn diagnose_openai_cli_tool_choice(
tool_choice: Option<&Value>,
) -> Option<RequestBodyBuildDiagnostic> {
let Some(Value::Object(object)) = tool_choice else {
return None;
};
let is_cli_function_choice = object.get("function").is_none()
&& object
.get("type")
.and_then(Value::as_str)
.is_some_and(|value| value.eq_ignore_ascii_case("function"));
if !is_cli_function_choice {
return None;
}
let valid_name = object
.get("name")
.and_then(Value::as_str)
.map(str::trim)
.is_some_and(|value| !value.is_empty());
if valid_name {
None
} else {
Some(diagnostic(
"$.tool_choice.name",
"OpenAI CLI tool_choice 指定 function 时必须包含非空 name",
))
}
}
fn diagnose_openai_tool_choice(tool_choice: Option<&Value>) -> Option<RequestBodyBuildDiagnostic> {
let Some(Value::Object(object)) = tool_choice else {
return None;
};
let valid_name = object
.get("function")
.and_then(Value::as_object)
.and_then(|function| function.get("name"))
.and_then(Value::as_str)
.map(str::trim)
.is_some_and(|value| !value.is_empty());
if valid_name {
None
} else {
Some(diagnostic(
"$.tool_choice.function.name",
"tool_choice 指定具体工具时必须包含非空 function.name",
))
}
}
fn image_part_url(part_object: &serde_json::Map<String, Value>) -> Option<&str> {
part_object
.get("image_url")
.and_then(|value| {
value.as_str().or_else(|| {
value
.as_object()
.and_then(|object| object.get("url"))
.and_then(Value::as_str)
})
})
.or_else(|| part_object.get("url").and_then(Value::as_str))
.map(str::trim)
.filter(|value| !value.is_empty())
}
fn diagnostic(path: impl Into<String>, message: impl Into<String>) -> RequestBodyBuildDiagnostic {
CandidateFailureDiagnostic::new(
CandidateFailureDiagnosticKind::RequestBodyBuild,
path,
message,
)
}
fn request_body_build_source(client_api_format: &str, provider_api_format: &str) -> String {
format!("{client_api_format}_to_{provider_api_format}")
}
#[cfg(test)]
mod tests {
use serde_json::json;
use super::request_body_build_failure_extra_data;
#[test]
fn openai_chat_to_claude_reports_claude_native_tool_shape() {
let body = json!({
"model": "gpt-5.4",
"messages": [{ "role": "user", "content": "hello" }],
"tools": [{
"name": "read_file",
"description": "Read a file",
"input_schema": { "type": "object" }
}]
});
let diagnostic = request_body_build_failure_extra_data(&body, "openai:chat", "claude:chat")
.expect("diagnostic");
assert_eq!(
diagnostic["request_body_build_error"]["path"],
"$.tools[0].function"
);
assert_eq!(
diagnostic["failure_diagnostic"]["kind"],
"request_body_build"
);
assert_eq!(
diagnostic["failure_diagnostic"]["source"],
"openai:chat_to_claude:chat"
);
assert!(diagnostic["request_body_build_error"]["message"]
.as_str()
.expect("message")
.contains("Claude 原生 tool"));
}
#[test]
fn openai_chat_to_claude_reports_invalid_message_content_part() {
let body = json!({
"model": "gpt-5.4",
"messages": [{
"role": "user",
"content": ["not-an-object"]
}]
});
let diagnostic = request_body_build_failure_extra_data(&body, "openai:chat", "claude:chat")
.expect("diagnostic");
assert_eq!(
diagnostic["request_body_build_error"]["path"],
"$.messages[0].content[0]"
);
}
#[test]
fn openai_chat_to_gemini_reports_gemini_native_tool_shape() {
let body = json!({
"model": "gpt-5.4",
"messages": [{ "role": "user", "content": "hello" }],
"tools": [{
"functionDeclarations": [{
"name": "search",
"parameters": { "type": "object" }
}]
}]
});
let diagnostic = request_body_build_failure_extra_data(&body, "openai:chat", "gemini:chat")
.expect("diagnostic");
assert_eq!(
diagnostic["request_body_build_error"]["path"],
"$.tools[0].function"
);
assert!(diagnostic["request_body_build_error"]["message"]
.as_str()
.expect("message")
.contains("Gemini 原生 tool"));
}
#[test]
fn openai_cli_reports_invalid_function_call_name() {
let body = json!({
"model": "gpt-5.4",
"input": [{
"type": "function_call",
"arguments": "{}"
}]
});
let diagnostic = request_body_build_failure_extra_data(&body, "openai:cli", "claude:chat")
.expect("diagnostic");
assert_eq!(
diagnostic["request_body_build_error"]["path"],
"$.input[0].name"
);
}
#[test]
fn openai_cli_reports_invalid_tool_choice_name() {
let body = json!({
"model": "gpt-5.4",
"input": "hello",
"tool_choice": { "type": "function" }
});
let diagnostic = request_body_build_failure_extra_data(&body, "openai:cli", "gemini:chat")
.expect("diagnostic");
assert_eq!(
diagnostic["request_body_build_error"]["path"],
"$.tool_choice.name"
);
}
#[test]
fn same_format_provider_reports_non_object_body() {
let diagnostic = super::same_format_provider_request_body_failure_extra_data(
&json!("raw"),
"openai:chat",
None,
"same_format",
)
.expect("diagnostic");
assert_eq!(diagnostic["request_body_build_error"]["path"], "$");
assert!(diagnostic["request_body_build_error"]["message"]
.as_str()
.expect("message")
.contains("反代请求体必须是 JSON object"));
}
#[test]
fn same_format_provider_reports_invalid_body_rules_shape() {
let diagnostic = super::same_format_provider_request_body_failure_extra_data(
&json!({ "model": "gpt-5.4" }),
"openai:chat",
Some(&json!({ "action": "set" })),
"same_format",
)
.expect("diagnostic");
assert_eq!(
diagnostic["request_body_build_error"]["path"],
"$.endpoint.body_rules"
);
}
}