fix(kiro): 修复 Claude CLI 跨格式 Responses 流式转换 (#329)

This commit is contained in:
Entropy.Xu
2026-04-24 20:23:02 +08:00
committed by GitHub
parent 657fcd595c
commit 343345529d
10 changed files with 956 additions and 97 deletions

View File

@@ -12,6 +12,10 @@ enum RewriteMode {
OpenAiImage(OpenAiImageStreamState),
Standard(StreamingStandardConversionState),
KiroToClaudeCli(KiroToClaudeCliStreamState),
KiroToClaudeCliThenStandard {
kiro: KiroToClaudeCliStreamState,
standard: StreamingStandardConversionState,
},
}
pub(crate) struct LocalStreamRewriter<'a> {
@@ -35,6 +39,12 @@ pub(crate) fn maybe_build_local_stream_rewriter<'a>(
FinalizeStreamRewriteMode::KiroToClaudeCli => {
RewriteMode::KiroToClaudeCli(KiroToClaudeCliStreamState::new(report_context))
}
FinalizeStreamRewriteMode::KiroToClaudeCliThenStandard => {
RewriteMode::KiroToClaudeCliThenStandard {
kiro: KiroToClaudeCliStreamState::new(report_context),
standard: StreamingStandardConversionState::default(),
}
}
};
Some(LocalStreamRewriter {
@@ -52,6 +62,10 @@ impl LocalStreamRewriter<'_> {
if let RewriteMode::KiroToClaudeCli(state) = &mut self.mode {
return state.push_chunk(self.report_context, chunk);
}
if let RewriteMode::KiroToClaudeCliThenStandard { kiro, standard } = &mut self.mode {
let claude_bytes = kiro.push_chunk(self.report_context, chunk)?;
return transform_standard_bytes(standard, self.report_context, claude_bytes);
}
self.buffered.extend_from_slice(chunk);
let mut output = Vec::new();
while let Some(line_end) = self.buffered.iter().position(|byte| *byte == b'\n') {
@@ -68,11 +82,21 @@ impl LocalStreamRewriter<'_> {
if let RewriteMode::KiroToClaudeCli(state) = &mut self.mode {
return state.finish(self.report_context);
}
if let RewriteMode::KiroToClaudeCliThenStandard { kiro, standard } = &mut self.mode {
let mut output = transform_standard_bytes(
standard,
self.report_context,
kiro.finish(self.report_context)?,
)?;
output.extend(standard.finish(self.report_context)?);
return Ok(output);
}
if self.buffered.is_empty() {
match &mut self.mode {
RewriteMode::Standard(state) => return state.finish(self.report_context),
RewriteMode::OpenAiImage(_) => {}
RewriteMode::KiroToClaudeCli(_) => {}
RewriteMode::KiroToClaudeCliThenStandard { .. } => {}
RewriteMode::EnvelopeUnwrap => {}
}
return Ok(Vec::new());
@@ -85,6 +109,7 @@ impl LocalStreamRewriter<'_> {
}
RewriteMode::OpenAiImage(_) => {}
RewriteMode::KiroToClaudeCli(_) => {}
RewriteMode::KiroToClaudeCliThenStandard { .. } => {}
RewriteMode::EnvelopeUnwrap => {}
}
Ok(output)
@@ -97,10 +122,26 @@ impl LocalStreamRewriter<'_> {
RewriteMode::OpenAiImage(_) => Ok(Vec::new()),
RewriteMode::Standard(state) => state.transform_line(self.report_context, line),
RewriteMode::KiroToClaudeCli(_) => Ok(Vec::new()),
RewriteMode::KiroToClaudeCliThenStandard { .. } => Ok(Vec::new()),
}
}
}
fn transform_standard_bytes(
standard: &mut StreamingStandardConversionState,
report_context: &Value,
bytes: Vec<u8>,
) -> Result<Vec<u8>, GatewayError> {
if bytes.is_empty() {
return Ok(Vec::new());
}
let mut output = Vec::new();
for line in bytes.split_inclusive(|byte| *byte == b'\n') {
output.extend(standard.transform_line(report_context, line.to_vec())?);
}
Ok(output)
}
#[derive(Default)]
struct OpenAiImageStreamState {
buffered: Vec<u8>,

View File

@@ -52,6 +52,12 @@ pub(super) async fn maybe_build_local_standard_decision_payload_for_candidate(
{
extra_fields.insert("proxy".to_string(), proxy_value);
}
if let Some(envelope_name) = resolved.envelope_name {
extra_fields.insert(
"envelope_name".to_string(),
serde_json::Value::String(envelope_name.to_string()),
);
}
let report_context = append_local_failover_policy_to_value(
append_execution_contract_fields_to_value(
build_local_execution_report_context(LocalExecutionReportContextParts {
@@ -83,7 +89,7 @@ pub(super) async fn maybe_build_local_standard_decision_payload_for_candidate(
.and_then(serde_json::Value::as_bool)
.unwrap_or(false),
upstream_is_stream: resolved.upstream_is_stream,
has_envelope: false,
has_envelope: resolved.envelope_name.is_some(),
needs_conversion: true,
extra_fields,
}),
@@ -105,6 +111,7 @@ pub(super) async fn maybe_build_local_standard_decision_payload_for_candidate(
provider_request_headers,
upstream_url,
upstream_is_stream,
envelope_name: _,
transport,
} = resolved;

View File

@@ -13,8 +13,12 @@ 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,
};
use crate::ai_pipeline::transport::kiro::{
build_kiro_provider_headers, build_kiro_provider_request_body, KiroProviderHeadersInput,
KiroRequestAuth, KIRO_ENVELOPE_NAME,
};
use crate::ai_pipeline::transport::vertex::uses_vertex_api_key_query_auth;
use crate::ai_pipeline::GatewayProviderTransportSnapshot;
use crate::ai_pipeline::{GatewayProviderTransportSnapshot, LocalResolvedOAuthRequestAuth};
use crate::AppState;
use super::payload::mark_skipped_local_standard_candidate;
@@ -29,6 +33,7 @@ pub(crate) struct LocalStandardCandidatePayloadParts {
pub(super) provider_request_headers: BTreeMap<String, String>,
pub(super) upstream_url: String,
pub(super) upstream_is_stream: bool,
pub(super) envelope_name: Option<&'static str>,
pub(super) transport: Arc<GatewayProviderTransportSnapshot>,
}
@@ -46,6 +51,12 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
let candidate = &attempt.eligible.candidate;
let transport = &attempt.eligible.transport;
let provider_api_format = attempt.eligible.provider_api_format.as_str();
let is_kiro_claude_cli = transport
.provider
.provider_type
.trim()
.eq_ignore_ascii_case("kiro")
&& provider_api_format.eq_ignore_ascii_case("claude:cli");
let Some(conversion_kind) = crate::ai_pipeline::conversion::request_conversion_kind(
spec_metadata.api_format,
provider_api_format,
@@ -72,32 +83,90 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
return None;
}
let prepared_candidate = match prepare_header_authenticated_candidate(
planner_state,
transport,
candidate,
crate::ai_pipeline::conversion::request_conversion_direct_auth(transport, conversion_kind),
OauthPreparationContext {
trace_id,
api_format: provider_api_format,
operation: "standard_family_cross_format",
},
)
.await
{
Ok(prepared) => prepared,
Err(skip_reason) => {
mark_skipped_local_standard_candidate(
state,
input,
trace_id,
let oauth_context = OauthPreparationContext {
trace_id,
api_format: provider_api_format,
operation: "standard_family_cross_format",
};
let kiro_auth = if is_kiro_claude_cli {
match crate::ai_pipeline::planner::candidate_preparation::resolve_candidate_oauth_auth(
planner_state,
transport,
oauth_context,
)
.await
{
Some(LocalResolvedOAuthRequestAuth::Kiro(auth)) => Some(auth),
_ => {
mark_skipped_local_standard_candidate(
state,
input,
trace_id,
candidate,
attempt.candidate_index,
&attempt.candidate_id,
"transport_auth_unavailable",
)
.await;
return None;
}
}
} else {
None
};
let prepared_candidate = if let Some(kiro_auth) = kiro_auth.as_ref() {
let mapped_model =
match crate::ai_pipeline::planner::candidate_preparation::resolve_candidate_mapped_model(
candidate,
attempt.candidate_index,
&attempt.candidate_id,
skip_reason,
)
.await;
return None;
) {
Ok(mapped_model) => mapped_model,
Err(skip_reason) => {
mark_skipped_local_standard_candidate(
state,
input,
trace_id,
candidate,
attempt.candidate_index,
&attempt.candidate_id,
skip_reason,
)
.await;
return None;
}
};
crate::ai_pipeline::planner::candidate_preparation::PreparedHeaderAuthenticatedCandidate {
auth_header: kiro_auth.name.to_string(),
auth_value: kiro_auth.value.clone(),
mapped_model,
}
} else {
match prepare_header_authenticated_candidate(
planner_state,
transport,
candidate,
crate::ai_pipeline::conversion::request_conversion_direct_auth(
transport,
conversion_kind,
),
oauth_context,
)
.await
{
Ok(prepared) => prepared,
Err(skip_reason) => {
mark_skipped_local_standard_candidate(
state,
input,
trace_id,
candidate,
attempt.candidate_index,
&attempt.candidate_id,
skip_reason,
)
.await;
return None;
}
}
};
@@ -115,7 +184,11 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
provider_api_format,
parts.uri.path(),
upstream_is_stream,
transport.endpoint.body_rules.as_ref(),
if is_kiro_claude_cli {
None
} else {
transport.endpoint.body_rules.as_ref()
},
Some(input.auth_context.api_key_id.as_str()),
) {
Some(body) => body,
@@ -134,6 +207,26 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
}
};
if let Some(kiro_auth) = kiro_auth.as_ref() {
return build_kiro_cross_format_payload_parts(
state,
parts,
trace_id,
body_json,
input,
attempt,
transport,
provider_api_format,
prepared_candidate.mapped_model,
prepared_candidate.auth_header,
prepared_candidate.auth_value,
provider_request_body,
upstream_is_stream,
kiro_auth,
)
.await;
}
let upstream_url = match crate::ai_pipeline::planner::standard::build_standard_upstream_url(
parts,
transport,
@@ -237,6 +330,109 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
provider_request_headers,
upstream_url,
upstream_is_stream,
envelope_name: None,
transport: Arc::clone(transport),
})
}
#[allow(clippy::too_many_arguments)]
async fn build_kiro_cross_format_payload_parts(
state: &AppState,
parts: &http::request::Parts,
trace_id: &str,
original_body_json: &serde_json::Value,
input: &LocalStandardDecisionInput,
attempt: &LocalStandardCandidateAttempt,
transport: &Arc<GatewayProviderTransportSnapshot>,
provider_api_format: &str,
mapped_model: String,
auth_header: String,
auth_value: String,
claude_request_body: Value,
upstream_is_stream: bool,
kiro_auth: &KiroRequestAuth,
) -> Option<LocalStandardCandidatePayloadParts> {
let candidate = &attempt.eligible.candidate;
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(body) => body,
None => {
mark_skipped_local_standard_candidate(
state,
input,
trace_id,
candidate,
attempt.candidate_index,
&attempt.candidate_id,
"provider_request_body_missing",
)
.await;
return None;
}
};
let upstream_url = match crate::ai_pipeline::build_provider_transport_request_url(
transport,
provider_api_format,
Some(&mapped_model),
upstream_is_stream,
parts.uri.query(),
Some(kiro_auth.auth_config.effective_api_region()),
) {
Some(url) => url,
None => {
mark_skipped_local_standard_candidate(
state,
input,
trace_id,
candidate,
attempt.candidate_index,
&attempt.candidate_id,
"upstream_url_missing",
)
.await;
return None;
}
};
let provider_request_headers = match build_kiro_provider_headers(KiroProviderHeadersInput {
headers: &parts.headers,
provider_request_body: &provider_request_body,
original_request_body: original_body_json,
header_rules: transport.endpoint.header_rules.as_ref(),
auth_header: &auth_header,
auth_value: &auth_value,
auth_config: &kiro_auth.auth_config,
machine_id: kiro_auth.machine_id.as_str(),
}) {
Some(headers) => headers,
None => {
mark_skipped_local_standard_candidate(
state,
input,
trace_id,
candidate,
attempt.candidate_index,
&attempt.candidate_id,
"transport_header_rules_apply_failed",
)
.await;
return None;
}
};
Some(LocalStandardCandidatePayloadParts {
auth_header,
auth_value,
mapped_model,
provider_api_format: provider_api_format.to_string(),
provider_request_body,
provider_request_headers,
upstream_url,
upstream_is_stream,
envelope_name: Some(KIRO_ENVELOPE_NAME),
transport: Arc::clone(transport),
})
}

View File

@@ -69,6 +69,12 @@ pub(crate) async fn maybe_build_local_openai_chat_decision_payload_for_candidate
{
extra_fields.insert("proxy".to_string(), proxy_value);
}
if let Some(envelope_name) = resolved.envelope_name {
extra_fields.insert(
"envelope_name".to_string(),
serde_json::Value::String(envelope_name.to_string()),
);
}
let report_context = append_local_failover_policy_to_value(
append_execution_contract_fields_to_value(
build_local_execution_report_context(LocalExecutionReportContextParts {
@@ -100,7 +106,7 @@ pub(crate) async fn maybe_build_local_openai_chat_decision_payload_for_candidate
.and_then(serde_json::Value::as_bool)
.unwrap_or(false),
upstream_is_stream,
has_envelope: false,
has_envelope: resolved.envelope_name.is_some(),
needs_conversion: matches!(
resolved.conversion_mode,
crate::ai_pipeline::ConversionMode::Bidirectional
@@ -125,6 +131,7 @@ pub(crate) async fn maybe_build_local_openai_chat_decision_payload_for_candidate
execution_strategy,
conversion_mode,
report_kind,
envelope_name: _,
transport,
} = resolved;

View File

@@ -20,9 +20,16 @@ use crate::ai_pipeline::transport::auth::{
build_openai_passthrough_headers, ensure_upstream_auth_header,
resolve_local_openai_bearer_auth,
};
use crate::ai_pipeline::transport::kiro::{
build_kiro_provider_headers, build_kiro_provider_request_body, KiroProviderHeadersInput,
KiroRequestAuth, KIRO_ENVELOPE_NAME,
};
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};
use crate::ai_pipeline::{
ConversionMode, ExecutionStrategy, GatewayProviderTransportSnapshot,
LocalResolvedOAuthRequestAuth,
};
use crate::AppState;
use super::support::{mark_skipped_local_openai_chat_candidate, LocalOpenAiChatDecisionInput};
@@ -38,6 +45,7 @@ pub(crate) struct LocalOpenAiChatCandidatePayloadParts {
pub(super) execution_strategy: ExecutionStrategy,
pub(super) conversion_mode: ConversionMode,
pub(super) report_kind: String,
pub(super) envelope_name: Option<&'static str>,
pub(super) transport: Arc<GatewayProviderTransportSnapshot>,
}
@@ -194,6 +202,7 @@ pub(crate) async fn resolve_local_openai_chat_candidate_payload_parts(
execution_strategy: ExecutionStrategy::LocalSameFormat,
conversion_mode: ConversionMode::None,
report_kind: report_kind.to_string(),
envelope_name: None,
transport: Arc::clone(transport),
});
}
@@ -232,32 +241,92 @@ pub(crate) async fn resolve_local_openai_chat_candidate_payload_parts(
.await;
return None;
}
let prepared_candidate = match prepare_header_authenticated_candidate(
planner_state,
transport,
candidate,
request_conversion_direct_auth(transport, conversion_kind),
OauthPreparationContext {
trace_id,
api_format: provider_api_format.as_str(),
operation: "openai_chat_cross_format",
},
)
.await
{
Ok(prepared) => prepared,
Err(skip_reason) => {
mark_skipped_local_openai_chat_candidate(
state,
input,
trace_id,
let is_kiro_claude_cli = transport
.provider
.provider_type
.trim()
.eq_ignore_ascii_case("kiro")
&& provider_api_format.eq_ignore_ascii_case("claude:cli");
let oauth_context = OauthPreparationContext {
trace_id,
api_format: provider_api_format.as_str(),
operation: "openai_chat_cross_format",
};
let kiro_auth = if is_kiro_claude_cli {
match crate::ai_pipeline::planner::candidate_preparation::resolve_candidate_oauth_auth(
planner_state,
transport,
oauth_context,
)
.await
{
Some(LocalResolvedOAuthRequestAuth::Kiro(auth)) => Some(auth),
_ => {
mark_skipped_local_openai_chat_candidate(
state,
input,
trace_id,
candidate,
candidate_index,
candidate_id,
"transport_auth_unavailable",
)
.await;
return None;
}
}
} else {
None
};
let prepared_candidate = if let Some(kiro_auth) = kiro_auth.as_ref() {
let mapped_model =
match crate::ai_pipeline::planner::candidate_preparation::resolve_candidate_mapped_model(
candidate,
candidate_index,
candidate_id,
skip_reason,
)
.await;
return None;
) {
Ok(mapped_model) => mapped_model,
Err(skip_reason) => {
mark_skipped_local_openai_chat_candidate(
state,
input,
trace_id,
candidate,
candidate_index,
candidate_id,
skip_reason,
)
.await;
return None;
}
};
crate::ai_pipeline::planner::candidate_preparation::PreparedHeaderAuthenticatedCandidate {
auth_header: kiro_auth.name.to_string(),
auth_value: kiro_auth.value.clone(),
mapped_model,
}
} else {
match prepare_header_authenticated_candidate(
planner_state,
transport,
candidate,
request_conversion_direct_auth(transport, conversion_kind),
oauth_context,
)
.await
{
Ok(prepared) => prepared,
Err(skip_reason) => {
mark_skipped_local_openai_chat_candidate(
state,
input,
trace_id,
candidate,
candidate_index,
candidate_id,
skip_reason,
)
.await;
return None;
}
}
};
@@ -267,7 +336,11 @@ pub(crate) async fn resolve_local_openai_chat_candidate_payload_parts(
transport.provider.provider_type.as_str(),
provider_api_format.as_str(),
upstream_is_stream,
transport.endpoint.body_rules.as_ref(),
if is_kiro_claude_cli {
None
} else {
transport.endpoint.body_rules.as_ref()
},
Some(input.auth_context.api_key_id.as_str()),
) else {
mark_skipped_local_openai_chat_candidate(
@@ -283,6 +356,29 @@ pub(crate) async fn resolve_local_openai_chat_candidate_payload_parts(
return None;
};
if let Some(kiro_auth) = kiro_auth.as_ref() {
return build_kiro_openai_chat_cross_format_payload_parts(
state,
parts,
trace_id,
body_json,
input,
eligible,
candidate_index,
candidate_id,
decision_kind,
transport,
provider_api_format.as_str(),
prepared_candidate.mapped_model,
prepared_candidate.auth_header,
prepared_candidate.auth_value,
provider_request_body,
upstream_is_stream,
kiro_auth,
)
.await;
}
let Some(upstream_url) = build_cross_format_openai_chat_upstream_url(
parts,
transport,
@@ -392,6 +488,119 @@ pub(crate) async fn resolve_local_openai_chat_candidate_payload_parts(
execution_strategy: ExecutionStrategy::LocalCrossFormat,
conversion_mode: ConversionMode::Bidirectional,
report_kind: resolved_report_kind,
envelope_name: None,
transport: Arc::clone(transport),
})
}
#[allow(clippy::too_many_arguments)]
async fn build_kiro_openai_chat_cross_format_payload_parts(
state: &AppState,
parts: &http::request::Parts,
trace_id: &str,
original_body_json: &serde_json::Value,
input: &LocalOpenAiChatDecisionInput,
eligible: &EligibleLocalExecutionCandidate,
candidate_index: u32,
candidate_id: &str,
decision_kind: &str,
transport: &Arc<GatewayProviderTransportSnapshot>,
provider_api_format: &str,
mapped_model: String,
auth_header: String,
auth_value: String,
claude_request_body: Value,
upstream_is_stream: bool,
kiro_auth: &KiroRequestAuth,
) -> Option<LocalOpenAiChatCandidatePayloadParts> {
let candidate = &eligible.candidate;
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(body) => body,
None => {
mark_skipped_local_openai_chat_candidate(
state,
input,
trace_id,
candidate,
candidate_index,
candidate_id,
"provider_request_body_missing",
)
.await;
return None;
}
};
let upstream_url = match crate::ai_pipeline::build_provider_transport_request_url(
transport,
provider_api_format,
Some(&mapped_model),
upstream_is_stream,
parts.uri.query(),
Some(kiro_auth.auth_config.effective_api_region()),
) {
Some(url) => url,
None => {
mark_skipped_local_openai_chat_candidate(
state,
input,
trace_id,
candidate,
candidate_index,
candidate_id,
"upstream_url_missing",
)
.await;
return None;
}
};
let provider_request_headers = match build_kiro_provider_headers(KiroProviderHeadersInput {
headers: &parts.headers,
provider_request_body: &provider_request_body,
original_request_body: original_body_json,
header_rules: transport.endpoint.header_rules.as_ref(),
auth_header: &auth_header,
auth_value: &auth_value,
auth_config: &kiro_auth.auth_config,
machine_id: kiro_auth.machine_id.as_str(),
}) {
Some(headers) => headers,
None => {
mark_skipped_local_openai_chat_candidate(
state,
input,
trace_id,
candidate,
candidate_index,
candidate_id,
"transport_header_rules_apply_failed",
)
.await;
return None;
}
};
let resolved_report_kind = if decision_kind == OPENAI_CHAT_STREAM_PLAN_KIND {
"openai_chat_stream_success".to_string()
} else {
"openai_chat_sync_finalize".to_string()
};
Some(LocalOpenAiChatCandidatePayloadParts {
auth_header,
auth_value,
mapped_model,
provider_api_format: provider_api_format.to_string(),
provider_request_body,
provider_request_headers,
upstream_url,
execution_strategy: ExecutionStrategy::LocalCrossFormat,
conversion_mode: ConversionMode::Bidirectional,
report_kind: resolved_report_kind,
envelope_name: Some(KIRO_ENVELOPE_NAME),
transport: Arc::clone(transport),
})
}

View File

@@ -70,8 +70,8 @@ pub(crate) async fn maybe_build_local_openai_cli_decision_payload_for_candidate(
{
extra_fields.insert("proxy".to_string(), proxy_value);
}
if resolved.is_antigravity {
extra_fields.insert("envelope_name".to_string(), json!("antigravity:v1internal"));
if let Some(envelope_name) = resolved.envelope_name {
extra_fields.insert("envelope_name".to_string(), json!(envelope_name));
}
let report_context = append_local_failover_policy_to_value(
append_execution_contract_fields_to_value(
@@ -104,7 +104,7 @@ pub(crate) async fn maybe_build_local_openai_cli_decision_payload_for_candidate(
.and_then(serde_json::Value::as_bool)
.unwrap_or(false),
upstream_is_stream: resolved.upstream_is_stream,
has_envelope: resolved.is_antigravity,
has_envelope: resolved.envelope_name.is_some(),
needs_conversion: matches!(
resolved.conversion_mode,
crate::ai_pipeline::ConversionMode::Bidirectional
@@ -139,7 +139,7 @@ pub(crate) async fn maybe_build_local_openai_cli_decision_payload_for_candidate(
upstream_base_url = %resolved.transport.endpoint.base_url,
upstream_url = %resolved.upstream_url,
upstream_is_stream = resolved.upstream_is_stream,
has_envelope = resolved.is_antigravity,
has_envelope = resolved.envelope_name.is_some(),
"gateway built local openai cli decision payload"
);
let super::request::LocalOpenAiCliCandidatePayloadParts {
@@ -153,6 +153,7 @@ pub(crate) async fn maybe_build_local_openai_cli_decision_payload_for_candidate(
execution_strategy,
conversion_mode,
is_antigravity: _,
envelope_name: _,
upstream_is_stream,
transport,
} = resolved;

View File

@@ -27,10 +27,17 @@ use crate::ai_pipeline::transport::auth::{
build_openai_passthrough_headers, ensure_upstream_auth_header, resolve_local_gemini_auth,
resolve_local_openai_bearer_auth, resolve_local_standard_auth,
};
use crate::ai_pipeline::transport::kiro::{
build_kiro_provider_headers, build_kiro_provider_request_body,
local_kiro_request_transport_unsupported_reason_with_network, KiroProviderHeadersInput,
KiroRequestAuth, KIRO_ENVELOPE_NAME,
};
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::{GatewayProviderTransportSnapshot, PlannerAppState};
use crate::ai_pipeline::{
GatewayProviderTransportSnapshot, LocalResolvedOAuthRequestAuth, PlannerAppState,
};
use crate::AppState;
use super::support::{mark_skipped_local_openai_cli_candidate, LocalOpenAiCliDecisionInput};
@@ -49,6 +56,7 @@ pub(crate) struct LocalOpenAiCliCandidatePayloadParts {
pub(super) execution_strategy: ExecutionStrategy,
pub(super) conversion_mode: ConversionMode,
pub(super) is_antigravity: bool,
pub(super) envelope_name: Option<&'static str>,
pub(super) upstream_is_stream: bool,
pub(super) transport: Arc<GatewayProviderTransportSnapshot>,
}
@@ -76,10 +84,18 @@ pub(crate) async fn resolve_local_openai_cli_candidate_payload_parts(
.provider_type
.trim()
.eq_ignore_ascii_case("antigravity");
let is_kiro_claude_cli = transport
.provider
.provider_type
.trim()
.eq_ignore_ascii_case("kiro")
&& provider_api_format.eq_ignore_ascii_case("claude:cli");
let same_format = provider_api_format == client_api_format;
let conversion_kind = request_conversion_kind(spec_metadata.api_format, provider_api_format);
let transport_unsupported_reason = if same_format {
let transport_unsupported_reason = if same_format && is_kiro_claude_cli {
local_kiro_request_transport_unsupported_reason_with_network(transport)
} else if same_format {
local_standard_transport_unsupported_reason_with_network(transport, provider_api_format)
} else {
match conversion_kind {
@@ -106,7 +122,41 @@ pub(crate) async fn resolve_local_openai_cli_candidate_payload_parts(
return None;
}
let direct_auth = if same_format {
let oauth_context = OauthPreparationContext {
trace_id,
api_format: provider_api_format,
operation: "openai_cli_candidate_request",
};
let kiro_auth = if is_kiro_claude_cli {
match crate::ai_pipeline::planner::candidate_preparation::resolve_candidate_oauth_auth(
planner_state,
transport,
oauth_context,
)
.await
{
Some(LocalResolvedOAuthRequestAuth::Kiro(auth)) => Some(auth),
_ => {
mark_skipped_local_openai_cli_candidate(
state,
input,
trace_id,
candidate,
candidate_index,
candidate_id,
"transport_auth_unavailable",
)
.await;
return None;
}
}
} else {
None
};
let direct_auth = if kiro_auth.is_some() {
None
} else if same_format {
match provider_api_format {
"gemini:cli" => resolve_local_gemini_auth(transport),
"claude:cli" => resolve_local_standard_auth(transport),
@@ -116,32 +166,55 @@ pub(crate) async fn resolve_local_openai_cli_candidate_payload_parts(
} else {
conversion_kind.and_then(|kind| request_conversion_direct_auth(transport, kind))
};
let prepared_candidate = match prepare_header_authenticated_candidate(
planner_state,
transport,
candidate,
direct_auth,
OauthPreparationContext {
trace_id,
api_format: provider_api_format,
operation: "openai_cli_candidate_request",
},
)
.await
{
Ok(prepared) => prepared,
Err(skip_reason) => {
mark_skipped_local_openai_cli_candidate(
state,
input,
trace_id,
let prepared_candidate = if let Some(kiro_auth) = kiro_auth.as_ref() {
let mapped_model =
match crate::ai_pipeline::planner::candidate_preparation::resolve_candidate_mapped_model(
candidate,
candidate_index,
candidate_id,
skip_reason,
)
.await;
return None;
) {
Ok(mapped_model) => mapped_model,
Err(skip_reason) => {
mark_skipped_local_openai_cli_candidate(
state,
input,
trace_id,
candidate,
candidate_index,
candidate_id,
skip_reason,
)
.await;
return None;
}
};
crate::ai_pipeline::planner::candidate_preparation::PreparedHeaderAuthenticatedCandidate {
auth_header: kiro_auth.name.to_string(),
auth_value: kiro_auth.value.clone(),
mapped_model,
}
} else {
match prepare_header_authenticated_candidate(
planner_state,
transport,
candidate,
direct_auth,
oauth_context,
)
.await
{
Ok(prepared) => prepared,
Err(skip_reason) => {
mark_skipped_local_openai_cli_candidate(
state,
input,
trace_id,
candidate,
candidate_index,
candidate_id,
skip_reason,
)
.await;
return None;
}
}
};
let auth_header = prepared_candidate.auth_header;
@@ -163,7 +236,11 @@ pub(crate) async fn resolve_local_openai_cli_candidate_payload_parts(
provider_api_format,
upstream_is_stream,
transport.provider.provider_type.as_str(),
transport.endpoint.body_rules.as_ref(),
if is_kiro_claude_cli {
None
} else {
transport.endpoint.body_rules.as_ref()
},
Some(input.auth_context.api_key_id.as_str()),
)
} else {
@@ -173,7 +250,11 @@ pub(crate) async fn resolve_local_openai_cli_candidate_payload_parts(
upstream_is_stream,
transport.provider.provider_type.as_str(),
provider_api_format,
transport.endpoint.body_rules.as_ref(),
if is_kiro_claude_cli {
None
} else {
transport.endpoint.body_rules.as_ref()
},
Some(input.auth_context.api_key_id.as_str()),
)
}) else {
@@ -240,6 +321,30 @@ pub(crate) async fn resolve_local_openai_cli_candidate_payload_parts(
base_provider_request_body
};
if let Some(kiro_auth) = kiro_auth.as_ref() {
return build_kiro_openai_cli_payload_parts(
state,
parts,
trace_id,
body_json,
input,
eligible,
candidate_index,
candidate_id,
spec_metadata.api_format,
transport,
provider_api_format,
mapped_model,
auth_header,
auth_value,
provider_request_body,
upstream_is_stream,
needs_bidirectional_conversion,
kiro_auth,
)
.await;
}
let Some(upstream_url) = (if needs_bidirectional_conversion {
build_cross_format_openai_cli_upstream_url(
parts,
@@ -392,6 +497,149 @@ pub(crate) async fn resolve_local_openai_cli_candidate_payload_parts(
conversion_mode,
is_antigravity: is_antigravity
|| antigravity_auth.is_some() && ANTIGRAVITY_ENVELOPE_NAME == "antigravity:v1internal",
envelope_name: if is_antigravity || antigravity_auth.is_some() {
Some(ANTIGRAVITY_ENVELOPE_NAME)
} else {
None
},
upstream_is_stream,
transport: Arc::clone(transport),
})
}
#[allow(clippy::too_many_arguments)]
async fn build_kiro_openai_cli_payload_parts(
state: &AppState,
parts: &http::request::Parts,
trace_id: &str,
original_body_json: &serde_json::Value,
input: &LocalOpenAiCliDecisionInput,
eligible: &EligibleLocalExecutionCandidate,
candidate_index: u32,
candidate_id: &str,
client_api_format: &str,
transport: &Arc<GatewayProviderTransportSnapshot>,
provider_api_format: &str,
mapped_model: String,
auth_header: String,
auth_value: String,
claude_request_body: Value,
upstream_is_stream: bool,
needs_bidirectional_conversion: bool,
kiro_auth: &KiroRequestAuth,
) -> Option<LocalOpenAiCliCandidatePayloadParts> {
let candidate = &eligible.candidate;
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(body) => body,
None => {
mark_skipped_local_openai_cli_candidate(
state,
input,
trace_id,
candidate,
candidate_index,
candidate_id,
"provider_request_body_missing",
)
.await;
return None;
}
};
let upstream_url = match crate::ai_pipeline::build_provider_transport_request_url(
transport,
provider_api_format,
Some(&mapped_model),
upstream_is_stream,
parts.uri.query(),
Some(kiro_auth.auth_config.effective_api_region()),
) {
Some(url) => url,
None => {
mark_skipped_local_openai_cli_candidate(
state,
input,
trace_id,
candidate,
candidate_index,
candidate_id,
"upstream_url_missing",
)
.await;
return None;
}
};
let provider_request_headers = match build_kiro_provider_headers(KiroProviderHeadersInput {
headers: &parts.headers,
provider_request_body: &provider_request_body,
original_request_body: original_body_json,
header_rules: transport.endpoint.header_rules.as_ref(),
auth_header: &auth_header,
auth_value: &auth_value,
auth_config: &kiro_auth.auth_config,
machine_id: kiro_auth.machine_id.as_str(),
}) {
Some(headers) => headers,
None => {
mark_skipped_local_openai_cli_candidate(
state,
input,
trace_id,
candidate,
candidate_index,
candidate_id,
"transport_header_rules_apply_failed",
)
.await;
return None;
}
};
let execution_strategy = if needs_bidirectional_conversion {
ExecutionStrategy::LocalCrossFormat
} else {
ExecutionStrategy::LocalSameFormat
};
let conversion_mode = if needs_bidirectional_conversion {
ConversionMode::Bidirectional
} else {
ConversionMode::None
};
debug!(
event_name = "local_openai_cli_kiro_upstream_url_resolved",
log_type = "debug",
trace_id = %trace_id,
candidate_id = %candidate_id,
candidate_index,
provider_id = %candidate.provider_id,
endpoint_id = %candidate.endpoint_id,
key_id = %candidate.key_id,
provider_type = %transport.provider.provider_type,
client_api_format = client_api_format,
provider_api_format = %provider_api_format,
execution_strategy = execution_strategy.as_str(),
conversion_mode = conversion_mode.as_str(),
upstream_url = %upstream_url,
upstream_is_stream,
"gateway resolved local openai cli kiro upstream url"
);
Some(LocalOpenAiCliCandidatePayloadParts {
auth_header,
auth_value,
mapped_model,
provider_api_format: provider_api_format.to_string(),
provider_request_body,
provider_request_headers,
upstream_url,
execution_strategy,
conversion_mode,
is_antigravity: false,
envelope_name: Some(KIRO_ENVELOPE_NAME),
upstream_is_stream,
transport: Arc::clone(transport),
})